Distributed systems — the basics
Be able to reason about consistency, failures and queues in systems with several machines.
Prerequisites
- EAsynchronous programmingrequired
- ENetworking — the basicsrequired
Intuition
As soon as the system consists of more than one machine, different rules apply. Three you notice immediately:
1. The network is unreliable. A call that does not answer may have failed — or succeeded with a response that was lost. You cannot know which. So operations have to be idempotent: the same call twice should give the same result as once.
2. CAP: under a network partition you have to choose between consistency (refuse to answer with possibly stale data) and availability (answer with what you have). You cannot have both. A payment system chooses C; a recommendation feed chooses A.
3. Queues decouple. Put slow or unreliable steps (LLM calls, training, email) behind a queue. Then the receiver can be down without the sender noticing.
Formal
Delivery guarantees — none is free:
| Guarantee | Means | The cost |
|---|---|---|
| At-most-once | sent once, can be lost | the simplest |
| At-least-once | delivered, can be duplicated | requires an idempotent receiver |
| Exactly-once | delivered exactly once | requires transactions or deduplication by id; expensive and often an illusion |
In practice you build at-least-once plus idempotence, which gives exactly-once behaviour at a reasonable cost.
Patterns that solve most of it:
- An idempotency key: the client sends an id; the server remembers the result per id.
- The outbox: write the event in the same database transaction as the data, then publish asynchronously. It solves «saved in the DB but crashed before publishing».
- A circuit breaker: stop calling a broken dependency and answer in a degraded mode directly.
- Exponential backoff with jitter: otherwise every client retries at the same moment and knocks the service over again.
- A dead letter queue: messages that have failed n times are set aside for review instead of blocking the queue.
AI-grafen's lab runs use exactly this: a run has an id, the result is written once, and a re-run with the same id does not create a second record.
Code
import hashlib, time, random
# An idempotent write: the same key → the same result, no double work
async def create_run(db, user_id: str, lab: str, files: dict, idem_key: str | None = None):
key = idem_key or hashlib.sha256(
f"{user_id}:{lab}:{sorted(files.items())}".encode()).hexdigest()
if existing := await db.fetch_one("SELECT id FROM labs.lab_run WHERE idem_key=%s", (key,)):
return existing["id"] # return the same result
return await db.fetch_one(
"INSERT INTO labs.lab_run (user_id, lab, idem_key) VALUES (%s,%s,%s) "
"ON CONFLICT (idem_key) DO UPDATE SET idem_key=EXCLUDED.idem_key RETURNING id",
(user_id, lab, key))["id"]
# Backoff with jitter — without jitter every client comes back at the same moment
async def call_with_backoff(fn, attempts=5, base=0.5, cap=30):
for i in range(attempts):
try:
return await fn()
except TransientError:
if i == attempts - 1:
raise
wait = min(cap, base * 2 ** i) * (0.5 + random.random())
await asyncio.sleep(wait)
Mastery means
- Reasons about consistency and availability under partition
- Designs for idempotence and retries
- Uses queues to decouple services
Sign in to do the exercises and build your mastery up.
Sources
- Wikipedia — CAP theorem (CC BY-SA 4.0) — CC BY-SA 4.0
- AWS — Exponential backoff and jitter — free to read