Event sourcing and event streaming are confused constantly because both put events in an append-only log. They answer different questions. Event sourcing is a persistence decision: the authoritative state of an entity is the ordered sequence of events that happened to it, and current state is a fold over that sequence. Event streaming is a transport and processing decision: facts flow through a durable, partitioned log so many consumers can react to them, replay them and derive new data from them.
The distinction matters most at write time. An event-sourced system must reject a write that was decided against stale state, so its store needs a conditional append. A streaming log is designed to accept writes as fast as possible from many producers, and it has no notion of "append only if this entity is still at version 7". This article builds both from first principles, shows exactly where they differ, walks through an order service that uses both, and ends with a checklist for deciding where each belongs in your system.
Two ideas that share a word
In an event-sourced service, every aggregate (an order, an account, a shipment) has its own stream: order-42 holds OrderPlaced, ItemAdded, PaymentCaptured in the order they were accepted. To handle a command the service reads that one stream, folds it into a state object, decides which new events the command produces, and appends them. The stream is the source of truth; tables and caches are disposable projections rebuilt from it.
In an event-streaming platform, a topic is split into partitions, each an ordered log. Producers write records with a key; the key hashes to a partition; consumers in a group split the partitions between them and track their offset. Retention is time or size based, or compacted to the latest value per key. The log is optimised for throughput and fan-out: one write, many independent readers, each at its own position.
Both are logs of immutable facts, so the vocabulary overlaps. The guarantees do not. An event store promises that a stream's history cannot change underneath a decision. A streaming log promises that a record, once acknowledged, will be delivered in partition order to every consumer group that reads it. For grounding in event sourcing itself, see event sourcing architecture; for the table-and-changelog view of a log, see stream-table duality.
The architecture in one picture
Notice the two arrows leaving the event store. Projections that belong to the same service read domain events directly. Other services receive integration events through a streaming log. The relay between them is where most production bugs live, so it gets its own section below.
What an event store must guarantee
Four properties make an event store fit to be a system of record:
- Per-stream ordering with versions. Each event in a stream carries a dense version number: 0, 1, 2. Version gaps are bugs.
- Conditional append. An append names the version the writer last saw. If another writer appended first, the store rejects the write. This is optimistic concurrency, and it is how two concurrent commands on the same order cannot both succeed against the same stale state.
- Cheap single-stream reads. Loading an aggregate means reading one stream from version 0 (or from a snapshot) quickly, without scanning unrelated events.
- An ordered all-stream feed. Projections and relays need to read every event in a stable global order and resume from a checkpoint.
Purpose-built stores provide all four; EventStoreDB, now renamed KurrentDB, is the best-known example, with expected-version checks on append. A relational database provides them with one table and a unique constraint, which is often the pragmatic choice:
CREATE TABLE events (
global_pos BIGSERIAL PRIMARY KEY,
stream_id TEXT NOT NULL,
version INT NOT NULL,
event_id UUID NOT NULL UNIQUE,
type TEXT NOT NULL,
data JSONB NOT NULL,
metadata JSONB NOT NULL,
recorded_at TIMESTAMPTZ NOT NULL DEFAULT now(),
UNIQUE (stream_id, version) -- the concurrency check lives here
);class ConcurrencyConflict(Exception): pass
def append(conn, stream_id, expected_version, new_events):
"""Append new_events only if the stream is still at expected_version."""
try:
with conn.transaction():
for i, e in enumerate(new_events, start=1):
conn.execute(
"INSERT INTO events (stream_id, version, event_id, type, data, metadata)"
" VALUES (%s, %s, %s, %s, %s, %s)",
(stream_id, expected_version + i, e.id, e.type, e.data_json, e.meta_json))
except UniqueViolation:
raise ConcurrencyConflict(stream_id, expected_version)
return expected_version + len(new_events)
def handle(conn, cmd, decide, fold, retries=3):
for _ in range(retries):
history = load_stream(conn, cmd.stream_id) # ordered by version
state, version = fold(history), len(history) - 1 # -1 means "no stream yet"
new_events = decide(state, cmd) # pure function, may raise
try:
return append(conn, cmd.stream_id, version, new_events)
except ConcurrencyConflict:
continue # someone else won; re-decide
raise ConcurrencyConflict(cmd.stream_id, version)The retry re-reads and re-decides. It never blindly re-appends the old events, because the decision might now be different: the item may already be out of stock, or the order cancelled.
What a streaming log guarantees, and what it does not
A partitioned log such as Kafka gives you ordering within a partition, durable replicated storage, consumer groups with independent offsets, and very high write and read throughput. Its idempotent producer removes duplicates created when one producer retries a send after a lost acknowledgement. Its transactions make a set of writes across partitions, plus consumer offsets, visible atomically or not at all.
None of that is a conditional append. Two producers that each read order 42 at version 7 and each decide to append version 8 will both succeed; the partition now holds two events that each believed they were the eighth. Idempotence deduplicates retries of the same send, not conflicting sends from different writers. Transactions guarantee atomic visibility, not that the state you decided against is still current. Reading one entity's history is also expensive: its events are interleaved with every other key that hashes to the same partition, so loading an aggregate means scanning the partition or maintaining a separate index.
There is a legitimate way to event-source on a streaming log: make each aggregate have exactly one writer. Route commands to a topic keyed by aggregate id, so every command for order 42 lands on the same partition, and let exactly one stream-processing task own that partition. The task keeps current state per key in a local state store, validates each command against it, and emits events transactionally together with its state change and input offset. Concurrency control becomes serial processing per partition, and fencing of zombie instances (via the transactional id) replaces the version check. It works, but it is a different design with its own costs: rebalances pause ownership, the state store must be restored on failover, and you still need an index or a compacted snapshot topic to answer "what happened to order 42?" quickly.
Domain events versus integration events
Even when both systems exist, the events in them should usually differ.
| Aspect | Domain event (event store) | Integration event (streaming log) |
|---|---|---|
| Audience | The owning service and its projections | Other teams and services |
| Granularity | Fine: every state transition the aggregate cares about | Coarse: facts others need, often a summary |
| Schema change | Upcast old versions on read; private to the team | Contract with compatibility rules and a schema registry |
| Payload | May contain internal fields and PII | Minimised; references instead of sensitive data |
| Retention | Forever: it is the system of record | Bounded by time, size or compaction |
| Ordering scope | Per stream, enforced on write | Per partition key, best effort across keys |
Publishing raw domain events to the whole company couples every consumer to your internal model. Translate in the relay: map ItemAdded plus ItemRemoved into a single OrderContentsChanged if that is what downstream services actually need.
Bridging the two: the relay
The relay reads the event store's all-stream feed from a checkpoint and produces integration events to the log. Because the append and the publish are two systems, there is no single transaction across them. The safe pattern is the same as the transactional outbox: the event store itself is the outbox, and the relay is at-least-once.
def relay_loop(conn, producer, checkpoint):
while True:
batch = conn.query(
"SELECT global_pos, stream_id, event_id, type, data FROM events"
" WHERE global_pos > %s AND global_pos <= safe_watermark()"
" ORDER BY global_pos LIMIT 500", (checkpoint.load(),))
for row in batch:
out = translate(row) # domain -> integration, or None to drop
if out is not None:
producer.send("orders.v1", key=row.stream_id, value=out,
headers={"event_id": str(row.event_id)})
producer.flush() # wait for acks before advancing
if batch:
checkpoint.save(batch[-1].global_pos)Two details carry the correctness. First, the key is the stream id, so all integration events for one order land in one partition and keep their relative order. Second, the checkpoint advances only after the broker acknowledges the batch. A crash between flush and save replays the batch, so consumers must deduplicate by event_id. If the store is a relational database, logical-replication change data capture can replace the polling loop with the same semantics.
A worked example: two concurrent edits to one order
Order 42 has stream order-42 at version 2: OrderPlaced (v0), ItemAdded sku=A (v1), ItemAdded sku=B (v2). A customer clicks "remove item B" while a warehouse job issues "ship order". Both handlers load the stream at version 2.
The ship handler decides OrderShipped and appends with expected version 2; it becomes v3. The remove handler decides ItemRemoved sku=B and appends with expected version 2; the unique constraint on (order-42, 3) rejects it. Its retry reloads four events, folds them into a shipped order, and decide now raises "cannot modify a shipped order". The customer sees a clear error instead of a shipment that silently disagrees with the order.
Had the same two writes gone straight to a Kafka topic, both would have been accepted. Downstream, billing would see a removal after shipment and have to guess which one was real. The relay later publishes OrderShipped to orders.v1 with key order-42; billing, search and analytics each consume it at their own pace and can replay from retention if they rebuild.
Failure modes
- Global-position gaps skip events. With a sequence-backed
global_pos, a transaction that took position 101 can commit after one that took 102. A relay that read 102 and checkpointed it never sees 101. Read only up to a watermark below the oldest in-flight transaction, or serialise appends to the global feed. - Lost-update via the log. Using a plain topic as the store with several writers per key produces conflicting histories. Either use a store with conditional append or enforce one writer per partition.
- Unbounded streams. An aggregate with millions of events loads slowly. Snapshot every N events and load from the latest snapshot plus the tail, or split the aggregate along its real consistency boundary.
- Schema drift. Old domain events must stay readable forever. Version event types and upcast on read; never rewrite history in place.
- Leaking internal events. Consumers bind to fields you meant to change. Publish translated integration events with a registered schema.
- Erasure requests. An immutable log and a right to erasure conflict. Keep personal data out of events where possible, or encrypt it per subject and delete the key (crypto-shredding). Remember that log retention and compaction on the streaming side do not erase backups.
- Relay lag hidden by healthy brokers. Alert on the age of the oldest unrelayed event, not only on consumer lag in the log.
Choosing, in practice
Use event sourcing where the history is the domain or where you must explain how state was reached: ledgers, orders, entitlements, workflows with audit needs. Do not event-source a CRUD settings table; the conditional append, snapshots and upcasting are real ongoing cost. Use event streaming wherever several systems need the same facts, need to replay them, or need to compute over them continuously, whether or not the producer is event-sourced. A service that stores state in ordinary tables can publish to a log through an outbox or CDC and be a perfectly good streaming citizen.
The common healthy shape is the one in the diagram: event-sourced (or table-backed) services inside, a streaming log between them, and a relay with explicit translation, checkpoints and deduplication at the boundary. For long-retained integration topics, log compaction keeps the latest fact per key without keeping every intermediate event.
What to do next
- For each service, write down whether its state is event-sourced, table-backed, or derived; most should not be event-sourced.
- For any event-sourced aggregate, confirm the store rejects an append at a stale expected version, and write a test that races two commands.
- If you are using a topic as an event store, prove there is exactly one writer per key, or move the write path to a store with conditional append.
- Separate domain events from integration events and register the integration schemas with compatibility rules.
- Build the relay with a checkpoint that advances only after broker acknowledgement, keyed by stream id, with an event id header for deduplication.
- Add alerts on oldest-unrelayed-event age and on concurrency-conflict rate per aggregate type.
- Decide your erasure strategy before the first personal field is written into an immutable stream.