การรวมกลุ่มและ join ข้อมูลอย่างมีประสิทธิภาพ

การแปลงข้อมูลด้วย Spark SQL ใน Databricks

Disha Mukherjee

Lead Data Engineer

สามคำถาม หนึ่งชุดข้อมูลที่สะอาด

recraft: half: นักวิเคราะห์ข้อมูลกำลังดูกราฟและรายงานธุรกิจหลากสีบนหลายจอ บนพื้นหลังโปร่งใส

 

$$

  • หมวดหมู่ไหนสร้างรายได้มากที่สุด?
  • ใครคือลูกค้าชั้นนำ?
  • จะเสริมบริบทด้วยตาราง department อย่างไร?
การแปลงข้อมูลด้วย Spark SQL ใน Databricks

`groupBy()` และ `agg()` ทำงานอย่างไร

 

$$

  • groupBy() - แบ่งแถวตามคอลัมน์คีย์
  • agg() - ใช้ฟังก์ชันกับแต่ละกลุ่มแบบ parallel
  • ทั้งคู่เป็น lazy — action เป็นตัวกระตุ้นการคำนวณ
  • เรียก agg() ครั้งเดียว = สแกนข้อมูลรอบเดียว

recraft: half: แถวข้อมูลหลากสีถูกแบ่งออกเป็นสามกลุ่มที่มีป้ายกำกับ พร้อมค่ารวมแสดงใต้แต่ละกลุ่ม บนพื้นหลังโปร่งใส

การแปลงข้อมูลด้วย Spark SQL ใน Databricks

รายได้ตามหมวดหมู่ - โค้ด

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())
)
การแปลงข้อมูลด้วย Spark SQL ใน Databricks

รายได้ตามหมวดหมู่ - ผลลัพธ์

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       |
+-----------+-----------------+-----------------+---------------+
การแปลงข้อมูลด้วย Spark SQL ใน Databricks

เสริมข้อมูลด้วย dimension table

recraft: half: ตารางฐานข้อมูลสองตารางถูกเชื่อมต่อกันด้วยลิงก์หรือสะพานที่เรืองแสง บนพื้นหลังโปร่งใส

$$

+-----------+----------+
|Category   |Department|
+-----------+----------+
|Clothing   |Retail    |
|Dining     |Food      |
|Electronics|Tech      |
|Groceries  |Food      |
|Savings    |Finance   |
+-----------+----------+
การแปลงข้อมูลด้วย Spark SQL ใน Databricks

Standard left join

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          |
+-----------+-----------+----------+------------------+
การแปลงข้อมูลด้วย Spark SQL ใน Databricks

ต้นทุนซ่อนเร้น - shuffle

recraft: half: แพ็กเก็ตข้อมูลไหลผ่านเครือข่ายของเซิร์ฟเวอร์และโหนดที่เชื่อมต่อกัน พร้อมลูกศรแสดงทิศทางการเคลื่อนที่ บนพื้นหลังโปร่งใส

 

  • คีย์ที่ตรงกันต้องอยู่บนเครื่องเดียวกัน
  • Spark ส่งแถวข้ามเครือข่ายเพื่อจัดเรียงข้อมูล
  • Shuffle = คอขวดเมื่อข้อมูลขนาดใหญ่
การแปลงข้อมูลด้วย Spark SQL ใน Databricks

Spark UI

recraft: full: แดชบอร์ดติดตามการวิเคราะห์โทนสีเข้มที่มีแผนภูมิแท่งแนวนอนหลากสี แสดงขั้นตอน job แถบความคืบหน้า และเมตริกปริมาณข้อมูล บนพื้นหลังโปร่งใส

 

  • Jobs และ Stages - ติดตามสิ่งที่รันไปและระยะเวลา
  • Shuffle Read/Write - ดูปริมาณข้อมูลที่ถูกย้าย
  • ใช้วินิจฉัยคอขวดของ pipeline
การแปลงข้อมูลด้วย Spark SQL ใน Databricks

อ่าน query plan ด้วย `.explain()`

df_joined.explain(mode="formatted")
...
+- PhotonBroadcastHashJoin LeftOuter (18)
   :- ...
   +- PhotonShuffleExchangeSource (17)
      +- PhotonShuffleMapStage (16)
         +- PhotonShuffleExchangeSink (15)
            +- LocalTableScan (13)
...
  • .explain() - ตรวจสอบ execution plan
  • Node 15 & 16 - ไม่มี routing Arguments ถูกบันทึก
การแปลงข้อมูลด้วย Spark SQL ใน Databricks

Broadcast join - วิธีแก้ปัญหา

# Wrap the small table in 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
การแปลงข้อมูลด้วย Spark SQL ใน Databricks

ก่อนและหลัง - ตรวจสอบ plan

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]
  • ตาราง dim มีคำสั่ง routing แล้ว
  • ตารางขนาดใหญ่ df_valid ไม่ถูกย้ายเลย
  • เปลี่ยนแค่บรรทัดเดียว = ลดเวลาจากนาทีเหลือวินาที
การแปลงข้อมูลด้วย Spark SQL ใน Databricks

มาฝึกกันเถอะ!

การแปลงข้อมูลด้วย Spark SQL ใน Databricks

Preparing Video For Download...