13.2 Unbundling Databases
The two philosophies, unresolved after 50 years:
| Unix (early 1970s) | Relational databases (early 1970s) |
|---|---|
| Presenting programmers with a logical but fairly low-level hardware abstraction. | Giving application programmers a high-level abstraction that would hide the complexities of data structures on disk, concurrency and crash recovery. |
| Pipes and files that are just sequences of bytes. | SQL and transactions. |
| Simpler in the sense that it is a fairly thin wrapper around hardware resources. | Simpler in the sense that a short declarative query can draw on a lot of powerful infrastructure — query optimization, indexes, join methods, concurrency control, replication — without the author needing to understand the implementation details. |
The NoSQL movement could be interpreted as wanting to apply a Unix-esque approach of low-level abstractions to the domain of distributed OLTP storage.
2.1 CREATE INDEX is a batch job
Think about what happens when you run
CREATE INDEX: the database must
- SCAN over a CONSISTENT SNAPSHOT of the table
- Pick out the field values, SORT them, write out the index
- PROCESS THE BACKLOG OF WRITES made since the snapshot was taken
- CONTINUE to keep the index up to date whenever a transaction writes
"This process is REMARKABLY SIMILAR TO SETTING UP A NEW FOLLOWER REPLICA, and also VERY SIMILAR TO BOOTSTRAPPING CDC in a streaming system. Whenever you run
CREATE INDEX, THE DATABASE ESSENTIALLY REPROCESSES THE EXISTING DATASET AND DERIVES THE INDEX AS A NEW VIEW ONTO THE EXISTING DATA."
2.2 The meta-database of everything
"In this light, THE DATAFLOW ACROSS AN ENTIRE ORGANIZATION STARTS LOOKING LIKE ONE HUGE DATABASE. Whenever a batch, stream, or ETL process transports data from one place and form to another, IT IS ACTING LIKE THE DATABASE SUBSYSTEM THAT KEEPS INDEXES OR MATERIALIZED VIEWS UP TO DATE."
"Viewed like this, BATCH AND STREAM PROCESSORS ARE LIKE ELABORATE IMPLEMENTATIONS OF TRIGGERS, STORED PROCEDURES, AND MATERIALIZED VIEW MAINTENANCE ALGORITHMS. The derived data systems they maintain are LIKE DIFFERENT INDEX TYPES. Instead of implementing those facilities as features of a SINGLE INTEGRATED DATABASE PRODUCT, they are provided by VARIOUS PIECES OF SOFTWARE, RUNNING ON DIFFERENT MACHINES, ADMINISTERED BY DIFFERENT TEAMS."
Two avenues, and they are complementary:
| Federated databases — unifying reads | Unbundled databases — unifying writes |
|---|---|
| A unified query interface to a wide variety of underlying storage engines and processing methods — a federated database or polystore. | Ensure that all data changes end up in all the right places, even in the face of faults. Making it easier to reliably plug together storage systems (through CDC and event logs) is like unbundling a database's index maintenance features in a way that can synchronize writes across disparate technologies. |
| PostgreSQL foreign data wrappers; Trino, Hoptimator, Xorq. | Debezium, Kafka, IVM engines. |
| Applications that need a specialized data model can still access the underlying engines directly, while users who want to combine data can do so through the federated interface. | Events for a given record flow through one path, so the pieces stay in step without any component having to know about all the others. |
| Follows the relational tradition: a single integrated system with a high-level query language and elegant semantics, but a complicated implementation. | Follows the Unix tradition: small tools that do one thing well, that communicate through a uniform low-level API (pipes), and that can be composed using a higher-level language (the shell). |
"Federated read-only querying requires MAPPING ONE DATA MODEL INTO ANOTHER, which takes some thought BUT IS ULTIMATELY QUITE A MANAGEABLE PROBLEM. KEEPING THE WRITES TO SEVERAL STORAGE SYSTEMS IN SYNC IS THE HARDER ENGINEERING PROBLEM."
2.3 Why log-based integration beats distributed transactions
"Transactions WITHIN a single storage or stream processing system are feasible, BUT WHEN DATA CROSSES THE BOUNDARY BETWEEN DIFFERENT TECHNOLOGIES, AN ASYNCHRONOUS EVENT LOG WITH IDEMPOTENT WRITES IS A MUCH MORE ROBUST AND PRACTICABLE APPROACH."
"Distributed transactions ARE used within some stream processors to achieve exactly-once semantics, and this can work quite well. HOWEVER, WHEN A TRANSACTION WOULD NEED TO INVOLVE SYSTEMS WRITTEN BY DIFFERENT GROUPS OF PEOPLE, THE LACK OF A STANDARDIZED TRANSACTION PROTOCOL MAKES INTEGRATION MUCH HARDER."
Loose coupling manifests in two ways:
| Level | Benefit |
|---|---|
| System | "Asynchronous event streams make the system AS A WHOLE more robust to outages or performance degradation of individual components. If a consumer runs slow or fails, THE EVENT LOG CAN BUFFER MESSAGES, allowing the producer and other consumers to CONTINUE RUNNING UNAFFECTED. The faulty consumer CAN CATCH UP WHEN IT IS FIXED, so it doesn't miss any data, AND THE FAULT IS CONTAINED. By contrast, THE SYNCHRONOUS INTERACTION OF DISTRIBUTED TRANSACTIONS TENDS TO ESCALATE LOCAL FAULTS INTO LARGE-SCALE FAILURES" |
| Human | "Unbundling allows software components and services to be DEVELOPED, IMPROVED, AND MAINTAINED INDEPENDENTLY BY DIFFERENT TEAMS. Specialization allows each team to FOCUS ON DOING ONE THING WELL, with well-defined interfaces. Event logs provide an interface POWERFUL ENOUGH to capture fairly strong consistency properties (because of durability and ordering), BUT ALSO GENERAL ENOUGH to be applicable to almost any kind of data" |
2.4 The honest caveat — don't unbundle prematurely
"If unbundling does become the way of the future, IT WILL NOT REPLACE DATABASES IN THEIR CURRENT FORM. They will still be needed AS MUCH AS EVER, for maintaining state in stream processors and to SERVE QUERIES for the output of batch and stream processors."
"The COMPLEXITY OF RUNNING SEVERAL PIECES OF INFRASTRUCTURE CAN BE A PROBLEM. Each piece has a LEARNING CURVE and its own CONFIGURATION ISSUES AND OPERATIONAL QUIRKS, so IT IS WORTH DEPLOYING AS FEW MOVING PARTS AS POSSIBLE. A single integrated product may also achieve BETTER AND MORE PREDICTABLE PERFORMANCE on the workloads it's designed for."
"BUILDING FOR SCALE THAT YOU DON'T NEED IS WASTED EFFORT AND MAY LOCK YOU INTO AN INFLEXIBLE DESIGN. IN EFFECT, IT IS A FORM OF PREMATURE OPTIMIZATION."
"THE GOAL OF UNBUNDLING IS NOT TO COMPETE WITH INDIVIDUAL DATABASES ON PERFORMANCE FOR PARTICULAR WORKLOADS; THE GOAL IS TO ALLOW YOU TO COMBINE SEVERAL DATABASES IN ORDER TO ACHIEVE GOOD PERFORMANCE FOR A MUCH WIDER RANGE OF WORKLOADS. IT'S ABOUT BREADTH, NOT DEPTH."
"Thus, IF A SINGLE TECHNOLOGY DOES EVERYTHING YOU NEED, YOU'RE MOST LIKELY BEST OFF SIMPLY USING THAT PRODUCT. The advantages of unbundling come into the picture ONLY WHEN NO SINGLE PIECE OF SOFTWARE SATISFIES ALL YOUR REQUIREMENTS."