13.1 Data Integration
The starting observation:
"If you have a problem such as 'I want to store some data and look it up again later,' THERE IS NO ONE RIGHT SOLUTION, but many approaches each appropriate in different circumstances. A software implementation typically has to PICK ONE. It's hard enough to get ONE code path robust and performing well; TRYING TO SATISFY TOO MANY USE CASES WITH MANY FEATURES IS LIKELY TO LEAD TO POOR IMPLEMENTATIONS OF THOSE FEATURES compared to specialized tools."
"EVERY PIECE OF SOFTWARE — EVEN A SO-CALLED 'GENERAL-PURPOSE' DATABASE — IS DESIGNED FOR A PARTICULAR USAGE PATTERN."
Two challenges, in order:
- Map software products to the circumstances they fit. "Vendors are UNDERSTANDABLY RELUCTANT to tell you about the kinds of workloads for which their software is POORLY SUITED, but hopefully the previous chapters have equipped you with QUESTIONS TO ASK that will help you READ BETWEEN THE LINES."
- Even with a perfect map, "in complex applications data is used in various ways, and ONE PIECE OF SOFTWARE IS UNLIKELY TO BE SUITABLE FOR ALL OF THEM. Therefore YOU INEVITABLY END UP HAVING TO COBBLE TOGETHER SEVERAL PIECES OF SOFTWARE."
(The canonical example: an OLTP database plus a full-text search index. "Some databases include full-text indexing, which can be sufficient for simple applications, but more sophisticated search requires specialist information-retrieval tools. Conversely, SEARCH INDEXES ARE GENERALLY NOT VERY SUITABLE AS A DURABLE SYSTEM OF RECORD.")
1.1 Reasoning about dataflows
"When copies of the same data must be maintained in several storage systems, YOU NEED TO BE VERY CLEAR ABOUT THE INPUTS AND OUTPUTS. WHERE IS DATA WRITTEN FIRST, AND WHICH REPRESENTATIONS ARE DERIVED FROM WHICH SOURCES?"
If CDC is the only way of updating the index, you can be confident it is entirely derived from the system of record and therefore consistent with it (barring bugs). Writing to the database is the only way of supplying new input.
With dual writes, neither the database nor the search index is in charge of determining the order of writes, so they may make contradictory decisions and become permanently inconsistent.
"If it is possible to FUNNEL ALL USER INPUT THROUGH A SINGLE SYSTEM THAT DECIDES ON AN ORDERING FOR ALL WRITES, it becomes much easier to derive other representations by PROCESSING THE WRITES IN THE SAME ORDER. This is an application of STATE MACHINE REPLICATION. WHETHER YOU USE CDC OR AN EVENT SOURCING LOG IS LESS IMPORTANT THAN THE PRINCIPLE OF DECIDING ON A TOTAL ORDER."
And: "updating a derived data system based on an event log can often be made DETERMINISTIC AND IDEMPOTENT, MAKING IT QUITE EASY TO RECOVER FROM FAULTS."
1.2 Derived data vs distributed transactions
| Distributed transactions | Log-based derived data | |
|---|---|---|
| Mechanism | An ATOMIC COMMIT PROTOCOL ensures changes are applied atomically | Correctness through DETERMINISTIC RETRY AND IDEMPOTENCE |
| The biggest difference | "After a value is written, YOU CAN IMMEDIATELY READ THE UP-TO-DATE VALUE" | "Updated ASYNCHRONOUSLY, so they DO NOT BY DEFAULT GUARANTEE THAT READS ARE UP TO DATE" |
| Verdict | "Used successfully in environments willing to absorb their performance and operational costs. HOWEVER, XA HAS POOR FAULT TOLERANCE AND PERFORMANCE CHARACTERISTICS, WHICH SEVERELY LIMIT ITS USEFULNESS. It might be possible to create a better protocol, BUT GETTING IT WIDELY ADOPTED AND INTEGRATED WITH EXISTING TOOLS WOULD BE CHALLENGING, AND IT IS UNLIKELY TO HAPPEN SOON" | "In the absence of widespread support for a good distributed transaction protocol, LOG-BASED DERIVED DATA IS THE MOST PROMISING APPROACH for integrating different data systems" |
The intellectual honesty worth quoting: "Guarantees such as reading your own writes ARE USEFUL, and IT IS NOT PRODUCTIVE TO TELL EVERYONE 'EVENTUAL CONSISTENCY IS INEVITABLE — SUCK IT UP AND LEARN TO DEAL WITH IT' (at least not without good guidance on how to deal with it)."
1.3 The limits of total ordering
Total order is feasible in small systems — "as demonstrated by the popularity of databases with single-leader replication, which construct precisely such a log." But four limits emerge at scale:
- Throughput. “Constructing a totally ordered log requires all events to pass through A single leader node that decides on the ordering. If throughput exceeds what one machine can handle, You need to shard the log — and the order of events in two shards is then ambiguous.”
- Geography. “If servers are spread across multiple regions, you typically have A separate leader in each datacenter, because network delays make synchronous cross-datacenter coordination inefficient. This implies an undefined ordering of events that originate in two different datacenters.”
- Microservices. “A common design choice is to deploy each service and its durable state as An independent unit, with no durable state shared between services. when two events originate in different services, those events have no defined order.”
- Client-side state. “Some applications maintain client-side state updated Immediately on user input (without waiting for server confirmation), and even Continue to work offline. Clients and servers are very likely to see events in different orders.”
★ “Deciding on a total order is Total order broadcast, which is Equivalent to consensus. Most consensus algorithms are designed for situations where The throughput of a single node is sufficient to process the entire stream, and They do not provide a mechanism for multiple nodes to share the work of ordering.”
1.4 Ordering events to capture causality — the unfriending problem
"If NO CAUSAL LINK exists between events, the lack of a total order IS NOT A BIG PROBLEM, since concurrent events can be ordered arbitrarily. Some cases are easy — multiple updates of the SAME OBJECT can be totally ordered by ROUTING ALL UPDATES FOR THAT OBJECT ID TO THE SAME LOG SHARD. HOWEVER, CAUSAL DEPENDENCIES SOMETIMES ARISE IN MORE SUBTLE WAYS."
Two users were in a relationship and have just broken up.
The user's intention is that their ex-partner should not see the rude message, since it was sent after that individual's friend status was revoked. But friendship status is stored in one place and messages in another, so the ordering dependency may be lost: a notification service may process the message-send event before the unfriend event, and so it sends a notification to the ex-partner.
The notifications are effectively a join between the messages and the friend list — the time-dependence-of-joins problem of Ch 12 §4.3, in disguise.
"Unfortunately, THIS PROBLEM DOESN'T SEEM TO HAVE A SIMPLE ANSWER." Three starting points:
- LOGICAL TIMESTAMPS "provide total ordering WITHOUT COORDINATION, so they may help when total order broadcast is not feasible. HOWEVER, THEY STILL REQUIRE RECIPIENTS TO HANDLE EVENTS DELIVERED OUT OF ORDER, and they require ADDITIONAL METADATA to be passed around"
- "If you can LOG AN EVENT TO RECORD THE STATE OF THE SYSTEM THAT THE USER SAW BEFORE MAKING A DECISION and give it a unique identifier, then ANY LATER EVENTS CAN REFERENCE THAT EVENT IDENTIFIER in order to record the causal dependency"
- CONFLICT RESOLUTION ALGORITHMS "help with processing events delivered in an unexpected order. They are USEFUL FOR MAINTAINING STATE, BUT THEY DO NOT HELP IF ACTIONS HAVE EXTERNAL SIDE EFFECTS (such as sending a notification)"
"Perhaps patterns will emerge in the future that allow causal dependencies to be captured efficiently WITHOUT FORCING ALL EVENTS THROUGH THE BOTTLENECK OF TOTAL ORDER BROADCAST."
1.5 Batch and stream, together
"The MAIN FUNDAMENTAL DIFFERENCE is that stream processors operate on UNBOUNDED datasets, whereas batch inputs are of a KNOWN, FINITE SIZE."
"Batch processing has A QUITE STRONG FUNCTIONAL FLAVOR (even if the code is not written in a functional language). It encourages DETERMINISTIC, PURE FUNCTIONS whose output depends only on the input and that have NO SIDE EFFECTS other than the explicit outputs, treating INPUTS AS IMMUTABLE AND OUTPUTS AS APPEND-ONLY. Stream processing is similar, but it EXTENDS OPERATORS TO ALLOW A MANAGED, FAULT-TOLERANT STATE."
Why asynchrony is the point, not a compromise:
"In principle, derived data systems COULD be maintained SYNCHRONOUSLY, just as a relational database updates secondary indexes synchronously within the same transaction. HOWEVER, ASYNCHRONY IS WHAT MAKES SYSTEMS BASED ON EVENT LOGS ROBUST. IT ALLOWS A FAULT IN ONE PART OF THE SYSTEM TO BE CONTAINED LOCALLY, WHEREAS DISTRIBUTED TRANSACTIONS ABORT IF ANY ONE PARTICIPANT FAILS, SO THEY TEND TO AMPLIFY FAILURES BY SPREADING THEM TO THE REST OF THE SYSTEM."
(And Ch 7's point again: "a sharded system with secondary indexes needs to either send writes to multiple shards (term-partitioned) or send reads to all shards (document-partitioned). SUCH CROSS-SHARD COMMUNICATION IS ALSO MOST RELIABLE AND SCALABLE IF THE INDEX IS MAINTAINED ASYNCHRONOUSLY.")
Reprocessing for application evolution — and the railway analogy:
"WITHOUT REPROCESSING, SCHEMA EVOLUTION IS LIMITED TO SIMPLE CHANGES like adding a new optional field. WITH REPROCESSING, IT IS POSSIBLE TO RESTRUCTURE A DATASET INTO A COMPLETELY DIFFERENT MODEL to better serve new requirements."
19th-century English railways faced the question of how to change the gauge without shutting down the line for months or years.
- Convert the track to dual gauge (mixed gauge) by adding a third rail. This conversion can be done gradually.
- Trains of both gauges now run on the line, using two of the three rails.
- Once all trains are converted, remove the rail providing the nonstandard gauge.
“Reprocessing” the existing tracks in this way, and allowing the old and new versions to exist side by side, makes it possible to change the gauge gradually over the course of years. Nevertheless, the undertaking is expensive, which is why nonstandard gauges still exist today — BART uses a different gauge from the majority of the US.
The data equivalent is to maintain the old and new schema side by side as two independently derived views onto the same underlying data:
- Shift a small number of users to the new view to test performance and find bugs.
- Gradually increase the proportion.
- Eventually drop the old view.
The beauty of such a gradual migration is that every stage is easily reversible if something goes wrong; you always have a working system to go back to. Reducing the risk of irreversible damage allows you to be more confident about going ahead and thus to move faster — Ch 2's evolvability, operationalized.
Unifying batch and stream — the kappa architecture:
"An early proposal was the LAMBDA ARCHITECTURE, which HAD A NUMBER OF PROBLEMS AND HAS FALLEN OUT OF USE. More recent systems allow batch computations (reprocessing historical data) and stream computations (processing events as they arrive) to be implemented IN THE SAME SYSTEM — sometimes known as the KAPPA ARCHITECTURE."
Three features required:
- "The ability to REPLAY HISTORICAL EVENTS THROUGH THE SAME PROCESSING ENGINE that handles the stream of recent events" — log-based brokers can replay; some stream processors can read from a DFS or object store
- "EXACTLY-ONCE SEMANTICS — ensuring the output is the same as if no faults had occurred. As with batch processing, this requires DISCARDING THE PARTIAL OUTPUTS OF ANY FAILED TASKS"
- "Tools for WINDOWING BY EVENT TIME, NOT BY PROCESSING TIME, SINCE PROCESSING TIME IS MEANINGLESS WHEN REPROCESSING HISTORICAL EVENTS" (Apache Beam provides such an API, runnable on Flink or Google Cloud Dataflow)