Distributed systems · communication
Talking across the wire
Processes on different machines share no memory and no clock. Everything they coordinate, they coordinate by sending bytes over a network that can delay, drop, duplicate, or reorder them. This page covers the abstractions built on top of that: sockets, remote procedure calls, and message passing. It also covers the guarantees each one can and cannot give.
Local vs. remote IPC
Inter-process communication (IPC) is any mechanism that lets two processes exchange data and coordinate. On one machine, the kernel mediates: pipes, signals, and shared memory are fast and fail together with the machine. Across machines, the only channel is the network, and the two ends can fail independently.
Three things change when the call leaves the machine. Latency rises by three to six orders of magnitude. Memory is no longer shared, so pointers mean nothing on the other side and data must be copied into messages. Failure becomes partial: a silent peer may be crashed, slow, or cut off, and the sender cannot tell which.
TutorialPython Socket Programming HOWTO · Beej's Guide to Network Programming
Transport basics
Almost every distributed protocol rides on one of two transports. TCP gives a connection, an ordered byte stream, retransmission, and flow control. UDP gives independent datagrams with no setup, no ordering, and no retransmission. TCP costs a handshake and can stall a whole stream behind one lost packet. UDP is cheaper but leaves reliability to the application.
"Reliable" in TCP means bytes arrive in order or the connection reports an error. It does not mean the remote application processed them. A TCP acknowledgment comes from the remote kernel, not the remote program. That gap is why the end-to-end argument holds that correctness checks belong at the endpoints: the application must confirm its own work.
PapersSaltzer, Reed & Clark, "End-to-End Arguments in System Design," ACM TOCS 2(4), 1984 · Cerf & Kahn, "A Protocol for Packet Network Intercommunication," IEEE Trans. Comm. 22(5), 1974
ToolWireshark user guide (watch a real handshake)
Marshalling
Before data can cross the wire, it must be flattened into bytes (marshalling or serialization) and rebuilt on the other side (unmarshalling). Both ends must agree on field order, types, character encoding, and byte order. One machine may be little-endian and the other big-endian; network protocols conventionally use big-endian "network order".
Text formats (JSON, XML) are readable and forgiving. Binary formats (Protocol Buffers, Thrift, Avro, MessagePack) are smaller and faster to parse. Schema-based formats also handle evolution: new fields get new tag numbers, and old readers skip tags they don't know. That lets clients and servers upgrade at different times, which is the normal case in a live distributed system.
TutorialProtocol Buffers wire format · Apache Avro docs · Kleppmann, schema evolution in Avro, Protobuf and Thrift
Synchronous vs. asynchronous
In synchronous communication the sender blocks until a reply arrives. The code is easy to read and reason about, but latency stacks up and a slow dependency stalls everything behind it. In asynchronous communication the sender continues immediately and handles the reply later through a callback, future, or event. Throughput and resilience improve, but error handling and state become harder to follow.
Two related ideas matter as soon as calls are async. Deadlines bound how long anyone waits, and should travel with the request so downstream services stop work the caller has abandoned. Backpressure lets a slow consumer tell a fast producer to slow down, instead of letting queues grow until memory runs out. Most production systems mix both styles: synchronous where the user is waiting for an answer, asynchronous for everything that can happen afterward.
TutorialPython asyncio: coroutines and tasks · Node.js event loop · gRPC deadlines
Anatomy of an RPC
A remote procedure call makes a function on another machine look like a local one. The programmer writes balance = bank.get_balance(acct) and generated code handles the network. The pieces:
- Interface definition (IDL): a language-neutral contract listing procedures and message types.
- Stubs: code generated from the IDL. The client stub marshals arguments and sends them. The server stub (skeleton) unmarshals them and calls the real function.
- Runtime: manages connections, timeouts, retries, and dispatch.
- Binding: how the client finds the server, through a fixed address, DNS, or a service registry.
syntax = "proto3"; service Bank { rpc GetBalance (BalanceRequest) returns (BalanceReply); } message BalanceRequest { string account_id = 1; } message BalanceReply { int64 cents = 1; }
The transparency is useful and also dangerous. Waldo and colleagues argued that a remote call can never be truly local: it differs in latency, memory access (no pointers), concurrency, and above all partial failure. Interfaces that hide those differences push the problems to where they are hardest to handle.
PapersBirrell & Nelson, "Implementing Remote Procedure Calls," ACM TOCS 2(1), 1984 · Waldo, Wyant, Wollrath & Kendall, "A Note on Distributed Computing," 1994
TutorialgRPC Python quick start · Apache Thrift tutorial
When calls fail
A client that sends a request and hears nothing back faces three possibilities. From the client's side they look identical.
If the client retries, case 2 runs the operation twice. If it doesn't, case 1 never runs it. The RPC system's call semantics define which risk you take:
| Semantics | Mechanism | Risk | Safe for |
|---|---|---|---|
| At-most-once | Never retry, or server filters duplicates by request ID | Operation may be lost | Non-repeatable actions where loss is acceptable |
| At-least-once | Retry until acknowledged | Operation may run more than once | Idempotent operations |
| Exactly-once (effective) | At-least-once delivery plus idempotency keys or deduplication at the receiver | Extra state and storage | Payments, orders, transfers |
An operation is idempotent if repeating it has the same effect as doing it once. "Set balance to 50" is idempotent; "add 10 to balance" is not. A common fix is to attach a unique idempotency key to each request so the server can recognize and ignore a duplicate. Retries also need exponential backoff with jitter. Without it, every client retries at the same moment and the recovering server is flooded again.
ReadTimeouts, retries, and backoff with jitter (Amazon Builders' Library) · Designing robust APIs with idempotency (Stripe)
Message-oriented middleware
Instead of calling a procedure, a sender can hand a message to a broker that stores it and delivers it later. This decouples the two sides in time (the receiver need not be running when the message is sent) and in space (the sender need not know who or where the receiver is).
Queues spread work: each message goes to exactly one of several competing consumers, which makes scaling as easy as adding workers. Topics broadcast events: every subscriber receives its own copy, so new consumers can be added without touching the producer. Request–reply over messaging pairs a request queue with a reply queue and a correlation ID, recovering RPC-like behavior with the broker's buffering.
Brokers differ in their storage model. Traditional brokers (RabbitMQ, ActiveMQ) delete a message once it is acknowledged. Log-based brokers (Kafka, Pulsar) append messages to a durable, partitioned log. Each consumer tracks its own offset, so history can be replayed.
TutorialRabbitMQ tutorials · Kafka design docs · ZeroMQ guide (brokerless messaging)
PaperEugster, Felber, Guerraoui & Kermarrec, "The Many Faces of Publish/Subscribe," ACM Computing Surveys 35(2), 2003
Delivery and ordering guarantees
Every messaging system states, explicitly or not, what it promises about each message. Delivery promises use the same three words as RPC:
- At-most-once
- Send and forget. Fast, but messages can vanish. Fine for metrics and telemetry.
- At-least-once
- Persist, acknowledge after processing, redeliver on timeout. Duplicates are possible, so consumers must be idempotent.
- Exactly-once
- At-least-once plus deduplication or transactional commit of output and offset together. It holds within one system and needs application help at the edges.
- Durability
- Whether a message survives a broker crash: memory only, written to disk, or replicated across brokers.
Ordering is a separate promise. FIFO order preserves each sender's sequence. Causal order also preserves cause and effect across senders. Total order makes every receiver see every message in the same sequence, which is what replicated state machines need.
PapersLamport, "Time, Clocks, and the Ordering of Events," CACM 21(7), 1978 · Défago, Schiper & Urbán, "Total Order Broadcast and Multicast Algorithms: Taxonomy and Survey," ACM Computing Surveys 36(4), 2004
ReadKafka message delivery semantics
Group communication
Many tasks need one message to reach many processes: replicating an update, announcing membership, spreading configuration. Three approaches trade control for scale:
- Repeated unicast: send N copies. Simple, but the sender becomes a bottleneck and its failure midway leaves some members updated and others not.
- Reliable multicast: a protocol guarantees all correct members receive the message, often with an ordering guarantee. Virtual synchrony adds consistent views of who is in the group.
- Gossip (epidemic) protocols: each node periodically forwards what it knows to a few random peers. No coordinator, tolerant of churn, and information reaches everyone in roughly log N rounds.
PapersDemers et al., "Epidemic Algorithms for Replicated Database Maintenance," PODC 1987 · Birman & Joseph, "Exploiting Virtual Synchrony in Distributed Systems," SOSP 1987
Choosing a model
No model is best overall. Start from what the caller needs: an answer now, a guarantee the work will happen eventually, or a stream of events to react to.
| Model | Coupling | Latency | On receiver failure | Fits |
|---|---|---|---|---|
| REST over HTTP | Sync, both online | Low to moderate; text payloads | Caller gets an error | Public APIs, browser clients |
| Binary RPC (gRPC, Thrift) | Sync or streaming | Low; compact binary, multiplexed | Caller gets an error or deadline | Internal service-to-service calls |
| Work queue | Async, decoupled in time | Higher, but absorbs bursts | Messages wait in the queue | Background jobs, load leveling |
| Pub/sub topic | Async, decoupled in time and identity | Moderate | Each subscriber catches up independently | Event notification, fan-out |
| Replayable log | Async, consumer-paced | Moderate; very high throughput | Consumer resumes from its offset | Event sourcing, data pipelines |
| Gossip | Async, peer-to-peer | Seconds; eventual | Routed around automatically | Membership, failure detection |
Latency and throughput pull in different directions. Batching many messages per network write raises throughput but makes each message wait. Measure both, and watch tail latency (p99), not only the average.
PaperFielding, Architectural Styles and the Design of Network-based Software Architectures (REST), 2000
In practice
gRPC
RPCHTTP/2Protobuf
Google open-sourced gRPC from its internal RPC system. It uses Protocol Buffers for the IDL and payloads, and HTTP/2 to multiplex many calls over one connection. It supports four call shapes: unary, server streaming, client streaming, and bidirectional streaming. Deadlines and cancellation propagate across service hops.
Apache Kafka
log-based messagingat-least-once / exactly-once
Built at LinkedIn to move activity data between systems. Topics are split into partitions, each an append-only log replicated across brokers. Order is guaranteed within a partition, not across partitions. Consumers in a group divide partitions among themselves and commit offsets, so a crashed consumer's replacement resumes where it left off.
Kreps, Narkhede & Rao, "Kafka: a Distributed Messaging System for Log Processing," NetDB 2011
Mixed designs
sync + async
A typical ride-hailing or checkout flow answers the user synchronously ("your request was accepted") and does everything else asynchronously: payment capture, notifications, analytics, and continuous location updates flow as events. The synchronous path stays short, and a slow downstream system delays an email instead of a purchase.
Queue-based load leveling pattern · Asynchronous request-reply pattern
Tools and reading
Try it
- Python socket and asyncio streamsWrite a multi-client chat server two ways and compare.
- gRPC basics tutorialDefine a service, generate stubs, and try all four streaming modes.
- RabbitMQ work queuesCompeting consumers, acknowledgments, and redelivery.
- Kafka quick startProduce, consume, then replay from offset zero.
- WiresharkCapture TCP handshakes and gRPC frames on loopback.
- ToxiproxyInject latency and dropped connections to test retries and timeouts.
Read it
- van Steen & Tanenbaum, Distributed Systems, ch. 4Free textbook chapter on communication.
- Kleppmann, Designing Data-Intensive Applications, ch. 4 and 11Encoding, dataflow, and stream processing.
- Birrell & Nelson, 1984The original RPC design.
- Waldo et al., 1994Why remote is not local.
- Saltzer, Reed & Clark, 1984Where reliability belongs.
- Eugster et al., 2003Taxonomy of publish/subscribe.