Tổng hợp và join dữ liệu hiệu quả

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

Disha Mukherjee

Lead Data Engineer

Ba câu hỏi, một bộ dữ liệu sạch

recraft: half: Một nhà phân tích dữ liệu đang xem các biểu đồ đầy màu sắc và báo cáo kinh doanh trên nhiều màn hình, nền trong suốt

 

$$

  • Những danh mục nào tạo ra doanh thu nhiều nhất?
  • Ai là khách hàng hàng đầu?
  • Làm sao làm giàu kết quả với ngữ cảnh phòng ban?
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Cách groupBy() và agg() hoạt động

 

$$

  • groupBy() - phân chia dòng theo cột khóa
  • agg() - áp dụng hàm cho từng nhóm song song
  • Cả hai đều lười - cần action để chạy tính toán
  • Một lần gọi agg() = quét dữ liệu một lần

recraft: half: Các hàng dữ liệu đầy màu sắc được chia thành ba nhóm có nhãn riêng, kèm giá trị tổng dưới mỗi nhóm, nền trong suốt

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

Doanh thu theo danh mục - mã nguồn

category_revenue = (
    df_valid
    .groupBy("Category")

.agg( F.round(F.sum("Transaction_Amount"), 2).alias("total_revenue"), F.count("Transaction_Amount").alias("transaction_count"), F.round(F.avg("Transaction_Amount"), 2).alias("avg_transaction"), )
.orderBy(F.col("total_revenue").desc())
)
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Doanh thu theo danh mục - đầu ra

category_revenue.show(truncate=False)
+-----------+-----------------+-----------------+---------------+
|Category   |total_revenue    |transaction_count|avg_transaction|
+-----------+-----------------+-----------------+---------------+
|Clothing   |338266605.24     |6742             |50173.04       |
|Dining     |337631718.64     |6719             |50250.29       |
|Electronics|331102843.72     |6614             |50060.91       |
|Savings    |329346609.07     |6596             |49931.26       |
|Groceries  |329003451.57     |6537             |50329.43       |
|Unknown    |880591.27        |15               |58706.08       |
+-----------+-----------------+-----------------+---------------+
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Làm giàu bằng bảng chiều (dimension)

recraft: half: Hai bảng cơ sở dữ liệu được nối với nhau bằng một liên kết phát sáng, nền trong suốt

$$

+-----------+----------+
|Category   |Department|
+-----------+----------+
|Clothing   |Retail    |
|Dining     |Food      |
|Electronics|Tech      |
|Groceries  |Food      |
|Savings    |Finance   |
+-----------+----------+
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Left join chuẩn

df_joined = df_valid.join(df_dim, on="Category", how="left")

df_joined.select( "Customer_ID", "Category", "Department", "Transaction_Amount" ).show(5, truncate=False)
+-----------+-----------+----------+------------------+
|Customer_ID|Category   |Department|Transaction_Amount|
+-----------+-----------+----------+------------------+
|CUST003    |Electronics|Tech      |5752.36           |
|CUST009    |Clothing   |Retail    |28959.12          |
|CUST010    |Savings    |Finance   |72098.18          |
|CUST011    |Savings    |Finance   |49771.05          |
|CUST020    |Groceries  |Food      |69825.14          |
+-----------+-----------+----------+------------------+
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Chi phí ẩn - shuffle

recraft: half: Các gói dữ liệu di chuyển qua mạng các máy chủ và nút được kết nối với mũi tên chỉ hướng, nền trong suốt

 

  • Khóa khớp phải ở trên cùng một máy
  • Spark chuyển các dòng qua mạng để căn thẳng hàng
  • Shuffle = nút thắt cổ chai ở quy mô lớn
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Spark UI

recraft: full: Bảng điều khiển giám sát phân tích nền tối với biểu đồ thanh ngang nhiều màu thể hiện các giai đoạn job, thanh tiến độ và chỉ số thông lượng dữ liệu, nền trong suốt

 

  • JobsStages - theo dõi cái gì chạy và mất bao lâu
  • Shuffle Read/Write - xem lượng dữ liệu đã di chuyển
  • Hữu ích để chẩn đoán nút thắt đường ống
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Đọc kế hoạch truy vấn với .explain()

df_joined.explain(mode="formatted")
...
+- PhotonBroadcastHashJoin LeftOuter (18)
   :- ...
   +- PhotonShuffleExchangeSource (17)
      +- PhotonShuffleMapStage (16)
         +- PhotonShuffleExchangeSink (15)
            +- LocalTableScan (13)
...
  • .explain() - xem kế hoạch thực thi
  • Nút 15 & 16 - không có tham số định tuyến trong log
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Broadcast join - cách khắc phục

# Bọc bảng nhỏ bằng F.broadcast()
df_broadcast = df_valid.join(
    F.broadcast(df_dim),
    on="Category",
    how="left"
)

print(f"Broadcast joined rows: {df_broadcast.count():,}")
Broadcast joined rows: 33,223
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Trước và sau - xác minh kế hoạch

df_broadcast.explain(mode="formatted")
+- PhotonBroadcastHashJoin LeftOuter (18)
   :- ...
   +- PhotonShuffleExchangeSource (17)
      +- PhotonShuffleMapStage (16)
         +- PhotonShuffleExchangeSink (15)
            +- LocalTableScan (13)
...            
(15) PhotonShuffleExchangeSink
Arguments: SinglePartition
(16) PhotonShuffleMapStage
Arguments: EXECUTOR_BROADCAST, [id=#11839]
  • Bảng chiều có hướng dẫn định tuyến
  • Bảng lớn df_valid không bao giờ di chuyển
  • Một thay đổi = từ phút xuống giây ở quy mô lớn
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...