Transport is not processing
A broker retains and distributes events. A processor interprets them, groups them, joins reference data and maintains results. An analytics store serves queries. These boundaries have different failure and retention contracts.
For website activity, the pipeline might be event receiver -> durable log -> window processor -> aggregate store -> dashboard. Define whether the displayed result is provisional or final and how old it may be. A broker accepting an event does not mean the dashboard already includes it.
Three clocks
| Time | Meaning | Example failure |
|---|---|---|
| Event time | When the recorded activity happened | Device clocks may be wrong |
| Ingestion time | When infrastructure received it | Offline uploads arrive hours late |
| Processing time | When a worker handles it | A backlog shifts apparent activity |
An event-time chart can describe when users acted; a processing-time chart can describe worker throughput. Neither is universally better. Store enough metadata to distinguish them and reject or quarantine implausible timestamps according to policy.
Window families
A tumbling window partitions time into nonoverlapping buckets. A hopping/sliding window has a fixed width and smaller advance, so one event contributes to multiple windows. A session window groups activity separated by a configured inactivity gap; a late event can bridge two previously separate sessions.
For ten-minute windows advancing every minute, an event may affect ten aggregates. This extra state and output work belongs in the capacity estimate. Specify boundaries such as [10:00,10:10) so an event at exactly 10:10 is assigned consistently.
Watermarks describe progress assumptions
A watermark expresses event-time progress under the processor's lateness assumptions. It is not proof that no older event can arrive. Multi-input processing may be held back by a slow input; idleness handling can prevent a silent partition from stalling progress indefinitely, but resuming that partition introduces late data.
Choose a bounded out-of-order allowance, watch watermark lag and distinguish it from broker consumer lag. A processor can keep up with bytes while waiting for event-time progress.
Worked late-event timeline
A purchase occurred at 10:02 and arrives at 10:09. The current watermark is 10:08, so it is older than the watermark. Whether it changes the 10:00–10:05 aggregate depends on the chosen window and allowed-lateness/output policy.
Options are to discard with an explicit late-event count, write to a correction stream, or update/retract the previously emitted result. For financial reconciliation, silently discarding valid activity is usually unacceptable. A live dashboard may tolerate provisional totals followed by corrections.
State size and expiration
Ten million active user-window keys at an assumed 80 bytes of application state occupy about 800 MB before index, serialization, allocator, checkpoint and runtime overhead. Overlapping windows and joins increase that total.
Bound state by event-time cleanup and business retention. Do not expire deduplication keys earlier than the retry/replay horizon if duplicate suppression depends on them. Hot keys can concentrate work even with many partitions; split compatible aggregates into partials, then merge. A non-associative computation may not admit this optimization.
Joins over streams
A stream-table join enriches each event with reference data. Decide whether it uses the latest profile or the profile valid when the event occurred. Changing the reference data can otherwise change the meaning of replayed results.
A stream-stream join buffers both sides and needs a time range and expiry rule. An unbounded join can retain unmatched records forever. Missing counterparts require a policy: wait, emit incomplete output, dead-letter, or reconcile later. Model duplicate inputs so a replay does not multiply the joined result.
Checkpoints and external outputs
A recoverable processing checkpoint coordinates source positions and operator state. After restart, the processor restores state and rereads events from the appropriate position. An external sink needs a compatible transactional protocol or idempotent writes to avoid duplicating already-visible effects.
A useful aggregate sink key is metric plus window plus dimensions, with a versioned upsert. For corrections, compare the result version atomically so a delayed old output does not overwrite a new one. Exactly-once claims must identify the covered boundary; sending email from inside a processing callback remains an external side effect.
Approximate analytics
HyperLogLog estimates distinct counts using bounded state; it does not tell you which exact users were present. Count-Min Sketch estimates frequency with a defined error model; collisions can overestimate. Top-k algorithms trade memory against completeness and update behavior.
Use approximations for dashboards where error is acceptable, not for billing or authorization. Label the result and measure error against a smaller exact sample. Merging sketch state requires compatible parameters and hashes. Counting distinct users separately in each partition and summing is wrong when the same user appears in several partitions.
OLAP serving and freshness
Columnar stores are effective for scanning selected columns and aggregating many rows. Preaggregate common queries but retain a rebuild path. Partition by useful time ranges, choose sorting keys for actual filters and avoid creating an unmanageable partition per user.
Expose ingestion freshness and correction status. A dashboard timestamp should mean something precise: newest event processed, last successful batch or query execution time are different indicators. Late events can change historical buckets after the latest ingestion time advances.
Backfill and replay
Replay retained source data through a versioned processor into a separate result namespace. Isolate its resource budget from live traffic. Compare representative windows, catch up current events and switch queries only after validation. Replaying old data through a newly changed reference table may produce different results; snapshot/reference-version policy makes that intentional or exposes the discrepancy.
Exercise
Design daily active learners and five-minute popular-course rankings. Include event identity, time selection, late-arrival policy, distinct-count semantics, state expiry, sink keys and a full rebuild. Inject a checkpoint failure after writing output and an event from yesterday. Continue with the Kafka lab.