5.7 Modes of Dataflow
Clients connect to servers; the API exposed by the server is a service.
Compatibility is a relationship between one process that encodes data and another that decodes it. Who encodes, who decodes? Four modes:
| Mode | Shape |
|---|---|
| 1 · Databases | Writer encodes → storage → reader decodes. The reader may be a later version of the same process — “storing something in the database is sending a message to your future self.” |
| 2 · Services (REST / RPC) | Client encodes the request; server decodes it and encodes a response; client decodes that. |
| 3 · Workflows / durable execution | Orchestrator drives tasks, with all RPCs and state changes logged. |
| 4 · Event-driven (brokers, actors) | Sender encodes → broker → recipient decodes, asynchronously. |
7.1 Dataflow through databases
Backward compatibility is clearly necessary — otherwise your future self can't decode what you previously wrote.
But forward compatibility is ALSO often required, because several processes access the database at once: different applications/services, or multiple instances of the same service, and during a rolling upgrade some instances run new code and some run old. So a value written by newer code may be read by older code still running.
"Data outlives code."
A database allows any value to be updated at any time — you may have values written five milliseconds ago and others written five years ago. When you deploy a new version of a server-side application, you may entirely replace the old version within a few minutes. The same is NOT true of database contents: the five-year-old data is still there, in the original encoding, unless you have explicitly rewritten it.
Migration is possible but expensive on a large dataset, so most databases defer it — asynchronously, best-effort:
| System | How it defers |
|---|---|
| LSM-tree storage engines | Rewrite data using the latest format during COMPACTION (Ch 4 §2.3) |
| Most relational databases | Simple schema changes (add a column with a null default) without rewriting existing data. When an old row is read, the database fills in nulls for columns missing from the encoded data on disk |
Schema evolution thus allows the entire database to APPEAR as if it were encoded with a single schema, even though the underlying storage contains records encoded with various historical schema versions.
More complex changes still require rewriting, often at the application level — changing a single-valued attribute to multivalued, or moving data into a separate table. Maintaining forward and backward compatibility across such migrations remains a research problem.
Archival storage. A periodic snapshot (for backup or warehouse loading) is typically encoded using the LATEST schema, even if the source contained a mixture of eras — since you're copying anyway, you might as well encode consistently. Because the dump is written in one go and thereafter immutable, Avro object container files are a good fit, and it's a good opportunity to write an analytics-friendly column-oriented format like Parquet.
7.2 Dataflow through services: REST and RPC
Clients connect to servers; the API exposed by the server is a service.
The web works this way — browsers GET HTML/CSS/JS/images and POST data; the API is a standardized set of protocols and formats (HTTP, URLs, SSL/TLS, HTML). Because browsers, servers, and site authors mostly agree on these standards, you can use any browser to access any website (at least in theory!).
But browsers aren't the only client — native mobile/desktop apps and client-side JavaScript also make HTTP requests, where the response is not HTML for a human but data convenient for further processing (most often JSON), and the API on top of HTTP is application-specific.
Services vs databases:
Both allow clients to submit and query data. But databases allow ARBITRARY queries using query languages; services expose an APPLICATION-SPECIFIC API allowing only inputs and outputs predetermined by the business logic. This restriction provides encapsulation — services can impose fine-grained restrictions on what clients can and cannot do.
Why compatibility matters here: the design goal of microservices is independently deployable and evolvable services, each owned by one team able to release frequently without coordinating with other teams. So we must expect old and new versions of servers and clients running at the same time.
As long as APIs remain compatible, teams are free to modify their systems any way they'd like — this property makes internal migrations of data, services, or even entire systems much easier.
Three contexts for web services:
- A client app on a user's device → a service over HTTP, typically over the public internet
- One service → another service in the same organization, often within the same private network
- One service → a service owned by a DIFFERENT organization, usually via the internet — data exchange between organizations' backends, e.g. credit card processing, or OAuth for shared access to user data
REST is the most popular design philosophy — building on HTTP's principles: simple data formats, URLs for identifying resources, and HTTP features for cache control, authentication, and content type negotiation.
IDLs for services: clients must know which endpoint to query and what data format to send/expect. The two most popular service IDLs: OpenAPI (Swagger) for JSON web services, and Protocol Buffers for gRPC.
openapi: 3.0.0
info: { title: "Ping, Pong", version: 1.0.0 }
servers: [ { url: http://localhost:8080 } ]
paths:
/ping:
get:
summary: Given a ping, returns a pong message
responses:
'200':
description: A pong
content:
application/json:
schema:
type: object
properties:
message: { type: string, example: "Pong!" }from fastapi import FastAPI
from pydantic import BaseModel
app = FastAPI(title="Ping, Pong", version="1.0.0")
class PongResponse(BaseModel):
message: str = "Pong!"
@app.get("/ping", response_model=PongResponse,
summary="Given a ping, returns a pong message")
async def ping():
return PongResponse()Two directions of coupling between definition and code:
Service frameworks (Spring Boot, FastAPI, gRPC) let developers focus on business logic while the framework handles routing, metrics, caching, authentication.
7.3 The problems with RPC — the classic critique
A long line of predecessors, all with serious problems:
| Technology | Problem |
|---|---|
| EJB, Java RMI | Limited to Java |
| DCOM | Limited to Microsoft platforms |
| CORBA | Excessively complex; does not provide backward or forward compatibility |
| SOAP / WS-* | Aims at cross-vendor interoperability but plagued by complexity and compatibility problems |
All are based on RPC (introduced in the 1970s), which tries to make a remote network service call look the same as a local function call — "location transparency."
Although this seems convenient at first, the approach is FUNDAMENTALLY FLAWED. A network request is very different from a local function call:
| # | Local function call | Network request |
|---|---|---|
| 1 | Predictable — succeeds or fails based on parameters under your control | Unpredictable, for reasons entirely outside your control: request/response lost, remote machine slow or unavailable. Network problems are common, so applications must anticipate them (e.g. retries) |
| 2 | Returns a result, throws an exception, or never returns (infinite loop / crash) | A third outcome: returns WITHOUT A RESULT because of a TIMEOUT. You simply don't know what happened — no way of knowing whether the request got through |
| 3 | No such problem | Retrying a "failed" request may perform the action MULTIPLE TIMES, if the request got through and only the response was lost — unless you build deduplication (idempotence) into the protocol |
| 4 | Takes about the same time each call | Much slower, and latency is wildly variable — <1 ms at good times, many seconds when the network is congested or the service overloaded, for exactly the same work |
| 5 | Can efficiently pass references (pointers) to local memory | All parameters must be encoded into bytes. Fine for immutable primitives; quickly problematic with larger amounts of data and mutable objects |
| 6 | Single process, single language — no type translation | Client and service may be in different languages, so the framework must translate datatypes. This can end up ugly — recall JavaScript's 2⁵³ problem |
There's no point trying to make a remote service look too much like a local object, because it's a FUNDAMENTALLY DIFFERENT THING. Part of the appeal of REST is that it treats state transfer over a network as a process DISTINCT from a function call.
7.4 Load balancing, service discovery, and service meshes
The problem: a client must know the address of the service it's connecting to — service discovery. The simplest approach (hardcode IP and port) works until the server goes offline, moves, or becomes overloaded — then the client must be manually reconfigured.
Load balancing = spreading requests across the multiple instances usually running for availability and scalability.
| Solution | How it works | Trade-off |
|---|---|---|
| Hardware load balancers | Specialized equipment in datacenters. Clients connect to a single host:port; connections routed to one of the servers. Detects network failures to a downstream server and shifts traffic | Requires an appliance |
| Software load balancers (NGINX, HAProxy) | Same behavior, as applications on a standard machine | — |
| DNS | Multiple IP addresses per domain name; the client's network layer picks which to use | DNS is designed to propagate changes over LONGER periods and to CACHE entries. If servers start/stop/move frequently, clients might see stale IPs with no server running on them |
| Service discovery systems (etcd, ZooKeeper) | A centralized registry. A new instance registers itself (host, port) with metadata: shard ownership, datacenter location, and periodically sends a HEARTBEAT to signal it's still available. The client queries the registry for endpoints, then connects directly | Supports a much more dynamic environment than DNS, and the extra metadata enables smarter load-balancing decisions by clients |
| Service meshes (Istio, Linkerd) | Combines software load balancing and service discovery. Deployed as an in-process client library or a "sidecar" container on BOTH client and server. The client app connects to its own local load balancer, which connects to the server's load balancer, which routes to the local server process | Complicated — but: because everything is routed through local connections, connection encryption is handled entirely at the load-balancer level, shielding clients and servers from SSL certificates and TLS. Also sophisticated observability: which services call each other in real time, failure detection, traffic load |
Choosing: very dynamic environments with Kubernetes → often a service mesh. Specialized infrastructure (databases, messaging) might require purpose-built load balancers. Simpler deployments are best served with software load balancers.
7.5 Encoding and evolution for RPC
A simplifying assumption vs databases: it is reasonable to assume ALL SERVERS WILL BE UPDATED FIRST and all clients second. Thus you need BACKWARD compatibility only on REQUESTS, and FORWARD compatibility on RESPONSES.
Compatibility properties are inherited from the encoding:
- gRPC (Protocol Buffers) and Avro RPC evolve per their format's rules
- RESTful APIs most commonly use JSON responses and JSON/URI-encoded/form-encoded requests. Adding optional request parameters and adding new fields to response objects are usually considered compatible changes
The hard part is organizational:
RPC is often used across ORGANIZATIONAL BOUNDARIES, so the service provider often has NO CONTROL over its clients and cannot force them to upgrade. Compatibility needs to be maintained for a long time, PERHAPS INDEFINITELY. If a breaking change is required, the provider frequently ends up maintaining multiple versions of the API side by side.
There is no agreement on how API versioning should work. Common approaches for REST:
- Version number in the URL (
/v2/users) - Version in the HTTP
Acceptheader - For services using API keys: store the client's requested version on the server, updatable through a separate administrative interface
7.6 Durable execution and workflows
The setup: a payment processor charges a credit card and deposits funds into a bank account. Services: fraud detection, credit card integration, bank integration.
A workflow engine is an orchestrator plus an executor. The orchestrator decides when and on which machine each task runs, what to do if a task fails — for instance the machine crashing mid-task — and how many tasks may run in parallel. The executor actually runs them. Triggers are a time-based schedule, an external service, or a human.
Workflow definitions may be written in a general-purpose language, a DSL, or a markup language such as BPEL; graphical notations like BPMN exist.
Three families of workflow engine:
| Family | Examples | Purpose |
|---|---|---|
| Data orchestration | Airflow, Dagster, Prefect | Integrate with data systems, orchestrate ETL tasks |
| Graphical / business | Camunda, Orkes | BPMN notation so non-engineers can define and execute workflows |
| Durable execution | Temporal, Restate | Exactly-once semantics for workflows |
Why durable execution exists:
We want to process each payment exactly once. A failure mid-workflow could result in a credit card charge but no corresponding bank deposit. In a service-based architecture you can't simply wrap the two tasks in a database transaction. Moreover, you might be interacting with third-party payment gateways you have limited control over.
How it works:
If a task fails, the framework re-executes it — but SKIPS any RPC calls or state changes that succeeded before the failure. It will PRETEND to make the call, but instead RETURN THE RESULTS FROM THE PREVIOUS CALL.
This is possible because durable execution frameworks LOG ALL RPCs AND STATE CHANGES TO DURABLE STORAGE, like a write-ahead log.
@workflow.defn
class PaymentWorkflow:
@workflow.run
async def run(self, payment: PaymentRequest) -> PaymentResult:
is_fraud = await workflow.execute_activity(
check_fraud, payment, start_to_close_timeout=timedelta(seconds=15))
if is_fraud:
return PaymentResultFraudulent
credit_card_response = await workflow.execute_activity(
debit_credit_card, payment, start_to_close_timeout=timedelta(seconds=15))
# ...Three real challenges — all consequences of "replay must be deterministic":
- External services must STILL provide an idempotent API. Developers must remember to use unique IDs for these APIs to prevent duplicate execution.
- Code changes are brittle. Because the framework logs each RPC call in order, it expects subsequent executions to make the same RPC calls in the same order. You might introduce undefined behavior simply by REORDERING FUNCTION CALLS.
Instead of modifying an existing workflow's code, it is safer to deploy a NEW VERSION separately, so re-executions of existing invocations continue using the old version and only new invocations use the new code.
- Nondeterministic code is problematic — random number generators, system clocks. Frameworks provide their own deterministic implementations, but you have to remember to use them. Some offer static analysis (Temporal's Workflow Check) to detect introduced nondeterminism.
Making code deterministic is a powerful idea but tricky to do robustly. (Returns in Ch 9.)
Note how this is the same constraint as event-sourcing projections in Ch 3 §3.3 — replayability demands determinism, everywhere it appears.
7.7 Event-driven architectures
A request is called an event or message. Two structural differences from RPC:
- The sender usually does not wait for the recipient to process the event
- Events are typically not sent via a direct network connection but via an intermediary — a MESSAGE BROKER (event broker, message queue, message-oriented middleware) which stores the message temporarily
Five advantages over direct RPC:
| Advantage | Why |
|---|---|
| Buffering | Acts as a buffer if the recipient is unavailable or overloaded — improving system reliability |
| Redelivery | Automatically redelivers messages to a process that has crashed, preventing message loss |
| No service discovery needed | Senders don't need to connect directly to the recipient's IP address |
| Fan-out | The same message can be sent to several recipients |
| Logical decoupling | The sender just publishes and doesn't care who consumes |
Communication is asynchronous — the sender sends and forgets. You can implement a synchronous RPC-like model by having the sender wait for a response on a separate channel.
The landscape: formerly commercial enterprise software (TIBCO, IBM WebSphere, webMethods), then open source (RabbitMQ, ActiveMQ, HornetQ, NATS, Redpanda, Apache Kafka), more recently cloud services (Amazon Kinesis, Azure Service Bus, Google Cloud Pub/Sub).
Two distribution patterns:
Encoding: brokers typically don't enforce any data model — a message is just bytes with metadata, so any encoding works. Common approach: Protocol Buffers, Avro, or JSON, with a SCHEMA REGISTRY deployed alongside the broker to store valid schema versions and check compatibility. AsyncAPI is the messaging equivalent of OpenAPI.
Durability varies. Many write messages to disk so they survive a broker crash/restart. Unlike databases, many brokers automatically DELETE messages after consumption. Some can be configured to store messages indefinitely — which you'd require for event sourcing (Ch 3 §3).
⚠️ If a consumer REPUBLISHES messages to another topic, be careful to PRESERVE UNKNOWN FIELDS — otherwise you reproduce the data-loss problem of §0 in the message pipeline.
7.8 Distributed actor frameworks
The actor model is a concurrency model for a single process. Rather than dealing directly with threads (and race conditions, locking, deadlock), logic is encapsulated in ACTORS. Each actor typically represents one client or entity, may have local state not shared with any other actor, and communicates by sending and receiving asynchronous messages.
Message delivery is NOT GUARANTEED — in certain error scenarios, messages will be lost. Since each actor processes only one message at a time, it needn't worry about threads, and each actor can be scheduled independently.
Distributed actor frameworks — Akka, Orleans, Erlang/OTP — use this model to scale across multiple nodes. The same message-passing mechanism is used whether sender and recipient are on the same node or different nodes; across nodes the message is transparently encoded, sent, and decoded.
Location transparency works BETTER in the actor model than in RPC, because the actor model ALREADY ASSUMES messages may be lost, even within a single process. Although network latency is higher than in-process, there is less of a fundamental mismatch between local and remote communication.
(This is the sharp counterpoint to §7.3: location transparency isn't inherently wrong — it's wrong when the local model promises more reliability than the network can deliver. Actors work because they promise less.)
But: a distributed actor framework is essentially a message broker + the actor programming model in one framework — and if you want rolling upgrades, you still have to worry about forward and backward compatibility, since messages may flow from a new-version node to an old-version node and vice versa.
5.6 The merits of schemas
Protocol Buffers' and Avro's schema languages are much simpler than XML Schema or JSON Schema, which support detailed validation rules ("must match this regex", "must be between 0…
5.8 Technology deep dives
gRPC adds: HTTP/2 transport (multiplexed streams, header compression), four call types (unary, server-streaming, client-streaming, bidirectional), deadlines propagated through the…