ESC
Type to search guides, tutorials, and reference documentation.

Real-Time Analytics and Stream Processing

Partitioned logs, event time versus processing time, watermarks, windows, state, delivery semantics, consumer lag and backpressure — and when streaming is not worth its cost.

Stream processing is the practice of computing over data as it arrives, as an unbounded sequence of events rather than a finite table. The defining difference from batch is not speed — it is that there is no end of input. A batch job can sort, count or join because it knows when it has seen everything; a streaming job never knows, so it must decide when a result is complete enough to emit, what to do when a contradicting event arrives afterwards, and how much state it is willing to hold while it waits. Real-time analytics is what you get when those decisions are made deliberately.

The log, partitions and consumer groups

Modern streaming platforms are built on a partitioned, append-only log. Producers append records; the log retains them for a configured period; consumers read forward at their own pace, tracking an offset. Three consequences follow, and they explain most of the behaviour you will observe.

Ordering is per partition, not global. Records are ordered within a partition and unordered across partitions. Since the partition is chosen by a key, all events for the same key travel in order — which is exactly the guarantee that per-entity processing needs, and the reason key selection is the most consequential design decision in a streaming system.

Reading is non-destructive. Consuming does not remove the record. Retention is time or size based, so several independent consumers can read the same stream, and a consumer can rewind and reprocess history. Replay-driven recovery is the single biggest operational advantage of a log over a queue.

Parallelism is bounded by partition count. One partition is consumed by at most one member of a consumer group, so the partition count sets the ceiling on parallelism. Raising it later redistributes keys and therefore breaks the co-location that stateful operators depend on, which is why partition count deserves thought before launch.

Event time, processing time and watermarks

Every streaming system deals with two clocks. Event time is when the thing happened, carried in the record. Processing time is when your operator saw it. They diverge because of network delay, retries, buffering, offline clients and outages — and the divergence is unbounded, not bounded by a nice constant.

Computing on processing time is trivial and produces results that are not reproducible: replaying the same input on a different day gives different windows. Computing on event time is correct and forces you to answer the completeness question. A watermark is the system's assertion that no further events with a timestamp earlier than some point are expected, and it is the mechanism that makes event-time windows emit. Watermarks are heuristics. Set them too tight and late events are dropped or force expensive retractions; set them too loose and every window waits, adding latency and holding state. This tension is the heart of streaming: latency, completeness and cost — choose two.

Because watermarks are heuristics, a serious design also specifies allowed lateness (how long after a window closes a straggler is still accepted) and what happens beyond it: dropped and counted, or routed to a side output for later correction. Silently dropping late data is how a streaming total quietly diverges from the batch total computed on the same source.

Windows, joins and state

Aggregation over an unbounded stream requires a window. Tumbling windows are fixed and non-overlapping, so each event lands in exactly one. Sliding (hopping) windows overlap, so each event lands in several — which multiplies both output volume and state. Session windows are defined by a gap of inactivity, so their boundaries depend on the data; they are the right model for user behaviour and the most expensive to maintain, because a session's state must be held open until the gap elapses.

Anything beyond a stateless map holds state: aggregates, join buffers, deduplication sets, session accumulators. State is the resource that actually constrains a streaming application. It must be kept locally to be fast, checkpointed durably to survive failure, restored on restart, and — crucially — bounded. Any keyspace that grows forever (user identifiers, session tokens, request ids) will exhaust the operator unless entries expire. Unbounded state growth is the most common cause of a streaming job that runs fine for weeks and then fails repeatedly at the same point.

Joins deserve particular caution. A stream-to-stream join is a join over two buffers held for a bounded interval, and it only produces matches that occur within that interval — a fact that must be explained to whoever consumes the output, because the missing rows look like data loss. A stream-to-table join against a slowly changing dataset is usually easier to reason about and far cheaper.

Delivery semantics and what "exactly once" means

At-most-once drops on failure. At-least-once is what you get by default from any system that retries: a record may be processed more than once, so duplicates reach your output. What is marketed as exactly-once is more precisely effectively-once processing, and it is achieved by combining three things — idempotent producers that deduplicate retried writes, a transactional commit that couples the output write with the consumer's offset advance, and a sink that can participate in that commit or is itself idempotent.

The part worth internalising is that the guarantee ends at the boundary of the system that provides it. If a job writes to an external store that cannot take part in the transaction, the effective guarantee reverts to at-least-once, and correctness depends on the write being idempotent — an upsert on a key, not an increment of a counter. Any side effect that is not idempotent, such as sending an email or charging a card, must be made idempotent with a deduplication key of your own. No platform setting can fix that for you.

Backpressure, lag and rebalancing

Consumer lag — the distance between the latest offset and the one a consumer has processed — is the primary health metric. Steady lag means keeping up; growing lag means the consumer is slower than the producer and will never catch up on its own. Lag that grows during a spike and drains afterwards is normal and is exactly what the log's buffering is for. Lag that never drains is a capacity or a skew problem.

Backpressure is the mechanism by which a slow stage slows the stages above it rather than accumulating an unbounded queue. A pipeline without backpressure does not stay fast; it fails later and harder, having buffered memory it could not release.

Rebalancing happens when a consumer joins or leaves a group and partitions are reassigned. During a rebalance, processing pauses and stateful operators may have to restore state on a new host. Frequent rebalances — usually caused by a processing loop exceeding the group's timeout, or by aggressive autoscaling — can consume more time than the actual work, producing a job that appears busy while making no progress.

Failure modes

  • Hot key. One key carries a large share of traffic, so one partition and one consumer saturate while the rest idle. Partition count cannot fix it; the key must be salted or the aggregation restructured.
  • Unbounded state. No expiry on a growing keyspace; the job dies at a size that only appears after weeks in production.
  • Silent late-data loss. Watermarks too tight, stragglers dropped without a counter, totals drift from the source and nobody can explain the gap.
  • Replay double-counting. Reprocessing history into a sink that appends rather than upserts. Every replay inflates the result, which is discovered only after the incident.
  • Retention shorter than recovery time. A consumer down longer than the retention window resumes at a truncated offset and has a permanent hole in its output.
  • Schema change without compatibility rules. A producer changes a field and every consumer fails at once; a schema registry with enforced compatibility exists precisely to make this impossible.
  • Two implementations of one definition. The streaming path and the batch path compute the same metric in different code, and they diverge. This is the standing cost of a dual-path architecture.

When not to stream

Streaming is a permanent increase in operational complexity: always-on jobs, state to manage, replay procedures to rehearse, and a class of bugs that only appear under real-world event-time disorder. It earns that cost when a decision is made automatically and immediately — fraud interception, alerting, live inventory, dynamic limits — or when the volume is high enough that repeated batch recomputation is more expensive than incremental maintenance.

It does not earn it when the consumer is a human looking at a dashboard on a cadence, when the upstream source only updates periodically anyway, or when "real time" was requested without anyone naming the decision that the freshness would change. Frequent micro-batches capture much of the benefit at a fraction of the complexity, and are the right default for most analytics. Compare with scheduled data pipelines, see big data architecture for where the streaming path sits in a platform, and data warehousing for where streamed data usually lands for historical analysis.