Data Streaming and Event-Driven Systems
Data streaming is the continuous movement and processing of events as they occur or soon after they occur. Event-driven systems organise software and data around meaningful changes in state—an order was placed, a payment was authorised, a sensor crossed a threshold, a student submitted an assessment, a document was updated.
An event is useful when it records a meaningful change with enough identity, time and context that future consumers can interpret it correctly.
Streaming is often described as “real time”, but speed is only one part of the problem. Reliable streaming also requires ordering, duplication control, event-time reasoning, replay, schema evolution, state management, backpressure, security and clear contracts between producers and consumers.
ARTICLE ID: DATA.MANAGEMENT.030
Canonical function: continuous event movement, replay and stateful processing
Series route: Data Engineering and Pipelines → Data Streaming and Event-Driven Systems.
The Simple Answer
A stream is an ordered or partially ordered sequence of events. A common route is:
Producer → Event → Topic or Stream → Partition → Consumer → State → Output → Replay
The architecture works when consumers can process events consistently even when events arrive late, are duplicated, are retried or must be replayed after logic changes.
What Is an Event?
An event is a record that something happened. Good events usually include:
- a stable event identifier;
- the entity or aggregate affected;
- event type;
- event time;
- payload;
- schema version;
- producer identity;
- relevant correlation or trace identifiers.
An event should describe a fact about state change rather than an instruction disguised as history.
Events vs Commands
A command says, “Do this.” An event says, “This happened.”
Commands may fail or be rejected. Events should represent accepted facts within the producer’s authority. Mixing the two creates ambiguity about whether a message describes intent or completed state.
Topics and Streams
Topics or streams provide named channels through which producers publish events and consumers subscribe to them.
A good topic boundary aligns with coherent meaning. One giant topic for every event creates operational and semantic confusion. Too many tiny topics create management overhead.
Partitions
Partitions divide a stream so processing can scale horizontally. Events with the same key are often routed to the same partition when relative ordering matters.
Partition keys should follow the ordering boundary. If order matters per customer, customer identity may be an appropriate key. If order matters per account, account identity may be better.
Ordering Is Usually Local, Not Global
Distributed systems can preserve order within a partition more easily than across the entire world. Demanding one total global order can reduce scalability and still fail to represent meaningful causality.
The stronger question is: which events must be ordered relative to which other events?
Event Time vs Processing Time
Event time is when the real-world or source event occurred. Processing time is when the streaming system handled it.
These can differ because networks delay, devices go offline, queues build up or systems retry. Analytical logic should use event time when the question concerns what happened in the world, and processing time when the question concerns system behaviour.
Late Events
Late events arrive after a consumer has already processed later event-time windows. A mature stream processor defines how long it will wait, whether closed windows can be updated, and how corrections are surfaced.
Late does not mean invalid. It means the system must reconcile world time with arrival time.
Watermarks
A watermark is a system estimate of how far event-time processing has progressed. It helps decide when a time window can be considered sufficiently complete for computation.
Watermarks are policies about lateness, not guarantees that no older event will ever appear.
Windows
Streaming systems often aggregate events over windows:
- fixed windows;
- sliding windows;
- session windows;
- custom business windows.
The window should match the business question. A technical one-minute window is not automatically meaningful to the receiver.
At-Most-Once, At-Least-Once and Effectively-Once Processing
Distributed messaging systems make trade-offs around delivery:
- at-most-once: an event may be lost but is not redelivered;
- at-least-once: events are retried, so duplicates are possible;
- effectively-once: the overall processing design makes repeated delivery produce one logical outcome.
Marketing language around “exactly once” should always be interpreted in the context of the full end-to-end system. Broker guarantees alone do not automatically prevent duplicated business effects in downstream systems.
Idempotency
An idempotent consumer can safely process the same logical event more than once without duplicating the intended outcome.
Stable event IDs, deduplication stores, deterministic upserts and version checks are common techniques.
Replay
Replay allows consumers to process retained events again. This is one of the strongest properties of event logs.
Replay supports:
- recovery after consumer failure;
- building a new downstream system;
- recomputing derived state after a bug fix;
- auditing event history;
- reconstructing analytical views.
Replay only works safely when consumers are designed for repeatability and historical schemas remain interpretable.
Retention
Streams often retain events for a defined period or volume. Longer retention increases replay capability and historical value but also increases storage, security and privacy responsibilities.
Retention should follow purpose rather than an assumption that every event log should be permanent.
Stateful Processing
Some stream operations need memory of prior events: running totals, sessions, deduplication sets, joins or fraud patterns.
This is stateful processing. State must be checkpointed, recoverable and aligned with event progress so failures do not silently create inconsistent results.
Checkpoints
A checkpoint records enough processing state that a consumer can resume after failure. Good checkpoint design coordinates input progress with stored state.
Checkpointing is part of correctness, not merely performance.
Backpressure
Backpressure occurs when consumers cannot keep up with producers. Lag grows, memory can fill and latency increases.
Systems need a strategy: scale consumers, slow producers, buffer safely, degrade non-critical work or shed load according to explicit priorities.
Consumer Lag
Consumer lag measures how far a consumer is behind the stream. It is a useful operational signal, but the business meaning depends on event rate and criticality.
Ten thousand events of lag can be trivial in one system and catastrophic in another.
Change Data Capture
Change Data Capture (CDC) turns database inserts, updates and deletes into a stream of changes. It is useful for synchronisation, analytics and migration.
CDC reflects changes in a database representation. It is not always the same as a domain event explaining what happened in business terms.
CDC vs Domain Events
A database row changing from PENDING to PAID is a technical change record. A domain event such as PaymentSettled expresses a business fact.
Both are useful, but they serve different consumers and should not be confused.
Transactional Outbox Pattern
One common design challenge is ensuring that a database update and its corresponding event are not separated by failure. An outbox pattern records the event intention in the same transactional boundary as the operational change, then publishes it asynchronously.
The principle is broader than any implementation: preserve one authoritative handoff between operational state and emitted event.
Schema Evolution
Event schemas change while old events remain in retained logs. Consumers may replay events produced under older schemas years later.
Safe streaming systems use explicit schema versions, compatibility rules and deprecation policies.
See Data Contracts and Data Products.
Schema Registries
A schema registry can store event schemas and validate compatibility between versions. It supports technical contract enforcement but does not decide whether the business meaning of an event changed.
Consumer Groups
Consumer groups allow several workers to share processing of one logical subscription. This provides scalability while ensuring each partition is processed by one worker in the group at a time.
Different consumer groups can independently read the same events for different purposes.
Fan-Out
One event can support many downstream consumers: analytics, notifications, fraud detection, inventory, auditing and machine learning.
This decoupling is powerful, but it increases governance responsibility because one producer can create many hidden dependencies.
Event Contracts
Event contracts should define:
- event meaning;
- key and ordering boundary;
- schema;
- required and optional fields;
- units;
- event-time semantics;
- sensitivity;
- retention;
- compatibility expectations;
- producer ownership.
Without a contract, an event stream can become an undocumented API with permanent downstream consequences.
Dead-Letter Paths
Events that cannot be processed safely may be routed into a dead-letter or quarantine path.
The failed event should preserve source identity, error reason, schema version and repair status. Dead-letter queues should not become permanent graveyards that nobody reviews.
Observability
Useful streaming signals include:
- producer rate;
- consumer lag;
- processing latency;
- error rate;
- late-event rate;
- duplicate rate;
- dead-letter volume;
- schema failures;
- checkpoint health;
- partition imbalance.
See Data Observability and Monitoring.
Streaming Quality
Streaming quality needs both record-level and temporal checks. A record can be structurally valid while arriving far too late for its intended use.
- schema validity;
- key completeness;
- event-time plausibility;
- reference-data validity;
- duplicate detection;
- ordering checks where meaningful;
- arrival-time service levels;
- business invariants.
Security
Streams can carry sensitive operational data at high velocity. Access should be controlled at producer, topic and consumer level, with encryption and strong identity where appropriate.
Logs and debugging tools should not expose full payloads unnecessarily.
Privacy
Event logs can retain detailed behavioural histories. A stream designed for operational recovery can become a long-term surveillance dataset if retention and reuse are not governed.
See Data Security and Privacy.
Streaming and Data Lakes
Streams often feed lakehouses for long-term analytical storage. The streaming layer captures continuous change; the lakehouse preserves historical state in queryable forms.
See Data Lakes and Lakehouses.
Streaming and Operational Databases
Operational databases remain useful for authoritative transactional state. Streams complement them by broadcasting change and enabling decoupled processing.
See Database Management and Transactional Integrity.
Event Sourcing
In event-sourced designs, the event history itself is the primary record from which current state can be reconstructed.
This can provide strong auditability and temporal history, but it also increases the burden of schema evolution, replay, event correction and long-term interpretation. Event sourcing should be chosen because the domain benefits from historical event truth, not because event logs are fashionable.
Snapshots
Long event histories can make full reconstruction expensive. Snapshots capture state at a known point, allowing later replay to begin from that state rather than from the first event ever recorded.
Snapshots should remain linked to the event position and schema version from which they were produced.
Education Example
An education platform may emit events such as StudentEnrolled, LessonAttended and AssessmentSubmitted. Operational systems consume them for immediate workflows; analytics consumers build near-real-time attendance views; the lakehouse retains history.
If a mobile device uploads yesterday’s attendance late, event-time logic ensures the record belongs to yesterday’s learning history rather than today’s merely because it arrived today.
Commerce Example
An e-commerce system may emit order, payment, shipment and return events. Consumers include inventory, fraud detection, customer communications and analytics.
Ordering by order identity protects local sequence without requiring every order in the world to share one global clock.
IoT Example
Sensors may produce readings continuously but intermittently lose connectivity. When connectivity returns, older events can arrive in bursts. Event-time windows and late-data policies prevent those readings from being misclassified simply because the network was delayed.
AI Example
Streaming can update AI retrieval indexes when documents change. A DocumentUpdated event can trigger re-chunking and re-indexing, while a DocumentRevoked event removes access.
The AI system should preserve source identity and process events idempotently so retries do not create duplicate index entries.
Common Failure Modes
- Real-time theatre: streaming is chosen when daily batch would serve the receiver.
- No event identity: duplicates cannot be detected reliably.
- Global-order fantasy: the system pays huge cost for ordering that the domain does not need.
- Processing time as world time: delayed events rewrite the wrong period.
- Replay without idempotency: recovery duplicates business effects.
- CDC equals domain semantics: technical row changes are mistaken for business events.
- Schema registry equals contract: business meaning is left undocumented.
- Dead-letter graveyard: failures accumulate without repair.
- Infinite retention: event history becomes unmanaged privacy and cost exposure.
- Lag without context: technical thresholds ignore receiver consequence.
A Streaming Design Checklist
- Why does the receiver need streaming rather than batch?
- What is the event’s business meaning?
- What stable event ID exists?
- Which entity defines the ordering boundary?
- What is event time?
- How are late events handled?
- Which delivery semantics apply?
- Are consumers idempotent?
- How long are events retained?
- Can events be replayed safely?
- How are schemas versioned?
- What quality and contract checks run?
- How is state checkpointed?
- What happens under backpressure?
- Which observability signals reveal receiver impact?
A Maturity Ladder
- Transported: events move from producers to consumers.
- Identified: events have stable IDs, keys and timestamps.
- Contracted: schemas and semantics are governed.
- Replayable: retained history can reconstruct downstream state.
- Stateful: checkpoints and event-time logic are reliable.
- Observable: lag, lateness, errors and dead-letter paths are visible.
- Governed: retention, privacy and consumer dependencies are controlled.
- Adaptive: incidents and receiver outcomes improve event and processing design.
The Deeper Principle: Streaming Is a History of Change, Not Merely Faster Data
The real power of streaming is not that information moves quickly. It is that systems can preserve a durable sequence of meaningful changes and allow many consumers to react independently while retaining a route back through time.
That power only becomes trustworthy when identity, event time, ordering, replay and contracts are designed deliberately.
Data Management Series
- Data Streaming and Event-Driven Systems
- Data Engineering and Pipelines
- Data Contracts and Data Products
- Data Lakes and Lakehouses
- Data Testing and Reliability Engineering
Final idea: a stream is a living record of change. Reliable event-driven architecture preserves enough order, identity, time and replay capability that many downstream systems can act quickly without losing the evidence needed to explain what happened.
