Learn Labs
13. A Philosophy of Streaming Systems

13.3 Designing Applications Around Dataflow

Derived state

update
12 ms
write p99
200 ms
staleness
appendasyncwriteordered logsystem of recordsearch indexcachewarehouse
write latency12 ms
write availability100%
Safe

Writes stay fast and available at the cost of 200 ms of staleness in the derived views. Because the log fixes a total order, every view converges to the same state — that is what makes asynchrony tolerable here rather than a source of permanent divergence.

Pick one system of record and derive everything else from its log. The only question left is whether the derived views update inside the write or after it — and that choice is the whole trade.

3.1 The spreadsheet that data systems still haven't caught up to

"Spreadsheets have POWERFUL DATAFLOW PROGRAMMING CAPABILITIES: you put a formula in one cell, and WHENEVER ANY INPUT CHANGES, THE RESULT IS AUTOMATICALLY RECALCULATED. THIS IS EXACTLY WHAT WE WANT AT A DATA SYSTEM LEVEL."

"MOST DATA SYSTEMS STILL HAVE SOMETHING TO LEARN FROM THE FEATURES THAT VISICALC ALREADY HAD IN 1979."

The difference: "today's data systems need to be FAULT-TOLERANT, SCALABLE, AND CAPABLE OF STORING DATA DURABLY. They also need to INTEGRATE DISPARATE TECHNOLOGIES written by different groups of people over time. IT IS UNREALISTIC TO EXPECT ALL SOFTWARE TO BE DEVELOPED USING ONE PARTICULAR LANGUAGE, FRAMEWORK, OR TOOL."

3.2 Application code as a derivation function

Derived datasetThe transformation function
Secondary index"For each row, PICK OUT THE VALUES in the columns being indexed and SORT by those values." (Built into databases — you invoke it by "merely running CREATE INDEX")
Full-text search index"LANGUAGE DETECTION, WORD SEGMENTATION, STEMMING OR LEMMATIZATION, SPELLING CORRECTION, AND SYNONYM IDENTIFICATION, followed by building a data structure for efficient lookups." (Basic linguistic features may be built in, but "MORE SOPHISTICATED FEATURES OFTEN REQUIRE DOMAIN-SPECIFIC TUNING")
ML model"Derived from the TRAINING DATA by applying various FEATURE EXTRACTION AND STATISTICAL ANALYSIS functions. When applied to new input, its output is derived from that input AND its learned parameters (and hence, INDIRECTLY, from the training data)." ("Feature engineering is NOTORIOUSLY APPLICATION-SPECIFIC")
Cache"An aggregation of data IN THE FORM IN WHICH IT IS GOING TO BE DISPLAYED IN A UI. Populating it thus requires KNOWLEDGE OF WHAT FIELDS ARE REFERENCED IN THE UI; CHANGES IN THE UI MAY REQUIRE UPDATING THE DEFINITION OF HOW THE CACHE IS POPULATED AND REBUILDING IT."

"When the function is NOT a standard cookie-cutter function, CUSTOM CODE IS REQUIRED. THIS CUSTOM CODE IS WHERE MANY DATABASES STRUGGLE. Although relational databases support triggers, stored procedures, and UDFs, THEY HAVE BEEN SOMEWHAT OF AN AFTERTHOUGHT IN DATABASE DESIGN."

3.3 The separation of Church and state

"In theory, databases COULD be deployment environments for arbitrary application code, LIKE AN OPERATING SYSTEM. However, in practice they have turned out to be POORLY SUITED. They do not fit well with the requirements of modern application development: DEPENDENCY AND PACKAGE MANAGEMENT, VERSION CONTROL, ROLLING UPGRADES, EVOLVABILITY, MONITORING, METRICS, CALLS TO NETWORK SERVICES, AND INTEGRATION WITH EXTERNAL SYSTEMS."

"On the other hand, KUBERNETES, DOCKER, MESOS, YARN are designed SPECIFICALLY for running application code. BY FOCUSING ON DOING ONE THING WELL, they are able to do it MUCH BETTER than a database that provides execution of UDFs as one of its many features."

