PySpark और SQL से डेटा को आकार दें

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

Disha Mukherjee

Lead Data Engineer

अपने डेटा को आकार दें

path = "/Volumes/.../default/transactions.csv"
df = spark.read.csv(path, header=True, inferSchema=True)
df.show(5, truncate=False)

Databricks में आउटपुट का प्रीव्यू

  • select() → कॉलम चुनें
  • filter() → रो चुनें
Databricks में Spark SQL के साथ Data Transformation

कॉलम चुनना

df.select("ID", "Customer_ID", "Transaction_Amount", "Category").show(5)
+---+-----------+------------------+-----------+
| ID|Customer_ID|Transaction_Amount|   Category|
+---+-----------+------------------+-----------+
|  1|    CUST001|          27337.49|Electronics|
|  2|    CUST002|           97716.6|Electronics|
|  3|    CUST003|           5752.36|Electronics|
|  4|    CUST004|          93443.22|    Savings|
|  5|    CUST005|          15109.98|Electronics|
+---+-----------+------------------+-----------+
Databricks में Spark SQL के साथ Data Transformation

रो को फ़िल्टर करना

from pyspark.sql import functions as F


filtered_df = df.filter(F.col("Transaction_Status") == "Completed") filtered_df.show(5)
+---+-------------------+-----------+------------------+------------------+
| ID|               Date|Customer_ID|Transaction_Amount|Transaction_Status|
+---+-------------------+-----------+------------------+------------------+
|  3|2023-01-03 00:00:00|    CUST003|           5752.36|         Completed|
|  9|2023-01-09 00:00:00|    CUST009|          28959.12|         Completed|
| 10|2023-01-10 00:00:00|    CUST010|          72098.18|         Completed|
| 11|2023-01-11 00:00:00|    CUST011|          49771.05|         Completed|
| 20|2023-01-20 00:00:00|    CUST020|          69825.14|         Completed|
+---+-------------------+-----------+------------------+------------------+
Databricks में Spark SQL के साथ Data Transformation

फ़िल्टर शर्तें जोड़ना

high_value_df = df.filter(

(F.col("Transaction_Amount") > 20000) &
(F.col("Transaction_Status") == "Completed")
) high_value_df.show(3)
+---+-------------------+-----------+------------------+------------------+
| ID|               Date|Customer_ID|Transaction_Amount|Transaction_Status|
+---+-------------------+-----------+------------------+------------------+
|  9|2023-01-09 00:00:00|    CUST009|          28959.12|         Completed|
| 10|2023-01-10 00:00:00|    CUST010|          72098.18|         Completed|
| 11|2023-01-11 00:00:00|    CUST011|          49771.05|         Completed|
+---+-------------------+-----------+------------------+------------------+
Databricks में Spark SQL के साथ Data Transformation

Python methods या SQL

 

  • Python methods: select(), filter()
    • ट्रांसफॉर्मेशन को कदम-दर-कदम बनाएं

 

$$

$$

  • SQL queries: Standard SQL सिंटैक्स
    • एग्रीगेशन और ग्रुपिंग के लिए अधिक पढ़ने योग्य

Python बनाम SQL

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

टेम्पररी व्यू बनाना

$$

df.createOrReplaceTempView("transactions")

 

  • Temporary view = DataFrame को इंगित करने वाला नाम
  • केवल वर्तमान Spark सत्र तक ही रहता है
Databricks में Spark SQL के साथ Data Transformation

SQL से क्वेरी करना

%sql
SELECT Category, COUNT(*) AS transaction_count

FROM transactions
WHERE Transaction_Status = 'Completed' GROUP BY Category ORDER BY transaction_count DESC
+-----------+-----------------+
|   Category|transaction_count|
+-----------+-----------------+
|   Clothing|             6764|
|     Dining|             6732|
|Electronics|             6640|
|    Savings|             6629|
|  Groceries|             6557|
+-----------+-----------------+
Databricks में Spark SQL के साथ Data Transformation

display() से वैलिडेट करना

display(filtered_df)

सॉर्टिंग और सर्चिंग के साथ इंटरएक्टिव टेबल का GIF

$$

$$

$$

  • परिणामों को सॉर्ट करें, पेजिनेट करें
  • nulls और outliers पहचानें
Databricks में Spark SQL के साथ Data Transformation

क्वेरी प्लान देखना

filtered_df.explain()
== Physical Plan ==
*(1) Filter (Transaction_Status = Completed)
+- FileScan csv [ID, Date, Customer_ID, ...]

$$

$$

  • SQL के साथ भी काम करता है: EXPLAIN SELECT ...
  • परफॉर्मेंस समस्याएँ ठीक करने में मददगार
Databricks में Spark SQL के साथ Data Transformation

Spark त्रुटियाँ पढ़ना

df.select("Account").show()
[UNRESOLVED_COLUMN.WITH_SUGGESTION] A column with name `Account` 
cannot be resolved. Did you mean one of the following? 
[`Date`, `Location`, `Customer_ID`, `Transaction_Amount`]
Databricks में Spark SQL के साथ Data Transformation

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

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

Preparing Video For Download...