Streaming vs Batch — কখন কোনটি
এই পাঠে যা শিখবেন
- Batch ও streaming-এর মৌলিক পার্থক্য — latency বনাম throughput
- Micro-batch বনাম true streaming — Spark Streaming বনাম Flink-এর architectural choice
- Event time, processing time, ও watermark — streaming-এর তিনটি timing concept
- কোন ব্যবহারে কোনটি — bKash, Daraz, Pathao, BTRC-র বাস্তব pattern দিয়ে
১ · মূল প্রশ্ন — ডেটা কত দ্রুত দরকার?
ডেটা ইঞ্জিনিয়ারিং-এর প্রতিটি pipeline-এর প্রথম প্রশ্ন: "এই ডেটা business-এ কত দ্রুত দরকার?" উত্তর "কাল সকালে রিপোর্টে" হলে batch যথেষ্ট। উত্তর "এক সেকেন্ডের ভিতর fraud আটকাতে হবে" হলে streaming লাগবে। এই দু'টির মাঝে সম্পূর্ণ ভিন্ন infrastructure, ভিন্ন framework, এমনকি ভিন্ন mindset।
Latency — একটি event আসা থেকে result বের হওয়া পর্যন্ত সময় (যেমন ১০০ ms)।
Throughput — প্রতি সেকেন্ডে কত event process করা যায় (যেমন ১০ লাখ events/sec)।
Batch latency বেশি কিন্তু throughput অসাধারণ। Streaming latency কম কিন্তু throughput batch-এর চেয়ে কম (একই hardware-এ)।
এই pair-কে engineer-রা ছোটবেলা থেকেই trade-off হিসেবে চিনে — গাড়ি বনাম ট্রাক, ব্যাংক ATM বনাম মাস-শেষের payroll।
২ · Batch processing — কাঁচা শক্তি
Batch processingBatch Processingএকটি নির্দিষ্ট সময়ের ডেটা একত্র করে একসাথে process করা। সাধারণত প্রতি ঘণ্টা বা দিন-ভিত্তিক চলে। উচ্চ throughput, simple code, কম খরচ। পুরনো এবং সবচেয়ে বেশি ব্যবহৃত প্যাটার্ন। গত ২৪ ঘণ্টার সব transaction → একটি বড় file/table → Spark/Hive job → আজ সকালের dashboard। ৬০ বছর ধরে enterprise-এর মেরুদণ্ড।
কেন এত popular:
- Idempotent পুনরায় চালানো সহজ: job ব্যর্থ হলে সম্পূর্ণ batch আবার চালান — state management ছাড়াই।
- Distributed compute optimal: Spark/MapReduce — পুরো ডেটা scan করে একসাথে। CPU/memory utilization ১০০%-এর কাছে।
- Cost কম: spot/preemptible VM ব্যবহার করা যায়। AWS-এ batch job-এ ৭০% খরচ সাশ্রয়।
- Debug সহজ: বিকেলে fail করলে — সকালে engineer দেখে fix করে আবার চালান।
৩ · Stream processing — সেকেন্ডে সিদ্ধান্ত
Stream processingStream Processingevents আসার সাথে সাথে process করা — কোনো batch জমানো ছাড়া। Kafka, Flink, Spark Structured Streaming এই কাজ করে। সাধারণ latency: ১০ ms থেকে ১ সেকেন্ড।-এ ডেটা একটি চলমান নদী। events একটার পর একটা আসছে — প্রতিটি event পৃথকভাবে বা ছোট window-এ process হয়। output কোনো final report নয় — এটি আরেকটি continuously-updating stream।
উদাহরণ: bKash-এ একটি transaction → fraud check → result ১০০ ms-এর ভেতর। পুরো batch জমানোর সুযোগ নেই — fraud already হয়ে গেছে।
Streaming-এর challenge:
- State management: "এই user-এর শেষ ১ ঘণ্টায় কতবার transaction?" — এই counter মেমরিতে রাখতে হয়, machine ক্র্যাশ করলে restore।
- Out-of-order events: Network-এর কারণে event ১৪:০২-এর পরে ১৪:০১-এর event আসতে পারে।
- Exactly-once delivery: একটি event দু'বার process হলে duplicate balance deduct। কঠিন কিন্তু critical।
- Always-on cluster: ২৪×৭ চলে, batch-এর মতো শেষ হয় না। Cost বেশি।
৪ · Micro-batch — মাঝামাঝি পথ
Spark Streaming জনপ্রিয় করেছে এই ধারণা: "আমি প্রতি ১ সেকেন্ডে একটি ছোট batch চালাই।" Latency ৫০০ ms-১০ s, কিন্তু code-এ batch-এর simplicity (DataFrame API)। বেশিরভাগ business-এ "near real-time" enough — আর Spark team-কে নতুন framework শিখতে হয় না।
Micro-batch (Spark): ছোট batch-এ — latency ১-১০ s সাধারণত।
Industry-তে ৮০% real-time use cases-এ micro-batch enough।
৫ · Event time বনাম Processing time
Streaming-এর সবচেয়ে subtle ধারণা — দু'টি ভিন্ন time-এর মধ্যে পার্থক্য:
- Event time: event আসলে কখন ঘটেছে (mobile-এ tap-এর সময়, যেমন ১৪:০০:০৫)।
- Processing time: server এ event পৌঁছানোর সময় (যেমন ১৪:০০:২৩ — ১৮ সেকেন্ড পর, network delay-এর কারণে)।
এই দু'য়ের পার্থক্য বোঝা critical। "গত ১ মিনিটে কত transaction?" — উত্তর processing time-এ ভুল হবে যদি network slow থাকে। Event time-এ সঠিক হবে।
৬ · Watermark — কখন বলা যায় "সব এসেছে?"
WatermarkWatermarkStreaming-এ একটি timestamp boundary যা বলে — "এর আগের event সব এসে গেছে।" Late event এর পর এলে drop বা separate handling। Flink ও Spark-এ critical concept। হলো streaming framework-এর "উপলব্ধি" — কোন event time পর্যন্ত সব data এসেছে। যদি বর্তমান watermark ১৪:০০:০৫ হয় — মানে ১৪:০০:০৫-এর আগের সব event আমরা পেয়েছি (probably)। এর পরে আসা ১৩:৫৯-এর event "late" — drop করা যায় বা পৃথক handle।
Watermark = max(event_time observed) − allowed_lateness। যেমন allowed_lateness = ১০ s হলে — সর্বশেষ ১৪:০০:১৫-এর event দেখার পর watermark ১৪:০০:০৫।
৭ · Bangladesh-এর use cases — কোনটি কোথায়
Streaming অপরিহার্য:
- bKash fraud detection: একটি transaction ১০০ ms-এ block করতে হবে। batch-এ ফেরা মানে টাকা চলে গেছে।
- Pathao ride matching: rider অনুরোধ → ১০ সেকেন্ডের ভেতর nearest driver খোঁজা।
- BTRC traffic monitoring: আক্রমণাত্মক traffic pattern detect — DDoS রোধ।
- Daraz cart abandonment alert: cart-এ যোগ করে ১৫ মিনিট চুপচাপ → push notification।
Batch যথেষ্ট:
- রাত্রে settlement report: bKash মার্চেন্ট-দের প্রতি রাতে commission হিসাব।
- মাসিক tax report: NBR-এ জমা — মাসে একবার।
- ML training: প্রতি সপ্তাহে নতুন model — historical ডেটায়।
- Daily KPI dashboard: CEO সকাল ৯টায় দেখেন — গতকালের sales।
৮ · Lambda ও Kappa architecture
বাস্তব production-এ অনেক system দু'টোই ব্যবহার করে। Nathan Marz (২০১১) প্রস্তাব করেছিলেন Lambda architecture: একই ডেটা একই সময়ে batch (accuracy-র জন্য) ও speed layer (latency-র জন্য) দু'টিতে যায়। Jay Kreps (২০১৪) সরলীকরণ করে Kappa: শুধু streaming, batch হলো streaming-এর special case (replay)।
আজকাল Kappa-ই বেশি জনপ্রিয় — Kafka + Flink/Spark Structured Streaming দিয়ে দু'টোই করা যায় same code-এ।
৯ · কোডে দেখা — Spark batch vs streaming
একই DataFrame API — পার্থক্য শুধু source ও sink-এ:
from pyspark.sql import SparkSession
from pyspark.sql.functions import sum as _sum, col
spark = SparkSession.builder.appName("bkash_batch").getOrCreate()
# Batch: গতকালের সব transaction পড়ুন
df = spark.read.parquet("s3://bkash-data/transactions/2026-05-08/")
# মার্চেন্ট-ভিত্তিক মোট হিসাব
result = (df.filter(col("status") == "success")
.groupBy("merchant_id")
.agg(_sum("amount").alias("daily_total")))
result.write.parquet("s3://bkash-reports/daily/2026-05-08/")
spark.stop()
from pyspark.sql import SparkSession
from pyspark.sql.functions import sum as _sum, col, window
spark = SparkSession.builder.appName("bkash_stream").getOrCreate()
# Streaming: Kafka থেকে continuous read
df = (spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "kafka-1:9092,kafka-2:9092")
.option("subscribe", "transactions")
.load())
# প্রতি ১-মিনিট window-এ মার্চেন্ট-ভিত্তিক যোগ
result = (df.filter(col("status") == "success")
.withWatermark("event_time", "10 seconds")
.groupBy(window("event_time", "1 minute"), "merchant_id")
.agg(_sum("amount").alias("minute_total")))
# Console-এ continuously print, প্রতি ১০ s checkpoint
query = (result.writeStream
.outputMode("update")
.format("console")
.option("checkpointLocation", "/tmp/ckpt")
.trigger(processingTime="10 seconds")
.start())
query.awaitTermination()
readStream ও writeStream — মূল পার্থক্য। Job কখনো শেষ হয় না; প্রতি ১০ সেকেন্ডে নতুন output। withWatermark ১০ সেকেন্ড late event allow করে; পরে drop।
ভাবনার প্রশ্ন
প্রতিটি প্রশ্ন নিজে কিছুক্ষণ ভাবুন — তারপর "→ উত্তর" চাপুন।
প্র ০১ আপনি bKash-এ junior data engineer। Manager বলছেন — "সব dashboard real-time হোক।" আপনি কী যুক্তি দেবেন? কোন metric streaming দরকার, কোনটা batch enough?
এই scenario বাস্তব — অনেক senior leader streaming-এর "চমক"-এ মুগ্ধ হয়ে সর্বত্র বাস্তবায়ন চান। Engineer-এর দায়িত্ব trade-off ব্যাখ্যা করা।
প্রথম যুক্তি — খরচ: bKash-এর daily transaction ~৪ কোটি। এই ভলিউমে Kafka cluster (৩+ broker, ZooKeeper, monitoring), Flink/Spark Streaming cluster (২৪×৭ চলবে), checkpoint storage — মাসে ১৫-৩০ লাখ টাকা infrastructure cost। এর তুলনায় batch (Airflow + Spark on-demand) ৩-৫ লাখ টাকায় শেষ।
দ্বিতীয় যুক্তি — complexity: Streaming-এ debugging কঠিন। একটি event late এলে — কোথায় drop হলো? Watermark কেন stuck? State store full কেন? Engineer team-কে নতুন skill set শিখতে হবে।
Streaming dরকার:
- Fraud detection: < ১ সেকেন্ডে block — batch-এ অসম্ভব। Direct revenue impact — ROI প্রমাণিত।
- OTP delivery monitoring: SMS gateway down হলে ৩০ সেকেন্ডের ভেতর alert — না হলে millions of users login ব্যর্থ।
- Cash-in/cash-out matching: agent ও user-এর leg মিলতে হবে real-time — নাহলে disputed transaction।
- Transaction success rate: banking partner down হলে — minutes-এ detect, fail-over।
Batch enough:
- CEO daily KPI: গতকাল কত revenue — সকাল ৬টায় ready। Real-time-এ কোনো actionable insight নেই।
- Merchant settlement: মাসের শেষে commission — batch-এ accuracy বেশি।
- Compliance reports (BFIU): দৈনিক submission — batch।
- Marketing analysis: কোন campaign কাজ করল — সাপ্তাহিক analysis।
- Customer 360: ML feature store — দিনে ১বার update enough।
Hybrid approach: "Live transaction count" dashboard streaming, কিন্তু "top performing merchants" batch — different freshness, different cost. এটাই Lambda philosophy।
মূল উপলব্ধি: "Real-time" একটি business question, technical default নয়। প্রতিটি metric-এ আলাদা প্রশ্ন: "৫ সেকেন্ড দেরিতে এলে কি কোনো decision বদলাবে?" উত্তর "না" হলে — batch।
প্র ০২ Pathao-এর একটি ride request — driver mobile-এ event time ১৪:০০:০০, কিন্তু network slow-এ server-এ আসে ১৪:০০:০৮। Watermark ৫ সেকেন্ড late allow করে। এই event কী হবে? Late events নিয়ে strategy কী?
এটি streaming-এর সবচেয়ে subtle problem। Bangladesh-এর mobile network reality — ৪G-তে p99 latency ১-৩ সেকেন্ড সাধারণ, gramophone area-এ আরও বেশি।
সংখ্যাটি বিশ্লেষণ:
- Event time: ১৪:০০:০০ (mobile-এ tap হলো)
- Processing time: ১৪:০০:০৮ (server-এ পৌঁছাল, ৮ সেকেন্ড পর)
- Watermark = max(observed event_time) − ৫ সেকেন্ড
- যদি এর মধ্যে আরও event এসে থাকে যেগুলোর event_time ১৪:০০:১০ — তাহলে current watermark = ১৪:০০:০৫
- আমাদের event-এর event_time (১৪:০০:০০) < watermark (১৪:০০:০৫) → late event
Spark-এ কী হবে: By default, late event drop। যদি window ১৪:০০:০০-১৪:০১:০০-এর জন্য aggregation ইতিমধ্যে emit হয়ে থাকে — এই late event সেই window-এ যোগ হবে না।
Strategy ১ — Acceptable loss: বেশিরভাগ analytics-এ ০.১% events drop চলে। Daily total revenue থেকে ০.১% হারিয়ে গেলে — কেউ দেখবে না।
Strategy ২ — Wider watermark: ৫ সেকেন্ডের জায়গায় ১ মিনিট। অনেক বেশি event capture, কিন্তু latency বেড়ে যায় (output ১ মিনিট পরে আসবে)।
Strategy ৩ — Allowed lateness with retraction (Flink): Flink window allow করে — late event এলে আগের aggregate retract করে নতুন emit। Downstream system (database) idempotent হতে হবে।
Strategy ৪ — Side output: Late events আলাদা stream-এ পাঠান — পরে batch-এ reprocess। Production-grade architecture।
Strategy ৫ — Ride request specific: Pathao-এর ক্ষেত্রে — যদি ride request ১০ সেকেন্ড late আসে, rider already cancel করেছে (UI 30s timeout)। তাই business-logic-এ event ignore করা সঠিক।
Trade-off matrix:
- কম watermark → কম latency, বেশি drop
- বেশি watermark → বেশি accuracy, বেশি latency, বেশি memory
- Allowed lateness → accurate কিন্তু complex
- Side output → সবচেয়ে robust, কিন্তু extra pipeline
মূল উপলব্ধি: Streaming-এ "correctness" ও "freshness" trade-off — দু'টি একসাথে সর্বোচ্চ পাওয়া যায় না। Business context define করে balance — ride matching-এ freshness বেশি গুরুত্বপূর্ণ; revenue counting-এ correctness।
প্র ০৩ Micro-batch (Spark, ১ s window) বনাম true streaming (Flink, event-by-event) — আজকের Bangladesh fintech startup কোনটি বাছবে? কী trade-offs?
এই প্রশ্ন প্রায় প্রতিটি tech lead interview-এ আসে। উত্তর nuanced — black-and-white নয়।
Spark Structured Streaming-এর pros:
- একই API batch ও stream-এ: DataFrame জানলেই হলো। Team upskilling সহজ।
- Massive ecosystem: Delta Lake, MLlib, GraphX integration native।
- Cloud support সর্বত্র: Databricks, EMR, Dataproc, Synapse — সব platform-এ।
- Bangladesh-এ engineer pool: Spark জানে অনেক — Flink জানে কম।
- Resource efficiency: Micro-batch-এ CPU utilization ভাল।
Spark-এর cons:
- Latency floor ৫০০ ms-১ s — sub-100ms সম্ভব না।
- Complex stateful workload (large state, many joins) — Flink-এর তুলনায় memory inefficient।
- Continuous processing mode (১ ms latency) এখনো experimental।
Flink-এর pros:
- True event-at-a-time: latency ১০ ms-এর নিচে সম্ভব।
- Sophisticated state management: RocksDB-backed, TB-level state।
- Better watermark handling: CEP (Complex Event Processing), session windows।
- Exactly-once প্রাকৃতিক: two-phase commit সরাসরি।
- Backpressure handling: automatic — Spark-এ manually tune।
Flink-এর cons:
- Engineer খুঁজে পাওয়া কঠিন। Bangladesh-এ Flink expert হাতেগোনা।
- Cloud managed offering কম (AWS Kinesis Data Analytics, Confluent Cloud)।
- Java/Scala-centric। PyFlink ক্রমে ভালো হচ্ছে কিন্তু feature gap আছে।
- Steep learning curve।
Bangladesh fintech startup-এর জন্য recommendation:
- Phase 1 (০-২ বছর): Spark Structured Streaming. Team-এর সাথে fit, ecosystem সমৃদ্ধ, hiring সহজ। ৯৫% use case cover।
- Phase 2 (২+ বছর): যদি sub-100ms latency দরকার হয় (HFT-জাতীয়, complex CEP) — তাহলে Flink-এ specific workload migrate।
- Hybrid is fine: বেশিরভাগ pipeline Spark, কিছু critical low-latency Flink — Kafka common backbone।
Bangladesh-specific consideration: AWS Mumbai region থেকে Dhaka latency ~৩০-৫০ ms। কোনো streaming-এ এটাই floor। তাই Flink-এর ১০ ms vs Spark-এর ৫০০ ms — user-perceived পার্থক্য সাধারণত ৫০-১৫০ ms (network dominates)। Real differentiation শুধু extreme low-latency বা complex stateful workload-এ।
মূল উপলব্ধি: Tooling-এর choice business need ও team capability-এর function — pure technical superiority না। ৯৫% Bangladesh fintech-এ Spark Structured Streaming + Kafka enough। Flink optimization, premature adoption নয়।
প্র ০৪ Daraz-এর checkout flow — ব্যবহারকারী cart-এ পণ্য যোগ করে, payment করে, order confirmed হয়। এই journey-তে কোন event streaming-এ যাবে, কোনটি batch-এ? একটি concrete pipeline design করুন।
E-commerce checkout — multi-step business process। প্রতিটি step-এর latency requirement আলাদা। ভাল design মানে — সঠিক step-এ সঠিক tool।
Events ও তাদের nature:
cart_add: user pn যোগ করল → streaming (recommendation, abandonment alert)।checkout_initiated: user "Pay" চাপলো → streaming (fraud check)।payment_success: bKash/card confirmed → streaming (instant order confirmation)।inventory_decrement: stock কমলো → streaming (overselling রোধ)।order_packed: warehouse staff scan → streaming (delivery ETA)।delivery_attempted: rider feedback → streaming (Pathao হস্তান্তর)।review_submitted: ৭ দিন পর — batch (recommendation model retrain)।
Pipeline design:
Layer 1 — Ingestion:
- Mobile app/web → backend API → Kafka topic (একটি event type = একটি topic)।
- Topics:
cart-events,checkout-events,order-events,delivery-events। - Partition key:
user_id— একই user-এর events একই partition-এ (ordering রক্ষা)। - Retention: ৭ দিন (batch reprocess-এর জন্য)।
Layer 2 — Streaming:
- Fraud detection job (Flink):
checkout-eventsconsume → user history check (Redis lookup) → ১০০ ms-এ allow/block। - Inventory job (Spark Streaming):
order-events→ inventory table update। Idempotent (order_id key)। - Cart abandonment (Spark Streaming):
cart-events→ ১৫ মিনিট inactivity window → notification topic। - Live ops dashboard (Spark Streaming): minute-level aggregation → Druid/ClickHouse → Grafana।
Layer 3 — Batch (নাইটলি):
- Data lake archive: সব Kafka topic → S3 (Parquet, partitioned by date) — Airflow DAG।
- dbt transformations: staging → marts → reporting tables।
- ML feature engineering: user purchase history → feature store।
- Recommendation model retrain: সাপ্তাহিক — Spark MLlib।
- Finance report: per-merchant settlement, tax — daily batch।
Layer 4 — Serving:
- Real-time: Redis (user session, recent activity)।
- Operational analytics: ClickHouse (live dashboard)।
- Reporting: Snowflake/BigQuery (BI team)।
- ML serving: Feature store + model server।
Why this split:
- Customer-facing actions (checkout, fraud, inventory) — streaming। Latency = revenue।
- Background analytics (recommendation, finance) — batch। Latency > ১ ঘণ্টা acceptable।
- Same Kafka backbone — Lambda-style dual access (real-time consumer + S3 archive)।
Cost consideration (Bangladesh context): Kafka cluster (3 broker, MSK Mumbai) ~১২০ লাখ/বছর; Flink (Kinesis Analytics) ~৬০ লাখ; Spark batch on-demand ~২৪ লাখ। মোট ~২ কোটি/বছর — Daraz-এর scale-এ justifiable।
মূল উপলব্ধি: Production data architecture কখনো pure streaming বা pure batch নয় — সবসময় hybrid। Kafka backbone, streaming-এ live decision, batch-এ deep analysis — এই pattern-এ আজকের সব unicorn (Stripe, Uber, Airbnb) দাঁড়িয়ে।
অনুশীলন
-
Classify করুন: নিচের use case গুলো streaming না batch — যুক্তিসহ:
- (ক) BTRC-র মাসিক operator-ভিত্তিক traffic report
- (খ) bKash agent-এর জন্য low balance alert
- (গ) Daraz-এর "trending now" carousel
- (ঘ) NBR-এ বার্ষিক ট্যাক্স রিটার্ন
- (ক) Batch। মাসিক report — latency hours acceptable, throughput বিশাল (TBs of CDR)।
- (খ) Streaming। Agent balance < threshold → ১ মিনিটের ভেতর SMS — না হলে cash-out fail।
- (গ) Streaming (micro-batch enough)। ১৫ মিনিট freshness ভালো; user behavior real-time প্রতিফলিত।
- (ঘ) Batch। বছরে একবার, accuracy critical, latency irrelevant।
-
Latency budget হিসাব: bKash transaction-এ user "Send" চাপার পর "Success" দেখাতে কত সময়? End-to-end latency কোথায় কোথায় খরচ হয়?
সাধারণ budget ~২-৩ সেকেন্ড user perception। Breakdown:
- Mobile → API gateway: ১০০-৩০০ ms (4G network)
- API → fraud check (Flink + Redis): ৫০-১০০ ms
- API → bank/MFS settlement: ৫০০-১৫০০ ms (external dependency)
- Database commit: ২০-৫০ ms
- SMS/push notification: ১০০-৫০০ ms
- Mobile UI update: ৫০ ms
Streaming layer-এ ১০০ ms budget — তাই Spark micro-batch (৫০০ ms+) বদলে Flink বা in-memory check।
-
Watermark চিন্তা: আপনি Pathao ride request stream design করছেন। allowed_lateness কত রাখবেন? কেন? trade-off কী?
Recommendation: ১০-৩০ সেকেন্ড।
যুক্তি:
- Bangladesh 4G p99 latency ~৩-৫ সেকেন্ড — most events arrive within।
- ৩০ s+ late event business-এ অর্থহীন (rider cancelled)।
- Memory cost: প্রতি partition × ৩০ s × event rate = state size — manageable।
Trade-offs:
- ৫ s — কম memory, বেশি drop (rural area-এর events miss)।
- ৬০ s — accurate, কিন্তু window output ১ minute delayed। Live "demand heatmap" stale।
- Side output strategy — late events alert-এ পাঠান, separately handle।
আরও পড়ুন · ABCL TECH-এ আপনার পরবর্তী পদক্ষেপ
- পাঠ ১৭ · Apache Kafka পরিচিতি পরবর্তী পাঠ Streaming-এর backbone — distributed log। আজ যা শিখলেন, কাল tools।
- পাঠ ১৫ · dbt — analytics engineering আগের পাঠ Batch-এর জগৎ থেকে — dbt transformation, modular SQL।
- পাঠ ১৯ · Spark Structured Streaming এই পাঠের সাথে সম্পর্কিত Micro-batch বাস্তবায়ন — DataFrame API-তে streaming।
- সব AI Courses দেখুন ABCL TECH Python, ML, DL, NLP, CV, GenAI, RL, MLOps — সব AI কোর্স একসাথে।