Every database is secretly a stream processor. Its write-ahead log is a sequence of changes, and its tables are what you get by applying those changes in order. Turn that around and you have the core idea of stream processing: a table is a stream folded up by key, and a stream is a table's history unfolded. This is the stream-table duality, and once you see it, Kafka Streams' KStream and KTable, Flink's dynamic tables, change data capture and materialized views stop being separate features and become one model.
The duality is easy to state and easy to get wrong in production. Treating a changelog as if it were a stream of independent events double-counts. Forgetting that a null value means delete leaks state. Joining against "the current table" when events arrive late gives answers that change on replay. This article builds the model from first principles, maps it onto Kafka Streams and Flink, works an example end to end, and catalogs the failure modes.
From first principles: fold and unfold
Start with a stream of events about bank accounts. Each event is a fact that happened: a deposit or a withdrawal. Folding the stream by key, keeping a running total per account, produces a table of balances. That table at any moment is a pure function of the stream prefix seen so far. Now record every change to the table as it happens: account 17's balance became 120, then 95. That record is a new stream, a changelog. Replaying the changelog into an empty map rebuilds the table exactly. The two directions are inverse views of the same state, with one important asymmetry: from the event stream you can derive any table, but from a changelog you can reliably rebuild only the table it came from.
from collections import defaultdict
events = [("acct-17", +100), ("acct-42", +50), ("acct-17", +20), ("acct-17", -25)]
def fold(stream):
'''Stream -> table, emitting the changelog as a side effect.'''
table, changelog = defaultdict(int), []
for key, amount in stream:
table[key] += amount
changelog.append((key, table[key])) # upsert: new value for the key
return dict(table), changelog
def replay(changelog):
'''Changelog -> table. Last write per key wins.'''
table = {}
for key, value in changelog:
if value is None:
table.pop(key, None) # tombstone: delete
else:
table[key] = value
return table
table, log = fold(events)
assert replay(log) == table # {'acct-17': 95, 'acct-42': 50}A complete changelog does let you recover the net changes by subtraction: 100, 120, 95 implies +20 then -25. What it never carries is the event's own attributes, such as type, source or reason. And in practice changelogs are rarely complete: compaction keeps only 95, and caches that coalesce updates drop intermediate values. That loss is the asymmetry, and treating a changelog as events is the root of the most common bug described below.
Three kinds of stream
Every stream in a pipeline is one of three kinds, and the consumer must know which. An event stream carries independent facts; every record counts. An upsert or changelog stream carries the new value for a key; each record replaces the previous one for that key, and a null value deletes. A retraction stream carries explicit inserts and deletes of whole rows, so an update is a delete of the old row followed by an insert of the new one.
| Concept | Kafka Streams | Flink SQL / Table API |
|---|---|---|
| Event stream | KStream | Append-only table (only +I rows) |
| Changelog by key | KTable (partitioned) / GlobalKTable (replicated) | Upsert stream: +I, +U, -D keyed by primary key |
| Retraction stream | Old and new values passed internally to KTable aggregations | Retract stream: +I, -U, +U, -D |
| Delete | Record with null value (tombstone) | -D row |
Flink makes the row kinds explicit: +I insert, -U retract the old version of an updated row, +U the new version, -D delete. An upsert sink needs a primary key and can drop -U rows, because the next +U for that key overwrites. A sink without a key needs the full retract stream to stay correct. Kafka Streams hides the same machinery inside KTable operations but exposes the result: a KTable's output topic is an upsert stream keyed by the table key.
Where the duality lives in Kafka
Kafka Streams stores each table in a local state store and backs it with a changelog topic. On restart or rebalance, a task restores its store by replaying that topic, which is the replay function above at scale. Because only the latest value per key is needed for restore, changelog topics are compacted, and compaction removes superseded values while keeping each key's last record, including tombstones until a retention period passes; the mechanics are in Kafka log compaction architecture. The topology-level view of stores, tasks and exactly-once processing is in Kafka Streams in depth.
Two consequences follow directly. A compacted topic is a table, not a history: you cannot recompute a sum of deposits from it. And a tombstone is a real delete instruction: reading a topic as a KTable turns a null value into removal of the key, while reading the same topic as a KStream yields a record with a null value, which aggregations will typically skip or crash on.
Worked example: the double-counting bug
A team computes balances as a KTable, then wants the total money held across all accounts. The first attempt reads the balance changelog as a stream and sums it.
// WRONG: consumes an upsert stream as if each record were an event.
KTable<String, Long> balances = deposits
.groupByKey()
.aggregate(() -> 0L, (acct, amount, bal) -> bal + amount,
Materialized.with(Serdes.String(), Serdes.Long()));
balances.toStream()
.groupBy((acct, bal) -> "ALL")
.reduce(Long::sum); // adds 100, 120 and 95 for acct-17With the events above, the true total is 145: account 17 holds 95 and account 42 holds 50. The wrong pipeline adds every intermediate balance, so it reports 100 + 50 + 120 + 95 = 365, and it gets worse with every update. The bug is invisible in a test with one event per key.
The fix is to stay in the table world, where an update carries both the old and the new value. Re-grouping a KTable produces a KGroupedTable whose aggregate takes an adder and a subtractor: when account 17 changes from 120 to 95, the framework calls the subtractor with 120 and the adder with 95.
// RIGHT: table-to-table aggregation with retraction.
KTable<String, Long> total = balances
.groupBy((acct, bal) -> KeyValue.pair("ALL", bal),
Grouped.with(Serdes.String(), Serdes.Long()))
.aggregate(
() -> 0L,
(key, bal, agg) -> agg + bal, // adder: new value
(key, bal, agg) -> agg - bal, // subtractor: old value
Materialized.with(Serdes.String(), Serdes.Long()));In Flink SQL the same query, SELECT SUM(balance) FROM balances over a table that receives updates, is correct automatically, because the planner propagates -U rows. The rule generalizes: aggregate events with a stream aggregation, aggregate tables with a table aggregation, and never cross the line without converting deliberately. Converting a table to a stream is legitimate when downstream wants "the latest balance, every time it changes", for notifications, for example, but then it must not be summed.
Joins: lookups now, or lookups as of then
A stream-table join enriches each event with the table's value for its key: orders joined with customers, payments joined with exchange rates. Semantically it is a lookup, triggered only by the stream side; a table update produces no output by itself. A table-table join is different: its result is itself a table, re-emitted whenever either side changes, and it is the streaming equivalent of a materialized join view such as those described in Materialize for streaming SQL.
The subtle part is time. With an ordinary store, the lookup sees the table as it is when the processor handles the event. If a payment from 10:00 arrives late, at 10:05, after the rate changed at 10:02, it is converted at the wrong rate, and a replay may produce a different answer from the live run. Kafka 3.5 added versioned state stores (KIP-889) with DSL semantics for them (KIP-914): if the table is materialized in a versioned store, the join looks up the value as of the stream record's timestamp. Lookups older than the store's history retention return nothing, which an inner join treats as no match and a left join as a null table side.
KTable<String, Rate> rates = builder.table(
"fx-rates",
Materialized.<String, Rate>as(
Stores.persistentVersionedKeyValueStore("fx-rates-versioned", Duration.ofHours(24)))
.withKeySerde(Serdes.String())
.withValueSerde(rateSerde));
KStream<String, Payment> payments = builder.stream("payments"); // keyed by currency
payments
.join(rates, (pay, rate) -> pay.convertWith(rate)) // rate as of pay's timestamp
.to("payments-usd");Set history retention longer than the worst lateness you accept, and monitor joins that miss because they fell outside it. Flink SQL expresses the same thing as an event-time temporal join with FOR SYSTEM_TIME AS OF against a versioned table, whose state is described in Flink state architecture.
Emitting tables downstream
A table changes on every input, and naively each change is emitted. A key updated 1,000 times a second becomes 1,000 output records, most immediately superseded. Kafka Streams places a record cache in front of stores that coalesces updates per key before forwarding, so the emission rate depends on cache size and commit interval; that makes intermediate values appear or vanish depending on configuration, which is correct for an upsert stream but surprising in tests. When downstream needs fewer, final-ish results, use suppress, for example Suppressed.untilTimeLimit(Duration.ofSeconds(5), BufferConfig.maxRecords(10_000)) for rate limiting, or untilWindowCloses for windowed aggregates that should emit once.
Failure modes
| Symptom | Cause | Fix |
|---|---|---|
| Totals grow with every update | Changelog summed as events | Table aggregation with subtractor, or retract stream |
| Deleted entities reappear | Tombstones dropped by a filter or serializer that cannot write null | Forward nulls; test delete paths |
| Cannot rebuild a metric after a bug fix | Only the compacted changelog was kept | Retain the source event stream with time-based retention |
| Join results differ on replay | Lookup against current table, late events | Versioned store or temporal join |
| Joins silently drop records | Keys or partition counts differ between sides | Co-partition: same key, same partitioner, same count |
| Wrong value wins for a key | Updates for one key written to different partitions | Key the changelog by the table key, always |
Trade-offs
Choosing a representation is choosing what downstream consumers must know. Upsert streams are compact and easy to sink into key-value stores, but they lose history and need keys. Retraction streams are general, work with keyless sinks and arbitrary SQL, but roughly double the traffic for updates. Versioned stores make joins deterministic at the cost of storage proportional to update rate times retention. Keeping the raw event stream is the insurance policy that lets you derive new tables later; compacted changelogs alone cannot do that.
There is also a placement trade-off for lookup tables. A KTable is partitioned, so a stream joined against it must be co-partitioned with it, and a key change before the join forces a repartition topic and an extra network hop. A GlobalKTable is replicated in full to every instance, so the join needs no co-partitioning and can use a foreign key extracted from the event, but every instance pays the memory, disk and restore time for the whole table, and global tables are bootstrapped before processing starts. Use global tables for small, slowly changing reference data such as currencies or product categories, and partitioned tables for anything that grows with your users or entities.
Finally, consider where the table should live at all. If the only consumer is a service that needs current values by key, a stream processor that materializes the table and serves queries from its local store avoids a second database, but it couples query availability to rebalances. Sinking the upsert stream into an external key-value store or database decouples them, at the cost of one more system whose writes must be idempotent per key, which upsert streams make natural.
What to do next
- Label every topic in your pipeline as event, upsert or retraction stream, and record the key.
- Search for
toStream()followed by an aggregation; each one is a double-counting candidate. - Write a test with several updates and a delete per key, not one event per key.
- Verify that tombstones survive every filter, mapper and serializer on the path.
- For stream-table joins with late data, move the table to a versioned store and size history retention.
- Keep the source events with time-based retention, separate from compacted changelogs.