Testing distributed systems and chaos engineering
How to gain confidence that a distributed system survives real failures: unit and contract tests, integration tests with real dependencies, load and soak tests, fault injection, chaos engineering, game days, deterministic simulation and testing in production safely.
Reading is half of it. See this used in a real interview: walk through Design a Distributed Key-Value Store →
Distributed systems fail in ways unit tests never see: a network partition during a leader election, a slow disk on one replica, a retry storm after a dependency recovers, clocks that disagree. "How would you know this design actually works?" is a fair interview question, especially for infrastructure designs like key-value stores and queues. A credible answer combines conventional tests with deliberate failure injection.
The layers
| Layer | What it catches | Notes |
|---|---|---|
| Unit tests | logic errors in one component | fast; mock nothing you do not own |
| Property-based tests | edge cases you did not think of | generate random inputs and check invariants |
| Contract tests | breaking API or event changes between services | each consumer's expectations checked against the provider |
| Integration tests | wiring with real databases, queues, caches | run real dependencies in containers |
| End-to-end tests | whole user journeys | few, slow, flaky; keep them for critical paths |
| Load and soak tests | capacity limits, leaks, slow degradation | realistic traffic for hours or days |
| Fault injection and chaos | behaviour under failures | the focus of this guide |
Invariants first
Before testing failures, write down what must always hold:
- No acknowledged write is ever lost.
- A payment is never captured twice. See payments and ledgers.
- Every message is delivered at least once, in order per key. See Design a Message Queue.
- Reads after a successful write return that write (if you promise read-your-writes).
Tests and checkers then verify invariants, not specific outputs.
Load testing
- Generate traffic that matches production in shape: the mix of endpoints, key distribution (including hot keys), payload sizes and burstiness.
- Find the breaking point and how the system fails there: gracefully shedding load, or collapsing.
- Soak tests run for hours to reveal memory leaks, growing queues and slow disk fill.
- Replay recorded production traffic (with sensitive data removed) for realism.
See back-of-envelope estimation for the targets to test against.
Fault injection
Inject the failures that really happen:
- Kill processes and nodes; restart them.
- Partition the network between some nodes, in one direction only.
- Add latency, packet loss and slow disks on a single replica (a slow node is often worse than a dead one).
- Skew clocks.
- Fill disks; exhaust file descriptors and connection pools.
- Make a dependency return errors or hang.
Tools range from simple proxies that add latency and errors to platform-level chaos tools. Jepsen-style testing runs a database under partitions and checks histories of operations for consistency violations; it has found bugs in many well-known databases. Use it as inspiration for testing your own key-value store.
Chaos engineering
Chaos engineering runs controlled experiments in production-like environments, or production itself:
- Define steady state with metrics (success rate, latency).
- Form a hypothesis: "if one cache node dies, success rate stays above 99.9 %."
- Inject the failure with a limited blast radius (one zone, a small percentage of traffic) and an abort switch.
- Compare with steady state; fix what broke.
Netflix's Chaos Monkey (randomly terminating instances) popularised the approach; the point is to make failure routine, so systems are built to survive it. See Design a Distributed Cache.
Game days
Teams rehearse major incidents: fail over a region, restore a database from backup, lose a critical dependency. Game days test runbooks, alerts and people, not just software, and they often reveal that a restore takes far longer than assumed. See backups and disaster recovery.
Deterministic simulation
Some teams (FoundationDB is the well-known example) run the entire distributed system inside a deterministic simulator: a single-threaded process that simulates network, disks and clocks, injects faults from a random seed, and checks invariants. Any failure reproduces exactly from its seed. It is a large investment but finds rare concurrency bugs that no other method catches.
Testing in production, safely
Some behaviour only shows up with real traffic:
- Canary releases and gradual rollouts with automatic rollback. See feature flags and A/B testing.
- Shadow traffic: mirror real requests to a new version and compare responses, discarding its output.
- Synthetic monitoring: scripted probes exercising critical journeys around the clock.
In the interview
When asked how you would validate a design: state the invariants, load test to the estimated peak with realistic skew, inject failures (node loss, partitions, slow replicas) and verify invariants still hold, run chaos experiments with small blast radius, and roll out with canaries and shadow traffic.
Checklist
- Invariants written down and checked automatically.
- Contract tests between services; integration tests with real dependencies.
- Load and soak tests with realistic traffic shape.
- Fault injection for crashes, partitions, slow nodes and clock skew.
- Chaos experiments with steady-state metrics and limited blast radius.
- Game days for failover and restore; canaries and shadow traffic for releases.