पाठ 19 / 25
Structured Streaming मॉडल
उसी 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 को क्रमिक रूप से चलाता है।
छोटा streaming aggregation (उदाहरण)
Built-in rate source (परीक्षण generator) और in-memory sink उपयोग करता है। मैंने यह run दर्ज नहीं किया, इसलिए कोई output नहीं दिखाया; batches आने पर गिनती बढ़ती रहती है। असली jobs में आप checkpointLocation के साथ Kafka या Delta में लिखेंगे।
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 दें, और उसे यूँ ही न हटाएँ या साझा न करें।
त्वरित जाँच: Checkpoint स्थान क्या रखता है?
- सिर्फ़ अंतिम output files
- आपके passwords
- Spark binaries
- Query की प्रगति और state ताकि restart आगे बढ़ सके
Answer
Query की प्रगति और state ताकि restart आगे बढ़ सके — Checkpoint के offsets और state दोबारा शुरू हुई query को वहीं से जारी रखने देते हैं जहाँ वह रुकी थी।