윈도우 함수와 스트리밍 쿼리

Databricks에서 Spark SQL로 데이터 변환하기

Disha Mukherjee

Lead Data Engineer

`groupBy()`만으로는 부족할 때

recraft: half: A magnifying glass over rows of spreadsheet data with highlighted calculated values

 

  • groupBy() → 그룹당 한 행 반환 (합계, 평균)
  • 윈도우 함수 → 열을 추가하고 모든 행 유지
  • 활용 사례: 누계, 순위, 행 비교
Databricks에서 Spark SQL로 데이터 변환하기

누계 계산

from pyspark.sql.window import Window

window_spec = (
    Window.partitionBy("Customer_ID")
    .orderBy("Date")
    .rowsBetween(Window.unboundedPreceding, Window.currentRow)
)

df_running = df_valid.withColumn( "running_total", F.round(F.sum("Transaction_Amount").over(window_spec), 2) )
Databricks에서 Spark SQL로 데이터 변환하기

누계 계산

+-----------+-------------------+------------------+-------------+
|Customer_ID|Date               |Transaction_Amount|running_total|
+-----------+-------------------+------------------+-------------+
|CUST003    |2023-01-03 00:00:00|5752.36           |5752.36      |
|CUST003    |2023-02-14 00:00:00|12340.00          |18092.36     |
+-----------+-------------------+------------------+-------------+
Databricks에서 Spark SQL로 데이터 변환하기

고객 순위 매기기

customer_totals = (
    df_valid.groupBy("Customer_ID")
    .agg(F.round(F.sum("Transaction_Amount"), 2).alias("total_revenue"))
)


rank_window = Window.orderBy(F.col("total_revenue").desc())
df_ranked = customer_totals.withColumn("revenue_rank", F.rank().over(rank_window))
+-----------+-------------+------------+
|Customer_ID|total_revenue|revenue_rank|
+-----------+-------------+------------+
|CUST1469   |99996.20     |1           |
|CUST19129  |99990.14     |2           |
|CUST39417  |99982.21     |3           |
+-----------+-------------+------------+
Databricks에서 Spark SQL로 데이터 변환하기

스트리밍이란?

 

  • 배치(Batch) - 고정된 데이터셋을 한 번에 처리
  • 스트리밍(Streaming) - 데이터가 도착하는 대로 점진적으로 처리
  • 트랜잭션 피드, 로그, 이벤트 데이터에 적합

recraft: half: Data flowing continuously as a stream into a processing engine with arrows

Databricks에서 Spark SQL로 데이터 변환하기

파일 기반 스트리밍

$$

Stream directory:
  day_1.csv  (108 KB)
  day_2.csv  (108 KB)
  day_3.csv  (108 KB)
  day_4.csv  (108 KB)
  day_5.csv  (108 KB)

 

  • CSV 파일 디렉터리 = 스트리밍 소스
  • 새 파일 하나 = 마이크로 배치 하나
  • 새 파일은 자동으로 감지
  • 스키마는 명시적으로 정의해야 함
Databricks에서 Spark SQL로 데이터 변환하기

스트림 읽기

df_stream = (
    spark.readStream.format("csv")
    .option("header", "true")
    .schema(streaming_schema)
    .load(STREAM_DIR)
)

print(df_stream.isStreaming)
True
Databricks에서 Spark SQL로 데이터 변환하기

체크포인트란?

recraft: half: A bookmark or save point marker in a flowing data pipeline with progress saved

 

$$

  • 체크포인트(Checkpoint) = 디스크의 메타데이터 디렉터리
  • 처리된 파일 목록을 기록
  • 재시작 시 Spark가 중단된 지점부터 재개
  • 중복 처리 방지
Databricks에서 Spark SQL로 데이터 변환하기

스트림 쓰기

query = (
    df_stream.writeStream.format("delta")
    .outputMode("append")
    .option("checkpointLocation", CHECKPOINT_DIR)
    .option("path", DELTA_PATH)
    .trigger(availableNow=True)
    .start()
)
query.awaitTermination()
Status:       Stopped
Rows written: 5,000
Databricks에서 Spark SQL로 데이터 변환하기

모니터링: status와 lastProgress

print(query.status)

progress = query.lastProgress
print(f"Rows processed: {progress['numInputRows']}")
print(f"Rows/sec: {progress['processedRowsPerSecond']:.0f}")
{'message': 'Stopped', 'isDataAvailable': False, 'isTriggerActive': False}
Rows processed: 5,000
Rows/sec:       752
Databricks에서 Spark SQL로 데이터 변환하기

체크포인트 복구

query_restart = (
    df_stream.writeStream.format("delta")
    .option("checkpointLocation", CHECKPOINT_DIR)
    .option("path", DELTA_PATH)
    .trigger(availableNow=True)
    .start()
)
query_restart.awaitTermination()
Rows on restart: 0
Databricks에서 Spark SQL로 데이터 변환하기

연습해 봅시다!

Databricks에서 Spark SQL로 데이터 변환하기

Preparing Video For Download...