Learn Labs
13. A Philosophy of Streaming Systems

13.4 Observing Derived State — the write path and the read path

index updatedread the indexdocument writtenbatch / stream processingthe derived datasetthe indexuser querymore processingresponseWrite path — precomputed, done eagerly as soon as the data comes in,regardless of whether anyone has asked to see it.Read path — happens only when someone asks for it.

The derived dataset is where the two paths meet: a trade-off between the amount of work done at write time and the amount done at read time. The write path is similar to eager evaluation and the read path is similar to lazy evaluation.

Figure 13.4.1Observing Derived State — the write path and the read path

The full-text search spectrum, which shows the boundary is a dial:

No indexAn indexPrecompute all results
WriteNothing to do.Update entries for all terms in the document.Impossible. The set of possible queries is infinite (or at least exponential in the number of terms), thus precomputing all possible results would not be possible.
ReadScan all documents, like grep — very expensive if you had a large number of documents.Search each word and apply Boolean logic (AND / OR).—

The middle ground: precompute results for only a fixed set of the most common queries, so they can be served quickly. The uncommon queries can still be served from the index. This would generally be called a cache of common queries — although we could also call it a materialized view, as it would need to be updated when new documents appear.

So the role of caches, indexes and materialized views is simple: they shift the boundary between the read path and the write path. They allow us to do more work on the write path, by precomputing results, in order to save effort on the read path.

"Shifting the boundary was in fact the topic of the social networking example in Ch 2. In that example we also saw how THE BOUNDARY MIGHT BE DRAWN DIFFERENTLY FOR CELEBRITIES COMPARED TO ORDINARY USERS. AFTER 500 PAGES, WE HAVE COME FULL CIRCLE!"

4.1 Extending the write path to the end-user device

"In the past, web browsers were STATELESS CLIENTS that could do useful things only when you had an internet connection (just about the only thing you could do offline was SCROLL UP AND DOWN in a page you had previously loaded). However, single-page JavaScript apps now have A LOT OF STATEFUL CAPABILITIES."

"When we move away from the assumption of stateless clients talking to a central database and toward state maintained ON END-USER DEVICES, A WORLD OF NEW OPPORTUNITIES OPENS UP. In particular, WE CAN THINK OF THE ON-DEVICE STATE AS A CACHE OF STATE ON THE SERVER. THE PIXELS ON THE SCREEN ARE A MATERIALIZED VIEW OF MODEL OBJECTS IN THE CLIENT APP; THE MODEL OBJECTS ARE A LOCAL REPLICA OF STATE IN A REMOTE DATACENTER."

Pushing state changes to clients:

"If you load a typical web page and the data subsequently changes on the server, THE BROWSER DOES NOT FIND OUT UNTIL YOU RELOAD. The browser reads the data at ONLY ONE POINT IN TIME, ASSUMING IT IS STATIC. Thus THE STATE IN THE BROWSER IS A STALE CACHE that is not updated unless you explicitly poll. (HTTP-based feed subscription protocols LIKE RSS ARE REALLY JUST A BASIC FORM OF POLLING.)"

"SERVER-SENT EVENTS (the EventSource API) and WEBSOCKETS provide channels by which a browser can keep AN OPEN TCP CONNECTION and the server can ACTIVELY PUSH MESSAGES."

"Actively pushing state changes all the way to client devices means EXTENDING THE WRITE PATH ALL THE WAY TO THE END USER. When a client is first initialized, it will STILL NEED TO USE A READ PATH TO GET ITS INITIAL STATE, BUT THEREAFTER IT CAN RELY ON A STREAM OF STATE CHANGES. The ideas around stream processing are thus NOT RESTRICTED TO RUNNING IN A DATACENTER; WE CAN EXTEND THEM ALL THE WAY TO END-USER DEVICES."

And offline is already solved: "The devices will be offline some of the time. BUT WE ALREADY SOLVED THAT PROBLEM: in Ch 12 we discussed how a consumer of a log-based broker can RECONNECT AFTER FAILING and ENSURE IT DOESN'T MISS ANY MESSAGES. THE SAME TECHNIQUE WORKS FOR INDIVIDUAL USERS, WHERE EACH DEVICE IS A SMALL SUBSCRIBER TO A SMALL STREAM OF EVENTS."

