Learn Labs
5. Encoding and Evolution

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:

ModeShape
1 · DatabasesWriter 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 executionOrchestrator 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:

SystemHow it defers
LSM-tree storage enginesRewrite data using the latest format during COMPACTION (Ch 4 §2.3)
Most relational databasesSimple 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:

  1. A client app on a user's device → a service over HTTP, typically over the public internet
  2. One service → another service in the same organization, often within the same private network
  3. 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:

code-firstFastAPI — write the server firstschema-firstgRPC — write the definition firstOpenAPIgenerated.protoscaffolds the servergenerated clients & SDKsBoth also give documentation generation, schema-change compatibility verification, and a GUI for testing.
Figure 5.7.2Two 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:

TechnologyProblem
EJB, Java RMILimited to Java
DCOMLimited to Microsoft platforms
CORBAExcessively 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 callNetwork request
1Predictable — succeeds or fails based on parameters under your controlUnpredictable, 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)
2Returns 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
3No such problemRetrying 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
4Takes about the same time each callMuch 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
5Can efficiently pass references (pointers) to local memoryAll parameters must be encoded into bytes. Fine for immutable primitives; quickly problematic with larger amounts of data and mutable objects
6Single process, single language — no type translationClient 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.

SolutionHow it worksTrade-off
Hardware load balancersSpecialized 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 trafficRequires an appliance
Software load balancers (NGINX, HAProxy)Same behavior, as applications on a standard machine—
DNSMultiple IP addresses per domain name; the client's network layer picks which to useDNS 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 directlySupports 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 processComplicated — 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
local, plaintextmTLSlocalappclient podsidecarsidecarserviceserver podThe app never sees TLS, and the mesh sees every call.Retries, timeouts and circuit breaking are configured centrally rather than in each app.
Figure 5.7.37.4 Load balancing, service discovery, and service meshes

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 Accept header
  • 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.

fraud?not fraudtriggercheck_fraudreturn Fraudulentdebit_credit_carddeposit_to_bankEach box is a task — an activity in Temporal, or a durable function. A workflow is a graph of tasks.

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.

Figure 5.7.47.6 Durable execution and workflows

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:

FamilyExamplesPurpose
Data orchestrationAirflow, Dagster, PrefectIntegrate with data systems, orchestrate ETL tasks
Graphical / businessCamunda, OrkesBPMN notation so non-engineers can define and execute workflows
Durable executionTemporal, RestateExactly-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":

  1. External services must STILL provide an idempotent API. Developers must remember to use unique IDs for these APIs to prevent duplicate execution.
  2. 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.

  3. 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:

  1. The sender usually does not wait for the recipient to process the event
  2. 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:

AdvantageWhy
BufferingActs as a buffer if the recipient is unavailable or overloaded — improving system reliability
RedeliveryAutomatically redelivers messages to a process that has crashed, preventing message loss
No service discovery neededSenders don't need to connect directly to the recipient's IP address
Fan-outThe same message can be sent to several recipients
Logical decouplingThe 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:

Queue — one consumer receives it
one of themproducernamed queueconsumer 1consumer 2
Topic — every subscriber receives it
all of themproducernamed topicsubscriber 1subscriber 2subscriber 3
Figure 5.7.5Two 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.


On this page