Stream Processing — Event Time, Windows, State & Exactly-Once
Studies in this cluster, in series order. Each one keeps its own URL.
Data engineering
Pipelines, sketches, approximate aggregations, object storage, and stream processing with event time, windows, and exactly-once sinks.
- 1.Stream Processing — Event Time, Windows, State & Exactly-OnceStream processing is continuous computation over unbounded data: you keep results fresh as events arrive instead of recomputing everything on a schedule. The hard part is not reading from Kafka fast. It is answering four questions when data arrives late, out of order, and forever: what you compute, where in event time, when you emit, and how refinements relate. A production engine must also keep large keyed state consistent across crashes and rescaling, and get results into sinks exactly once.
- 2.Windowing — Tumbling, Hopping, Session & Global Windows with TriggersA window turns an infinite stream into finite chunks you can aggregate. The window type is a product decision, not a tuning knob: tumbling windows answer per minute, hopping windows answer over the last five minutes updated every minute, session windows answer per visit, and a global window answers ever, until you say so. Triggers decide when a window emits, and the accumulation mode decides what each emission means downstream. Getting these wrong produces numbers that look plausible and are silently wrong.
- 3.Event Time & Watermarks — Allowed Lateness, Side Outputs & Idle SourcesIn event-time processing the engine never knows a window is complete. A phone that was in airplane mode can upload yesterday's clicks right now. A watermark is the engine's declaration that it believes no more events older than W will arrive. It drives on-time window firing, timers, and state cleanup. Too aggressive and you drop or mis-count late data. Too conservative and every result is delayed. Allowed lateness keeps window state around a bit longer for corrections, a side output captures anything later than that, and idleness handling stops one silent partition from freezing the whole job.
- 4.Stateful Streaming — Keyed State, RocksDB, Checkpoints & RescalingStateful streaming means the operator remembers things between events: running counts, open windows, join buffers, dedupe sets, feature rows. That state can be terabytes. It must survive crashes consistently with the source offsets, and it must be redistributable when you change parallelism. Flink's reference design is keyed state partitioned into key groups, stored in RocksDB, snapshotted by barrier checkpoints into object storage, and restorable through savepoints. Spark and Kafka Streams solve the same problem with different trade-offs.
- 5.End-to-End Exactly-Once — Replayable Sources, 2PC & Idempotent SinksExactly-once in streaming never means each record is physically processed once. After a crash, records are replayed. It means the effect of each record on state and on externally visible output happens once. That property is end-to-end only if three links hold: a replayable source with offsets you can rewind, checkpointed state consistent with those offsets, and a sink that is idempotent or transactional and commits only what a completed checkpoint covers. Break any link and you are back to at-least-once or at-most-once.
- 6.Streaming Joins, Backpressure & Operations — Interval, Temporal, Skew & BackfillJoins and flow control are where streaming jobs fall over in production. A stream-stream join must buffer both sides, so it needs a time bound or its state grows forever. A stream-table join must decide which version of the row applies: the latest value, or the version that was valid at the event's time. Backpressure is how a healthy engine slows a fast source instead of crashing. Skew is why one subtask is at 100 percent while the rest idle. Operating the job means watching lag, planning a backfill, and upgrading state safely.