# केस स्टडी: दैनिक Sales Pipeline — Apache Spark: DataFrames से Big Data Processing

Source: https://www.geekswithgeeks.com/hi/spark/wrap-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, निगरानी।](assets/figures/spark/section-8-map.svg) — चित्र 8.1 — Ingest, transform, publish और निगरानी।

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

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

```python
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 पर हाथ की सर्जरी से नहीं।

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

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

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