डेटा को कुशलता से aggregate करें और join करें

Databricks में Spark SQL के साथ Data Transformation

Disha Mukherjee

Lead Data Engineer

तीन सवाल, एक साफ़ डेटासेट

recraft: half: एक डेटा विश्लेषक कई स्क्रीन पर रंगीन चार्ट और बिज़नेस रिपोर्ट की समीक्षा कर रहा है, ट्रांसपेरेंट बैकग्राउंड पर अलग-थलग

 

$$

  • कौन-सी categories सबसे ज़्यादा revenue चलाती हैं?
  • Top customers कौन हैं?
  • Department संदर्भ से results को कैसे enrich करें?
Databricks में Spark SQL के साथ Data Transformation

groupBy() और agg() कैसे काम करते हैं

 

$$

  • groupBy() - key कॉलम से rows को विभाजित करता है
  • agg() - हर group पर functions parallel चलाता है
  • दोनों lazy हैं - action पर computation होता है
  • एक agg() कॉल = एक data scan

recraft: half: डेटा की रंगीन rows तीन अलग-अलग लेबल वाले groups में बँटी हुई हैं और हर group के नीचे total value निकली है, ट्रांसपेरेंट बैकग्राउंड पर अलग-थलग

Databricks में Spark SQL के साथ Data Transformation

श्रेणी अनुसार revenue - कोड

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())
)
Databricks में Spark SQL के साथ Data Transformation

श्रेणी अनुसार revenue - आउटपुट

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       |
+-----------+-----------------+-----------------+---------------+
Databricks में Spark SQL के साथ Data Transformation

डाइमेंशन टेबल से समृद्ध करना

recraft: half: दो डेटाबेस टेबल्स एक चमकती कड़ी या पुल से जुड़ रही हैं, ट्रांसपेरेंट बैकग्राउंड पर अलग-थलग

$$

+-----------+----------+
|Category   |Department|
+-----------+----------+
|Clothing   |Retail    |
|Dining     |Food      |
|Electronics|Tech      |
|Groceries  |Food      |
|Savings    |Finance   |
+-----------+----------+
Databricks में Spark SQL के साथ Data Transformation

स्टैंडर्ड 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          |
+-----------+-----------+----------+------------------+
Databricks में Spark SQL के साथ Data Transformation

छिपी लागत - shuffle

recraft: half: जुड़े हुए सर्वर और नोड्स के नेटवर्क में डेटा पैकेट्स तीरों के साथ बह रहे हैं, ट्रांसपेरेंट बैकग्राउंड पर अलग-थलग

 

  • Matching keys को एक ही मशीन पर पहुँचना चाहिए
  • Spark rows को align करने के लिए नेटवर्क में भेजता है
  • Shuffle = बड़े पैमाने पर bottleneck
Databricks में Spark SQL के साथ Data Transformation

Spark UI

recraft: full: एक डार्क एनालिटिक्स मॉनिटरिंग डैशबोर्ड, जिसमें जॉब स्टेजेज, प्रोग्रेस बार्स और डेटा थ्रूपुट मेट्रिक्स दिखाते रंगीन हॉरिज़ॉन्टल बार चार्ट हैं, ट्रांसपेरेंट बैकग्राउंड पर अलग-थलग

 

  • Jobs और Stages - क्या चला और कितनी देर, ट्रैक करें
  • Shuffle Read/Write - कितना डेटा मूव हुआ, देखें
  • Pipeline bottlenecks की जाँच में उपयोगी
Databricks में Spark SQL के साथ Data Transformation

.explain() से क्वेरी प्लान पढ़ना

df_joined.explain(mode="formatted")
...
+- PhotonBroadcastHashJoin LeftOuter (18)
   :- ...
   +- PhotonShuffleExchangeSource (17)
      +- PhotonShuffleMapStage (16)
         +- PhotonShuffleExchangeSink (15)
            +- LocalTableScan (13)
...
  • .explain() - execution plan जाँचें
  • Nodes 15 और 16 - कोई routing Arguments लॉग नहीं हुए
Databricks में Spark SQL के साथ Data Transformation

Broadcast join - समाधान

# छोटी टेबल को F.broadcast() में wrap करें
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
Databricks में Spark SQL के साथ Data Transformation

पहले और बाद में - प्लान की जाँच

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 कभी move नहीं होती
  • एक बदलाव = बड़े पैमाने पर minutes से seconds
Databricks में Spark SQL के साथ Data Transformation

अभ्यास करते हैं!

Databricks में Spark SQL के साथ Data Transformation

Preparing Video For Download...