पाठ 24 / 25

केस स्टडी: दैनिक Sales Pipeline

ऐसा Spark batch job डिज़ाइन करें जो कच्ची order files को partitioned revenue table में बदले।

डिज़ाइन

Job किसी दिए run_date के लिए दिन में एक बार चलता है। वह उस तारीख़ की कच्ची CSVs को स्पष्ट schema और fail-fast mode के साथ पढ़ता है, उन्हें साफ़ करता है (strings trim करना, types cast करना, nulls भरना या चिह्नित करना, दोहराए order_id हटाना), छोटी customer dimension से broadcast hint के साथ join करता है, प्रति शहर प्रति दिन revenue aggregate करता है, order_date से partitioned Parquet लिखता है, सिर्फ़ उसी partition के लिए overwrite mode में (idempotent: वही तारीख़ दोबारा चलाने पर वह तारीख़ बदल जाती है), और प्रकाशित करने से पहले डेटा-गुणवत्ता जाँचें चलाता है (row count न्यूनतम से ऊपर, कोई null key नहीं)। उसे Airflow जैसा orchestrator schedule करता है, अस्थायी विफलताओं पर retry करता है, row counts और अवधियाँ log करता है, और runtime बहकने पर उसका Spark UI देखा जाता है। Tests हर transformation function को छोटे DataFrames से ढकते हैं।

शुरू से अंत तक दैनिक pipeline

पढ़ें, साफ़ करें, join करें, aggregate करें, partitioned output लिखें और monitor करें, एक भरोसेमंद job के रूप में।

चार चरण: ingest, transform, publish, निगरानी।
चित्र 8.1 — Ingest, transform, publish और निगरानी।

कोड में job (उदाहरण)

Dynamic partition overwrite सिर्फ़ output में मौजूद partitions बदलता है। Paths, column नाम और run_date argument उदाहरण हैं।

spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")

raw = (spark.read.schema(SCHEMA).option("header", True).option("mode", "FAILFAST")
       .csv(f"s3a://lake/raw/orders/{run_date}/"))

clean = (raw.dropDuplicates(["order_id"])
            .withColumn("customer", F.trim("customer"))
            .withColumn("amount", F.col("amount").cast("double")))

revenue = (clean.join(F.broadcast(dim_customers), "customer", "left")
                .groupBy("order_date", "city")
                .agg(F.sum("amount").alias("revenue"), F.count("*").alias("orders")))

assert revenue.count() > 0, "no rows produced"
(revenue.write.mode("overwrite").partitionBy("order_date")
        .parquet("s3a://lake/gold/daily_revenue/"))

दोबारा चलाना सामान्य मामला बनाएँ

हर job ऐसा डिज़ाइन करें कि किसी तारीख़ को दोबारा चलाना सुरक्षित हो। तब रात 3 बजे की विफलता retry या backfill से ठीक होती है, tables पर हाथ की सर्जरी से नहीं।

त्वरित जाँच: इस डिज़ाइन में वही तारीख़ दोबारा चलाना क्या सुरक्षित बनाता है?

  • Random file नाम उपयोग करना
  • हर बार table में append करना
  • सिर्फ़ उस तारीख़ का partition overwrite करने से हर बार वही नतीजा मिलता है
  • डेटा-गुणवत्ता जाँचें छोड़ना
Answer

सिर्फ़ उस तारीख़ का partition overwrite करने से हर बार वही नतीजा मिलता है — Idempotent, partition-सीमित writes retries और backfills को हानिरहित बनाते हैं।