Engineering5 min read

Real-Time Data Pipelines: Freshness, Recovery, and Reporting

Oleksandr Melnychenko··
On this page

“Real time” can mean seconds to one team and minutes to another. Before choosing a streaming platform, decide which action depends on the data and how late an update can be before it becomes misleading. A shipment status, stock reservation, machine alarm and monthly report do not need the same delivery path.

Define the decision and its time limit

A dispatcher may need to know which deliveries are likely to miss a time window. The useful answer depends on order data, carrier updates, event timestamps and a rule for what counts as late. The system must show when a source last updated; an old location should not look current.

Write down the source of each event, its owner, expected frequency, format and identifier. Include what happens if a device is offline, an API rate-limits requests, or a partner sends the same event twice.

Separate ingestion from business meaning

An ingestion service receives data. The processing step validates it, links it to the correct business record and decides whether it changes the current state. Storage and reporting then make that state available to users. These may be separate services or a simpler application, depending on volume, freshness and recovery needs.

A broker or durable queue can help when producers and consumers run at different speeds. It is not an automatic requirement for every integration. Choose it after estimating peak volume, replay needs and the cost of operating it.

Expect late and repeated events

Networks, devices and third-party systems fail. Define how to handle duplicate messages, events that arrive out of order and corrections to previously accepted data. Keep the event identifier distinct from the business-record identifier: one shipment can have many legitimate updates.

An idempotent operation produces no additional business effect when the same operation is retried. Recognising a duplicate requires a stable operation key; preventing a second write also requires that the duplicate check and write cannot race with another worker.

Do not treat a queue’s delivery setting as a promise that the whole business workflow will process each event exactly once. Database updates, external API calls and notifications may need their own deduplication or reconciliation rules. Confluent’s Kafka delivery semantics documentation explains why retries can produce repeated delivery and where stronger guarantees apply.

Show freshness and quality to the user

Distinguish event time, when the source says something happened, from processing time, when a pipeline stage handles it. A late delivery confirmation may belong to yesterday's operational totals even though it was received today. Apache Beam's discussion of event time and late data explains this distinction.

An operational view should show the source observation time as well as the latest successful refresh. Refreshing a page does not make an old location current. A quiet source is not necessarily a failed source: define an expected update interval or a separate health signal. Agree what users see when data is delayed, missing or conflicting.

For reporting, define whether a measure uses event time or processing time, what happens when a late event changes a past period, and how totals are checked against the source. If an existing BI environment is kept, specify the exchange frequency and a reconciliation check.

Set targets you can test

Measure end-to-end delay from the source event timestamp to the time the update becomes available in the operational view. For an illustrative event observed at 10:00:00, received at 10:00:02 and visible at 10:00:15, that delay is 15 seconds; 13 seconds occur after receipt. This calculation assumes comparable clocks. Clock skew or an unreliable source timestamp must be recorded as a measurement limitation.

Report the median and a tail percentile, such as p95, over a stated interval and workload. A p95 of 15 seconds means approximately 95% of measured updates completed within that time; it says nothing about events that never arrived unless those failures are counted separately. Google's SRE monitoring guidance explains why an average alone can hide slow requests.

Check completeness independently: reconcile unique expected source events with accepted, rejected and still-pending events for the same cutoff. If the source cannot provide an expected count or sequence, describe the coverage you can verify instead of reporting an unsupported completeness percentage. Test normal load, a burst, a source outage and recovery; retain the event IDs and timestamps needed to reproduce a discrepancy.

Plan monitoring for the pipeline itself: source failures, rising consumer lag, invalid records, repeated retries and reports that no longer reconcile. AWS’s Kinesis monitoring guidance illustrates the kinds of stream and consumer signals that need watching; the choice of platform remains project-specific.

Design the first release around one flow

Connect one carrier feed to shipment records, show dispatchers delayed or missing updates and reconcile daily statuses. Define the interfaces, expected event volume, retention, alert rules, recovery procedure and the team responsible for a failed feed.

If your team is moving data between old and new systems, also separate migration of historical records from synchronisation of ongoing changes. Both need ownership and validation. Our modernisation guide explains that transition, and the ERP reporting example shows how to define a report from source records.

To discuss a pipeline, bring one operational decision, the data sources behind it and the time window in which the answer remains useful. We can scope the first connected workflow from there.