Skip to main content
PYSPARK • LESSON 275

Tumbling Time Window Aggregations

How can we aggregate real-time transaction counts into non-overlapping, contiguous 5-minute time windows?

Production3 Minutes840 XP
🤔 THE QUESTION

How can we aggregate real-time transaction counts into non-overlapping, contiguous 5-minute time windows?

💡 WHAT IS IT?

groupBy(window('event_timestamp', '5 minutes'), 'store_id') groups events into fixed, non-overlapping temporal buckets.

🎯 WHAT IS IT USED FOR?

Real-time store traffic monitoring, infrastructure error rate tracking, and 5-minute financial volume metrics.

💻 EXAMPLE
df_tumbling = df_watermarked.groupBy(window(col("event_timestamp"), "5 minutes"), col("store_id")).agg(sum("amount").alias("window_sales"), count("transaction_id").alias("tx_count"))

🎯 Mission Objectives

Practice typing production-grade PySpark code for Tumbling Time Window Aggregations.

  • window() tumbling specification
  • Non-overlapping time buckets
  • Streaming window aggregation