Hàm cửa sổ và truy vấn streaming

Chuyển đổi dữ liệu với Spark SQL trong Databricks

Disha Mukherjee

Lead Data Engineer

Khi groupBy là chưa đủ

recraft: half: Kính lúp soi các hàng dữ liệu bảng tính với các giá trị tính toán được tô sáng

 

  • groupBy() → một hàng cho mỗi nhóm (tổng, trung bình)
  • Hàm cửa sổ → thêm cột, giữ mọi hàng
  • Dùng để: tổng lũy kế, xếp hạng, so sánh hàng
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Tổng lũy kế

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) )
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Tổng lũy kế

+-----------+-------------------+------------------+-------------+
|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     |
+-----------+-------------------+------------------+-------------+
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Xếp hạng khách hàng

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           |
+-----------+-------------+------------+
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Streaming là gì?

 

  • Batch - xử lý trọn một tập dữ liệu cố định
  • Streaming - xử lý dần dữ liệu khi đến
  • Lý tưởng cho luồng giao dịch, log và dữ liệu sự kiện

recraft: half: Dữ liệu chảy liên tục như dòng suối vào một engine xử lý với các mũi tên

Chuyển đổi dữ liệu với Spark SQL trong Databricks

Streaming dựa trên file

$$

Thư mục stream:
  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)

 

  • Thư mục các file CSV = nguồn streaming
  • Mỗi file mới = một micro-batch
  • Tự động nhận file mới
  • Phải khai báo schema rõ ràng
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Đọc stream

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

print(df_stream.isStreaming)
True
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Checkpoint là gì?

recraft: half: Dấu trang hoặc điểm lưu trong một pipeline dữ liệu đang chảy với tiến độ được lưu

 

$$

  • Checkpoint = thư mục metadata trên đĩa
  • Ghi nhận file nào đã xử lý
  • Khi khởi động lại, Spark tiếp tục từ điểm dừng
  • Ngăn xử lý trùng lặp
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Ghi stream

query = (
    df_stream.writeStream.format("delta")
    .outputMode("append")
    .option("checkpointLocation", CHECKPOINT_DIR)
    .option("path", DELTA_PATH)
    .trigger(availableNow=True)
    .start()
)
query.awaitTermination()
Trạng thái:   Đã dừng
Số hàng ghi:  5,000
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Giám sát: status và 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
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Phục hồi từ checkpoint

query_restart = (
    df_stream.writeStream.format("delta")
    .option("checkpointLocation", CHECKPOINT_DIR)
    .option("path", DELTA_PATH)
    .trigger(availableNow=True)
    .start()
)
query_restart.awaitTermination()
Số hàng khi khởi động lại: 0
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Cùng luyện tập nào!

Chuyển đổi dữ liệu với Spark SQL trong Databricks

Preparing Video For Download...