Workflows के साथ प्रोडक्शन पाइपलाइन्स

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

Disha Mukherjee

Lead Data Engineer

Delta Lake क्यों?

सुरक्षित चमकता डिजिटल वॉल्ट, संगठित डेटा टेबल्स और सुरक्षा शील्ड, आधुनिक फ्लैट स्टाइल

 

$$

  • ACID transactions → असफल writes को रोल बैक करें
  • Schema enforcement → मिसमैच types को ब्लॉक करें
  • Versioning → किसी भी पिछले state को क्वेरी करें
Databricks में Spark SQL के साथ Data Transformation

Delta में लिखना

df_valid.write.format("delta") \
    .mode("overwrite") \
    .saveAsTable("transactions_clean")

print(f"Rows written: {df_valid.count():,}")
Rows written: 33,223

$$

$$

  • Unity Catalog में एक नई Delta table दिखेगी
Databricks में Spark SQL के साथ Data Transformation

Notebook tasks

task1_ingest: लोड और क्लीन

  • CSV लोड करता है → क्लीनिंग लगाता है → Delta table लिखता है
Databricks में Spark SQL के साथ Data Transformation

Notebook tasks

task2_metrics: श्रेणी अनुसार revenue

  • क्लीन की गई table पढ़ता है → metrics निकालता है → नई table में लिखता है
Databricks में Spark SQL के साथ Data Transformation

Notebook tasks

task3_customers: खर्च के आधार पर रैंक

  • क्लीन table पढ़ता है → customers को रैंक करता है → नई table में सेव करता है
Databricks में Spark SQL के साथ Data Transformation

जॉब बनाना

Jobs और Pipelines UI, तीन-task जॉब का पूरा DAG; task1_ingest से task2_metrics से task3_customers तक dependency तीर

  • आप मैन्युअल चला सकते हैं, jobs शेड्यूल करें, और triggers सेट करें
Databricks में Spark SQL के साथ Data Transformation

जॉब चलाना

जॉब रन DAG: task1_ingest और task2_metrics हरे (सफल), task3_customers लाल (फेल)

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

जॉब चलाना

तीसरे task में त्रुटि जाँचना

  • Error: customer-id को Customer_ID होना चाहिए
Databricks में Spark SQL के साथ Data Transformation

जॉब चलाना

सफलता के साथ ग्राफ़ व्यू

टाइमलाइन व्यू

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

Lakeflow क्या है?

$$

$$

तुलना: Imperative (Jobs) | Declarative (Lakeflow)

 

  • Jobs → हर स्टेप हम मैनेज करते हैं
  • Lakeflow → बस बताइए टेबल्स में क्या होना चाहिए
  • Databricks ऑर्डर, रिट्राई, और compute संभालता है
Databricks में Spark SQL के साथ Data Transformation

@dlt.table पैटर्न

@dlt.table(name="transactions_bronze")
def transactions_bronze():
    return spark.read.format("csv").schema(schema).load(FILE_PATH)

@dlt.table(name="transactions_silver") def transactions_silver(): return dlt.read("transactions_bronze").na.drop(...).filter(...)
@dlt.table(name="category_revenue_gold") def category_revenue_gold(): return dlt.read("transactions_silver").groupBy("Category").agg(...)
Databricks में Spark SQL के साथ Data Transformation

पाइपलाइन रन

Lakeflow पाइपलाइन DAG: transactions_bronze 100K rows, फिर transactions_silver 33K rows, फिर category_revenue_gold 6 rows, सभी green materialized views

  • Bronze (100K rows) → Silver (33K rows) → Gold (6 rows)
Databricks में Spark SQL के साथ Data Transformation

Notebooks, Jobs, या Lakeflow?

$$

लेयर्स: Notebooks, Databricks Jobs, Lakeflow Pipelines

 

$$

  • Notebooks → exploration, prototyping
  • Databricks Jobs → multi-step, scheduled pipelines
  • Lakeflow → पूरी तरह managed, declarative
Databricks में Spark SQL के साथ Data Transformation

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

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

Preparing Video For Download...