← notes

스파크 스트리밍 쿼리 101

2025-08-11 · spark

몇 년 동안이나 pyspark로 대부분의 데이터 처리를 해왔지만 스트리밍은 처음이라 기초적인 내용들 정리해봄

df.groupBy(F.window("event_time", "10 minutes")).agg(F.count("*").alias("cnt")) #tumbling
df.groupBy(F.window("event_time", "10 minutes", "5 minutes")).agg(F.count("*")) #sliding
df.groupBy(F.session_window("event_time", "10 minutes"), "user_id").agg(F.count("*")) #session
df.withWatermark(time_column, threshold)
#e.g. df.withWatermark("event_time", "5 minutes")

static이냐 stream이냐 누가 왼쪽이냐에 따라 지원여부 달라짐. steam-stream join은 제약이 강함..

https://spark.apache.org/docs/latest/streaming/index.html