Learn Labs
9. The Trouble with Distributed Systems

9.10 Decision cheat sheet

There is no correct constant. Measure the RTT distribution across many machines over an extended period, pick a target trade-off between detection delay and false-positive rate, a…

How do I set a timeout? There is no correct constant. Measure the RTT distribution across many machines over an extended period, pick a target trade-off between detection delay and false-positive rate, and prefer an adaptive detector (Phi Accrual) over a constant. Then tune the threshold per decision by the cost of being wrong — evicting from a read pool is cheap; failing over a leader is not.

Which clock do I use? Monotonic for every duration — timeouts, latency measurement, rate limiting, retry backoff. Time-of-day only for human-facing points in time. Never order events across machines by wall clock. For ordering, use logical clocks (Ch 10) or version vectors (Ch 6).

Can I trust LWW? Only if you never update existing records (Ch 6 §3.4). Otherwise it is a data-loss mechanism whose loss rate is proportional to your clock skew — and you will not get an error.

Do I need fencing? If two holders of the "same" lease could corrupt data or lose writes — yes, always. A lease alone is not mutual exclusion; the resource must reject stale tokens. If the wasted work from a double execution is merely inefficient (the third example in §5.2), you can skip it.

Do I need Byzantine fault tolerance? Almost certainly not. Datacenter nodes are yours; multitenancy is handled by isolation, not BFT; and BFT cannot protect against a shared software bug or a compromise that reaches all nodes. Do add the cheap "weak lying" defences: application-level checksums, input sanitization and size limits, and multiple NTP servers with outlier rejection.

What system model should I assume? Partially synchronous + crash-recovery, and design so that safety properties hold unconditionally while liveness may be conditioned on a majority surviving and the network eventually recovering. Then add explicit handling for fail-slow nodes, because the model doesn't cover them and they're the hardest case in practice.

Should I go distributed at all? "Distributed systems engineers will often regard a problem as trivial if it can be solved on a single computer — and indeed a single computer can do a lot nowadays. If you can avoid opening Pandora's box and simply keep things on a single machine, IT IS GENERALLY WORTH DOING SO." Go distributed for fault tolerance and low latency, which a single node cannot provide — not reflexively for scale.

How do I gain confidence in correctness? Combine all three: a TLA+ spec for the protocol, Jepsen against the real deployment, and DST if your architecture can support it. And fix what they find.