# Structured Streaming मॉडल — Apache Spark: DataFrames से Big Data Processing

Source: https://www.geekswithgeeks.com/hi/spark/st-model

> उसी DataFrame API से streaming queries लिखें और triggers व output modes समझें।

## वही API, असीमित input

Structured Streaming आपको streaming job को ऐसे लिखने देता है जैसे वह उस table पर batch job हो जिसमें **नई rows जुड़ती रहती हैं**। आप source (Kafka, files, परीक्षण के लिए `rate` source) से `spark.readStream` से पढ़ते हैं, सामान्य DataFrame ऑपरेशन लगाते हैं और `writeStream` से **sink** (Kafka, files/Delta, परीक्षण के लिए console या memory table) में लिखते हैं। **Trigger** तय करता है कि micro-batches कितनी बार चलें (`processingTime="10 seconds"`, या जो मौजूद है उसे प्रोसेस करके रुकने को `availableNow=True`)। **Output mode** बताता है कि हर बार क्या लिखा जाए: `append` (सिर्फ़ नई अंतिम rows), `update` (बदली rows) या `complete` (पूरी नतीजा table, छोटे aggregations के लिए)। **Checkpoint स्थान** प्रगति और state रखता है ताकि दोबारा शुरू हुई query बिना डेटा खोए या दोहराए आगे बढ़े (replay योग्य sources और idempotent या transactional sinks के साथ)।

## लगातार बढ़ती table

Structured Streaming live stream को असीमित table मानता है और आपकी query को क्रमिक रूप से चलाता है।

![तीन हिस्से: source, query, sink।](assets/figures/spark/section-6-map.svg) — चित्र 6.1 — Source, query और sink।

## छोटा streaming aggregation (उदाहरण)

Built-in `rate` source (परीक्षण generator) और in-memory sink उपयोग करता है। मैंने यह run दर्ज नहीं किया, इसलिए कोई output नहीं दिखाया; batches आने पर गिनती बढ़ती रहती है। असली jobs में आप `checkpointLocation` के साथ Kafka या Delta में लिखेंगे।

```python
import time
q = (spark.readStream.format("rate").option("rowsPerSecond", 20).load()
     .withColumn("bucket", F.col("value") % 3)
     .groupBy("bucket").count()
     .writeStream.outputMode("complete").format("memory").queryName("rates")
     .trigger(processingTime="1 second").start())

time.sleep(5)
spark.sql("SELECT bucket, count FROM rates ORDER BY bucket").show()
q.stop()
```

## Production में हमेशा checkpoint रखें

Checkpoint के बिना दोबारा शुरू हुई query फिर से शुरू होती है (या सही ढंग से आगे नहीं बढ़ पाती)। हर query को अपना टिकाऊ checkpoint folder दें, और उसे यूँ ही न हटाएँ या साझा न करें।

**Quiz:** Checkpoint स्थान क्या रखता है?

- [ ] सिर्फ़ अंतिम output files
- [ ] आपके passwords
- [ ] Spark binaries
- [x] Query की प्रगति और state ताकि restart आगे बढ़ सके

*Answer:* Query की प्रगति और state ताकि restart आगे बढ़ सके. Checkpoint के offsets और state दोबारा शुरू हुई query को वहीं से जारी रखने देते हैं जहाँ वह रुकी थी।
