A stream-scoring service must identify the event, the state it read and the model it used before a score is defensible.
Stream inference: event time, feature state and decision identity
Separate three clocks
A receipt event has an occurrence time, an arrival time and a scoring time. Those clocks can disagree after mobile disconnection, broker retry or upstream delay. A merchant-spend feature should include only events eligible at the receipt occurrence time if training used that definition. If the product intentionally scores from the latest available state instead, document that different contract and train accordingly. Training-serving parity fails quietly when one side uses event time and the other uses processing time.
Pin state with the decision
Record an event ID, source partition and offset, feature-state revision or window cutoff, model digest, preprocessing digest and policy revision with each output. A watermark indicates progress, not proof that no older event can arrive; late data needs an explicit policy. Do not overwrite an earlier decision merely because a later event corrected the aggregate. The inference chain names the runtime stages, while the state revision says which history those stages actually saw.
Define late-event behavior
Choose a bounded lateness period and one of three actions for events beyond it: ignore for live scoring, emit a revision with a new identity, or send to reconciliation. Decide whether late input updates future feature state even when the old decision is immutable. That choice changes later scores. Keep a dead-letter route for malformed timestamps and impossible clock skew. Late feature repair is useful only when its effect on decisions is stated.
Test the stream in event order and arrival order
Replay the same events in several arrival orders while preserving event timestamps, then compare final state and decisions under the declared lateness policy. Include duplicates, a delayed partition and a state restore. If the output differs, verify whether that difference is permitted and visible in the decision ledger. Replay controls make recovery testable; the project applies them to a burst of offline receipt uploads.
Implementation
def score_stream_event(event, state_cutoff, model_digest, allowed_lag_seconds):
if not event.get("event_id") or not model_digest:
raise ValueError("decision identity required")
if allowed_lag_seconds < 0:
raise ValueError("negative lag")
lag = state_cutoff - event["occurred_at"]
if lag > allowed_lag_seconds:
return {"route": "reconcile", "event_id": event["event_id"]}
return {"route": "score", "event_id": event["event_id"],
"state_cutoff": state_cutoff, "model_digest": model_digest}
receipt = {"event_id": "receipt-47", "occurred_at": 1_000}
assert score_stream_event(receipt, 1_020, "risk-r8", 82)["route"] == "score"
assert score_stream_event(receipt, 1_200, "risk-r8", 82)["route"] == "reconcile"
Performance and operating cost
The gate is O(1) time and output space per event. Maintaining keyed windows costs O(k) state for active keys and retained history; late-event tolerance enlarges that state and delays completeness. A tighter bound reduces memory but sends more valid delayed events to reconciliation. Choose it from observed delay and the cost of a late correction.
Common Mistakes
- Using broker arrival order as if it were business event order.
- Treating a watermark as an absolute guarantee against late data.
- Logging a model digest without the feature-state cutoff.
- Silently replacing a prior decision when old data arrives.
Read next
- Stream scoring replay: checkpoints, duplicates and output reconciliation
- Project: recover a receipt stream scorer after offline uploads
- Training-serving parity: compare feature values at one prediction clock
- Multi-stage inference: pin each stage and its contract
- Feature materialization lag: detect stuck updates and repair safely