End-to-end event streams — and why we don't build everything this way:

"Tools such as REACT and ELM already have the ability to UPDATE THE RENDERED UI IN RESPONSE TO CHANGES IN THE UNDERLYING STATE. It would be very natural to extend this programming model to ALSO ALLOW A SERVER TO PUSH STATE CHANGE EVENTS INTO THIS CLIENT-SIDE EVENT PIPELINE."

"State changes could then flow through an END-TO-END WRITE PATH: from the interaction on ONE DEVICE, through event logs and various derived data systems and stream processors, ALL THE WAY TO THE USER INTERFACE ON ANOTHER DEVICE — with fairly low delay, say UNDER ONE SECOND END TO END."

"Some applications — INSTANT MESSAGING AND ONLINE GAMES — ALREADY HAVE SUCH A 'REAL-TIME' ARCHITECTURE. WHY DON'T WE BUILD ALL APPLICATIONS THIS WAY?"

"THE CHALLENGE IS THAT THE ASSUMPTION OF STATELESS CLIENTS AND REQUEST/RESPONSE INTERACTIONS IS DEEPLY INGRAINED IN OUR DATABASES, LIBRARIES, FRAMEWORKS, AND PROTOCOLS. Many datastores support operations where a request returns A SINGLE RESPONSE; FAR FEWER SUPPORT OPERATIONS WHERE A REQUEST RETURNS A STREAM OF RESPONSES OVER TIME."

4.2 Reads are events too

"It is possible to represent READ REQUESTS AS STREAMS OF EVENTS and send BOTH the read events AND the write events through a stream processor. The processor responds to read events by EMITTING THE RESULT OF THE READ TO AN OUTPUT STREAM."

"When both writes and reads are represented as events and routed to the same operator, WE ARE IN FACT PERFORMING A STREAM-TABLE JOIN BETWEEN THE STREAM OF READ QUERIES AND THE DATABASE."

"This correspondence between SERVING REQUESTS and PERFORMING JOINS IS QUITE FUNDAMENTAL: A ONE-OFF READ REQUEST PASSES THROUGH THE JOIN OPERATOR, WHICH THEN IMMEDIATELY FORGETS THE REQUEST; A SUBSCRIBE REQUEST IS A PERSISTENT JOIN WITH PAST AND FUTURE EVENTS ON THE OTHER SIDE OF THE JOIN."

Why you might want a log of reads:

"Recording a log of read events potentially has benefits with regard to TRACKING CAUSAL DEPENDENCIES AND DATA PROVENANCE. The log would allow you to RECONSTRUCT WHAT THE USER SAW BEFORE THEY MADE A PARTICULAR DECISION. For example, in an online shop, THE PREDICTED SHIPPING DATE AND THE INVENTORY STATUS SHOWN TO A CUSTOMER LIKELY AFFECT WHETHER THEY CHOOSE TO BUY. To analyze this connection, YOU NEED TO RECORD THE RESULT OF THE USER'S QUERY." (This is §1.4's option 2, made concrete.)

The cost: "additional STORAGE AND I/O costs. Optimizing such systems is STILL AN OPEN RESEARCH PROBLEM — but IF YOU ALREADY LOG READ REQUESTS FOR OPERATIONAL PURPOSES, IT'S NOT A BIG CHANGE TO MAKE THAT LOG THE SOURCE OF THE REQUESTS INSTEAD."

Multishard data processing: "For single-shard queries this is perhaps OVERKILL. However, it opens the possibility of DISTRIBUTED EXECUTION OF COMPLEX QUERIES that need to combine data from several shards, TAKING ADVANTAGE OF THE INFRASTRUCTURE FOR MESSAGE ROUTING, SHARDING, AND JOINING THAT IS ALREADY PROVIDED." Examples: Storm's distributed RPC computing "the number of people who have seen a URL — the UNION OF THE FOLLOWER SETS of everyone who posted it"; and fraud prevention, where assessing a purchase requires "reputation scores of the user's IP address, email address, billing address, shipping address — each database itself sharded, so collecting the scores requires A SEQUENCE OF JOINS WITH DIFFERENTLY SHARDED DATASETS."


On this page