Lesson 3 of 4 · 25 min

Design a delayed-event alert pipeline

Define a streaming alert with event-time and false-positive controls.

Mechanism and reasoning

An alert pipeline changes the consequence of a data error. A duplicate dashboard row may confuse an analyst; a duplicate alert can wake an operator or contact a user twice. Begin with the event contract and the action's idempotency key. Separate detection from delivery so retries do not create repeated business effects.
Define what the alert means. A device that has not sent data is not necessarily stationary or unhealthy. Missing observations can reflect network loss. Use explicit states such as moving, stopped, unknown and alert-pending, with transition rules. The absence of a message needs its own interpretation.
Event time matters when telemetry arrives late. A delayed movement event may show that an apparent stop never lasted long enough. Decide whether the pipeline waits for a lateness allowance, issues a provisional alert, or sends a correction. The correct choice depends on urgency and the cost of false positives. No single watermark setting solves that business trade-off.
Use bounded state per entity and an expiry policy. Keep only the history needed to establish the condition and deduplicate delivery. A long offline period should not accumulate unlimited state. Track source completeness and clock plausibility; a device with a wildly wrong timestamp can distort windows.
In an interview, draw a timeline with at least one late event and a repeated event. Explain exactly when an alert becomes eligible and how a recovery or correction changes the state. Do not turn a vague phrase such as 'real time' into an unexamined promise. State the delay target and the uncertainty the pipeline accepts.

Turn the alert rule into a state transition table

For this exercise, a stopped observation starts an episode only if telemetry is fresh enough under a separate freshness rule. A moving observation ends that episode. Missing data can move the entity to unknown, depending on the defined freshness timeout. The state machine must not infer a physical stop solely from a lack of messages.
Current stateAccepted event or timerNext stateAction
MovingExplicit stopped observationStopped-pendingCreate episode identity
Stopped-pendingDuplicate same observationStopped-pendingNo new episode
Stopped-pendingMovement within intervalMovingCancel pending eligibility
Stopped-pendingDuration met and lateness policy permitsAlert-eligibleCreate one detection intent
Any observed stateFreshness timeoutUnknownApply separate telemetry-loss rule
AlertedAccepted late contradictionPolicy-defined corrected stateRetract or annotate under contract
The final row is a business decision. Some urgent alerts must be sent quickly and later corrected. Others can wait for a lateness allowance to reduce false positives. State the maximum added delay and what evidence can still change the result. “Exactly ten minutes” is not a complete delivery promise when event-time progress and network delay also matter.
code
1Illustrative event-time sequence209:00 stopped, episode E1309:06 moving, arrives at processing time 09:11409:08 stopped, episode E25Eligibility for E1: cancelled by accepted movement6Eligibility for E2: no earlier than event time 09:187Delivery: after the configured progress/lateness condition is met
The sequence separates event time from processing time. A timer firing at processing time 09:10 cannot establish that all relevant earlier movement observations have arrived. Conversely, waiting forever for complete certainty prevents useful alerts. The contract chooses how much uncertainty is acceptable and how later contradictions are handled.
A watermark is one way to express stream progress, but its exact relationship to timers and allowed lateness depends on the runner and windowing design. The phrase “two minutes of lateness allowance” should not be implemented as a universal wall-clock sleep without checking those semantics. For the workbook, treat the allowance as a stated acceptance/correction policy and explain delivery timing separately. This prevents a framework-specific default from becoming an unsupported business guarantee.
Detection identity and delivery identity should be stable. Machine ID plus episode ID plus alert-rule revision can identify the intended action. A delivery worker can retry that intent under a unique key without creating a new alert each time. If the rule changes during an episode, decide whether the existing episode is evaluated under its original revision or re-evaluated explicitly. Do not accidentally send two notifications simply because a deployment changed the rule label.
Bound state by the supported horizon. The processor may need recent observations, the current episode, pending timers and delivered-alert identity. Retaining all telemetry forever is unnecessary for this local state machine, though an authorized raw archive may serve separate analytics. Expiring delivery identities too early can let a replay resend old alerts. Align the detection replay horizon and deduplication retention.
Clock plausibility matters. A device timestamp far in the future can distort apparent duration or progress. Validate timestamp ranges and source identity before using events to advance state. A correction from an authoritative source may require replacing an earlier observation. Record why an event was rejected so that data loss is not confused with a quiet device.
The interview exercise is related to a first-person report of a stopped-cab safety-alert design, but delayed telemetry, state names, episode IDs and numerical fixtures here are our original extensions. The report does not confirm this exact formulation or answer. Keep that evidence boundary separate from the technical design.
Verification should include duplicated stop events, movement arriving late within policy, a contradiction beyond the automatic horizon, a telemetry outage and delivery retry after an uncertain acknowledgement. Inspect both the final episode state and the number of business notifications. Correct event processing can still produce duplicate alerts if the delivery path lacks its own idempotency contract.

