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…
8.1 Protocol Buffers / gRPC
Problem it solves. Compact, fast, cross-language encoding with explicitly defined forward/backward compatibility, plus generated type-safe client and server code.
Why wasn't JSON enough? ~2.5× larger on the wire, no schema by default (so no compatibility checking and no generated types), ambiguous number semantics, no binary strings, and slower to parse.
Why wasn't Avro enough? Protobuf's tag-based encoding means the reader doesn't need the writer's schema at all — each message is self-delimiting. That makes it far better suited to request/response over a network where negotiating a writer schema per call would be a burden.
How it works internally. Tag+wire-type packed into a key byte; varint integers; length-delimited strings/bytes/submessages; repeated = repeated occurrences of the same tag. Unknown fields are skippable because the wire type tells the parser the length — and proto3 (since 3.5) retains unknown fields rather than dropping them, which is what makes forward compatibility safe.
gRPC adds: HTTP/2 transport (multiplexed streams, header compression), four call types (unary, server-streaming, client-streaming, bidirectional), deadlines propagated through the call chain, and per-call metadata.
Deployment. .proto files in a shared repo or a schema registry, with generated code published as versioned artifacts per language. buf (or protolock) in CI to enforce compatibility rules — this is the single highest-value practice, and it's the mechanical equivalent of the rules in §4.2.
Monitoring. Per-method latency and status-code distribution (gRPC codes, not HTTP); deadline-exceeded rate (the clearest signal of a too-tight timeout or a slow dependency); message sizes (p99 — a 4 MB default max is a real limit people hit); unknown-field rate if you can instrument it, as a leading indicator of version skew; connection churn.
Scaling. HTTP/2 multiplexes over one long-lived TCP connection, which defeats connection-level L4 load balancers — you need L7/gRPC-aware balancing or client-side balancing with a service mesh. This surprises nearly everyone the first time.
Backup. N/A for the wire format; the .proto files and the tag-reservation history are the artifact to preserve. Losing the record of which tags were once used is how you get silent corruption years later.
What actually breaks in production.
- Reusing a retired field tag. The nightmare case: old persisted data has tag 7 as a string; new code declares tag 7 as an int64. Decoding produces garbage or crashes, and there's no error at deploy time.
reserved 7;exists precisely to prevent this — use it every single time you delete a field. int32→int64widening, then old code truncating a large value it reads back. Silent data corruption.- proto2
requiredfields — the field that can never be removed, on either side, forever. This is why proto3 removed the concept. - Default-value ambiguity in proto3: an unset field and a field explicitly set to
0/""/falseare indistinguishable on the wire. "Set the discount to 0" and "don't change the discount" become the same message. Requires wrapper types oroptional(reintroduced in proto3 3.15). - A gRPC deadline not propagated, so the server keeps working on a request whose caller gave up.
- Load balancing on L4 with HTTP/2 → all traffic pinned to whichever backend won the connection race.
8.2 Avro + a schema registry (Confluent Schema Registry, Apicurio)
Problem it solves. Encode millions of records compactly where the schema is known per-file or per-topic rather than per-message, and support dynamically generated schemas from a database or ETL job.
Why wasn't Protobuf enough? Tag numbers must be assigned by hand, which makes generating a schema from a relational table (or a changing CSV, or an inferred JSON shape) an administrative chore that's dangerous to automate — the generator must never reassign a previously used tag. Avro identifies fields by name, so regeneration is safe.
Why wasn't JSON-per-record enough? Field names repeated in every one of a billion records; no compatibility enforcement.
How it works internally. §5: values concatenated with no field identifiers; schema resolution matches writer↔reader by name; defaults fill absent fields; unions for nullability.
The registry's role: each Kafka message is prefixed with a magic byte + 4-byte schema ID; the consumer fetches the writer's schema by ID (and caches it), then resolves against its own reader schema.
The registry enforces a compatibility level per subject: BACKWARD (a new reader against old data — the default), FORWARD (an old reader against new data), FULL, the *_TRANSITIVE variants, or NONE.
Deployment. Registry as an HA service (Confluent's is itself backed by a Kafka topic); compatibility level set per subject; schema evolution gated in CI by a test-compatibility call before merge.
Monitoring. Registry availability and p99 latency (it is on the critical path of every consumer's cold start); cache hit rate in clients; schema version count per subject (rapid growth means something is generating schemas per-deploy, usually a bug); compatibility-check rejection rate.
Scaling. Clients cache schemas indefinitely, so steady-state load is near zero — the load is a thundering herd at consumer restart.
Backup. Back up the _schemas topic / registry store. Losing the registry makes every historical message on every topic undecodable. This is a genuine single point of catastrophic failure, and it's routinely un-backed-up.
What actually breaks.
- Adding a field without a default. Breaks backward compatibility; the registry rejects it if configured correctly, and silently poisons every consumer if it isn't.
- Registry down at consumer start → the whole consumer fleet can't decode anything. Cache warm-up and a fallback are worth designing.
NONEcompatibility level set "temporarily" during an incident and never restored.- Renaming a field — backward compatible via aliases, not forward compatible. Old readers break. Almost nobody remembers this asymmetry.
- A default of
nullnot first in the union, so it's rejected — or worse, a union reordered so branch indexes shift and old data decodes as the wrong type. - Two producers writing the same topic with divergent schema lineages.
8.3 REST/OpenAPI services
Problem it solves. Cross-organization, long-lived, human-inspectable APIs over ubiquitous HTTP infrastructure.
Why wasn't RPC enough? §7.3's six reasons, plus: REST treats state transfer over a network as a process distinct from a function call, which sets honest expectations. And HTTP brings caching, content negotiation, auth, and intermediaries for free.
Why wasn't gRPC enough for public APIs? Browsers can't speak raw gRPC without a proxy; binary payloads are undebuggable with curl; and external partners overwhelmingly expect JSON.
How it works internally. URLs identify resources; HTTP verbs and status codes carry semantics; Cache-Control/ETag for caching; Accept for content negotiation. OpenAPI describes it; tooling generates docs, SDKs, mocks, and compatibility checks.
Deployment. Gateway (auth, rate limiting, quotas) → service. Versioning strategy chosen once, deliberately — URL path, Accept header, or server-stored per-API-key version.
Monitoring. Per-endpoint p99 and error rate; usage per API version and per client — the only data that tells you when a deprecated version can actually be removed; deprecation-header adoption; payload sizes; Sunset/Deprecation header delivery.
Scaling. Stateless services behind L7 balancers; CDN for cacheable GETs; per-client rate limits.
What actually breaks.
- A "compatible" change that isn't. Tightening a validation rule, narrowing an enum, changing a nullable field to non-nullable, or changing a number to a string — all break clients while looking innocuous in a diff.
- The forgotten v1 you can never turn off, because three enterprise customers still call it and nobody instrumented per-version usage.
- Clients that break on unknown fields. Adding a response field is supposed to be safe — it isn't, if the client uses a strict deserializer. This is the §0 trap on the API surface, and it's why "additive changes are safe" is only true if clients are lenient.
- 2⁵³ in JSON — a bigint ID silently mangled by a JavaScript client.
- Breaking change shipped without a deprecation window, because the provider had no control over the clients and no way to force an upgrade.
8.4 Message brokers (Kafka, RabbitMQ, SQS, Pub/Sub)
Problem it solves. Decouple sender from recipient in time, in location, and in fan-out — buffering when the recipient is down, redelivering on crash, and letting many consumers read the same message.
Why wasn't direct RPC enough? If the recipient is down, the request fails; the sender must know the recipient's address; one message can't reach five consumers; and a slow consumer applies backpressure directly to the user-facing request path.
How it works internally. Ch 12 covers this properly. Briefly: RabbitMQ-style brokers track per-message acknowledgment and delete on ack; Kafka-style log brokers append to a partitioned, ordered, retained log and track a consumer offset — which is what makes replay and event sourcing possible.
Deployment. Broker cluster with replication; a schema registry alongside (§8.2); dead-letter queues; consumer groups sized to partition count.
Monitoring. Consumer lag per partition — the metric; broker disk usage vs retention; rebalance frequency (frequent rebalances mean consumers are timing out); DLQ depth and age; redelivery rate; end-to-end latency from produce to process.
Scaling. Partition count bounds consumer parallelism and is painful to increase (it changes key→partition mapping). Choose it with room to grow.
Backup. Topics with infinite retention are a system of record and need real backups (mirroring to a second cluster, or tiered storage to object storage). Topics with 7-day retention are a buffer, not a store — know which each of yours is.
What actually breaks.
- Unknown-field loss on republish — the §7.7 warning, made concrete: a consumer decodes into an old model, transforms, and republishes, silently dropping fields added by a newer producer.
- Poison messages looping forever without a DLQ, blocking the partition behind them.
- Retention expiry before a stalled consumer catches up → permanent data loss with no error, just a jump in the offset.
- Ordering assumptions across partitions — order is per-partition only.
- A consumer that isn't idempotent, meeting at-least-once delivery.
- Rebalance storms when processing time exceeds
max.poll.interval.ms, so the group never stabilizes and throughput goes to zero.
8.5 Durable execution engines (Temporal, Restate) — and workflow orchestrators (Airflow, Dagster)
Problem it solves. Exactly-once multi-step business processes spanning services and third parties, where you cannot wrap the steps in a database transaction.
Why weren't distributed transactions enough? Ch 8 shows 2PC is rarely usable across microservices and third-party APIs, and it runs counter to service independence.
Why wasn't a retry loop enough? Retrying step 3 must not re-charge the card in step 2. You need to know what already succeeded, durably, across process crashes.
How it works internally. Every RPC and state change is appended to a durable event history (a write-ahead log). On replay, the framework re-executes the workflow code from the top, but intercepts each activity call: if the history shows it already completed, it returns the recorded result instead of calling out again. The workflow function is therefore a deterministic reduction over its own history — which is precisely why nondeterminism is fatal.
The replayed tasks are read from history rather than re-executed, so the card is not charged twice — only the step that never completed actually runs.
Deployment. Worker fleet polling task queues + a service cluster backed by a database (Cassandra/MySQL/Postgres for Temporal). Workflow code versioned, never edited in place — deploy a new version so in-flight executions keep the old one.
Monitoring. Workflow-execution latency and failure rate; activity retry counts; stuck/blocked workflows (the ones that silently accumulate); task-queue backlog and worker poll-success rate; nondeterminism-error count — a nonzero value means a deploy broke replay for in-flight executions; workflow history length (there are hard limits, and a long-running loop hits them).
Scaling. Workers scale horizontally and statelessly; the state store is the bottleneck. Long-running workflows must use continueAsNew to bound history size.
Backup. The event history is the system of record — back it up like a database, because losing it strands every in-flight business process in an unknowable state.
What actually breaks.
- Nondeterminism introduced by an innocent edit. Reorder two activity calls, or add an
ifbefore an existing one, and replay of in-flight workflows diverges from history →NonDeterministicError. This is the #1 Temporal production incident and it's a deploy-time self-inflicted wound. datetime.now()/random()/uuid4()inside workflow code, producing different values on replay.- A non-idempotent third party. The framework guarantees it won't re-issue a recorded call — it cannot guarantee the gateway didn't process a call whose response was lost. You still need idempotency keys.
- Unbounded history on a long-lived workflow, hitting the event-count limit.
- Treating Airflow like Temporal. Airflow orchestrates idempotent batch tasks on a schedule; it does not give you exactly-once semantics for stateful business transactions. Conflating the two is a common and expensive design error.