0%
Data-Intensive Applications
Foundations of Data Systems
Distributed Data
Encoding and Evolution
Batch Processing
Data Quality and Governance
Operational Patterns
Event Streaming Fundamentals
Every database you have used stores the current state of the world. A user row says "balance = 150." But how did it get to 150? Was it a single deposit, or five deposits and three withdrawals? The row does not tell you. An event log takes a fundamentally different approach: it records every change as an immutable fact, in the order it happened.
An event is a record of something that occurred. "User 42 deposited $50 at 14:03:07." "Order 789 was cancelled at 14:03:09." Each event is appended to the end of a log and never modified or deleted. The log is append-only and ordered. This is not a new idea. Bank ledgers, ship logs, and accounting journals have worked this way for centuries. Databases took a shortcut by storing only the latest state, but event logs preserve the full history.
Event streaming takes this idea and makes it the backbone of distributed systems. Instead of services communicating through synchronous API calls, they communicate through events written to a shared log. A payment service does not call the notification service directly. It writes a "PaymentCompleted" event to the log. The notification service reads that event and sends an email. The analytics service reads the same event and updates a dashboard. The services are decoupled: they share data through the log but know nothing about each other.
This decoupling has a profound operational benefit. If the notification service goes down for 2 hours, no events are lost. They accumulate in the log. When the service recovers, it reads from where it left off and processes the backlog. In a synchronous architecture, those 2 hours of notifications would be permanently lost unless the payment service implemented retry logic, dead letter queues, and failure tracking, all complexity that the event log handles implicitly.
Event streaming also enables temporal decoupling: the producer and consumer do not need to be running at the same time. A batch analytics job can run every night at 2am, consuming a full day's worth of events in one pass. The producer ran throughout the day. The consumer runs once. Neither needs to coordinate timing with the other.
This is a fundamental shift from request-response architectures where both sides must be online simultaneously for communication to succeed.
Why Immutability Matters
Mutable state creates three problems that immutability solves.
Debugging: When a user's balance is wrong, you need to know what happened. With mutable state, you see "balance = 150" and nothing else. With an event log, you replay the sequence: deposit 100, deposit 100, withdraw 50. You find the missing event or the duplicate.
Auditing: Regulations in finance, healthcare, and government require a complete history of changes. An event log is a built-in audit trail. No separate audit table needed, no triggers, no application-level logging that someone forgets to add. The audit trail is not a secondary artifact bolted onto the system; it is the primary data structure from which everything else is derived.
Reprocessing: When you deploy a new version of your analytics pipeline, you can replay the entire event log through the new code and recompute results from scratch. With mutable state, you cannot go back. The old values are gone.
Temporal queries: With an event log, you can answer questions like "what was user 42's address on March 15?" by replaying events up to that date. This is impossible with mutable state that only reflects the current moment. Temporal queries are essential in domains like insurance (what policy was active when the claim was filed?), logistics (what warehouse held this item when the order was placed?), and compliance (what permissions did this employee have when they accessed this record?).
Append-Only Logs in Practice
An append-only log is a file where new records are written at the end and existing records are never overwritten. Each record gets a monotonically increasing sequence number called an offset. Offset 0 is the first event ever written. Offset 1,000,000 is the millionth. A consumer reads from any offset and moves forward. Two consumers can read the same log at different speeds, at different positions, without interfering with each other.
This is fundamentally different from a message queue where messages are deleted after consumption. In a log, messages persist. Consumer A can read message 500 today. Consumer B can read message 500 next week. The message is still there.
The offset is also a timestamp proxy. Because offsets are sequential and events arrive in time order, offset 500,000 was written before offset 1,000,000. You can convert between offsets and timestamps using Kafka's offsetsForTimes API, which lets a consumer say "start reading from events written after 2pm yesterday" without knowing the exact offset.
In interviews, emphasize that event logs decouple producers from consumers in both time and processing speed. A producer writes events without knowing or caring how many consumers exist, when they will read, or how fast they process. This temporal decoupling is what makes event-driven architectures resilient to downstream failures.
Events vs. Commands
A common confusion is between events and commands. An event says "OrderPlaced" -- it is a fact that already happened. A command says "PlaceOrder" -- it is a request that might be rejected. Events are past tense and immutable. Commands are imperative and may fail. This distinction matters because event logs store only events. The command validation happens upstream, and only successful outcomes become events in the log.
The naming convention reinforces this: events use past participle (OrderPlaced, PaymentProcessed, UserRegistered) while commands use imperative (PlaceOrder, ProcessPayment, RegisterUser). When you see "OrderPlaced" in a log, you know it happened. There is no ambiguity, no need to check whether it succeeded. This certainty is what makes event-driven architectures reliable: every consumer processes facts, not intentions.
Some systems blur this line with "fat events" that carry enough data for consumers to act without querying other services, versus "thin events" that carry only identifiers and require consumers to fetch details. Fat events increase event size and coupling to the producer's data model, but reduce the number of downstream API calls. Thin events are smaller and more stable but create runtime dependencies between consumers and source services. Most production systems land somewhere in between: include the data that most consumers need, omit rarely-used fields that consumers can fetch on demand.
Deriving State from Events
If your event log contains every state change, you can derive the current state at any point in time by replaying events from the beginning up to that moment. This is the foundation of event sourcing: the event log is the source of truth, and every "view" of the data (a database table, a search index, a cache) is a derived projection.
For example, replaying all OrderPlaced, OrderShipped, and OrderDelivered events produces a table of order statuses. Replaying only OrderPlaced events with their item quantities produces an aggregate revenue report. The same event stream feeds different views by applying different projection logic. This is why event logs pair naturally with CQRS (Command Query Responsibility Segregation): writes append events to the log, and reads query materialized views derived from those events.
The power of replay extends to disaster recovery. If a database becomes corrupted, you can rebuild it entirely from the event log. If a new microservice needs historical data, it consumes the log from offset 0 and bootstraps itself without any migration scripts or data dumps.
Event Schema Design
An event should be self-contained: it carries enough context for any consumer to process it without querying another service. A poor event says "OrderUpdated, orderId=789." A good event says "OrderShipped, orderId=789, customerId=42, trackingNumber=1Z999AA10123456784, carrier=UPS, shippedAt=2024-03-15T14:03:07Z." The consumer building the shipment notification has everything it needs in the event itself.
Events also need a schema that evolves without breaking consumers. Add new fields freely (old consumers ignore them). Never remove or rename fields (old consumers depend on them). Never change a field's type (an integer field becoming a string breaks deserialization). These rules are the same as Protobuf's backward compatibility rules. Tools like Apache Avro with a schema registry enforce these rules automatically: the registry rejects schema changes that would break existing consumers.
The Tradeoff: Storage
Immutability means you keep everything, and keeping everything costs storage. A busy e-commerce system generating 10,000 events per second produces 864 million events per day. At 500 bytes per event, that is 400 GB daily. This is why event streaming systems offer retention policies: keep events for 7 days, 30 days, or forever. Log compaction (covered later) offers a middle ground by keeping only the latest event per key while discarding older versions.
The storage cost is not just disk space. More data means longer replay times. Rebuilding a projection from 10 billion events takes hours, not minutes. Snapshots help: periodically save the current derived state, and on recovery, replay only events after the snapshot. This balances the benefits of immutability with the practical need for fast recovery.
There is also a cost to schema management. Every event in the log must remain parseable forever (or at least for the retention period). If you change the event format, old consumers must still be able to read old events. This is why event schema evolution follows strict rules: add fields, never remove them. Include a version field in each event so consumers can branch on the format. Use a schema registry to enforce compatibility at write time rather than discovering incompatibility at read time, when the damage is already done.
Finally, there is a cognitive cost. Teams accustomed to CRUD databases must shift their thinking from "what is the current state" to "what sequence of events produced the current state." Debugging requires reading event sequences instead of inspecting rows. Querying requires materialized views instead of direct table access. Testing requires event replay instead of fixture loading. These are not insurmountable challenges, but they represent a real learning curve that organizations should account for when adopting event-driven architectures.