Overview
Curated: · Written: · Reviewed:
Batch versus streaming is an outcome and time-semantics decision
Batch processing operates on a bounded collection—often a snapshot or interval—and publishes results after that input is known. Streaming processing continuously consumes an unbounded sequence and updates state or outputs as events arrive. Micro-batch engines execute frequent bounded increments while offering a streaming programming model. None is automatically faster, cheaper, simpler, or more correct for every use case. Start with the decision being served: how quickly does new information change an action, how much lateness and correction is acceptable, what ordering and state are required, and what happens during failure?
Time has several meanings
Event time is when the business event occurred; ingestion time is when the platform accepted it; processing time is when an operator handled it. Offline devices, source buffering, network partitions, retries, replays and clock problems make these diverge. Event-time windows support reproducible historical reasoning, while processing-time behavior can be appropriate for operational actions explicitly tied to arrival. A watermark is a system's progress assertion or heuristic about event time, not proof that no earlier event can ever arrive. The watermark policy trades result latency, state retention and correction. Late data may update prior results, go to a side output, trigger later recomputation, or be dropped only under an explicit business contract.
Windows group an unbounded stream into finite logical results: fixed/tumbling, sliding/hopping, session, or global with triggers. Triggers decide when panes are emitted; accumulation mode and sink semantics decide whether later panes replace, retract, or add. Consumers need to know whether a value is preliminary or final and how corrections arrive. Joining streams adds two-sided time bounds, watermarks, state expiry and missing-match semantics. One idle or skewed input partition can stall a downstream watermark unless idleness and alignment are handled deliberately.
Delivery and state are end-to-end concerns
At-most-once may lose records, at-least-once may replay them, and exactly-once engine state does not automatically make an external side effect exactly once. End-to-end outcomes depend on replayable sources, checkpoints, deterministic logic, transactional or idempotent sinks, stable event/operation identities and commit coordination. Checkpoint frequency trades runtime overhead, recovery work and state durability. State must be bounded with windows, TTL or business completion; large or hot keys cause skew, memory pressure and slow checkpoints.
Partition keys determine ordering and parallelism. Most systems preserve order only within a partition, not globally. More partitions raise throughput but increase coordination and can break assumptions during repartitioning. Backpressure is a correctness and reliability signal: when downstream capacity falls behind input, systems buffer, slow sources, spill, shed under an explicit policy, or fail. Hidden unbounded queues only turn latency into a later outage. Monitor input rate, processing rate, consumer lag, event-time lag, watermarks, state size, checkpoint time/failure, backpressure, late/drop/correction counts, sink commits and end-to-end decision freshness.
Batch remains valuable for complete-period reporting, large recomputation, model training, reconciliation and consumers whose useful latency is hours. Streaming suits fraud response, monitoring, personalization and control loops where seconds or minutes materially change outcomes. Hybrid architectures often keep one durable event log, produce low-latency provisional views, and reconcile or restate authoritative history in batch. Avoid two unrelated implementations of the same business logic when a shared semantic model or unified engine can reduce drift, but validate batch/stream equivalence with replayed representative history.
Engine mechanics: watermarks, checkpoints, and sinks
Flink assigns timestamps per record and generates watermarks from source-partition progress. An idle partition that never emits can stall a downstream watermark unless idleness timeouts and alignment are configured; a too-aggressive watermark closes windows while late events are still in flight. Checkpointing snapshots operator state at a barrier; exactly-once checkpoint mode coordinates with two-phase commit sinks where supported, but at-least-once checkpoints plus a non-idempotent HTTP call still duplicate effects. Spark Structured Streaming typically executes micro-batch increments with a checkpoint location that stores offsets and state; changing the query identity or deleting that directory is a new pipeline, not a repair.
Kafka preserves order only within a partition. The consumer group commits offsets independently of sink commits: committing offsets before a durable write creates at-most-once loss on crash; writing then failing before commit creates at-least-once replay. Increasing partitions raises throughput but does not globally order events, and key-based joins break if a producer changes the partitioner. Kinesis shards have analogous capacity and sequence-number semantics; a reshard changes the unit of parallelism and must be reflected in consumer state. Google Dataflow's exactly-once documentation is explicit about limitations: engine-side exactly-once does not automatically cover side effects outside the pipeline's sink protocol.
Apache Beam's programming model unifies bounded and unbounded PCollections with windows, watermarks, triggers, and accumulation. An after-watermark trigger plus discarding accumulation emits a pane once and drops later updates; accumulating or retracting modes require sinks that can retract. Iceberg atomic commits are a common analytical sink: a streaming job should commit a snapshot that is complete for the checkpoint, not a directory of files a reader can partially list. OpenLineage run events for streaming jobs need a stable job identity plus run or checkpoint identity so a restart is not graphed as an unrelated producer.
Figures and failure modes worth measuring
Track event-time lag (processing time minus watermark), watermark delay (max observed event time minus watermark), processing-time lag, consumer lag in offsets or shard iterators, state size bytes, checkpoint duration and failure rate, and late/drop/correction counts as first-class SLIs. A median latency of 200 ms with a 99th-percentile event-time lag of 45 minutes is a watermark or idle-source problem, not a “streaming is fast” success. Poison records that crash an operator on one key can stall a partition while others proceed; isolate them to a dead-letter path with the original offset retained. Hybrid designs that emit a low-latency provisional view and a nightly Iceberg restatement must declare which figure is authoritative for which decision, or two dashboards will both look fresh while disagreeing.
Operate either mode as a data product: version schemas and logic, protect sensitive events and metadata, test duplicates, gaps, disorder, late events, partition movement, source and sink outages, checkpoint corruption, poison records and replay, and publish coherent versions. Measure decision value, correctness, recovery and total operational burden—not only median processing latency.
Worked example: a $500 checkout delayed two minutes
Event time 10:00:01, $500 authorized. The producer buffer holds it until 10:02:05. Fraud wants a one-minute event-time window covering 10:00.
| clock | watermark / close | where the $500 lands | fraud decision |
|---|---|---|---|
| processing-time 1-minute windows | close on wall clock | 10:02 window | wrong minute; lookalike replay |
| event time, watermark +5s | window [10:00,10:01) closed at 10:00:06 | dropped as late | false negative, $500 unreviewed |
| event time, watermark +3 min | same window stays open | counted in 10:00 | correct, 3 minutes slower |
Streaming is not “faster.” That table is the interview: which clock, how late, and which figure is allowed to be wrong.
