9.2 Unreliable Networks
Shared-nothing systems: each machine has its own memory and disk, and one machine cannot access another's except by making requests over the network. (Even with shared object storage, machines communicate with it over the network.)
The internet and most datacenter networks (Ethernet) are ASYNCHRONOUS PACKET NETWORKS: one node can send a packet, but the network gives NO GUARANTEES as to WHEN it will arrive or WHETHER it will arrive at all.
2.1 The six things that could have happened
You send a request and get no response. Which of these happened?
- Your request was lost — someone unplugged a cable.
- Your request is queued and will arrive later — the network or the recipient is overloaded.
- The remote node failed — crashed or powered down.
- The remote node temporarily stopped responding but will resume — a long GC pause.
- The node processed your request, but the response was lost — a misconfigured switch.
- The node processed your request, but the response is delayed — overload.
These are indistinguishable in an asynchronous network. “The only information you have is that you haven't received a response yet.”
The sender can't even tell whether the packet was delivered. The only option is for the recipient to send a response message — which may in turn be lost or delayed.
The usual way of handling this is a TIMEOUT. However, when a timeout occurs, YOU STILL DON'T KNOW WHETHER THE REMOTE NODE GOT YOUR REQUEST — and if the request is still queued somewhere, it may still be delivered, even though you've given up.
2.2 Why TCP doesn't save you
What TCP genuinely provides: detects and retransmits dropped packets; detects reordered packets and restores order; detects corruption via a simple checksum; and figures out how fast it can send — congestion control / flow control / backpressure.
How a "send" actually works:
Four reasons TCP's "reliability" is not your reliability:
- TCP decides a packet must have been lost if no ACK arrives within a timeout — but it CAN'T TELL whether the outbound packet or the ACK was lost. It can resend, but can't guarantee the new packet gets through (if the network cable is unplugged, TCP can't plug it back in for you). Eventually it gives up and signals an error.
- TCP's deduplication and retransmission apply only to A SINGLE CONNECTION. If the application reconnects and retransmits, data could be duplicated.
- If a TCP connection closes with an error, you have NO WAY OF KNOWING HOW MUCH DATA WAS ACTUALLY PROCESSED by the remote node.
- Even if you receive an ACK, that means only that the OS KERNEL on the remote node received it; THE APPLICATION MAY HAVE CRASHED BEFORE HANDLING THAT DATA.
If you want to be sure a request was successful, you need A POSITIVE RESPONSE FROM THE APPLICATION ITSELF.
(Most of this applies equally to QUIC, SCTP (WebRTC), and BitTorrent uTP.)
2.3 Network faults in practice — the numbers
| Finding | Detail |
|---|---|
| Frequency | One study of a medium-sized datacenter: ~12 network faults per month — half disconnected a single machine, half disconnected an ENTIRE RACK |
| Redundancy doesn't help as much as you'd think | Adding redundant networking gear doesn't reduce faults as much as expected, since it DOESN'T GUARD AGAINST HUMAN ERROR (misconfigured switches), which is a MAJOR CAUSE OF OUTAGES |
| Physical causes | Wide-area fiber interruptions blamed on cows, beavers, and sharks (shark bites now rarer with better cable shielding). Humans too: accidental misconfiguration, scavenging, sabotage |
| Delay magnitude | Across cloud regions, round-trip times of UP TO SEVERAL MINUTES at high percentiles. Even within one datacenter, packet delay of more than a minute during a topology reconfiguration triggered by a switch software upgrade. We have to assume messages might be delayed ARBITRARILY |
| Asymmetric and partial faults | A and B can communicate, B and C can communicate, but A and C cannot. Or: a network interface that DROPS ALL INBOUND PACKETS BUT SENDS OUTBOUND PACKETS SUCCESSFULLY. "Just because a network link works in one direction doesn't guarantee it's also working in the opposite direction." |
| Aftershocks | Even a BRIEF network interruption can have repercussions lasting MUCH LONGER than the original issue |
If the error handling of network faults is not defined and tested, ARBITRARILY BAD THINGS could happen — the cluster could become DEADLOCKED and permanently unable to serve requests even when the network recovers, or it could POTENTIALLY DELETE ALL OF YOUR DATA. If software is put in an unanticipated situation, it may do arbitrary unexpected things.
Important nuance:
Handling network faults doesn't necessarily mean TOLERATING them. If your network is normally fairly reliable, a valid approach may be to SIMPLY SHOW AN ERROR MESSAGE to users. However, you DO need to know how your software reacts and ensure the system can RECOVER.
(Terminology: network partition / netsplit — one part of the network cut off from the rest. Not fundamentally different from other network interruptions, and NOT related to sharding, Ch 7.)
2.4 Fault detection
Systems that must detect faulty nodes: a load balancer taking a dead node out of rotation; a single-leader database promoting a follower.
Four cases where you get explicit feedback:
| Signal | Limitation |
|---|---|
| TCP RST / FIN — machine reachable, no process listening on the port | Only if the machine is up and the process crashed |
| A script notifying other nodes of a crash (HBase does this) — process died but the OS is running | Only if the OS survives |
| Querying network switch management interfaces for link failures | Ruled out if you're on the internet, in a shared datacenter with no switch access, or if a network problem blocks the management interface |
| ICMP Destination Unreachable from a router | The router doesn't have a magic failure detection capability either; it is SUBJECT TO THE SAME LIMITATIONS as other participants |
Rapid feedback is useful, but YOU CAN'T COUNT ON IT. In general you have to assume you will get NO RESPONSE AT ALL.
You need to strike a balance between FALSE POSITIVES and FALSE NEGATIVES: too short a timeout causes alive nodes to be incorrectly suspected dead; too long a timeout causes unnecessary delays waiting for dead nodes.
2.5 Timeouts and unbounded delays
The cost of a premature declaration of death:
- If the node is actually alive and in the middle of an action — sending an email — and another node takes over, the action may be performed twice.
- Its responsibilities transfer to other nodes, which means additional load on those nodes and on the network.
- If the system is already struggling with high load, this makes the problem worse.
- In particular, the node may not have been dead at all, only slow because of overload. Transferring its load can cause a cascading failure — in the extreme, all nodes declare each other dead and everything stops working.
The fictitious world where timeouts are easy:
If every packet is either delivered within time
dor lost, and a non-failed node always handles a request within timer, then every successful request receives a response within2d + r— a reasonable timeout.Unfortunately, most systems have NEITHER guarantee. Asynchronous networks have UNBOUNDED DELAYS (no upper limit on arrival time), and most server implementations cannot guarantee handling requests within a maximum time.
For failure detection, it's NOT SUFFICIENT for the system to be fast MOST of the time: if your timeout is low, it takes only A TRANSIENT SPIKE in round-trip times to throw the system off balance.
Network congestion and queueing — the four queues
- Switch queue. Several nodes send to the same destination, so the switch must queue them and feed them into the destination link one by one. If the queue fills up, the packet is dropped and must be resent, even though the network is functioning fine.
- OS queue at the destination. If all CPU cores and application threads are busy, the incoming request is queued by the operating system until the application is ready. This may take an arbitrary length of time depending on load.
- VM monitor queue. In virtualized environments an OS is often paused for tens of milliseconds while another VM uses a CPU core. During that time the VM cannot consume data from the network, so it is buffered by the VM monitor.
- Sender-side TCP queue. TCP limits its send rate to avoid overloading the network, which means additional queueing at the sender before the data even enters the network.
Plus: when TCP retransmits a lost packet, the application doesn't see the packet loss directly — it sees THE RESULTING DELAY (waiting for the timeout, then for the retransmitted packet's ACK).
TCP vs UDP: latency-sensitive applications (videoconferencing, VoIP) use UDP — a trade-off between reliability and variability of delays. UDP avoids some causes of variable delay (no flow control, no retransmission) though still susceptible to switch queues and scheduling delays.
UDP is a good choice WHEN DELAYED DATA IS WORTHLESS. In a VoIP call there isn't time to retransmit before the data is due to play; the application fills the missing slot with silence and moves on. THE RETRY HAPPENS AT THE HUMAN LAYER INSTEAD — "Could you repeat that please? The sound just cut out."
Two aggravating factors:
- Queueing delays have an especially wide range when a system is close to maximum capacity. A system with plenty of spare capacity can easily DRAIN queues; in a highly utilized system, LONG QUEUES BUILD UP VERY QUICKLY. (Same curve as Ch 2 §2.1.)
- Multitenancy: network links, switches, NICs, and CPUs are shared. Because you have no control over or insight into other customers' usage, network delays can be highly variable if a NOISY NEIGHBOUR is using a lot of resources.
Therefore, timeouts must be chosen experimentally:
Measure the distribution of round-trip times over an extended period and over many machines to determine the expected variability. Then determine an appropriate trade-off between failure-detection delay and risk of premature timeouts.
Even better: rather than configured constants, CONTINUALLY MEASURE response times and their variability (JITTER) and AUTOMATICALLY ADJUST timeouts. The Phi Accrual failure detector (used in Akka and Cassandra) does this. TCP retransmission timeouts work similarly.
2.6 Synchronous vs asynchronous networks — why we chose unreliability
The telephone network is the counter-example: a call establishes a CIRCUIT — a fixed, guaranteed amount of bandwidth allocated along the entire route, remaining in place until the call ends. ISDN runs at 4,000 frames/second; a call is allocated 16 bits within each frame in each direction — so each side is guaranteed to send exactly 16 bits of audio every 250 microseconds.
This network is SYNCHRONOUS: even passing through several routers, it does not suffer from queueing, BECAUSE THE SPACE HAS ALREADY BEEN RESERVED IN THE NEXT HOP. And because there is no queueing, the maximum end-to-end latency is fixed — a BOUNDED DELAY.
Why datacenters and the internet don't do this:
| Circuit | Packet (Ethernet, IP) |
|---|---|
| A fixed reserved bandwidth nobody else can use while the circuit is established. | Opportunistically uses whatever bandwidth is available. |
| Good for audio and video calls — a fairly constant number of bits per second. | Good for bursty traffic — a web page, an email, a file: no particular bandwidth requirement, we just want it as fast as possible. |
| To transfer a file over a circuit you would have to guess a bandwidth allocation: too low and it is unnecessarily slow with capacity left unused; too high and the circuit cannot be set up at all. | TCP dynamically adapts the transfer rate to the available capacity. |
The deeper framing — variable delay as dynamic resource partitioning:
A wire carrying 10,000 simultaneous calls divides the resource STATICALLY: even if you're the only call and 9,999 slots are unused, your circuit gets the same fixed bandwidth.
The internet shares bandwidth DYNAMICALLY. Senders push and jostle to get packets over the wire, and switches decide which packet to send from one moment to the next. The downside is queueing; the advantage is that IT MAXIMIZES UTILIZATION OF THE WIRE. The wire has a fixed cost, so if you utilize it better, EACH BYTE IS CHEAPER.
The same applies to CPUs: sharing a core dynamically among threads means a thread can be paused for varying lengths of time — but it utilizes the hardware better than a static allocation. Better utilization is also why cloud platforms run several VMs from different customers on the same physical machine.
Variable delays in networks are NOT A LAW OF NATURE but simply the result of a COST/BENEFIT TRADE-OFF.
Latency guarantees ARE achievable if resources are STATICALLY PARTITIONED — but at the cost of reduced utilization, i.e. MORE EXPENSIVE. Multitenancy with dynamic partitioning gives better utilization, so it is cheaper, with the downside of variable delays.
Hybrid attempts: ATM (an Ethernet competitor in the 1980s, little adoption outside telephone core switches); InfiniBand (end-to-end flow control at the link layer, reducing queueing — though it can still suffer link-congestion delays); QoS mechanisms (packet prioritization/scheduling, admission control/rate-limiting) which can emulate circuit switching or provide statistically bounded delay; L4S; Linux TC.
However, such QoS mechanisms are NOT currently enabled in multitenant datacenters and public clouds, or on the internet. Currently deployed technology does not allow us to make ANY guarantees about delays or reliability. CONSEQUENTLY, THERE'S NO "CORRECT" VALUE FOR TIMEOUTS — they need to be determined experimentally.