"The trend has been to keep STATELESS APPLICATION LOGIC SEPARATE FROM STATE MANAGEMENT: NOT PUTTING APPLICATION LOGIC IN THE DATABASE AND NOT PUTTING PERSISTENT STATE IN THE APPLICATION. As people in the functional programming community like to joke, 'WE BELIEVE IN THE SEPARATION OF CHURCH AND STATE.'"

(The joke explained: Alonzo Church created the lambda calculus, which has no mutable state — "so one could say that mutable state is separate from Church's work.")

The passivity problem:

"In this typical model, the database acts as A KIND OF MUTABLE SHARED VARIABLE that can be accessed synchronously over the network. HOWEVER, IN MOST PROGRAMMING LANGUAGES YOU CANNOT SUBSCRIBE TO CHANGES IN A MUTABLE VARIABLE — YOU CAN ONLY READ IT PERIODICALLY. Unlike in a spreadsheet, READERS OF THE VARIABLE DON'T GET NOTIFIED IF THE VALUE CHANGES." (You can implement it yourself — the observer pattern — but most languages don't have it built in.)

"DATABASES HAVE INHERITED THIS PASSIVE APPROACH TO MUTABLE DATA. If you want to find out whether the content has changed, OFTEN YOUR ONLY OPTION IS TO POLL. SUBSCRIBING TO CHANGES IS ONLY JUST BEGINNING TO EMERGE AS A FEATURE."

3.4 Dataflow: state changes and code, collaborating

"Instead of treating a database as A PASSIVE VARIABLE THAT IS MANIPULATED BY THE APPLICATION, we think much more about THE INTERPLAY AND COLLABORATION BETWEEN STATE, STATE CHANGES, AND CODE THAT PROCESSES THEM. APPLICATION CODE RESPONDS TO STATE CHANGES IN ONE PLACE BY TRIGGERING STATE CHANGES IN ANOTHER PLACE."

Two properties log-based brokers must provide:

  1. "THE ORDER of state changes is often important — if several views are derived from an event log, THEY NEED TO PROCESS THE EVENTS IN THE SAME ORDER so that they remain consistent with one another"
  2. "FAULT TOLERANCE is essential — LOSING JUST A SINGLE MESSAGE CAUSES THE DERIVED DATASET TO GO PERMANENTLY OUT OF SYNC with its data source"

"Stable message ordering and fault-tolerant processing are QUITE STRINGENT DEMANDS, BUT THEY ARE MUCH LESS EXPENSIVE AND MORE OPERATIONALLY ROBUST THAN DISTRIBUTED TRANSACTIONS."

"Like UNIX TOOLS CHAINED BY PIPES, stream operators can be composed to build large systems around dataflow. EACH OPERATOR TAKES STREAMS OF STATE CHANGES AS INPUT AND PRODUCES OTHER STREAMS OF STATE CHANGES AS OUTPUT."

3.5 Stream processors vs services — the currency example

The microservices way
synchronous network requestwaits for the responsepurchase processorexchange rate servicecontinue
The dataflow way
the query is localrecord the current rate whenever it changespurchase processorexchange rate updatessubscribed ahead of timelocal databasesame machine, maybe same processcontinue

The second approach replaces a synchronous network request with a query to a local database. Not only is it faster, it is also more robust to the failure of another service. The fastest and most reliable network request is no network request at all! Instead of RPC, we now have a stream join between purchase events and exchange rate update events.

In the microservices approach you could avoid the network request by caching the rate locally. However, to keep that cache fresh you would need to periodically poll or subscribe to a stream of changes — which is exactly what happens in the dataflow approach.

The join is time-dependent: if the purchase events are reprocessed at a later point, the exchange rate will have changed. To reconstruct the original output you need the historical rate at the original time of purchase. No matter whether you query a service or subscribe to a stream, you will need to handle this.

Figure 13.3.13.5 Stream processors vs services — the currency example

The organizational similarity: "Composing stream operators into dataflow systems has a lot of similar characteristics to microservices — the advantage is primarily ORGANIZATIONAL SCALABILITY THROUGH LOOSE COUPLING. However, the underlying communication mechanism is very different: ONE-DIRECTIONAL, ASYNCHRONOUS MESSAGE STREAMS rather than synchronous request/response."


On this page