Streams, Event Time and Late Arrivals
Introductions, exercises and summaries stay visible.
41.0 What this chapter gives you#
- A stream keeps bringing records while a question is being answered. There may be no natural final record telling the system that a day’s information is complete.
- You will separate the time an event claims to describe from the time a processor receives it. Then you will build windowed answers that can be revised when information arrives late, without counting a replay twice.
- The chapter uses an original deterministic time-window model. Its explicit watermarks, lateness limits and output revisions are teaching policies, not an implementation or performance test of Apache Beam or another distributed stream processor.
41.1 Continuous records#
41.1.1 PLAIN — in simple words#
- A stream is a continuing sequence of records. A till, delivery scanner or application can keep producing new information while yesterday’s records are still being processed.
- Continuous arrival does not guarantee continuous progress. A source can pause, a connection can fail, and a processor can fall behind while the producer keeps working.
- Before processing a stream, define what one record represents and how repetitions, corrections and gaps are recognised. Those responsibilities do not disappear when a system is described as real-time.
41.1.2 PLAIN — a picture in your head#
- Mira receives numbered slips through a slot while she is counting the previous slips. There is no sign saying that the last slip has arrived forever.
- She can still produce a useful interim answer, provided she says which slips it includes and what might change later.
- Where the comparison breaks: a distributed stream can have several independently ordered inputs rather than one slot. A position in one input is not automatically a position in a single global sequence.
41.1.3 PLAIN — a worked example#
- Suppose source A has processed records through position 50 and source B through position 80. These positions belong to different sequences and cannot be added to mean “global position 130.”
- An honest progress record is a pair: A through 50 and B through 80, together with any required gap checks or source-generation identifiers.
- If A stops sending, B may reach 500 while A remains at 50. A large maximum position does not establish that a report using both sources is complete.
- A report can either wait, publish a provisional result with the lag disclosed, or follow an explicit incomplete-data policy. The choice changes the product’s meaning, not just its speed.
41.1.4 PLAIN — what is really happening inside#
- A processor reads records, updates state and emits results. Durable progress and state recovery determine which work is repeated after a failure.
- Partitioning permits parallel work, but each partition may advance independently. Coordination is needed when one answer depends on several partitions.
- A bounded batch can be viewed as a finite collection. An unbounded stream needs continuing policies for time, state size, late input and result revision rather than assuming that all data will eventually fit in memory.
41.1.5 TECHNICAL — the engineer’s version#
- Apache Beam distinguishes bounded and unbounded collections and provides windowing and triggers for grouping and emitting results over continuing input. The runtime details depend on the selected runner and connectors. [S167]
- Progress is typically partition-scoped. Include stream identity, generation and offset semantics when designing recovery; a numeric offset alone can be ambiguous after recreation or migration.
- The local exercises use finite sequences to expose streaming rules reproducibly. They do not claim that a finite successful run establishes unbounded operational stability or distributed fault tolerance.
41.1.6 WORDS — remember these#
Stream: an ongoing supply of records — input processed as a continuing sequence rather than only as a completed finite batch. Unbounded input: no known final record — a collection whose completion is not generally available during normal operation. Partition progress: how far one input segment has advanced — a scoped position that must not be confused with universal completion.
41.2 Event and processing time#
41.2.1 PLAIN — in simple words#
- Event time is the time attached to the event being described. Processing time is when a particular system handles the record.
- They can differ because of offline devices, queues, retries, transmission delays or incorrect source clocks. A late arrival is not necessarily a late real-world event.
- Choose which time answers the question. “What happened before noon?” and “What did the server learn before noon?” are both useful, but they are different reports.
41.2.2 PLAIN — a picture in your head#
- A delivery driver writes 10:05 on a paper receipt and returns to the shop at 12:20. The delivery time and the time Mira receives the receipt are different facts.
- Filing the receipt under 12:20 may help explain office workload. Using it to claim the delivery happened at 12:20 would change the event’s meaning.
- Where the comparison breaks: a source timestamp may come from a poorly set or manipulated clock. Preserving it does not require treating it as unquestionably accurate.
41.2.3 PLAIN — a worked example#
- Consider a separate stream fixture with integer minutes from an agreed origin. E1 has event time 2 and value 100; E2 has event time 7 and value 200.
- The processor receives both before it announces completion progress for minute 10. Later it receives E3, whose event time is 6 and value is 50.
- By event time, all three belong to the first ten-minute interval. By arrival order, E3 follows the first published answer. Both descriptions can be stored without rewriting either fact.
- These values are invented abstract amounts, not additional sales in the canonical shop ledger. The example must not increase the earlier agreed total of 28,650 paise.
41.2.4 PLAIN — what is really happening inside#
- Timestamp assignment may come from source fields, transport metadata or a processor’s clock. The assignment rule is part of the pipeline contract.
- Conversion must preserve the relevant time zone and precision. A local date without a zone cannot always be mapped to one instant, and a clock reading is not a causal ordering proof.
- Implausible future timestamps can distort progress estimates if the system blindly uses the largest observed event time. Validating timestamps and recording clock uncertainty is therefore an operational concern, not only a presentation choice.
41.2.5 TECHNICAL — the engineer’s version#
- Event-time semantics require a defined timestamp extractor and temporal domain. Processing-time semantics instead follow the executing system’s clock and can change when the same input is replayed. [S167]
- Preserve original event time, ingestion time and processing metadata when their distinct meanings matter. Do not overwrite an inconvenient source timestamp with arrival time and keep the old label.
- Ordering by timestamp is not equivalent to happened-before. Clock error and concurrent events remain possible, so causal or per-entity ordering may require other identifiers or sequence information. [S41]
41.2.6 WORDS — remember these#
Event time: when the described event is said to have happened — a timestamp assigned under the source and pipeline’s semantic contract. Processing time: when a processor handles a record — time measured by the executing system, potentially different on replay. Ingestion time: when a system accepted the input — a separately useful observation that does not replace the original event time.
41.3 Windows#
41.3.1 PLAIN — in simple words#
- A window gives a bounded group over which a stream can answer a question. It might mean each ten-minute interval or each period of activity separated by a long gap.
- Window membership needs exact boundaries. Otherwise an event precisely at 10:00 can be counted twice or in neither interval.
- Windowing and publication are separate. Deciding which group contains a record does not decide when the result should be shown or whether it may later change.
41.3.2 PLAIN — a picture in your head#
- Mira places time-stamped slips into trays labelled 09:00 up to, but not including, 09:10; then 09:10 up to, but not including, 09:20.
- A slip at exactly 09:10 belongs to the second tray. The labels make the handover unambiguous.
- Where the comparison breaks: sliding windows can deliberately place one event into several groups. Session windows can merge after a late event connects periods that previously looked separate.
41.3.3 PLAIN — a worked example#
- With fixed windows of width 10, event times 2, 6 and 7 belong to
[0, 10). Time 10 belongs to[10, 20). The square bracket includes the start; the round bracket excludes the end. - For non-negative integer time t, the window start is
(t // 10) * 10. At t = 27, integer division gives 2, so the window is[20, 30). - With width 10 and a new sliding window every 5 minutes, time 7
belongs to both
[0, 10)and[5, 15). Summing the totals of overlapping windows would count some events more than once. - In a separate session rule, each event creates an interval
[t, t + 5), and overlapping intervals merge. Events at 1 and 4 form[1, 9). An event at 12 forms[12, 17). A late event at 8 bridges them into[1, 17).
41.3.4 PLAIN — what is really happening inside#
- A window function assigns each record to one or more windows. Grouping usually also includes an entity or report key, so one branch’s total is not accidentally combined with another’s.
- Fixed windows simplify alignment. Sliding windows answer overlapping recent-history questions. Session windows follow bursts of activity and require rules for merging state and revising previously emitted results.
- Calendar windows need additional definitions. A local day is not always a fixed duration in regions with clock changes. For a business report, the relevant zone and calendar policy belong in the metric definition.
41.3.5 TECHNICAL — the engineer’s version#
- Window assignment, trigger policy, accumulation mode and allowed lateness are distinct configuration dimensions. Matching only the window width does not make two streaming reports equivalent. [S167]
- The local model supports fixed ten-unit windows and explicitly non-negative integer event times. Its session example is a separate interval-union calculation, not a claim to implement every streaming session-merge protocol.
- Key the published result by report definition, entity scope and window boundary. A revision identifier can then distinguish successive answers to the same question without treating them as independent amounts to add.
41.3.6 WORDS — remember these#
Window: a defined group of time-related records — a membership rule used to bound aggregation over a stream. Sliding window: overlapping intervals — windows whose width exceeds their step, allowing one event to contribute to multiple results. Session window: a group linked by activity gaps — a window that may merge when an arriving event connects earlier groups.
41.4 Watermarks#
41.4.1 PLAIN — in simple words#
- A watermark is a progress statement about event time. It helps a processor decide that the ordinary waiting period for earlier records has passed.
- It is not a magical observation of every device in the world. Its strength depends on the source and the assumptions used to produce it.
- A record arriving behind the watermark is late under that policy. It may still describe a valid event, so “late” and “wrong” must not be used as synonyms.
41.4.2 PLAIN — a picture in your head#
- Mira announces that the normal courier has delivered slips through the morning interval. She publishes a provisional count while keeping a procedure for exceptional late envelopes.
- If a branch is offline, she must decide whether to wait for it or disclose that the count can still change.
- Where the comparison breaks: some sources can offer stronger progress guarantees than an estimated courier schedule. Other watermarks are heuristic estimates. Their meanings must be read from the actual source and framework contract.
41.4.3 PLAIN — a worked example#
- In our model, an external test driver advances a non-decreasing
watermark. When it reaches 10, the model publishes revision 1 for
[0, 10)using E1 and E2: 100 + 200 = 300. - Advancing the watermark is a declared input to the exercise. The model does not derive it from a network, a real clock or an unbounded source.
- Suppose two required input partitions have watermarks 20 and 7. A conservative combined frontier is 7, because the second partition has not provided the same progress evidence as the first.
- Marking the slow partition idle may permit progress under a chosen framework policy. When it resumes with older records, those records still need a late-data policy; idleness does not make them impossible.
41.4.4 PLAIN — what is really happening inside#
- Sources and operators maintain progress estimates or guarantees. Downstream stages combine them with outstanding work, timers and state dependencies.
- A watermark should not move backward within the same progress history. Recovered or recreated sources also need identity and state rules so that an old position is not mistaken for current progress.
- Monitoring should distinguish event-time lag, processing backlog and idle inputs. A processor may be fast at handling available records while a disconnected source makes its answer incomplete.
41.4.5 TECHNICAL — the engineer’s version#
- Beam describes a watermark as its notion of when data for earlier event times can be expected to have arrived. Its guide separates this from trigger configuration and allowed late input. [S167]
- A common estimate based on observed event time minus a delay is a policy, not a proof of source completeness. A forged future timestamp can advance such an estimate unsafely unless the system applies appropriate validation and source-specific safeguards.
- Our integer model rejects a decreasing watermark and records late dispositions. Those checks establish internal consistency of the chosen model, not the truth of an externally supplied progress claim.
41.4.6 WORDS — remember these#
Watermark: a statement of event-time progress — a source- and framework-dependent frontier used for window and timer decisions. Late data: a record behind the accepted progress frontier — input whose event time precedes the relevant watermark or closed-window boundary. Idle input: a source currently producing no records — a condition requiring a policy, not proof that earlier events can never return.
41.5 Late corrections#
41.5.1 PLAIN — in simple words#
- Publishing quickly and waiting for every possible late record are competing goals. A system must choose how long to accept revisions and what to do after that limit.
- A later result should say whether it replaces the earlier result or adds a difference to it. Mixing those meanings can double a report even when each emitted number is correct.
- Repeated delivery of the same event is different from a new late event. Recognising the difference requires identity and payload checks, not only comparing timestamps.
41.5.2 PLAIN — a picture in your head#
- Mira writes “morning total: 300, revision 1.” A late slip adds 50, so she writes “morning total: 350, revision 2.”
- Dev must replace the old total with 350. Adding the two displayed totals would produce 650, even though the actual included amounts sum to only 350.
- Where the comparison breaks: systems can instead emit a +50 change. Both replacement and delta designs can work, but the consumer must know which protocol it receives and how repeats are handled.
41.5.3 PLAIN — a worked example#
- The model permits late updates while the watermark is below window
end plus 5. For
[0, 10), the retention cutoff is therefore 15. - At watermark 11, E3 arrives with event time 6 and value 50. The model accepts it, changes the stored sum from 300 to 350 and publishes replacement revision 2.
- Repeating E3 with the same identity and content creates no revision and no additional 50. Reusing E3 with value 70 is a conflict, not a correction. A genuine correction needs its own explicit identity and relation to the earlier event.
- At watermark 15, a new E4 with event time 8 and value 25 is beyond the model’s acceptance boundary. It is recorded as too late and does not silently change 350 to 375. This equality-at-cutoff rule is our declared policy.
41.5.4 PLAIN — what is really happening inside#
- A trigger decides when to emit. Accumulating output may contain all accepted data so far; discarding output may contain only a new pane’s contribution. Consumers need the matching interpretation.
- Revisions need stable window identity and ordered version handling. Receiving revision 2 before revision 1 must not allow the older replacement to overwrite the newer answer.
- Records rejected by the streaming lateness policy may still enter a separately governed batch correction process. That process must preserve the difference between the earlier published view and a later restatement.
41.5.5 TECHNICAL — the engineer’s version#
- Beam’s trigger and accumulation settings affect when and how output panes are emitted. Allowed lateness provides a bounded opportunity to react to records after the watermark passes a window end. These settings must be coordinated with sink semantics. [S167]
- Our publisher uses
(window_start, revision, replacement_total)and a remembered event ID/payload mapping. Duplicate identity with identical content is a replay; conflicting content raises an error without mutating the accepted state. - The exercise retains final summaries and disposition evidence for inspection. It is not an unbounded deduplication store or a complete state-expiration implementation. Real retention must account for replay horizons and privacy requirements.
41.5.6 WORDS — remember these#
Trigger: the rule deciding when to publish — a condition for emitting a window result or pane. Replacement result: a new complete answer for the same key — output that supersedes an earlier revision rather than adding to it. Allowed lateness: the accepted revision interval — a configured policy governing how long late window data may affect retained computation.
41.6 State, checkpoints and replay#
41.6.1 PLAIN — in simple words#
- A streaming answer depends on remembered state: partial totals, seen identities, window boundaries and progress. Saving only the total can lose the information needed to resume safely.
- A restart may repeat records. If the restored deduplication state is older than the restored total, a repeated event can be added again.
- Recovery should restore a compatible combination of input progress, state and output publication evidence. Separate components must agree on what the checkpoint actually covers.
41.6.2 PLAIN — a picture in your head#
- Mira copies her total into a notebook but forgets to copy the list of slips already counted. After returning from a break, she cannot reliably distinguish counted slips from new ones.
- Saving the total and the counted-slip list together avoids that particular mismatch. She also needs to know which displayed revision customers may already have seen.
- Where the comparison breaks: distributed systems may checkpoint many operators and coordinate with transactional sinks. A single notebook illustrates the dependency but does not solve the distributed checkpoint protocol.
41.6.3 PLAIN — a worked example#
- After accepting E3, the local state contains sum 350, revision 2, seen identities E1/E2/E3 and watermark 11. Snapshot that complete model state.
- Restore it into a fresh model and replay E3 unchanged. The result remains 350 at revision 2. A new valid E5 in the same window can still be accepted while the watermark remains below 15.
- Now imagine restoring sum 350 but only seen identities E1/E2. E3 would look new and could incorrectly produce 400. This counterexample explains why state components must share a recovery boundary.
- If an external reader already received revision 2 before a crash, resending revision 2 should be recognised as the same result. Saving internal state alone does not prevent a remote consumer from adding a replacement twice.
41.6.4 PLAIN — what is really happening inside#
- Checkpointing serialises the state required to continue the computation. Source replay then reconstructs work after that checkpoint under the same rules.
- State schema evolution matters too. A newer program must either understand the older checkpoint format or use a reviewed migration/rebuild process.
- End-to-end correctness includes the sink. A framework can restore its own state consistently while an ordinary external API repeats a side effect that was acknowledged ambiguously.
41.6.5 TECHNICAL — the engineer’s version#
- The model snapshot includes its schema version, width, lateness policy, watermark, per-window state and accepted event identities. Restore validates the representation rather than executing arbitrary serialized objects.
- Use bounded, explicit data encodings for untrusted or persisted state. Python pickle is not a safe interchange format for untrusted inputs because loading it can execute code. [S62]
- The companion tests cover deterministic replay and revision handling on finite inputs. They do not exercise a distributed checkpoint coordinator, exactly-once external delivery, clock synchronization or recovery from physical power loss.
41.6.6 WORDS — remember these#
Operator state: remembered information used by a computation — partial aggregates, identities, timers and metadata needed to process later input. State checkpoint: a recoverable state boundary — a saved representation coordinated with the progress and effects it claims. Restatement: a later version of an earlier report — a deliberate correction whose relationship to previously published results remains visible.
41.97 Practice and worked answers#
- Question: Does a record received at 12:20 necessarily describe a 12:20 event? Answer: No. Event and receipt times are distinct, and the source clock’s reliability must also be considered.
- Question: Where does event time 10 belong with
fixed width 10? Answer: In
[10, 20), not[0, 10), under the declared half-open convention. - Question: Why should overlapping sliding-window totals not simply be added? Answer: The same event may belong to several windows, producing deliberate overlap.
- Question: The two input watermarks are 20 and 7. What is the conservative combined frontier? Answer: 7, when both inputs are required and no separate idleness policy changes the interpretation.
- Question: What replaces revision 1’s 300 after E3 adds 50? Answer: Revision 2’s complete total 350. Adding replacement totals would be incorrect.
- Question: What happens to E4 at watermark 15 in this model? Answer: It is recorded as too late for the first window and does not alter the published total.
- Question: Why save seen event IDs together with totals? Answer: Restoring inconsistent versions can cause already included events to be counted again.
- Question: Does a passing deterministic window test prove a source watermark is truthful? Answer: No. The test checks the model’s response to supplied progress, not the external source’s completeness.
41.98 Common wrong ideas#
- Wrong: Real-time means no delay. Right: Continuing systems still have queues, unavailable inputs and processing lag.
- Wrong: Arrival time and event time are interchangeable. Right: They answer different questions.
- Wrong: A watermark makes late records impossible. Right: Its meaning depends on the source’s guarantee or estimate.
- Wrong: Window membership determines publication timing. Right: Triggers and revision policy are separate choices.
- Wrong: Every result pane is an additional amount. Right: Some panes replace earlier results.
- Wrong: The same timestamp proves duplicate identity. Right: Different events can share a timestamp, and replays can arrive later.
- Wrong: Keeping a final sum is enough for recovery. Right: Progress, identities and publication state may also be required.
- Wrong: A framework’s state recovery automatically makes an external effect exactly once. Right: The sink and retry protocol must support the intended guarantee.
41.99 Chapter summary in 20 lines#
- Streams can continue without a natural final record.
- Define record identity and correction meaning before processing.
- Progress is scoped to its source and partition.
- Event time describes the event under a timestamp contract.
- Processing time describes the processor’s handling of a record.
- Source clocks can be wrong or manipulated.
- Windows assign records to bounded groups.
- Half-open boundaries prevent accidental double membership at an edge.
- Sliding windows deliberately overlap.
- Session windows may merge after a bridging arrival.
- Watermarks express event-time progress under stated assumptions.
- Slow or idle inputs need an explicit policy.
- Late does not necessarily mean invalid.
- Triggers decide when results are emitted.
- Consumers must distinguish replacements from deltas.
- Stable identities separate replays from new events.
- A lateness cutoff bounds revisions, not real-world truth.
- Recover totals, identities and progress consistently.
- External sinks have their own acknowledgement and replay limits.
- Preserve the history of published answers and later restatements.