Skip to main content
PYSPARK • LESSON 276

Sliding (Hopping) Time Window Aggregations

How can we calculate 10-minute moving window aggregations that update every 2 minutes (overlapping windows)?

Production3 Minutes850 XP
🤔 THE QUESTION

How can we calculate 10-minute moving window aggregations that update every 2 minutes (overlapping windows)?

💡 WHAT IS IT?

window('event_timestamp', '10 minutes', '2 minutes') specifies a 10-minute window length sliding forward every 2 minutes.

🎯 WHAT IS IT USED FOR?

Smoothed real-time anomaly detection, fraud velocity scoring, and continuous moving averages on streams.

💻 EXAMPLE
df_sliding = df_watermarked.groupBy(window(col("event_timestamp"), "10 minutes", "2 minutes"), col("device_id")).agg(avg("cpu_usage").alias("avg_cpu"))

🎯 Mission Objectives

Practice typing production-grade PySpark code for Sliding (Hopping) Time Window Aggregations.

  • window() sliding specification
  • Overlapping hopping windows
  • Streaming velocity monitoring