Worked example

Teaching condition: alert when a machine is explicitly stopped for ten event-time minutes, with two minutes of lateness allowance. At 09:00 it reports stopped. At processing time 09:11, a movement event with event time 09:06 arrives. The stated runner watermark and acceptance policy still accept this event; a five-minute event-to-arrival delay alone does not determine whether it is beyond allowed lateness. That event interrupts the apparent stop, so the ten-minute condition was never established. A timer based only on processing time could have issued a false alert at 09:10. A stable alert key could be machine ID plus stop-episode ID, preventing duplicate delivery if the same detection is retried. Qualification also requires observation freshness throughout the interval under the configured freshness timeout. Silence that crosses that timeout moves the device to unknown and prevents a claim of a continuously observed stop.

Exercise

A device reports stopped at 14:00, moving at 14:04 and stopped again at 14:05. The condition is ten uninterrupted minutes stopped. Assume all events are accepted in order and continued telemetry keeps the stopped state fresh until the eligibility time. What is the earliest event-time eligibility and which episode should the alert key use?

Model solution and rubric

The first stop lasts four minutes and does not qualify. The second episode can qualify at 14:15 if no accepted movement event interrupts it. Use the second stop episode's identity in the alert key. Delivery time may be later because of lateness allowance and processing. Verify duplicate telemetry does not reset or duplicate the episode and that missing telemetry enters the defined unknown state.
Score out of four: one point for the correct result, one for showing the intermediate reasoning, one for identifying the stated failure case, and one for a verification that could disprove the answer. Do not award the reasoning point for a tool name alone.

Failure modes and misconceptions

“A timer expiry proves the physical condition lasted.” It proves time passed in one clock domain. Accepted movement and freshness evidence can change the state.
“Deduplicating telemetry prevents duplicate notifications.” Detection and delivery are separate effects. Stable alert identity must survive delivery retries and replay.

Interview probe

Evidence class: recommended. Original practice.
How do you prevent delayed telemetry from generating a false incident?
Strong answer: Define event-time states and a lateness policy, distinguish missing data from explicit state, and make alert delivery idempotent. I would test a late movement event that invalidates an apparent stop and specify whether already-sent alerts need correction.
Follow-up: How does the design change when immediate action matters more than false-positive rate?
Weak answer indicators: Treating silence as proof of a stop; using a processing timer without lateness; deduplicating only messages but not alerts.

Sources

Technical references: Apache Beam programming guide; Kafka 4.1 design; Delhivery candidate report. Sources support the documented mechanisms. The numbers, decisions, rubrics and interview prompts in this lesson are original teaching examples, not measurements or employer question claims.
docsApache Beam programming guidebeam.apache.orgdocsKafka 4.1 designkafka.apache.orgdocsDelhivery candidate reportreddit.com

Checkpoint

No telemetry arrives for ten minutes. Without an explicit physical-state rule, what is known?

AThe machine is stoppedBThe machine is movingCA stop alert was deliveredDObservation is missing; physical state is not established
Sign up free to answer and see why

Checkpoint

A late accepted movement event falls inside a pending stop. What must happen?

ARe-evaluate the episode and cancel or correct under policyBIgnore it because a timer existsCCreate an unrelated duplicate stopDTreat movement as the same event as the original stop
Sign up free to answer and see why

Checkpoint

Stopped at 14:00, moving 14:04, stopped 14:05; ten-minute rule. Earliest event-time eligibility?

A14:10B14:15C14:14D14:05
Sign up free to answer and see why

Checkpoint

Detection is correct, but delivery retries after a lost acknowledgement. What prevents duplicate business alerts?

AOnly telemetry deduplicationBOnly a longer detection windowCStable alert identity and a duplicate-safe delivery contractDOnly a new random ID per send
Sign up free to answer and see why

Checkpoint

Urgency requires alerts before the lateness horizon closes. What must be explicit?

AThat all late events are impossibleBThat watermark equals wall clockCThat state need never expireDThe accepted false-positive risk and later correction behavior
Sign up free to answer and see why

Explain how you would define a streaming alert with event-time and false-positive controls without reading the solution. State one assumption that could change your answer, and one observation that would make you revise it.

Not yetGetting thereConfident

Wrap-up

  • Specify the state machine and correction policy before selecting a streaming tool. Detection and delivery need separate guarantees. The linked first-person Delhivery report names DE-1 and a stopped-cab safety-alert design. Employer confirmation and exact publication/interview dates are unverified. Late-telemetry rules, machine examples, episode identities and numerical cases here are original extensions, not a reported exact prompt.

Sources

Free to read · better with Enzo

Learn it with Enzo

Save your progress, answer the checkpoints, and let Enzo quiz you on what you just read.