Spark Structured Streaming
এই পাঠে যা শিখবেন
- Structured Streaming-এর "unbounded table" mental model
- Source (Kafka, file, socket) থেকে stream পড়া ও sink-এ লেখা
- Event time, watermark, এবং late data handling
- Output mode তিনটির পার্থক্য ও checkpoint-এর প্রয়োজনীয়তা
১ · কেন Structured Streaming, পুরোনো DStream নয়?
Spark-এর প্রথম streaming API ছিল DStreamDStream (Discretized Stream)Spark Streaming-এর পুরোনো API (২০১৩)। RDD-র উপর ভিত্তি, মূলত micro-batch। এখন legacy — নতুন code DStream-এ লেখা উচিত নয়। (২০১৩) — RDD-ভিত্তিক, low-level, time-windowing কঠিন। ২০১৬-তে Databricks Structured Streaming আনে — DataFrame/SQL API-র উপরে। আজ Daraz, Pathao, bKash সহ প্রায় সব production Spark streaming এই নতুন API ব্যবহার করে।
Streaming ডেটা = একটি টেবিল যা কখনো শেষ হয় না (unbounded table)। প্রতিটি নতুন ইভেন্ট = টেবিলে নতুন row। আপনি যেমন static টেবিলের উপর SQL/DataFrame query লেখেন, ঠিক সেই API দিয়েই streaming query লিখবেন — Spark বাকি কাজ (incremental execution, state, fault tolerance) সামলায়।
২ · Source — কোথা থেকে ডেটা আসে
Production-এ সবচেয়ে কমন source তিনটি:
- Kafka: bKash transaction stream, Daraz click event — ৯৫% real-time use case।
- File source: S3/HDFS-এ নতুন file আসলে auto pickup। Daraz-এর nightly partner upload-এর জন্য আদর্শ।
- Socket: শুধু দ্রুত demo/testing-এর জন্য — production-এ কখনই নয়।
আরও আছে Kinesis, Pub/Sub, Delta Lake (CDC stream), Rate (synthetic test stream)।
৩ · Sink — কোথায় ডেটা যায়
Sink হলো গন্তব্য — Kafka topic, Delta Lake, Cassandra, console (debug), memory (dashboard), foreachBatch (custom DB write)। একটি streaming query একটি sink-এ লেখে; multiple destination-এ লিখতে চাইলে foreachBatch-এ branch করুন বা fan-outFan-outএকটি stream-কে একাধিক destination-এ পাঠানো। Kafka topic-এর পর Spark একই stream-কে Snowflake, Elasticsearch ও Redis-এ লিখলে — এটা fan-out। Pathao-র ride event এক জায়গা থেকে ৪+ system-এ যায়। করুন।
৪ · Trigger — কখন micro-batch চালু হবে
- Default: previous batch শেষ হলেই পরের শুরু (~১০০ ms gap)।
- ProcessingTime("৩০ seconds"): ফিক্সড interval — bKash-এর per-minute fraud check-এ চমৎকার।
- Once: একবার চালিয়ে exit — Airflow scheduler থেকে nightly trigger করার জন্য।
- Continuous("১ second"): sub-ms latency, কিন্তু aggregation/join সাপোর্ট সীমিত।
৫ · Event time vs processing time — পার্থক্যটি গুরুত্বপূর্ণ
Event time = ইভেন্টটি যখন ঘটেছিল (ফোনে bKash send বাটন চাপার মুহূর্ত)। Processing time = সেই ইভেন্ট Spark-এ কখন পৌঁছাল। মোবাইল network খারাপ হলে — ১০ মিনিট পরেও event আসতে পারে। আপনি যদি "শেষ ১ মিনিটে কত transaction?" জিজ্ঞেস করেন — কোন time-এ? এটাই পার্থক্য।
৬ · Watermark — late data কতদূর সহ্য করব?
WatermarkWatermark"এই event-time-এর আগের ইভেন্ট আর গ্রহণ করব না" — এই threshold। সাধারণত maxEventTime − delay। Watermark সরে গেলে — সেই window-এর state memory থেকে drop হয়। নাহলে state অনন্তকাল বাড়বে। = "সর্বোচ্চ event time − allowed delay"। আপনি ১০ মিনিট watermark সেট করলে — ১০ মিনিট আগের ইভেন্ট grace period-এ accept; তার চেয়ে পুরোনো ইভেন্ট drop। এটা ছাড়া aggregation-এর state অসীমভাবে বাড়বে।
$$\text{watermark}(t) = \max_{\text{seen}} (\text{eventTime}) - \text{delayThreshold}$$
groupBy(window(...)) aggregation চালানো যায়, কিন্তু output mode complete ছাড়া আর কিছু কাজ করবে না — এবং state কখনো clean হবে না। Production-এ সবসময় watermark দিন।
৭ · Output mode — কী লিখব sink-এ?
- Append: শুধু নতুন row — কখনো update হবে না। Watermark-এর সাথে aggregation-এ কাজ করে (window expire হলে final row লেখা হয়)।
- Update: যে row বদলেছে শুধু সেগুলো। dashboard upsert-এ আদর্শ।
- Complete: পুরো result table প্রতি batch-এ। ছোট aggregation-এ ঠিক, বড় হলে ভয়াবহ।
৮ · Checkpointing — fault tolerance-এর মেরুদণ্ড
Spark cluster crash করলে কোথা থেকে শুরু করবে? CheckpointCheckpointS3/HDFS-এ Spark যে metadata রাখে — শেষ পড়া Kafka offset, aggregation state, query progress। Restart হলে এই থেকে exactly-once recovery। directory-তে Spark লেখে — শেষ পড়া Kafka offset, aggregation state, query metadata। restart করলে exact একই জায়গা থেকে শুরু — exactly-once guarantee (Kafka source + idempotent sink-এর সাথে)।
৯ · কোডে — Kafka থেকে real-time fraud check
একটি bKash-এর মতো scenario। Kafka topic bkash_txn থেকে পড়ে — প্রতি ৩০ সেকেন্ডে user-প্রতি transaction count window aggregate। ১০-এর বেশি হলে fraud alert।
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, from_json, window, count
from pyspark.sql.types import StructType, StringType, DoubleType, TimestampType
spark = (SparkSession.builder
.appName("bkash-fraud-check")
.getOrCreate())
# Kafka payload schema
schema = (StructType()
.add("user_id", StringType())
.add("amount", DoubleType())
.add("event_time", TimestampType()))
# 1) Kafka source — unbounded table
raw = (spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "kafka:9092")
.option("subscribe", "bkash_txn")
.option("startingOffsets", "latest")
.load())
txn = (raw.select(from_json(col("value").cast("string"), schema).alias("d"))
.select("d.*"))
# 2) Watermark + windowed aggregation
suspect = (txn
.withWatermark("event_time", "10 minutes")
.groupBy(window("event_time", "1 minute"), "user_id")
.agg(count("*").alias("txn_count"))
.where(col("txn_count") > 10))
# 3) Sink — append to Delta with checkpoint
query = (suspect.writeStream
.format("delta")
.outputMode("append")
.option("checkpointLocation", "s3://bkash-de/chk/fraud_v1")
.option("path", "s3://bkash-de/tables/fraud_alerts")
.trigger(processingTime="30 seconds")
.start())
query.awaitTermination()
readStream static read-এর সমান্তরাল — পার্থক্য হলো query কখনো শেষ হয় না। withWatermark ১০ মিনিটের চেয়ে পুরোনো event drop; window("1 minute") ১-মিনিট tumbling bucket; append mode নতুন finalized row লেখে। checkpointLocation ছাড়া start করলে warning।
১০ · Console sink — দ্রুত debug
# Rate source = synthetic 1 row/sec — শেখার জন্য চমৎকার
from pyspark.sql import SparkSession
from pyspark.sql.functions import col
spark = SparkSession.builder.appName("rate-demo").getOrCreate()
stream = (spark.readStream.format("rate")
.option("rowsPerSecond", 1).load())
(stream.where(col("value") % 2 == 0)
.writeStream
.format("console")
.outputMode("append")
.trigger(processingTime="2 seconds")
.start()
.awaitTermination())
rate source কোনো external dependency ছাড়াই Spark Streaming শেখার সবচেয়ে সহজ উপায়। প্রতি ২ সেকেন্ডে console-এ সম-সংখ্যার row print হবে।
ভাবনার প্রশ্ন
প্রতিটি প্রশ্ন নিজে কিছুক্ষণ ভাবুন — তারপর "→ উত্তর" চাপুন।
প্র ০১ Spark Structured Streaming "true streaming" না, এটি micro-batch। তবু Databricks বলে এটি "high-throughput streaming"। কেন? Flink-এর সাথে এই trade-off কোথায়?
এটি streaming জগতের সবচেয়ে আলোচিত architectural debate। Spark Structured Streaming default-এ micro-batch — প্রতি ৫০-৫০০ ms-এ একটি ছোট batch চালায়। এটি "true streaming" (event-by-event) নয়।
তবু high-throughput কেন?
- প্রতি micro-batch-এ Spark পুরো batch optimization (Catalyst, Tungsten) কাজে লাগায় — vectorized execution, code generation। এক-একটি ইভেন্ট করলে এসব overhead-এ বেশি ব্যয় হতো।
- JVM-এ object allocation ব্যয়বহুল। batch-এ allocate করলে — per-row cost কমে।
- Kafka-র batch fetch (poll = ৫০০ records) Spark-এর batch model-এর সাথে natural fit।
Trade-off — Flink-এর সাথে:
- Latency: Spark default ~৫০০ ms, tuned ~১০০ ms। Flink ১০-৫০ ms (event-driven, no batching)। bKash-এর fraud detection-এ ১ second যথেষ্ট, তাই Spark। কিন্তু stock trading-এ Flink।
- Throughput: Spark সাধারণত উঁচু — batch optimization-এর জন্য। Flink-এও ভাল কিন্তু per-event overhead-এ পিছিয়ে।
- State: Flink-এর native state backend (RocksDB) Spark-এর HDFS-state-এর চেয়ে অনেক দ্রুত — large state handling-এ Flink এগিয়ে।
- Ecosystem: Spark-এর batch ও ML সবকিছু। Flink-এ pure streaming। Daraz-এর মতো mixed workload — Spark। Pure event-driven Pathao live tracking — Flink।
Continuous mode: Spark ২.৩ থেকে continuous mode (sub-ms) যোগ হয়েছে — কিন্তু এখনো experimental, aggregation/join সাপোর্ট সীমিত। বেশিরভাগ production এখনো micro-batch।
মূল উপলব্ধি: "True streaming = better" — এটা naive thinking। বাস্তবে latency requirement অনুযায়ী tool বাছতে হয়। ১ সেকেন্ড latency-তে ১০০x throughput পাওয়া গেলে — সেটাই উত্তম।
প্র ০২ Watermark ১০ মিনিট সেট করলেন। একটি ইভেন্ট ১২ মিনিট দেরিতে এলো — drop হলো। গ্রাহক complain করল। কীভাবে design বদলাবেন?
এটি real-world streaming-এর সবচেয়ে কঠিন trade-off — correctness vs resource cost। Watermark যত বড় — তত বেশি state memory-তে রাখতে হয়; যত ছোট — তত বেশি late event drop।
প্রথমে প্রশ্ন করুন — late কেন?
- Mobile network খারাপ — গ্রামাঞ্চলে Pathao driver-এর GPS ping প্রায়ই ১৫-২০ মিনিট পরে আসে।
- Producer side buffering — bKash app offline মোডে transaction জমিয়ে রাখে।
- Kafka rebalance / consumer lag — অস্বাভাবিক।
Solution ১ — Watermark বাড়ান: ১০ → ৩০ মিনিট। সরল, কিন্তু state size বাড়বে। যদি state ১ GB থেকে ৫ GB হয়, RocksDB SSD plan করুন। Daraz checkout-এর মতো use case-এ এটি যথেষ্ট।
Solution ২ — Tiered processing (lambda-ish):
- Real-time path: ১০-মিনিট watermark, "approximate" dashboard-এর জন্য।
- Reconciliation path: প্রতি রাতে batch job — সব ইভেন্ট, late গুলো সহ — "true" daily report।
- Bank reconciliation, BTRC reporting-এর জন্য এটাই স্ট্যান্ডার্ড।
Solution ৩ — সাইড output of late events:
- Spark-এ direct API নেই, কিন্তু
foreachBatch-এ filter দিয়ে late rows আলাদা topic-এ পাঠান। - সেগুলোর জন্য আলাদা reprocessing pipeline (CDC-style update পুরোনো result-এ)।
Solution ৪ — Update mode + idempotent sink: append-এর বদলে update mode। Result store-এ upsert করলে — late event এসে existing row update করতে পারে। DynamoDB, Cassandra, Delta Lake-এ কাজ করে।
আসল মূল কথা: "Late event drop" পদ্ধতিগত ভুল নয় — প্রায় সব streaming system-এই কিছু drop করে। আপনার job হলো — এই drop-এর rate পরিমাপ করা (metric দিন), business stakeholder-এর সাথে SLA নির্ধারণ করা ("দিনে ০.১% drop OK?"), এবং tolerance-এর বাইরে গেলে design বদলানো। নীরব data loss-ই সবচেয়ে বিপজ্জনক।
প্র ০৩ Daraz-এ এক জুনিয়র engineer একটি streaming job লিখলেন — checkpoint দিলেন না। ২ দিন পরে cluster restart হলে সব data হারাল। কেন? কী পদক্ষেপ ছিল সঠিক?
Production Spark Streaming-এ checkpoint ছাড়া start করা = মৃত্যুদণ্ড। এটা শেখার একটা কঠিন কিন্তু চিরকালীন পাঠ।
কী হারাল ও কেন:
- Kafka offset: Spark কোন offset পর্যন্ত পড়েছিল — সেই metadata checkpoint-এ থাকে। ছাড়া দিলে — restart-এ
startingOffsets="latest"থেকে শুরু হবে। মাঝের ২ দিনের data — চিরতরে গেছে (Kafka retention কম থাকলে)। - Aggregation state: windowed count, running sum — সব memory-তে। Cluster বন্ধ হলেই hash table গায়েব। restart-এ count শূন্য থেকে শুরু — historical aggregate ভুল।
- Idempotency সম্ভাবনা শূন্য: exactly-once guarantee checkpoint-এর উপর নির্ভরশীল। ছাড়া দিলে — at-most-once বা at-least-once, যা business-এ duplicate বা missing data দু'টোই ঘটাতে পারে।
সঠিক pattern (production checklist):
- Checkpoint location S3/HDFS-এ:
s3://daraz-de/chk/click_v1। প্রতি query-র জন্য আলাদা path। - Versioned path:
v1, v2, …— schema বদলালে পুরোনো checkpoint বাতিল, নতুন path। - Monitoring:
StreamingQueryListenerদিয়ে batch duration, input rows, processed rows track। Datadog/Grafana alert। - Sink idempotency: Delta Lake / upsert-supporting DB। Kafka sink-এ
kafka.transactional.id। - Recovery test: staging-এ deliberately cluster kill করে দেখুন — restart-এ exact একই counts আসে কিনা।
- Schema lock: JSON parsing-এ explicit schema দিন;
inferSchemaproduction-এ নয়। - Backpressure:
maxOffsetsPerTriggerদিয়ে batch size cap — restart-এ huge backlog এলে memory blow হবে না।
Recovery action — যদি ইতিমধ্যে disaster:
- Kafka retention দেখুন — যদি ৭ দিন থাকে, ২ দিনের data এখনো reachable।
- Source DB থেকে batch backfill। CDC log থাকলে ভালো।
- Downstream সঙ্গী team-কে জানান — দু'দিনের metric "uncertain"।
- Postmortem — cause, blast radius, prevention।
বড় কথা: "এটা শুধু dev environment, পরে fix করব" — এই attitude থেকেই production accident হয়। Streaming system fault tolerance opt-in নয়, opt-out হওয়া উচিত — checkpoint দেওয়া যেন breathing-এর মতো reflexive।
প্র ০৪ Pathao-র live ride map — যেখানে ১,০০,০০০ driver-এর GPS প্রতি ৫ সেকেন্ডে আসে — আপনি কীভাবে design করবেন? Spark Streaming, Flink, না অন্য কিছু?
এটি একটি classic streaming architecture প্রশ্ন — Bangladeshi context-এ রিয়াল scale-এর। আসুন বিচ্ছিন্ন করে দেখি।
Workload analysis:
- 1,00,000 drivers × ১/৫ sec = ২০,০০০ events/sec। peak hour-এ ৫০,০০০।
- Latency requirement: end-user map-এ ৩-৫ সেকেন্ড লাগলেই OK। ১ মিনিট হলে বিরক্তিকর, ১০ সেকেন্ড সহনীয়।
- Stateful — প্রতি driver-এর last position কয়েক মিনিট রাখতে হবে।
Architecture proposal:
- Ingest: Kafka। ৬-১২ partition, key=driver_id (একই driver-এর update একই partition-এ → ordering)। Retention ১ দিন (live use case)।
- Processing: এই scale-এ Spark Structured Streaming যথেষ্ট। ৩-৫ সেকেন্ড latency-তে micro-batch জিতবে। Flink-এ গেলে ১০০ ms-এ পৌঁছানো যেত — কিন্তু extra ops complexity-র যোগ্য নয়।
- State:
mapGroupsWithStateবাflatMapGroupsWithState— driver-প্রতি last-known location, status (online/offline), idle timer। - Sink: Redis (TTL=১ মিনিট) — frontend মেপ-এর জন্য। সমান্তরালভাবে S3 raw archive (ML training-এর জন্য)।
- Frontend: Server-Sent Events / WebSocket দিয়ে map UI Redis থেকে pull/subscribe।
Trade-off আলোচনা:
- Spark বনাম Flink: Spark-এর pros — Daraz/Pathao team-এ Spark expertise, batch ও streaming একসাথে, mature। Flink-এর pros — ১০ ms latency, low memory state। এই use case-এ Spark জিতবে।
- Watermark: ১ মিনিট। মোবাইল network-এ delay আছে, কিন্তু "১ ঘণ্টা পুরোনো GPS" useless।
- Backpressure: peak hour traffic surge —
maxOffsetsPerTrigger=500000দিয়ে batch cap; না হলে memory crash। - Cost: ২০,০০০ events/sec ~১.৭ billion/day। Spark cluster ~৪-৮ executor (4 core, 16 GB) যথেষ্ট। মাসিক ~$৫০০-৮০০ AWS।
যে scenario-তে Flink লাগত:
- ৫,০০,০০০+ driver, sub-second latency requirement — financial trading-এর কাছাকাছি।
- Complex event-time joins (এক ride-এর সব event ১০-মিনিট window-এ correlate)।
- Massive state (TB-scale) — RocksDB native backend prefer।
মূল উপলব্ধি: "Best technology" বলে কিছু নেই — fit-for-purpose আছে। Pathao-র live map-এ Spark Streaming + Kafka + Redis সবচেয়ে cost-effective ও maintainable choice। Engineer-এর কাজ trade-off justify করা, novelty chase করা নয়।
অনুশীলন
-
Output mode বাছুন: "প্রতি ১-মিনিট window-এ user-প্রতি click count" — Delta Lake-এ append করতে চান। কোন output mode? কেন?
Append + watermark। কারণ — append mode-এ প্রতি window-এর result শুধু একবার final হলে লেখা হয় (watermark পার হলে)। এতে Delta Lake-এ duplicate row হবে না, এবং downstream batch query simple।
Update mode-এ একই window-এর row একাধিকবার লেখা হবে যতক্ষণ window open। Complete mode-এ পুরো aggregate পুনঃলিখিত — বড় state-এ অসহনীয়।
-
Watermark-এর প্রভাব: watermark ১ মিনিট। ১২:০০:০০-এর event ১২:০২:৩০-এ এলো। কী হবে?
Drop। Watermark = max(seenEventTime) − 1 min। ১২:০২:৩০-এ যদি কোনো event ১২:০২:০০-এ পৌঁছায়, watermark ~ ১২:০১:০০। ১২:০০:০০ event তার চেয়ে পুরোনো — drop।
সমাধান: watermark বাড়ান (৫ মিনিট), বা পুরো mobile network reality পরিমাপ করে SLA সেট করুন।
-
Mini code: Kafka topic
clicksথেকে পড়ে — প্রতি ৫ মিনিট window-এpage-প্রতি count, console-এ দেখান।from pyspark.sql import SparkSession from pyspark.sql.functions import from_json, col, window, count from pyspark.sql.types import StructType, StringType, TimestampType spark = SparkSession.builder.appName("clicks").getOrCreate() schema = StructType().add("page", StringType()).add("event_time", TimestampType()) raw = (spark.readStream.format("kafka") .option("kafka.bootstrap.servers", "kafka:9092") .option("subscribe", "clicks").load()) events = raw.select(from_json(col("value").cast("string"), schema).alias("d")).select("d.*") agg = (events.withWatermark("event_time", "10 minutes") .groupBy(window("event_time", "5 minutes"), "page") .agg(count("*").alias("clicks"))) (agg.writeStream.format("console").outputMode("update") .option("truncate", False).start().awaitTermination())
আরও পড়ুন · ABCL TECH-এ আপনার পরবর্তী পদক্ষেপ
- পাঠ ২০ · Apache Flink পরিচিতি পরবর্তী পাঠ True streaming engine — sub-ms latency, native state, event-time mastery।
- পাঠ ১৮ · Producer, Consumer, Topic আগের পাঠ Kafka-র core abstraction — Spark Streaming-এর source হিসেবে অপরিহার্য।
- পাঠ ১৬ · Streaming vs Batch এই পাঠের সাথে সম্পর্কিত কখন streaming, কখন batch — design decision-এর ভিত্তি।
- সব AI Courses দেখুন ABCL TECH Python, ML, DL, NLP, CV, GenAI, RL, MLOps — সব AI কোর্স একসাথে।
!pip install pyspark দিয়ে local mode-এ structured streaming পরীক্ষা সম্ভব।