12.3 Databases and Streams
Every write to a database is an event that can be captured, stored, and processed. The connection between databases and streams runs deeper than just the physical storage of logs on disk — IT IS QUITE FUNDAMENTAL.
- A REPLICATION LOG is a stream of database write events produced by the leader. Followers apply that stream and end up with an accurate copy.
- STATE MACHINE REPLICATION (Ch 10): if every event represents a write, and every replica processes the same events in the same order, all replicas end in the same state. (Processing is assumed DETERMINISTIC.) IT'S JUST ANOTHER CASE OF EVENT STREAMS.
3.1 The dual-write problem
The setup: an OLTP database for user requests, a cache, a full-text index, a warehouse — each with its own copy of the data, in its own representation optimized for its own purposes. They must be kept in sync.
Dual writes — application code explicitly writes to each system — has two serious problems:
“The two systems are now permanently inconsistent with each other, even though no error occurred.” And “unless you have an additional concurrency detection mechanism such as version vectors, you will not even notice that concurrent writes occurred. One value will simply silently overwrite another.”
Partial failure is the second problem: “one of the writes may fail while the other succeeds. This is a fault-tolerance problem rather than a concurrency problem, but it also has the effect of the two systems becoming inconsistent. Ensuring they both succeed or both fail is the atomic commit problem, which is expensive to solve.”
The root cause: "In Figure 12-4 THERE ISN'T A SINGLE LEADER. The database may have a leader and the search index may have a leader, BUT NEITHER FOLLOWS THE OTHER, so conflicts can occur" — it's accidental multi-leader replication (Ch 6).
The fix: "The situation would be better if there REALLY WAS ONLY ONE LEADER — the database — and if we could MAKE THE SEARCH INDEX A FOLLOWER OF THE DATABASE."
3.2 Change Data Capture (CDC)
The problem with most databases' replication logs is that they have long been considered AN INTERNAL IMPLEMENTATION DETAIL, NOT A PUBLIC API. For decades, many databases simply DID NOT HAVE A DOCUMENTED WAY of getting the log of changes.
CDC is the process of OBSERVING ALL DATA CHANGES written to a database and EXTRACTING THEM IN A FORM IN WHICH THEY CAN BE REPLICATED to other systems.
“Essentially, CDC makes one database the leader and turns the others into followers. A log-based message broker is well suited for transporting the change events, since it preserves the ordering of messages.”
Implementations: Debezium (source connectors for MySQL, PostgreSQL, Oracle, SQL Server, Db2, Cassandra, and more — attaching to replication logs and surfacing changes in a standard event schema), Kafka Connect, Maxwell (parses the MySQL binlog), GoldenGate (Oracle), pgcapture (PostgreSQL).
Like message brokers, CDC is usually ASYNCHRONOUS: the source database DOES NOT WAIT for a change to be applied to consumers before committing. This has the operational advantage that ADDING A SLOW CONSUMER DOES NOT AFFECT THE SYSTEM OF RECORD TOO MUCH, but the downside that ALL THE ISSUES OF REPLICATION LAG APPLY.
Initial snapshot
If you have the log of ALL changes ever made, you can reconstruct the entire state by replaying it. However, keeping all changes forever would require TOO MUCH DISK SPACE, and replaying would TAKE TOO LONG — so the log needs to be truncated.
Building a new full-text index requires A FULL COPY of the database. Applying only recent changes would MISS ITEMS THAT WERE NOT RECENTLY UPDATED.
The snapshot MUST CORRESPOND TO A KNOWN POSITION OR OFFSET IN THE CHANGE LOG, so you know where to start applying changes afterward. (Debezium uses Netflix's DBLog watermarking algorithm for INCREMENTAL SNAPSHOTS.)
Log compaction — the elegant alternative
Before compaction (key = cat video ID, value = play count)
[mew:1][purr:1][mew:2][scratch:1][purr:2][mew:3][yawn:1][purr:3][mew:4]After compaction — keep only the Most recent value for each key
[scratch:1][yawn:1][purr:3][mew:4]“The storage engine periodically looks for records with the Same key, throws away duplicates, and keeps only the most recent update. Segments may also be Merged as part of compaction. This process runs in the background.”
A Tombstone (a special null value) indicates deletion and causes removal.
★ “As long as a key is not overwritten or deleted, it stays in the log forever. the disk space required depends only on the current contents of the database, not the number of writes that have ever occurred.”
The payoff: "whenever you want to rebuild a derived data system such as a search index, START A NEW CONSUMER FROM OFFSET 0 of the log-compacted topic and sequentially scan all messages. THE LOG IS GUARANTEED TO CONTAIN THE MOST RECENT VALUE FOR EVERY KEY. You can use it to obtain A FULL COPY OF THE DATABASE CONTENTS WITHOUT HAVING TO TAKE ANOTHER SNAPSHOT."
This allows the message broker to be used FOR DURABLE STORAGE, NOT JUST FOR TRANSIENT MESSAGING.
API support today
Most popular databases now expose change streams AS A FIRST-CLASS INTERFACE, rather than the RETROFITTED AND REVERSE-ENGINEERED CDC efforts of the past. MySQL and PostgreSQL send changes through the same replication log they use for their own replicas.
The quorum-database challenge, and Cassandra's answer:
CDC support for QUORUM WRITES is challenging because THERE'S NO SINGLE SOURCE OF TRUTH TO SUBSCRIBE TO. Whether data is visible DEPENDS ON EACH READER'S CONSISTENCY PREFERENCES.
Cassandra SIDESTEPS this by exposing RAW LOG SEGMENTS FOR EACH NODE rather than a single stream of mutations. Systems that wish to consume must READ THE RAW SEGMENTS FOR EACH NODE AND DECIDE HOW BEST TO MERGE THEM — MUCH AS A QUORUM READER DOES.
3.3 CDC vs event sourcing
| CDC | Event sourcing | |
|---|---|---|
| Application's view | Uses the database in a MUTABLE way, updating and deleting at will | Application logic is EXPLICITLY BUILT on immutable events written to an event log; updates/deletes are DISCOURAGED OR PROHIBITED |
| Level of abstraction | Extracted at a LOW LEVEL (parsing the replication log) — which is what ensures the order of writes matches the order they were actually written, avoiding the dual-write race | Events reflect things that happened AT THE APPLICATION LEVEL rather than low-level state changes |
| Adoption cost | Can be added to an existing database WITH MINIMAL CHANGES; the application might NOT EVEN KNOW CDC IS OCCURRING | A BIG CHANGE for an application not already doing it |
| Log compaction | ✔ A CDC update event contains THE ENTIRE NEW VERSION of the record, so the current value is entirely determined by the most recent event ⇒ previous events CAN be discarded | ✗ Events express the INTENT of a user action, not the mechanics of the state update. LATER EVENTS TYPICALLY DO NOT OVERRIDE PRIOR ONES, so YOU NEED THE FULL HISTORY. LOG COMPACTION IS NOT POSSIBLE IN THE SAME WAY |
(Event-sourced apps typically store snapshots of current state — but this is ONLY A PERFORMANCE OPTIMIZATION; the intention is that the system can store all raw events forever and reprocess the full log whenever required.)
The schema-as-public-API problem — and the outbox pattern
In a microservices architecture, a database is typically accessed from ONLY ONE SERVICE, making it AN INTERNAL IMPLEMENTATION DETAIL that developers can change freely.
HOWEVER, CDC SYSTEMS TYPICALLY USE THE UPSTREAM DATABASE'S SCHEMA WHEN REPLICATING, WHICH TURNS THESE SCHEMAS INTO PUBLIC APIs THAT MUST BE MANAGED like the public API of the service. REMOVING A COLUMN WILL BREAK DOWNSTREAM CONSUMERS.
Such challenges always existed with data pipelines, but they typically impacted ONLY DATA WAREHOUSE ETL. SINCE CDC IS OFTEN A DATA STREAM, OTHER PRODUCTION SERVICES MIGHT BE CONSUMERS. BREAKING THEM CAN CAUSE A CUSTOMER-FACING OUTAGE. (Data contracts are used to prevent this.)
The outbox pattern:
-- ONE DATABASE TRANSACTION
UPDATE orders SET status = 'shipped' WHERE id = 42; -- internal domain model
INSERT INTO outbox (event_type, payload) VALUES (…); -- the PUBLIC schema“This might look like a dual write — it is. However, outboxes avoid the challenges above by keeping both writes in the same system (the database). This design allows both writes to appear in a single transaction.”
The trade-offs: “developers must still maintain the transformation between their internal and outbox schemas, which can be challenging. An outbox also increases the amount of data the database writes to its underlying storage, which might trigger performance problems.”
3.4 State, Streams, and Immutability — the philosophical core
Whenever you have state that changes, THAT STATE IS THE RESULT OF THE EVENTS THAT MUTATED IT OVER TIME. Your list of available seats is the result of reservations processed; the current balance is the result of credits and debits; the response-time graph is an aggregation of individual response times.
No matter how the state changes, THERE WAS ALWAYS A SEQUENCE OF EVENTS THAT CAUSED THOSE CHANGES. EVEN AS THINGS ARE DONE AND UNDONE, THE FACT REMAINS TRUE THAT THOSE EVENTS OCCURRED.
MUTABLE STATE AND AN APPEND-ONLY LOG OF IMMUTABLE EVENTS DO NOT CONTRADICT EACH OTHER; THEY ARE TWO SIDES OF THE SAME COIN.
“Application state is what you get when you integrate an event stream over time, and a change stream is what you get when you differentiate the state by time.” The analogy has limits — “the second derivative of state does not seem to be meaningful” — but it is a useful starting point.
Jim Gray and Andreas Reuter, 1992: "THERE IS NO FUNDAMENTAL NEED TO KEEP A DATABASE AT ALL; THE LOG CONTAINS ALL THE INFORMATION THERE IS. THE ONLY REASON FOR STORING THE DATABASE (i.e., the current end-of-the-log) IS PERFORMANCE OF RETRIEVAL OPERATIONS."
Advantages of immutable events:
① The accounting precedent — centuries old.
When a transaction occurs, it is recorded in an APPEND-ONLY LEDGER. The accounts — profit and loss, the balance sheet — are DERIVED from the transactions by adding them up.
IF A MISTAKE IS MADE, ACCOUNTANTS DON'T ERASE OR CHANGE THE INCORRECT TRANSACTION. Instead they ADD ANOTHER TRANSACTION THAT COMPENSATES for the mistake — e.g. refunding an incorrect charge. THE INCORRECT TRANSACTION REMAINS IN THE LEDGER FOREVER, because it might be important for auditing. If incorrect figures were already published, the NEXT ACCOUNTING PERIOD INCLUDES A CORRECTION. THIS PROCESS IS ENTIRELY NORMAL IN ACCOUNTING.
② Recovery from buggy code. "If you accidentally deploy buggy code that writes bad data, RECOVERY IS MUCH HARDER IF THE CODE IS ABLE TO DESTRUCTIVELY OVERWRITE DATA." Also: customer service can use an audit log to diagnose requests and complaints.
③ More information than the current state.
On a shopping site, a customer may add an item to their cart and then remove it. Although the second event cancels the first FROM THE POINT OF VIEW OF ORDER FULFILLMENT, IT MAY BE USEFUL FOR ANALYTICS to know THE CUSTOMER WAS CONSIDERING A PARTICULAR ITEM BUT THEN DECIDED AGAINST IT. Perhaps they will buy it in the future, or found a substitute. THIS INFORMATION WOULD BE LOST IN A DATABASE THAT DELETES ITEMS.
④ Deriving several views from the same log.
Having an explicit translation step makes it easier to EVOLVE YOUR APPLICATION. To introduce a new feature presenting existing data in a new way, USE THE EVENT LOG TO BUILD A SEPARATE READ-OPTIMIZED VIEW and RUN IT ALONGSIDE THE EXISTING SYSTEMS WITHOUT MODIFYING THEM. RUNNING OLD AND NEW SIDE BY SIDE IS OFTEN EASIER THAN PERFORMING A COMPLICATED SCHEMA MIGRATION. Once readers have switched, SHUT THE OLD ONE DOWN AND RECLAIM ITS RESOURCES.
"The traditional approach to database and schema design is based on the FALLACY THAT DATA MUST BE WRITTEN IN THE SAME FORM AS IT WILL BE QUERIED. DEBATES ABOUT NORMALIZATION AND DENORMALIZATION BECOME LARGELY IRRELEVANT if you can translate data from a write-optimized event log to read-optimized application state. IT IS ENTIRELY REASONABLE TO DENORMALIZE in the read-optimized views, as the translation process gives you a mechanism for KEEPING IT CONSISTENT WITH THE EVENT LOG."
(The Ch 2 home timeline is exactly this: highly denormalized read-optimized state, kept in sync by the fan-out service.)
Concurrency control — the two-sided effect:
✗ The downside: “consumers of the event log are usually Asynchronous, so a user could make a write, then Read from a derived view and find their write has not yet been reflected.” (Ch 6’s read-your-writes.) Updating the view Synchronously would require “either a Distributed transaction across the event log and the view, or some way of Waiting until an event has been reflected. Both approaches are usually impractical.”
✔ The upside: “Much of the need for Multi-object transactions stems from a single user action requiring data to be changed in Several places. With event sourcing, you can design an event to be a Self-contained description of a user action. The action then requires Only a single write in one place — appending to the log — Which is easy to make atomic.”
✔ And: “If the event log and the application state are Sharded the same way, then a straightforward Single-threaded log consumer needs no concurrency control for writes. By construction it processes only a single event at a time. The log removes the nondeterminism of concurrency by defining a serial order of events in A shard.” (= Ch 8’s actual serial execution.)
Limitations of immutability — three of them:
- Churn. "Some workloads MOSTLY ADD data and rarely update or delete; they are EASY to make immutable. Other workloads have a HIGH RATE OF UPDATES AND DELETES ON A COMPARATIVELY SMALL DATASET; in these cases THE IMMUTABLE HISTORY MAY GROW PROHIBITIVELY LARGE, FRAGMENTATION MAY BECOME AN ISSUE, and THE PERFORMANCE OF COMPACTION AND GARBAGE COLLECTION BECOMES CRUCIAL."
- Legal deletion. GDPR erasure, or containing an accidental leak. "It's NOT SUFFICIENT to just append another event indicating the data should be considered deleted — YOU ACTUALLY WANT TO REWRITE HISTORY AND PRETEND THE DATA WAS NEVER WRITTEN." (Datomic calls this excision; Fossil calls it shunning.)
- Truly deleting is surprisingly hard. "Copies can live in many places. Storage engines, filesystems, and SSDs often WRITE TO A NEW LOCATION RATHER THAN OVERWRITING IN PLACE, and BACKUPS ARE OFTEN DELIBERATELY IMMUTABLE to prevent accidental deletion."
Crypto-shredding, and its honest limits:
Store data you may want to delete ENCRYPTED; when you want to get rid of it, FORGET THE ENCRYPTION KEY. The encrypted data is still there, but nobody can use it.
In a sense, THIS ONLY MOVES THE PROBLEM AROUND; the actual data is still immutable, BUT YOUR KEY STORAGE IS MUTABLE. Moreover, YOU HAVE TO DECIDE UP FRONT which data is encrypted with the same key — an important decision, since YOU CAN LATER CRYPTO-SHRED EITHER ALL OR NONE OF THE DATA ENCRYPTED WITH A PARTICULAR KEY, BUT NOT SOME OF IT. Storing a separate key for every data item would get TOO UNWIELDY, as the key storage would get AS BIG AS THE PRIMARY DATA STORAGE. (Puncturable encryption allows selective revocation but is not yet widely used.)
"Overall, deletion is more a matter of MAKING IT HARDER TO RETRIEVE THE DATA than actually MAKING IT IMPOSSIBLE. Nevertheless, you sometimes have to try."