পাঠ ১৯ · ২৯-এর মধ্যে · মডিউল ৩
Home / AI Courses / Data Engineering / Spark Streaming

Spark Structured Streaming

Spark Streaming — unbounded tables, micro-batch & watermarks
৭ মিনিট পড়া মাঝারি · Intermediate PySpark কোডসহ

এই পাঠে যা শিখবেন

  • 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-এ? এটাই পার্থক্য।

ভাবুন ঢাকার একটি কুরিয়ার কোম্পানি — Sundarban। গ্রাহক চট্টগ্রাম থেকে ৫:০০টায় parcel পাঠালেন (event time), কিন্তু সেটা ঢাকায় ৭:৩০টায় পৌঁছাল (processing time)। "৫টায় কতটা parcel?" — event time চাই। "৭:৩০টায় warehouse-এ কী আছে?" — processing 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}$$

Watermark ছাড়া 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-এর সাথে)।

Checkpoint directory স্থায়ী, fault-tolerant storage-এ থাকতে হবে — S3, HDFS, GCS। Local disk দিলে — node মারা গেলে state হারাবেন। দু'টি query একই checkpoint share করলে — corruption।
Structured Streaming = Unbounded Table source → query → sink (incremental) 📥 Source Kafka topic bkash_txn t=12:00 · row1 t=12:01 · row2 t=12:02 · row3 … অসীম ⚙️ Streaming Query groupBy(window, user) withWatermark("10 min") micro-batch every 30s incremental · stateful 📤 Sink Delta Lake fraud_alerts append checkpoint → S3 path Daraz click stream · bKash transaction · Pathao ride event batch API + incremental engine = streaming, with the same code
Source → query → sink। প্রতিটি micro-batch-এ Spark unbounded table-এর নতুন rows process করে — checkpoint-এ state সংরক্ষণ।

৯ · কোডে — Kafka থেকে real-time fraud check

একটি bKash-এর মতো scenario। Kafka topic bkash_txn থেকে পড়ে — প্রতি ৩০ সেকেন্ডে user-প্রতি transaction count window aggregate। ১০-এর বেশি হলে fraud alert।

Python · PySpark Structured Streaming
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

Python · debug stream
# 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 দিন; inferSchema production-এ নয়।
  • 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 করা নয়।

অনুশীলন

  1. 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-এ অসহনীয়।

  2. Watermark-এর প্রভাব: watermark ১ মিনিট। ১২:০০:০০-এর event ১২:০২:৩০-এ এলো। কী হবে?

    Drop। Watermark = max(seenEventTime) − 1 min। ১২:০২:৩০-এ যদি কোনো event ১২:০২:০০-এ পৌঁছায়, watermark ~ ১২:০১:০০। ১২:০০:০০ event তার চেয়ে পুরোনো — drop।

    সমাধান: watermark বাড়ান (৫ মিনিট), বা পুরো mobile network reality পরিমাপ করে SLA সেট করুন।

  3. 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-এ আপনার পরবর্তী পদক্ষেপ

কোড রানার কাজ না করলে? Spark cluster ছাড়া PySpark চালাতে Google Colab ব্যবহার করুন — !pip install pyspark দিয়ে local mode-এ structured streaming পরীক্ষা সম্ভব।
পূর্ববর্তী পাঠ
পাঠ ১৮ · Producer, Consumer, Topic