पाठ 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 के रूप में।
कोड में 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 को हानिरहित बनाते हैं।