Distributed Communication Patterns#
The moment your code makes a call that leaves the process, everything you knew about function calls stops being true. This page is the set of patterns that exist because of that one fact — the ones you will reach for every week. Mental model first, code second. Press Play on any animation to watch the idea move.
Every pattern on this page answers one question: how do you build something dependable out of components that can be slow, can be down, and can never tell you which? Timeouts, retries, idempotency keys, circuit breakers, queues, outboxes and sagas are not seven unrelated tools. They are seven answers to that single question, and each one buys reliability with a currency you must be willing to spend: latency, duplication, staleness, or complexity.
The Third Outcome#
Theory#
Before any pattern makes sense you need one idea. A local function call has two outcomes: it returns, or it throws. A remote call has three: it returns, it fails, or you never find out. That third outcome does not exist inside a process, it cannot be eliminated, and every single pattern in this course is machinery for surviving it.
You post a cheque and hear nothing for a week. Did it get lost on the way there? Did it arrive and get lost in their office? Did they cash it and the receipt got lost on the way back? You cannot tell. And the decision you have to make — send another cheque or not — is exactly the decision your HTTP client makes on every timeout, dozens of times a second.
Notice what went wrong there. Nobody wrote a bug. The server did its job, the network dropped one packet out of billions, and the client did the textbook thing. The double charge is an emergent property of a correct client talking to a correct server over an imperfect wire.
- A timeout is not an errorAn HTTP
500tells you the work failed. A timeout tells you nothing — the request may never have arrived, may have failed, or may have completely succeeded. Code that treats them identically is code that double-charges people. - Slow is worse than downA dead service gives you a fast connection-refused, which is easy to handle. A service answering in 30 seconds holds your threads, your connections and your memory hostage, and takes your service down with it. Grey failure is the hard case.
What the network does to your numbers#
Those numbers are six orders of magnitude apart, and the spread is why a refactor that “just” moves a method into another service can make a page ten times slower. But latency is only half of it. The other half is arithmetic:
Availability multiplies and latency adds. Every synchronous dependency you put in a request path makes both worse, permanently. This is the real cost of a microservice boundary, and it is why the interesting question is never “should these be separate services?” but “does this call have to be in the request path at all?”
Two rules of thumb worth memorising before anything else. First: the network is not reliable, not fast, not secure, and its topology changes — these are four of the famous eight fallacies of distributed computing, and every outage you will ever debug is someone having assumed one of them. Second: you cannot make a distributed call as safe as a local one; you can only make its failure cheap.
Calling Another Service#
Choosing whether to wait for an answer, picking an API style, and finding a healthy instance to send the call to.
Synchronous or Asynchronous#
Theory#
This is the first and biggest fork in the road, and it is not really a technology choice. It is a choice about who waits. In a synchronous call the caller holds a connection open until the work is finished. In an asynchronous one the caller hands the work to a broker, gets an acknowledgement, and leaves.
Synchronous is the barista taking your order, making it, and handing it over before serving the next person. Simple, and the queue is only as fast as the slowest drink. Asynchronous is the cashier writing your name on a cup and putting it on the rail: the till keeps moving, three baristas work the rail in parallel, and you find out your drink is ready later. The second design serves far more people — and introduces a state the first one never had: paid for, not yet made.
- Choose synchronous when the caller needs the answerA price, an authorisation, a seat availability, a login. If the next line of the caller's code cannot run without the result, a queue only adds machinery around a wait you still have to do.
- Choose asynchronous when the caller needs an acknowledgementSend the email, update the search index, recalculate the recommendations, emit the audit record. Nobody is watching, so do not make a customer wait for it — and do not let its outage become yours.
- The sync trap“It is only one more call” is how a 200 ms endpoint becomes a 2 s endpoint over eighteen months, one reasonable pull request at a time.
- The async trapEventual consistency is a product problem, not just an engineering one. Somebody has to decide what the screen says while the order exists but is not yet paid for.
The answer interviewers are listening for is not “use a queue”. It is the sentence “asynchronous buys availability and pays for it in consistency, so the real question is what the user sees in between”. Naming the intermediate state is what separates someone who has read about queues from someone who has shipped one.
Questions#
A user uploads a profile photo. What happens synchronously, and what happens asynchronously?
Split by what the user must see before the request can return. They must see that the upload succeeded and, ideally, their new photo — so storing the original and returning its URL is synchronous. Everything derived from it is not: generating five thumbnail sizes, running content moderation, stripping EXIF, invalidating CDN paths, and notifying their followers. Each of those is slow, each can fail independently, and none of them should be able to make the upload button return an error.
The interesting part is the in-between state. Thumbnails are not ready yet, so the answer is not “do it later and hope” — it is to design the intermediate state: serve the original scaled by the browser until the thumbnail exists, or show a shimmer. That decision is the actual work of going asynchronous.
@app.post("/profile/photo", status_code=202)
async def upload_photo(user_id: str, file: UploadFile):
# Synchronous: only what the user must see before the button stops spinning.
key = f"users/{user_id}/original/{uuid.uuid4()}.jpg"
await s3.put_object(Bucket=BUCKET, Key=key, Body=await file.read())
await db.execute(
"update users set photo_key = $1, photo_state = 'processing' where id = $2",
key, user_id,
)
# Asynchronous: one event, four independent consumers, none in the request path.
await outbox.append(
topic="photo.uploaded",
key=user_id,
payload={"userId": user_id, "objectKey": key, "uploadedAt": now_iso()},
)
return {"state": "processing", "url": cdn_url(key)}
app.post("/profile/photo", async (req, res) => {
// Synchronous: only what the user must see before the button stops spinning.
const key = `users/${req.userId}/original/${randomUUID()}.jpg`;
await s3.send(new PutObjectCommand({ Bucket: BUCKET, Key: key, Body: req.file.buffer }));
await db.query(
"update users set photo_key = $1, photo_state = 'processing' where id = $2",
[key, req.userId],
);
// Asynchronous: one event, four independent consumers, none in the request path.
await outbox.append({
topic: "photo.uploaded",
key: req.userId,
payload: { userId: req.userId, objectKey: key, uploadedAt: new Date().toISOString() },
});
res.status(202).json({ state: "processing", url: cdnUrl(key) });
});
REST, gRPC & GraphQL#
Theory#
Once you have decided to make a synchronous call, you pick a style. All three of these are request/response over HTTP; they differ in who decides the shape of the response and how much the two ends must agree in advance.
The animation shows gRPC's four call shapes, because those are the ones people have not met. The encoding difference matters just as much, and it is where most of the performance argument lives:
syntax = "proto3";
package orders.v1;
service Orders {
rpc GetOrder(GetOrderRequest) returns (Order);
rpc WatchOrder(GetOrderRequest) returns (stream OrderEvent);
}
message GetOrderRequest {
string order_id = 1;
}
message Order {
string order_id = 1;
string status = 2; // never renumber a tag
int64 total_minor = 3;
string currency = 4;
reserved 5; // 5 was 'legacy_total', do not reuse
}
GET /v1/orders/A-4471
{
"id": "A-4471",
"status": "paid",
"amount": 8210,
"currency": "GBP",
"_links": { "cancel": "/v1/orders/A-4471/cancellation" }
}
# GraphQL: the CLIENT chooses the fields, in one round trip
query {
order(id: "A-4471") {
status
amount
customer { name tier } # would be a second REST call
items { sku qty } # would be a third
}
}
- REST wins at the edgeCacheable by URL, debuggable with
curl, understood by every proxy and every developer, and it survives a receiver who has never seen your schema. For a public API this is usually the end of the discussion. - gRPC wins between your own servicesGenerated clients in every language, a schema that CI can check for breaking changes, HTTP/2 multiplexing, streaming, and 4× smaller payloads. Deadlines and cancellation are built into the protocol rather than bolted on.
- GraphQL wins for varied clientsWhen an iOS screen, an Android screen and a web dashboard all need different slices of the same graph, letting the client specify the slice beats maintaining three bespoke endpoints.
- And the costsgRPC is awkward from a browser (you need grpc-web or a gateway) and uncacheable. GraphQL moves query cost to runtime — an innocent query can fan out into an N+1 storm, so you need dataloaders, depth limits and cost analysis before you expose it publicly.
Whatever you pick, your real contract is the schema, and it will outlive your code. You never get to upgrade both sides at once — during any rollout v1 and v2 are both live and talking to each other. Add optional fields with defaults, rename freely, and never reuse or retype a field number.
Finding & Choosing an Instance#
Theory#
“Call the orders service” hides two questions. Where is it? — its instances are ephemeral, rescheduled constantly, and their addresses change every deploy. And which one? — because there are eleven of them and they are not equally healthy.
Discovery is not a lookup problem, it is a freshness problem. Every design —
DNS with short TTLs, a sidecar watching a registry, a client library that subscribes — is
trading how stale the caller's address list may be against how much load the registry can take.
The window between “instance died” and “callers know” is
probe interval × failure threshold, and during it, callers send
traffic into a black hole.
Then you pick one. The default everyone inherits is round robin, and it fails in a specific, common way that is worth seeing once:
Use the picker above to compare all three. Round robin distributes requests equally, which only helps if every request costs the same and every server is identical. Least-connections distributes work, so a slow server naturally receives less. Consistent hashing distributes keys, which is what you want when the server has a cache worth keeping warm.
Questions#
Your p50 is healthy but your p99 tripled after a deploy. All instances pass their health checks. What is your first hypothesis?
One instance is grey failing — up, answering, passing a shallow /healthz that
only proves the process can return 200, but slow. With round robin it keeps receiving its full
share, so a fixed fraction of requests is slow and the tail moves while the median does not. The
tell is the shape: if p99 tripled and p50 did not, look for a subset of instances, not a
global regression.
Three fixes, in increasing order of value: make health checks deep (check the dependencies the endpoint actually needs); switch to least-request load balancing so slowness is self-limiting; and add outlier detection so an instance whose error or latency profile deviates from its peers is ejected automatically for a cool-down period.
Surviving Failure#
Timeouts, retries, idempotency and circuit breakers — the tools that turn the third outcome from an outage into a blip.
Timeouts & Deadlines#
Theory#
A timeout is the only way to convert the third outcome into the second one: it turns “I do not know” into a definite failure you can act on. Every remote call needs one. But a timeout set independently at each hop produces a subtle, expensive bug.
A table orders, waits twenty minutes, gives up and walks out. Nobody tells the kitchen. The kitchen finishes the order, plates it, and the waiter carries it to an empty table — while eleven new tables wait for the same stove. The wasted work is not the meal; it is the stove time. That is exactly what a missing deadline costs you: capacity, spent precisely when you have none to spare.
# The caller sets ONE budget for the whole request tree.
async def handle_request(request):
deadline = time.monotonic() + 3.0 # absolute, not a duration
price = await get_price(request.sku, deadline)
stock = await get_stock(request.sku, deadline)
return render(price, stock)
async def get_price(sku, deadline):
remaining = deadline - time.monotonic()
if remaining <= 0.05: # not enough time to be useful
raise DeadlineExceeded() # fail now, do not start the call
# Pass the remainder down; the next hop subtracts its own elapsed time.
async with httpx.AsyncClient() as client:
return await client.get(
f"{PRICING}/price/{sku}",
timeout=remaining,
headers={"X-Request-Deadline-Ms": str(int(remaining * 1000))},
)
// The caller sets ONE budget for the whole request tree.
async function handleRequest(req) {
const deadline = Date.now() + 3000; // absolute, not a duration
const price = await getPrice(req.sku, deadline);
const stock = await getStock(req.sku, deadline);
return render(price, stock);
}
async function getPrice(sku, deadline) {
const remaining = deadline - Date.now();
if (remaining <= 50) throw new DeadlineExceeded(); // fail now, do not start
// AbortController is how you actually cancel in-flight work in JS.
const ac = new AbortController();
const timer = setTimeout(() => ac.abort(), remaining);
try {
return await fetch(`${PRICING}/price/${sku}`, {
signal: ac.signal,
headers: { "X-Request-Deadline-Ms": String(remaining) },
});
} finally {
clearTimeout(timer);
}
}
Three timeout bugs that are nearly universal. One: the default is often
infinite — Python's requests, many JDBC drivers and plenty of gRPC setups
will wait forever unless told otherwise. Two: connect timeout and read timeout are different
settings, and setting only one leaves the other unbounded. Three: an inner timeout
larger than an outer one is always a bug; the inner call can never win, so it is pure
wasted capacity. Timeouts must shrink as you go down the stack.
Retries, Backoff & Jitter#
Theory#
Retrying is the obvious response to a failed call and the easiest way to turn a small problem into an outage. The difference between the two is entirely in when you retry.
Step through all three variants above. The first is a load amplifier aimed at the thing that is already struggling. The second fixes the volume and leaves every client synchronised. Only the third — adding randomness — produces the smooth ramp a recovering service can actually survive.
RETRYABLE = {408, 429, 500, 502, 503, 504}
async def call_with_retries(fn, *, attempts=4, base=0.2, cap=10.0):
for attempt in range(attempts):
try:
return await fn()
except HttpError as e:
# Never retry a 400 or a 422 - the request itself is wrong,
# and it will be exactly as wrong the second time.
if e.status not in RETRYABLE or attempt == attempts - 1:
raise
# Honour the server if it told us when to come back.
wait = e.retry_after or random.uniform(0, min(cap, base * 2 ** attempt))
except (TimeoutError, ConnectionError):
if attempt == attempts - 1:
raise
wait = random.uniform(0, min(cap, base * 2 ** attempt))
await asyncio.sleep(wait)
const RETRYABLE = new Set([408, 429, 500, 502, 503, 504]);
async function callWithRetries(fn, { attempts = 4, base = 200, cap = 10_000 } = {}) {
for (let attempt = 0; attempt < attempts; attempt += 1) {
try {
return await fn();
} catch (err) {
const last = attempt === attempts - 1;
const retryable = err.status === undefined || RETRYABLE.has(err.status);
// Never retry a 400 or a 422 - the request itself is wrong.
if (last || !retryable) throw err;
// Full jitter: a uniform sample from the whole backoff window.
const window = Math.min(cap, base * 2 ** attempt);
const wait = err.retryAfterMs ?? Math.random() * window;
await new Promise((r) => setTimeout(r, wait));
}
}
}
The rule that people miss is that retries compose multiplicatively. Three layers each retrying three times is not three retries, it is twenty-seven:
Retry at exactly one layer — usually the one closest to the user that still understands the request — and cap it with a retry budget: if retries exceed ~10% of requests, stop retrying entirely until the ratio recovers. Every other layer fails fast and reports upward. Add a circuit breaker so that repeated failure stops producing traffic at all.
Idempotency#
Theory#
Retries and at-least-once delivery both mean the same thing: your handler will be called more than once with the same input, and not occasionally — routinely, under load, exactly when it matters. Idempotency is the property that makes that harmless, and it is the single highest-leverage idea in this course.
Press the call button once, or jab it eleven times: the lift comes once. The button is idempotent, and that is precisely why nobody has to think about how many times they pressed it. A vending machine's coin slot is not idempotent, which is why you watch it like a hawk. Design your endpoints to be lift buttons.
@app.post("/charges")
async def create_charge(body: ChargeRequest, idempotency_key: str = Header(...)):
async with db.transaction(): # ONE transaction
try:
# The unique index is the lock. Never "select then insert" -
# two concurrent retries will both pass that check.
await db.execute(
"insert into idempotency_keys (key, request_hash) values ($1, $2)",
idempotency_key, hash_of(body),
)
except UniqueViolation:
stored = await db.fetchrow(
"select response, request_hash from idempotency_keys where key = $1",
idempotency_key,
)
# Same key, different body = a client bug worth shouting about.
if stored["request_hash"] != hash_of(body):
raise HTTPException(422, "idempotency key reused with a different body")
return json.loads(stored["response"]) # replay the original answer
charge = await ledger.charge(body.amount, body.currency, body.card_token)
await db.execute(
"update idempotency_keys set response = $1 where key = $2",
json.dumps(charge), idempotency_key,
)
return charge
app.post("/charges", async (req, res) => {
const key = req.header("Idempotency-Key");
if (!key) return res.status(400).json({ error: "Idempotency-Key required" });
await db.transaction(async (tx) => { // ONE transaction
try {
// The unique index is the lock. Never "select then insert" -
// two concurrent retries will both pass that check.
await tx.query(
"insert into idempotency_keys (key, request_hash) values ($1, $2)",
[key, hashOf(req.body)],
);
} catch (err) {
if (err.code !== UNIQUE_VIOLATION) throw err;
const { rows } = await tx.query(
"select response, request_hash from idempotency_keys where key = $1", [key],
);
// Same key, different body = a client bug worth shouting about.
if (rows[0].request_hash !== hashOf(req.body)) {
return res.status(422).json({ error: "key reused with a different body" });
}
return res.json(rows[0].response); // replay the original answer
}
const charge = await ledger.charge(req.body);
await tx.query("update idempotency_keys set response = $1 where key = $2",
[charge, key]);
res.json(charge);
});
});
- Naturally idempotent — do this where you can
PUT /orders/44 {status:"shipped"}sets an absolute state, so running it twice is identical to running it once. Prefer absolute assignments (set balance = 90) over relative ones (balance = balance - 10) whenever the domain allows. - Needs an explicit keyAnything that appends, sends, charges or emits.
POST /charges,POST /emails,append to ledger. There is no way to make these safe without remembering that you already did them.
The three ways people get this wrong. Generating a new key per attempt
instead of per intent — the key must be created before the first try and reused by every
retry. Checking with a SELECT before inserting — two concurrent retries both see
nothing and both proceed; only a unique constraint is atomic. And writing the key in a different
transaction from the effect — which reintroduces exactly the dual-write problem you will meet in
section 11.
Circuit Breakers & Bulkheads#
Theory#
Timeouts, retries and idempotency make a single call survivable. These two patterns stop a failing dependency from consuming the calling service itself. A circuit breaker decides whether to call at all; a bulkhead decides how much of you can be tied up when the answer is yes.
A short circuit in the kitchen trips one breaker. Your lights stay on, your router stays on, and the house does not burn down. The breaker does not repair the fault — it contains it, and it makes the fault obvious and the recovery deliberate. A software circuit breaker has exactly this job, and like the electrical one it must be scoped: one per dependency, never a single breaker for everything.
The breaker keeps you from calling a broken thing. The bulkhead keeps the calls you do make from eating every resource you have:
Resource exhaustion is how one failure becomes all failures. A shared thread pool, connection pool or event loop is a single point of contagion: the slow dependency does not have to break you directly, it just has to hold your workers long enough that everything else queues behind it. Partition by dependency, and pair every breaker with a fallback — a cached value, a degraded feature, a queued write, or an honest error.
Questions#
Your recommendations service is down. The product page is timing out. Walk through what should have prevented the outage.
Four layers, each of which would have been enough on its own. A short timeout (200 ms, not 30 s) means a hanging call releases its worker almost immediately. A bulkhead caps recommendations at, say, 20% of the pool, so even with a long timeout it cannot consume all of it. A circuit breaker notices the failure rate and stops calling entirely, which also removes load from a service trying to restart. A fallback — last week's popular items, or simply omitting the carousel — turns the whole thing from an outage into a slightly less personalised page.
The framing to say out loud: recommendations is a non-critical dependency, so the design fault was letting a non-critical dependency fail a critical path at all. Classify every dependency as critical or not, and make the non-critical ones structurally incapable of taking you down.
Asynchronous Messaging#
Putting a broker between services: how queues and topics decouple them, what delivery guarantees really promise, and what to do when consumers fall behind.
Queues & Topics#
Theory#
A broker gives you three things at once: it decouples in time (the receiver need not be up), it decouples in identity (the sender need not know who receives), and it absorbs load (a spike becomes a backlog instead of an error page). What you choose between is not really a product but a delivery semantic.
The two numbers to alarm on for any queue are depth (is it growing without bound?) and oldest message age (are we already late?). Depth alone lies to you: a queue with a steady depth of 40 looks identical whether messages take 50 ms or 50 minutes to get through it.
Now the same infrastructure with the opposite semantic. A queue splits work between consumers; a topic copies it to all of them:
Getting this backwards is a classic production incident. Deploy three replicas of a service that reads a queue expecting each replica to see every message, and each replica sees a third of them — two thirds of your orders are never billed, silently. Say the semantic out loud when you create the subscription: “competing consumers” or “fan-out”.
The third shape is the log — Kafka, Pulsar, Kinesis. It looks like a topic but it does not delete what it delivers, and that one difference changes everything downstream of it:
- Choose a queuewhen work is disposable once done, consumers are interchangeable, and you want per-message retries and a dead letter queue for free. Job processing, emails, thumbnails.
- Choose a logwhen the history is the product: multiple independent readers, replay after a bug, ordering within a key, or feeding both a real-time consumer and a nightly analytics job from the same data.
- Log costsParallelism is capped by partition count, which you must choose up front. Rebalances produce duplicates. And there is no per-message retry — one poison record blocks its whole partition unless you handle it yourself.
Delivery Guarantees & Ordering#
Theory#
Every broker's guarantee comes down to one decision in the consumer: do you acknowledge before or after doing the work? That is the whole thing. Everything else is marketing.
There is no exactly-once delivery. It is provably impossible over an unreliable network — the classic two-generals result. What you can have is exactly-once effect: at-least-once delivery plus a dedupe key committed in the same transaction as the side effect. When a vendor says “exactly once”, they mean either that, or exactly-once within their own system only. Your call to Stripe is still yours to protect.
Ordering has the same character — a guarantee that is narrower than people assume, and silently lost by a default setting:
# 1. Same entity -> same key -> same partition -> ordered.
await producer.send(
topic="order-events",
key=order_id.encode(), # NOT a random uuid, NOT round robin
value=event.to_json(),
headers=[("event-id", event.id.encode()), ("traceparent", ctx.encode())],
)
# 2. Belt and braces: reject stale events even if they do arrive in order.
async def handle(event):
updated = await db.execute(
"""update orders
set status = $1, version = $2
where id = $3 and version < $2""", # monotonic guard
event.status, event.version, event.order_id,
)
if updated == 0:
log.info("stale or duplicate event ignored", event_id=event.id)
// 1. Same entity -> same key -> same partition -> ordered.
await producer.send({
topic: "order-events",
messages: [{
key: orderId, // NOT a random uuid, NOT round robin
value: JSON.stringify(event),
headers: { "event-id": event.id, traceparent: ctx },
}],
});
// 2. Belt and braces: reject stale events even if they do arrive in order.
async function handle(event) {
const { rowCount } = await db.query(
`update orders
set status = $1, version = $2
where id = $3 and version < $2`, // monotonic guard
[event.status, event.version, event.orderId],
);
if (rowCount === 0) log.info({ eventId: event.id }, "stale or duplicate ignored");
}
Never ask for global ordering. It means one partition, one consumer, and no horizontal scale — you have bought a distributed system and configured it to be a single-threaded one. Ask instead: which entity's events must not overtake each other? That entity's id is your partition key, and everything else is free to be reordered.
Dead Letters & Backpressure#
Theory#
Retries assume failure is transient. Some failures are not — a malformed payload will fail identically forever, and while it is being retried it blocks everything behind it.
A dead letter queue is a bug report with the payload attached, not a graveyard.
Alarm on depth > 0, keep the original message plus the exception and stack trace,
and put a one-click redrive in the runbook. The most common failure in practice is
not a missing DLQ — it is a DLQ nobody has looked at since March.
The other broker failure mode is quieter, and it is the one that kills processes at 3 a.m.: a producer that is simply faster than its consumer.
Every queue must be bounded, and the bound must come from a latency target rather than from
available memory. “Unbounded” always means “bounded by RAM,
discovered during an incident”. Once bounded you have three real options when it fills:
push back (stop reading, so the producer slows down), shed (reject with
503 plus Retry-After), or drop by policy (discard the oldest,
which is right for live data whose value decays).
Consistency Across Services#
Keeping data correct when one business action spans a database, a broker and several services that can each fail halfway through.
The Dual Write & the Outbox#
Theory#
Here is the most common bug in event-driven systems, and it is in code that looks completely reasonable: save to the database, then publish an event. Two systems, two commits, and no transaction that spans them.
A dual write has no correct implementation. Retrying narrows the window; it never
closes it, because the process can die inside the retry loop. Reversing the order just swaps
“an order nobody was told about” for “an event about an order that does not
exist”. If you catch yourself writing db.save(); broker.publish();, you have
this bug.
async def place_order(order):
async with db.transaction(): # ONE transaction, ONE system
await db.execute(
"insert into orders (id, customer_id, total_minor) values ($1, $2, $3)",
order.id, order.customer_id, order.total_minor,
)
await db.execute(
"""insert into outbox (id, topic, partition_key, payload)
values ($1, $2, $3, $4)""",
uuid.uuid4(), "order-events", order.id, order.to_event_json(),
)
# No publish here. Nothing to fail. Nothing to retry. Nothing to get wrong.
# A separate process, restartable at any moment, publishes what committed.
async def relay():
while True:
rows = await db.fetch(
"select * from outbox where published_at is null order by id limit 100"
)
for row in rows:
await broker.publish(row["topic"], key=row["partition_key"],
value=row["payload"], headers={"event-id": row["id"]})
# A crash here republishes on restart - consumers dedupe on event-id.
await db.execute("update outbox set published_at = now() where id = $1",
row["id"])
await asyncio.sleep(0.2)
async function placeOrder(order) {
await db.transaction(async (tx) => { // ONE transaction, ONE system
await tx.query(
"insert into orders (id, customer_id, total_minor) values ($1, $2, $3)",
[order.id, order.customerId, order.totalMinor],
);
await tx.query(
`insert into outbox (id, topic, partition_key, payload)
values ($1, $2, $3, $4)`,
[randomUUID(), "order-events", order.id, order.toEventJson()],
);
});
// No publish here. Nothing to fail. Nothing to retry. Nothing to get wrong.
}
// A separate process, restartable at any moment, publishes what committed.
async function relay() {
for (;;) {
const { rows } = await db.query(
"select * from outbox where published_at is null order by id limit 100",
);
for (const row of rows) {
await broker.publish(row.topic, {
key: row.partition_key, value: row.payload, headers: { "event-id": row.id },
});
// A crash here republishes on restart - consumers dedupe on event-id.
await db.query("update outbox set published_at = now() where id = $1", [row.id]);
}
await sleep(200);
}
}
The outbox converts an impossible problem into a solved one plus a tolerable one: atomicity within a single database, and at-least-once publishing that consumers deduplicate. The alternative, if you cannot change the application at all, is to let the database's own log be the source of events:
Sagas#
Theory#
One business operation, four services, four databases. You cannot wrap that in a transaction, so you give up atomicity and replace it with a sequence of local transactions, each with a matching compensation that semantically offsets it.
A database rollback leaves no trace: the write never existed. A compensation is a new forward transaction — the customer's statement shows the charge and the refund, the warehouse log shows the reservation and the release, and for a few seconds the money genuinely was gone. Compensation is bookkeeping, not time travel, and designing it is the real work of a saga.
- OrchestrationOne component owns the state machine and calls each step. Easy to see, easy to debug, easy to add a timeout or a manual intervention step. Use it for anything long, critical, or with more than three steps.
- ChoreographyEach service reacts to events and publishes its own. Maximum decoupling, and adding a participant needs no change anywhere else. Good for two or three stable steps.
- Orchestration's costThe orchestrator knows about everyone, so it becomes a place logic accumulates and a service every team must change.
- Choreography's costThe flow exists nowhere as a readable artefact. “Why did this order stall?” means correlating logs across five services, and cyclic subscriptions become loops nobody designed.
Two rules that make sagas work. Every step must be idempotent, because the orchestrator will retry on timeout. And design the compensations first — the happy path is the easy half, and discovering during an incident that a step has no compensation (an email is already sent, a warehouse has already shipped) is how sagas get a bad reputation. Where a step truly cannot be undone, put it last.
For contrast, the mechanism sagas replace — two-phase commit — is worth understanding precisely so you can explain why you did not use it:
Reaching Clients & Long-Running Work#
Pushing updates to browsers and partners, handing back results that take minutes, and the whole course on one page.
Polling, SSE & WebSockets#
Theory#
Everything so far has been servers talking to servers. The last hop is different, because the client is behind a NAT, on a flaky radio, and cannot be called. So the connection must start from their side, and the question becomes how long you keep it open: a fresh request every few seconds (polling), one request the server holds until it has news (long polling), one response that never ends (server-sent events), or a connection that stops being HTTP and carries messages both ways (WebSockets).
Polling is walking to the post office every hour to ask whether anything has come. Long polling is waiting at the counter until something arrives, then joining the back of the queue again. Server-sent events is a courier who rings your bell for each parcel — and when you were out, checks the number on your last receipt and brings everything since. A WebSocket is a phone line you keep open all day: either of you can talk at any time, but somebody has to pay for the line, and when it drops nobody redials for you.
- Strength — SSE is just HTTPIt carries your cookies, auth, CORS rules, tracing and HTTP/2 multiplexing unchanged, and
EventSourcereconnects by itself, sendingLast-Event-IDso the server can resume. - Strength — WebSockets are truly two-wayEither side sends whenever it likes with a few bytes of framing, which is what chat, collaborative editing and games need.
- Weakness — SSE is one-way textThe client still sends with ordinary requests, binary costs a base64 tax, and
EventSourcecannot set anAuthorizationheader. - Weakness — a WebSocket is a protocol you now ownAfter the
101there are no status codes, retries, acks or per-request auth. Envelopes, acks, heartbeats, resume and reconnect are all yours to build.
The detail that separates a demo from a production stream is what happens on reconnect. Connections will drop — every deploy drops all of them at once — so each event needs an id, and the server needs a backlog it can replay from:
Default to server-sent events. They ride on ordinary HTTP with your existing auth, reconnect automatically with Last-Event-ID resumption built in, and cover the common case of server-to-client updates. Reach for WebSockets only when the client needs to send frequently — chat, collaborative editing, games. Keep polling in your pocket for updates that are rare and not urgent; it is not a bad design, just a bad default.
Long-lived connections are state, and that is the real cost. Every deploy disconnects everyone at once (so you need jittered reconnect, or they all come back together). Autoscaling is driven by connection count, not CPU. And because any server may need to reach any client, you need a pub/sub backplane behind them — which means the fan-out problem you solved in section 8 reappears, one layer down.
“It works locally but updates arrive in bursts in production” is almost always buffering. nginx, gzip middleware, some CDNs and most serverless gateways hold a response until it is complete or a buffer fills — which for a stream that never ends means “much later”. Send X-Accel-Buffering: no, exclude text/event-stream from compression, and test through the real proxy chain before launch.
Questions#
A food-delivery app shows live order status — “preparing”, “picked up”, “two minutes away”. Two million users have the screen open at peak. How do you deliver the updates?
Start with direction: the server talks, the client only watches. That rules out WebSockets as unnecessary and points to SSE — ordinary HTTP, existing auth, and resumption built in. Polling every 5 seconds would be 400,000 requests per second of mostly “nothing new” for six status changes per order.
Then the parts that make it work at scale. Each status change gets a monotonic id, and the server keeps a short per-order backlog, so a phone that drives through a tunnel reconnects with Last-Event-ID and gets what it missed. Connections land on arbitrary nodes, so every node subscribes to a pub/sub backplane keyed by order id. Proxy buffering is switched off for the route, a : keepalive comment goes out every 15 seconds, and the retry: value is jittered so a deploy does not bring two million clients back in the same second. Finally, a 30-second poll of GET /orders/{id} stays as the fallback for networks that eat streams.
@app.get("/orders/{order_id}/events")
async def order_events(order_id: str, request: Request,
last_event_id: str | None = Header(None)):
async def gen():
last = int(last_event_id or 0)
live = bus.subscribe(f"order:{order_id}") # subscribe BEFORE replaying
try:
for ev in await backlog.since(order_id, last): # what the tunnel ate
yield f"id: {ev.seq}\ndata: {json.dumps(ev.data)}\n\n"
last = ev.seq
while not await request.is_disconnected():
try:
ev = await asyncio.wait_for(live.next(), timeout=15)
except asyncio.TimeoutError:
yield ": keepalive\n\n"
continue
if ev.seq > last: # skip what the replay sent
yield f"id: {ev.seq}\ndata: {json.dumps(ev.data)}\n\n"
last = ev.seq
finally:
await live.close()
return StreamingResponse(gen(), media_type="text/event-stream",
headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"})
app.get("/orders/:id/events", async (req, res) => {
res.writeHead(200, {
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache",
"X-Accel-Buffering": "no",
});
let last = Number(req.get("Last-Event-ID") ?? 0);
const send = (ev) => {
if (ev.seq <= last) return; // skip what the replay sent
last = ev.seq;
res.write(`id: ${ev.seq}\ndata: ${JSON.stringify(ev.data)}\n\n`);
};
const pending = [];
let replaying = true;
const live = await bus.subscribe(`order:${req.params.id}`, (ev) => // subscribe BEFORE replaying
replaying ? pending.push(ev) : send(ev));
for (const ev of await backlog.since(req.params.id, last)) send(ev); // what the tunnel ate
replaying = false;
pending.forEach(send);
const keepalive = setInterval(() => res.write(": keepalive\n\n"), 15_000);
req.on("close", () => { clearInterval(keepalive); live.unsubscribe(); });
});
Webhooks & Async Request-Reply#
Theory#
Two patterns for when the answer cannot come back on the connection that asked for it. A webhook is used when the receiver is another company's server: you make the HTTP call, to an endpoint you do not control and cannot debug. Asynchronous request-reply is used when the work outlasts any sane timeout: you return a handle immediately (202 Accepted plus a status URL) and deliver the result later — by the client polling, by a push, or by a webhook. Inside a system the same shape runs over queues, with a reply_to queue and a correlation id that matches each answer to its question.
You do not stand at the counter while your suit is cleaned. You get a numbered ticket and leave. You can come back and ask (polling the status URL), or leave your phone number so they call you (a webhook callback). The ticket number is the correlation id: it is the only thing that connects the suit on the rail to you. And if the shop loses the ticket book in a fire — an in-memory job table during a deploy — the suits are still there but nobody can claim them.
- Strength — nobody holds a connectionThe caller is released in milliseconds, the work runs on workers that scale independently, and a deploy in the middle costs nothing because the job lives in a database.
- Strength — the server sets the pace
Retry-Afteron every status response lets you slow every poller down during an incident without a client release. - Weakness — a resource you now ownThe job handle must survive restarts, be idempotent to create, expose progress and errors, and expire explicitly with
410 Gone. - Weakness — webhooks fail on their schedule, not yoursReceivers are down for hours, deliveries arrive out of order and more than once, and the callback URL is an SSRF vector you have to validate.
Never hold a request open for a long job. The load balancer's idle timeout (often 60 s) kills it, the client sees an error, the user clicks again — and now two jobs are running, because the first one never stopped. Return 202 fast, make the start idempotent with a key the client sends, and let the client wait on the status, not on the connection.
A timed-out reply is not a failed request. Whether the answer comes back over HTTP polling or a reply queue, the requester can always give up before the work finishes — so the work must be idempotent, a late reply with an unknown correlation id must be dropped quietly, and “I stopped waiting” must never be confused with “it did not happen”. It is the third outcome from section 0, wearing a ticket number.
Questions#
A “download my annual statement” endpoint takes 2–4 minutes to build the PDF. It times out at the load balancer, users click again, and support sees duplicate statements and angry tickets. Fix it.
Two bugs are hiding in one symptom. The first is shape: a synchronous request cannot outlive the shortest timeout on the path, so the work has to move off the request and the client needs a handle to come back with. That is async request-reply — POST returns 202 Accepted with Location: /jobs/{id} and Retry-After, a worker builds the PDF, and GET /jobs/{id} answers “running” until it redirects with 303 to the finished file.
The second is duplication, and the new shape does not fix it alone: a user who clicks twice still sends two POSTs. The client generates one idempotency key per intent — per click of “download 2025 statement”, reused on any retry — and the server enforces it with a unique constraint, so a repeated start returns the existing job. Finally, keep the status URL even if you also push a “ready” notification, because the notification can be missed and the status cannot.
@app.post("/statements", status_code=202)
async def start_statement(req: StatementRequest, response: Response,
idempotency_key: str = Header(...)):
job = await jobs.get_by_key(idempotency_key) # a double click returns the same job
if job is None:
async with db.transaction(): # unique constraint on key
job = await jobs.create(key=idempotency_key, params=req.model_dump())
await outbox.append("statement.requested", {"jobId": job.id})
response.headers["Location"] = f"/jobs/{job.id}"
response.headers["Retry-After"] = "5"
return {"jobId": job.id, "status": job.status}
@app.get("/jobs/{job_id}")
async def job_status(job_id: str, response: Response):
job = await jobs.get_or_404(job_id)
if job.status == "succeeded":
return RedirectResponse(f"/files/{job.file_id}", status_code=303)
response.headers["Retry-After"] = "10"
return {"status": job.status, "progress": job.progress, "error": job.error}
app.post("/statements", async (req, res) => {
const key = req.get("Idempotency-Key");
let job = await jobs.getByKey(key); // a double click returns the same job
if (!job) {
job = await db.transaction(async (tx) => { // unique constraint on key
const created = await jobs.create(tx, { key, params: req.body });
await outbox.append(tx, "statement.requested", { jobId: created.id });
return created;
});
}
res.status(202).location(`/jobs/${job.id}`).set("Retry-After", "5")
.json({ jobId: job.id, status: job.status });
});
app.get("/jobs/:id", async (req, res) => {
const job = await jobs.getOr404(req.params.id);
if (job.status === "succeeded") return res.redirect(303, `/files/${job.fileId}`);
res.set("Retry-After", "10").json({ status: job.status, progress: job.progress, error: job.error });
});
The Whole Thing on One Page#
Theory#
Everything above, compressed into the three decisions you actually make.
Pick a communication style#
| If the caller… | Use | You are accepting |
|---|---|---|
| needs the answer to continue | synchronous RPC (gRPC internally, REST at the edge) | latency adds, availability multiplies |
| needs only an acknowledgement | queue — competing consumers | eventual consistency, duplicate handling |
| is announcing a fact to nobody in particular | topic — pub/sub fan-out | no visible call graph; the event schema is now an API |
| wants replay, or many independent readers | log (Kafka, Kinesis, Pulsar) | partition-capped parallelism, rebalance duplicates |
| is a browser that needs updates | SSE, or WebSocket if it must send too | connections are state: deploys, scaling, backplane |
| is another company's server | webhook, signed and retried for hours | out-of-order delivery, endpoints you cannot debug |
| starts a job that outlasts any timeout | async request-reply: 202 + status URL (or a reply queue + correlation id) |
a job resource to store, poll traffic, results that expire |
Pick a reliability pattern#
| Symptom | Pattern | The one detail people get wrong |
|---|---|---|
| calls hang forever | timeouts | the default is often infinite; connect and read are separate settings |
| work continues after the caller left | deadline propagation | pass an absolute deadline, not a duration |
| transient errors reach users | retry with full jitter | retry at one layer only, under a budget |
| retries cause double effects | idempotency key | one key per intent; enforce with a unique constraint |
| a broken dependency is still being called | circuit breaker | scope per dependency, and always pair with a fallback |
| one dependency exhausts the whole service | bulkhead | partition the pool before you need to |
| a queue grows without limit | backpressure / load shedding | bound it from a latency target, not from RAM |
| one message fails forever | dead letter queue | alarm on depth > 0 and keep a redrive path |
| db write and event can diverge | transactional outbox (or CDC) | never save() then publish() |
| a multi-service operation half-failed | saga with compensations | design the compensations first; put un-undoable steps last |
| a fan-out is as slow as its worst shard | hedged requests, partial results | measure the p99 of the aggregate, not of the parts |
| cache expiry floods the origin | single-flight + jittered TTL | stale-while-revalidate beats an error every time |
The questions to ask of any design#
If you remember one sentence from this page: the network gives you three outcomes instead of two, so every pattern here is either making the unknown outcome rare (timeouts, breakers, bulkheads), making it harmless (idempotency, outbox, compensations), or making it visible (tracing, queue depth, DLQs). When you meet a pattern that is not on this page, ask which of those three it is doing — there is no fourth category.
Where to go next#
- 📘 The Distributed Communication Patterns Detailed Course — the same ground in 47 sections, plus transports and HTTP/3, polling, server-sent events, WebSockets and async request-reply in depth, service meshes, consensus and quorums, CQRS and event sourcing, scatter-gather and tail latency, distributed tracing, contract testing and a practice roadmap.
- 📚 Browse both guides in the Distributed Communication Courses catalog.
- 📖 AWS Builders' Library — Timeouts, retries, and backoff with jitter — the measured data behind section 5, from people running it at planetary scale.
- 🧩 microservices.io pattern catalog — Chris Richardson's reference for saga, outbox, API gateway and their variants.
- 🧠 Martin Fowler — What do you mean by “event-driven”? — the clearest short treatment of the four things people mean by that phrase.
"A distributed system is one in which the failure of a computer you didn't even know existed can render your own computer unusable." — Leslie Lamport