視窗函式與串流查詢

在 Databricks 使用 Spark SQL 進行資料轉換

Disha Mukherjee

Lead Data Engineer

當 groupBy 不夠用時

recraft: half: 放大鏡置於試算表列上,部分計算值已標示

 

  • groupBy() → 每組一列(總計、平均)
  • 視窗函式 → 加一欄,保留所有列
  • 用途: 累計、排名、列間比較
在 Databricks 使用 Spark SQL 進行資料轉換

累計總額(Running total)

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 進行資料轉換

累計總額(Running total)

+-----------+-------------------+------------------+-------------+
|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: 連續資料流以箭頭流入處理引擎

在 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 檔案目錄 = 串流來源
  • 每個新檔 = 一個微批次
  • 會自動擷取新檔
  • 結構(schema)需明確定義
在 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 進行資料轉換

什麼是檢查點(checkpoint)?

recraft: half: 流動資料管線中的書籤或儲存進度標記

 

$$

  • 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...