Hiring BarSupport

Distributed Coordination — Cross-Cutting Pattern

Pattern35 min read5 diagrams

Technologies that implement this pattern: ZooKeeper & etcd · Apache Kafka · PostgreSQL · DynamoDB · Redis

Why This Matters#

Coordination is the most expensive thing a distributed system can do. Every leader election, lock, lease and distributed transaction is a moment where independent machines must agree — across a network that drops packets, clocks that drift, and processes that pause for seconds in garbage collection without knowing it. The bugs this produces don't show up in tests. They show up as two workers both charging a card, two leaders both writing to the same shard, or an order that shipped without being paid for.

Most candidates treat coordination as a tooling question: "use ZooKeeper for leader election," "use Redis for the lock," "use 2PC for the transaction." Staff engineers treat it as a question of what happens when the coordinator is wrong. Every lease expires while its holder still thinks it's valid. Every lock can be held by two processes during a pause. Every cross-service transaction will have a step that fails after the previous step committed. The design isn't the coordination primitive; it's the fencing, idempotency and compensation that make the system correct when the primitive lies.

The third reframe — the one that separates Staff from Senior most reliably: the best coordination is the coordination you designed away. Partition ownership so each key has exactly one writer. Make operations idempotent so duplicate execution is harmless. Write the event and the state change in one local transaction (the outbox) so you never need a distributed one. Reach for consensus only for the small, slow-changing facts — who owns which partition, which node is leader — and keep it off the request path.

If you can explain why a lock without a fencing token is not a correctness mechanism, when a saga beats two-phase commit, and how the outbox removes the dual-write problem, you are answering at Staff level.

The 60-Second Version#

  • Locks are for efficiency; fencing is for correctness. A lease-based lock can be held by two processes at once (a 10s GC pause outlives a 5s lease). If double execution is unsafe, the resource must reject stale holders via a monotonically increasing fencing token (etcd revision, ZooKeeper zxid, a DB version column).
  • Use a consensus store; never build one. etcd/ZooKeeper clusters of 3 or 5 nodes (tolerating 1 or 2 failures) give linearizable leases with write latency of ~2–10ms in one region. They handle hundreds to low thousands of writes/sec — they are for metadata, not the data path.
  • Leases trade availability for safety with time. Typical TTLs are 10–30s (Kubernetes' default leader-election lease is 15s with a 10s renew deadline and 2s retry). Failover time ≈ TTL + detection; a shorter TTL means faster failover but more false failovers under GC or network blips.
  • Sagas over 2PC across services. 2PC blocks all participants if the coordinator dies between prepare and commit and requires every participant to support XA. Sagas commit locally step by step and undo with compensating actions — accept that intermediate states are visible for seconds to minutes.
  • Idempotency keys turn retries from bugs into features. Store (key → result) for at least the client retry window (24h is common) in the same transaction as the side effect. At-least-once delivery + idempotent consumers = effectively-once.
  • Outbox beats dual writes. Writing to the database and then publishing to Kafka loses or duplicates events on every crash between the two. Write the event to an outbox table in the same transaction; a relay or CDC publishes it with typical lag of 100ms–1s.

The Problem#

Several processes must act as if they were one: exactly one scheduler fires the 9am job, exactly one node owns partition 17, an order is either paid-and-reserved or neither, a message is processed once even though it was delivered three times. Networks partition, processes pause, clocks drift by tens to hundreds of milliseconds, and crashes happen between any two lines of code. Every coordination mechanism is a bet on some bound — lease duration, clock skew, timeout — and production eventually violates every bound. The job is to choose the cheapest mechanism whose failure is safe, and to put the correctness check where the side effect happens, not where the coordinator lives.


Case Studies That Use This Pattern#


The Four Intents#

"We need coordination here" hides four different problems. The mechanisms are not interchangeable.

IntentConstraintStrategyFailure ModeCorrectness Bar
Single active owner (leader, partition owner, scheduler)Exactly one actor mutates a resource at a timeLease from etcd/ZooKeeper + fencing token checked by the resourceTwo leaders during pause/partition; failover gap = TTLNo stale-owner write ever accepted
Mutual exclusion for efficiency (avoid duplicate work)Duplicate work is wasteful, not harmfulBest-effort lock (Redis SET NX PX), short TTLOccasional duplicate executionDuplicates tolerated, rare
Atomic business transaction across servicesOrder + payment + inventory must all happen or be undoneSaga (orchestrated) with compensations; 2PC only inside one DB/vendorPartial completion, failed compensationEventually consistent; every step undoable or retriable
Effectively-once side effectsRetries, redeliveries and replays must not double-applyIdempotency keys + outbox + dedup on consumerDedup window shorter than retry windowSame input → same effect, once

🎯 Staff Move: "Let me separate two questions: do I need a single owner, or do I need the effect to happen once? Most of the time it's the second, and idempotency solves it without any lock at all. I'll only bring in etcd for the part that truly needs one owner — and even then the database checks a fencing token."


The Core Tradeoff#

StrategyWhat WorksWhat BreaksWho Pays
Consensus-backed lease (etcd/ZooKeeper)Linearizable ownership; watches for fast handoffHolder can outlive its lease (pause); cluster is a shared dependencyFailover latency (TTL) paid by users; platform owns the cluster
Redis lock (SET NX PX, Redlock)Fast (~1ms), simpleNo safe fencing by default; async replication can lose the lock on failoverWhoever cleans up double-executions
Fencing tokens at the resourceMakes stale holders harmlessEvery protected resource must check the tokenResource owners must implement and test it
Two-phase commitAtomic across participantsBlocking on coordinator failure; locks held across network round trips; XA support neededThroughput; on-call resolving in-doubt transactions
Saga (orchestrated)No distributed locks; each step local; long-running friendlyIntermediate states visible; compensations can fail; design effort per stepProduct (visible pending states); domain teams writing compensations
Idempotency keysSafe retries end-to-endStorage for keys; key scope/retention mistakesService team; small storage cost
Transactional outboxNo dual-write; events match DB stateRelay/CDC lag; at-least-once publish → consumers must dedupePlatform (CDC pipeline); consumers (idempotency)

Staff Default Position#

Design coordination away first; when you can't, put the correctness check at the side effect, not at the coordinator.

Partition work so each key has a single owner and the owner assignment changes rarely. Make every externally visible operation idempotent with a client-supplied key stored transactionally with its effect. Publish events with the outbox pattern, never with a dual write. For the remaining genuine single-owner needs, take a lease from a managed consensus store (etcd, ZooKeeper, or the lease primitive of a strongly consistent database) and pass its fencing token to every write. Across service boundaries, use an orchestrated saga with explicit compensations; reserve 2PC for participants inside one database or one vendor's transactional boundary.


When to Deviate#

  • Duplicate work is only wasteful — A cache warmer or report generator that occasionally runs twice: a Redis lock with a short TTL is fine. Skip fencing; document that duplicates are tolerated.
  • Single database — If all participants live in one Postgres, use a local transaction. SELECT … FOR UPDATE beats any distributed mechanism. Don't invent a saga inside one database.
  • Strict atomicity inside one vendor — Spanner-style or DynamoDB transactional writes (up to 100 items per transaction) give you atomic multi-item commits without building 2PC. Use them when the scope fits.
  • Ultra-low-latency single leader — Exchanges and sequencers may use a primary/backup with sub-second failover and hardware-level fencing (STONITH, port shutdown) instead of a general-purpose lease service.

The L5 → L6 → L7 Contrast#

BehaviorSenior (L5)Staff (L6)Principal (L7)
First move"Use a distributed lock / ZooKeeper""Do I need a single owner, or an effect that happens once? Can partitioning or idempotency remove the need?""How many teams are running their own lock and election code, and which of those are correctness bugs waiting on a GC pause?"
LocksRedis lock with TTLLease + fencing token enforced at the resource; lock is efficiency, fencing is correctnessPaved-road coordination library; hand-rolled locks on caches banned for correctness use
Transactions2PC or "wrap it in a transaction"Orchestrated saga, compensations designed per step, pending states in the UXWorkflow engine as a platform; domain teams own compensations, platform owns durability and replay
RetriesRetry with backoffIdempotency keys stored with the effect, retention ≥ retry windowOrg-wide idempotency contract on every mutating API; enforced in the API gateway and linter
Failure"ZooKeeper handles failover"Names the split-brain window (TTL + pause) and what the resource does with a stale tokenTreats the consensus cluster as shared fate for N teams: cells, quotas, upgrade discipline, correlated-failure drills
OwnershipService teamService owns semantics; platform owns etcd/ZKCoordination platform with SLOs and quotas; decides what not to centralize
Why "First move" separates levels

The L5 answer reaches for a primitive. It's usually a reasonable primitive. The Staff answer questions whether the primitive is needed at all: most "we need a lock" cases are really "we need this to happen once," which idempotency solves with no coordination and no failover gap. The Principal answer looks at the org and finds that the most common source of double-execution incidents is fifteen teams' Redis locks — and fixes the class, not the instance.

Why "Locks" separates levels

A lock service can only tell you who held the lock when it granted it. It can't stop a paused process from waking up and writing after its lease expired. Only the resource being written can reject that write — by comparing a fencing token. Candidates who say "the lock has a TTL, so it's safe" are describing the exact bug. This is the single most reliable coordination question that separates L5 from L6.

Why "Transactions" separates levels

2PC is correct and blocking. In a microservice world it also requires every participant to hold locks across network round trips and support a prepare phase — which third-party APIs (payment processors, shipping carriers) don't. The Staff answer accepts eventual consistency with explicit compensations and designs what the user sees in between ("Payment pending"). The Principal answer notices every team is writing its own state machine with its own retry bugs and invests in one durable workflow platform.


The Five Fault Lines#

#Fault LineThe Tension
1Lease Duration: Fast Failover vs False FailoverShort TTL = quick recovery and more spurious leader changes
2Coordinator-Enforced vs Resource-Enforced SafetyTrusting the lock service vs making the resource reject stale actors
32PC vs SagaAtomicity with blocking vs availability with visible intermediate states
4Dual Write vs Outbox vs Event SourcingSimplicity vs consistency between state and events
5Centralized Coordination vs Partitioned OwnershipOne coordinator for everything vs coordination only at partition-assignment time

Fault Line 1: Lease Duration#

A lease TTL of 5s gives ~5–7s failover but a 6s GC pause or a network blip causes a spurious leader change, each of which costs warm-up and may cause a brief dual-leader window. A 30s TTL rarely flaps but leaves the role vacant for 30s after a real crash. Who pays: short TTL — the on-call chasing flapping and the data path absorbing extra handoffs; long TTL — users during the failover gap. Staff default: 10–15s TTL, renew at 1/3 TTL, step down voluntarily if renewal fails at 2/3 TTL, and fence every write. Deviate when: the role is latency-critical (sequencer) — use sub-second heartbeats with a dedicated network path and hard fencing.

Fault Line 2: Coordinator-Enforced vs Resource-Enforced Safety#

Diagram: Fault Line 2: Coordinator-Enforced vs Resource-Enforced Safety

Without the token check in the last step, both writes succeed and the data is corrupted — even though the lock service behaved perfectly. Who pays: without fencing, whoever reconciles corrupted state; with fencing, resource owners who must store and compare a token. Staff default: every write to a lease-protected resource carries the token; the resource performs WHERE token <= :my_token (or equivalent conditional write). Deviate when: duplicate execution is only wasteful — then fencing is optional.

Fault Line 3: 2PC vs Saga#

2PC gives atomic commit but a coordinator crash after PREPARE leaves participants holding locks, "in doubt," until the coordinator recovers. Sagas never hold cross-service locks but expose intermediate states (order created, payment pending) and require a compensation for every step that can't be retried forward. Who pays: 2PC — throughput and availability; sagas — product must design pending states, domain teams must write and test compensations. Staff default: saga across services, orchestrated (a durable workflow owns the state machine) rather than choreographed once there are more than ~3 steps. Deviate when: all participants are inside one database or a system that provides transactions natively.

Fault Line 4: Dual Write vs Outbox vs Event Sourcing#

"Update DB, then publish to Kafka" fails in two ways: crash after commit → event lost; publish succeeds, commit fails → phantom event. Who pays: downstream consumers and whoever reconciles. Staff default: transactional outbox, published via polling relay or CDC; consumers are idempotent because publishing is at-least-once. Deviate when: the domain is already event-sourced (the event log is the state) — then there's no second write to reconcile.

Fault Line 5: Centralized Coordination vs Partitioned Ownership#

Taking a lock per request puts the coordinator on the hot path: 10K req/sec × 2 round trips to etcd = more than a small etcd cluster should serve. Assigning partition ownership (partition → owner) via consensus and letting each owner act locally moves coordination to rebalance time only. Who pays: centralized — latency on every request and the platform team; partitioned — complexity of rebalancing and hand-off. Staff default: coordinate ownership, not operations. Deviate when: operations are rare (a nightly job) — a per-operation lease is fine.


Common Interview Mistakes#

What Candidates SayWhat Interviewers HearWhat Staff Engineers Say
"We'll take a Redis lock with a TTL""Hasn't considered pauses, failover, or fencing""The lock avoids duplicate work; correctness comes from a fencing token the database checks."
"Use 2PC to keep services consistent""Hasn't operated in-doubt transactions""Across services I'd use an orchestrated saga with compensations; 2PC only inside one database."
"Exactly-once delivery with Kafka""Confuses delivery with effect""Delivery is at-least-once. Effects are once because consumers dedupe on an idempotency key."
"Write to the DB then publish the event""Dual-write bug""Outbox: the event row commits with the state change; CDC publishes it."
"ZooKeeper guarantees only one leader""Doesn't know the old leader may not know it's deposed""ZooKeeper guarantees one session holds the node. The old leader may still be running — so writes carry the zxid as a fence."
"Retry until it succeeds""Duplicate side effects""Retries carry the same idempotency key; the server returns the stored result for a repeat."

Quick Reference#

Diagram: Quick Reference

Staff Sentence Templates#

"Before I pick a lock, I want to know whether duplicate execution of [operation] is harmful or just wasteful. If it's harmful, the lock is an optimization and the real guarantee is [idempotency key / fencing token] enforced by [the resource]."

"Leadership comes from a [15s] lease in etcd. Failover is roughly [TTL + detection ≈ 20s]; during that window [work queues up / reads continue]. Every write from the leader carries the lease revision, and [the database] rejects anything older than the highest it's seen."

"This spans [N] services, so I'll run it as an orchestrated saga. Each step is idempotent; [steps X, Y] have compensations; [step Z] is the pivot — after it, we only retry forward. The user sees [pending state] for up to [T seconds]."

"I won't publish the event separately from the write. It goes into an outbox row in the same transaction, CDC ships it within about a second, and consumers dedupe on [event ID]."


Implementation Deep Dive#

1. Leader Election with an etcd Lease + Fencing#

etcd's lease + transaction API gives a leader election whose revision number is a natural fencing token: it only ever increases.

LEASE_TTL = 15                             # seconds
RENEW_EVERY = 5                            # 1/3 of TTL
STEP_DOWN_AFTER = 10                       # renew failing for 2/3 TTL -> stop acting

function campaign(node_id):
    lease = etcd.grant(ttl=LEASE_TTL)
    # Atomic: create the key only if it doesn't exist
    txn = etcd.txn(
        if   = [ create_revision("/election/scheduler") == 0 ],
        then = [ put("/election/scheduler", node_id, lease=lease) ],
        else = [ get("/election/scheduler") ])
    if txn.succeeded:
        fence = txn.header.revision        # monotonically increasing token
        start_keepalive(lease, every=RENEW_EVERY, on_fail_for=STEP_DOWN_AFTER, then=step_down)
        return Leader(fence)
    watch("/election/scheduler", on_delete=campaign)   # re-campaign when it disappears

# Every side effect carries the fence; the resource enforces it
function leader_write(job_id, payload, fence):
    rows = db.exec("""
        UPDATE jobs SET state = $1, fence = $2
        WHERE id = $3 AND fence <= $2""", payload, fence, job_id)
    if rows == 0:
        metrics.incr("fence.rejected")
        step_down()                         # someone newer is leader

Numbers: Kubernetes' client-go leader election defaults (15s lease, 10s renew deadline, 2s retry) are a good reference point; failover lands around 15–20s. etcd recommends 3 or 5 members; a 5-member cluster tolerates 2 failures, and each write needs a majority fsync, typically a few milliseconds on SSD within one region.

🎯 Staff Insight: The step-down timer is as important as the fence. A leader that can't renew for 2/3 of its TTL should stop issuing side effects before the lease expires, so the gap between "old leader stops" and "new leader starts" is positive rather than overlapping. The fence covers the cases where the step-down timer itself is paused.

2. Idempotency Keys — Stored With the Effect#

CREATE TABLE idempotency_keys (
    tenant_id     bigint,
    key           text,
    request_hash  bytea,            -- detect same key reused with a different body
    status        text,             -- 'in_progress' | 'completed'
    response      jsonb,
    created_at    timestamptz default now(),
    PRIMARY KEY (tenant_id, key)
);
function handle(req):                          # POST /charges, Idempotency-Key: k
    BEGIN
      row = INSERT INTO idempotency_keys (tenant_id, key, request_hash, status)
            VALUES (req.tenant, req.key, hash(req.body), 'in_progress')
            ON CONFLICT DO NOTHING RETURNING *
      if row == null:                          # key seen before
          existing = SELECT ... FOR UPDATE
          if existing.request_hash != hash(req.body): return 422  # key reuse, different body
          if existing.status == 'completed':  return existing.response   # replay
          return 409, retry_after=1            # concurrent in-flight duplicate
      result = ledger.debit(req)               # side effect in the SAME transaction
      UPDATE idempotency_keys SET status='completed', response=result
    COMMIT
    return result

Retention: keep keys at least as long as any client may retry — Stripe documents that its idempotency keys can be pruned after 24 hours. Size: 10M mutating requests/day × ~300 bytes ≈ 3 GB/day, trivially partitioned by day and dropped.

When the side effect is external (a PSP call that can't join your transaction), split it: record intent, call the external API with your idempotency key passed through, then record the result. The external provider's own idempotency support is what makes the retry safe.

3. Transactional Outbox + CDC#

# Service code: one local transaction, no broker call
BEGIN
  INSERT INTO orders (id, status, ...) VALUES ($1, 'PLACED', ...)
  INSERT INTO outbox (id, aggregate_id, type, payload, created_at)
         VALUES (gen_uuid(), $1, 'OrderPlaced', $2, now())
COMMIT

# Relay option A: CDC (Debezium reading the WAL) -> Kafka topic 'orders.events'
#   - key = aggregate_id  (per-order ordering within a partition)
#   - lag: ~100ms-1s typical; alert on outbox.publish_lag_ms p99 > 5s

# Relay option B: polling publisher
every 200ms:
    rows = SELECT * FROM outbox WHERE published_at IS NULL
           ORDER BY created_at LIMIT 500 FOR UPDATE SKIP LOCKED
    for r in rows: kafka.send(topic, key=r.aggregate_id, value=r.payload, headers={event_id: r.id})
    kafka.flush()
    UPDATE outbox SET published_at = now() WHERE id IN (rows.ids)

# Consumers: dedupe on event_id (processed_events table or idempotent upsert)

Why it's safe: the event exists if and only if the state change committed. Publication is at-least-once (a relay crash after send and before UPDATE republishes), so consumers must be idempotent — which they need to be anyway, because Kafka redelivers on rebalance.

4. Orchestrated Saga with a Pivot Step#

Diagram: 4. Orchestrated Saga with a Pivot Step
saga PlaceOrder(order):                       # runs on a durable workflow engine
    step reserve_inventory   compensate release_inventory
    step authorize_payment   compensate void_authorization
    step capture_payment     # PIVOT: after this, never compensate — only retry forward
    step create_shipment     retry(forever, backoff=exp(1s..10m)), alert_after=1h

# Every step call carries idempotency key = saga_id + step_name
# Compensations are themselves idempotent and retried; a failed compensation pages a human

Design rules: order steps so the ones most likely to fail and cheapest to undo come first; put the irreversible step (capture, email, shipment) as late as possible — the pivot. Steps after the pivot must be retriable forward. Temporal (and Uber's Cadence, its predecessor) are the widely used open-source engines for this style.

Mechanism Comparison

MechanismLatency AddedThroughputFailure SemanticsCorrect Under Pause?
Local DB transaction1–5ms10K+ tx/s per nodeAtomicYes
etcd/ZooKeeper lease2–10ms per acquire/renewHundreds–low thousands writes/sLinearizable grant; holder may outlive leaseOnly with fencing
Redis lock~1ms100K+ ops/sLost on async failoverNo
2PC (XA)2+ round trips, locks held throughoutLow; bounded by slowest participantBlocks on coordinator failureYes, but blocking
Saga (orchestrated)Per step; total seconds–hoursHigh (no cross-service locks)Eventually consistent, compensationsYes, with idempotent steps
Outbox + CDC100ms–1s publish lagBounded by WAL/relay (~10K+ events/s)At-least-onceYes, with idempotent consumers

Architecture Diagram#

Diagram: Architecture Diagram

How to narrate it: consensus (etcd) appears in only two places — who leads the scheduler and who owns which partition — and never on the per-request path. Everything the request touches is a local transaction; cross-service consistency comes from the outbox, the saga and idempotency keys.


Failure Scenarios#

1. The GC Pause Double-Charge — Two Schedulers, One Billing Run#

A billing scheduler held leadership via a Redis lock (SET NX PX 30000). A full GC on the leader paused it for 38 seconds during the monthly billing run.

t=0      Leader A starts billing batch 1 of 40 (lock TTL 30s, renew every 10s).
t=+5s    Full GC pause begins on A (heap 28 GB, old-gen compaction).
t=+30s   Lock expires. Standby B acquires it and starts billing from batch 1.
t=+43s   A resumes, unaware it lost the lock; continues batch 2..
t=+45s   Both A and B issue charges for batches 1-9.
t=+8min  Payments anomaly alert: charge volume 1.9x expected for the hour.
t=+12min A killed manually. 11,400 customers double-charged.
t=+3d    Refunds complete; support and PSP fees ~$60K; trust damage unpriced.

Detection: billing.charges_per_customer_per_run > 1; leader.count from heartbeat metrics > 1; JVM gc.pause_ms > lease TTL/3. Blast radius: every customer in the overlapping batches. Mitigation: kill the stale leader; refund via ledger reconciliation. Prevention: idempotency key per (customer_id, billing_period) enforced by the ledger (the real fix — makes double execution harmless); lease from etcd with revision fencing on the billing-run table; step-down if renewal fails at 2/3 TTL; GC tuning to cap pauses. Owner: billing team owns idempotency; platform team owns the coordination library that made Redis locks look safe.

🎯 Staff Insight: The postmortem action item isn't "use ZooKeeper instead of Redis." ZooKeeper sessions time out during a 38s pause too. The fix is that charging a customer twice for the same period must be impossible at the ledger, regardless of how many schedulers think they are leader.

2. The Compensation That Couldn't — Saga Stuck Between States#

A travel booking saga reserved a flight, a hotel and a car. The car step failed; compensation to release the hotel called a partner API that had changed its cancellation endpoint.

t=0      Saga: flight reserved, hotel reserved, car reservation fails (sold out).
t=+1s    Compensate hotel: partner returns 404 on cancel endpoint (API v2 migration).
t=+1s    Orchestrator retries with backoff: 1s, 2s, 4s ... capped at 10 min.
t=+6h    1,900 sagas stuck in COMPENSATING. Hotels billing no-show fees.
t=+7h    Alert fires on saga age (threshold was 6h). Team patches client.
t=+9h    Stuck sagas drain; ~$95K in partner no-show fees disputed.

Detection: saga.state_age_seconds{state=COMPENSATING} p99 > 15 min; saga.compensation_failures rate. Blast radius: every saga that hit a car failure during the window. Mitigation: fix client; replay compensations (idempotent); manual cancel via partner portal for the oldest. Prevention: contract tests for compensation endpoints in CI; compensation failures page after 3 attempts rather than retrying silently for hours; car step (most likely to fail) moved before hotel. Owner: booking team owns compensations; partner-integration team owns contract tests.

3. The Outbox That Stopped — CDC Replication Slot Lag#

A Postgres outbox published through a logical replication slot. The CDC connector crashed on a malformed payload and stopped consuming; the slot retained WAL.

t=0      Connector fails on event with invalid UTF-8; enters crash loop.
t=+20min No events published. Downstream search index, emails, analytics stale.
t=+3h    Replication slot retains 180 GB of WAL; primary disk 85 percent.
t=+3.5h  Disk alert pages DB on-call (first page of the incident).
t=+4h    Poison event skipped to a DLQ; connector resumes; 3.1M events backfill.
t=+4.5h  Backlog drained; downstream consumers dedupe replays by event_id.

Detection: outbox.publish_lag_ms p99 > 5s (should have paged at t=+1min); pg_replication_slots retained WAL bytes > 20 GB. Blast radius: every downstream consumer, and nearly the primary database itself (disk full stops writes). Mitigation: DLQ for poison events; max_slot_wal_keep_size to cap WAL retention. Prevention: publish-lag SLO with paging; payload validation at write time; connector restart backoff with poison-message skip. Owner: data platform team owns CDC; order team owns payload schema.

Operational Reality Matrix#

FailureDetection SignalBlast RadiusMitigationOwner
Split-brain leader (pause, partition)leader.count > 1; fence.rejected > 0All writes by the roleFencing token at resource; step-down timerService team + platform library
Consensus cluster loss of quorumetcd.has_leader = 0; proposal failuresEvery team using leases/elections5-node cluster across 3 AZs; clients keep current role until lease expiryCoordination platform
Lease flappingleader.transitions > 3/hourWarm-up cost, brief dual-activityLonger TTL, GC tuning, dedicated networkService team
In-doubt 2PC transactionxa.prepared_age_seconds > 60Locked rows on all participantsCoordinator recovery, manual resolveDB team
Saga stuck compensatingsaga.state_age_seconds p99 > 15 minBusiness exposure per stuck sagaPage after 3 failures, manual runbookDomain team
Idempotency window too shortDuplicate effect with same key after purgeDouble side effectsRetention ≥ max client retry windowService team
Outbox/CDC stalledoutbox.publish_lag_ms > 5sAll downstream consumers; primary diskDLQ poison events, cap WAL retentionData platform

The Principal Lens#

Why L7 Sees This Problem Differently#

A Staff engineer makes one system's coordination correct. A Principal engineer sees that coordination bugs are an org-level defect class: fifteen teams have written leader election, nine use Redis locks as if they were correctness mechanisms, four maintain their own saga state machines with their own retry bugs, and two share an etcd cluster with Kubernetes' control plane. Each team's design review passed; the org still double-charges customers once a year. At L7, distributed coordination becomes a paved-road question — which primitives the org provides with SLOs, which patterns it mandates (idempotency on every mutating API, outbox for every event), and which it forbids — plus a shared-fate question, because one consensus cluster serving thirty teams is the single most correlated failure point in the company.

The Org-Level Fault Line#

A shared coordination platform vs per-team coordination.

OptionWhat WorksWhat BreaksWho Pays
Each team runs its own locks/elections/sagasAutonomy, isolated blast radiusSame correctness bugs rediscovered per team; Redis used as a lock service; no one owns the etcd upgradeCustomers (double effects); every team's on-call
One big shared etcd/ZooKeeper for everyoneOperational expertise concentratedNoisy tenants (a team writing 5K keys/s) degrade elections for all; one quorum loss is an org-wide eventEveryone at once
Coordination platform: library + cells + workflow engineCorrect-by-default leases with fencing, idempotency middleware, managed outbox/CDC, durable workflows; clusters per cell with quotasPlatform team is a dependency; migrations take quartersPlatform headcount (4–8), migration effort per team

The Principal default is the third, with one deliberate non-standardization: domain teams own their compensation logic and idempotency scopes. Centralizing the business semantics of "undo" would make the platform team a bottleneck on every product change.

Cost Model#

Assumptions: managed or self-run etcd/ZooKeeper on 3–5 small instances (~$150–400 each/month), workflow engine self-hosted on its own database, CDC on managed Kafka Connect, fully loaded engineer ~$25K/month. Incident cost of a double-effect event estimated from refunds, fees and support time.

ScaleCoordination NeedsInfra $/monthPeopleOn-call LoadTypical Incident Exposure
Startup (10 services)Postgres transactions, advisory locks, idempotency table, polling outbox~$0 incremental (lives in existing DB)0.2 FTEService team rotationLow; fix-forward
Growth (100 services)1 etcd cluster (5 nodes) for leases, CDC pipeline, workflow engine for 3–5 sagas~$6–12K (etcd, Kafka Connect, workflow DB + workers)2–3 FTE platformShared platform rotation, ~2 pages/month~$50–100K per double-charge class incident
Large (1,000+ services)etcd per cell/region (6–12 clusters), org-wide idempotency middleware, multi-tenant workflow platform, CDC for hundreds of tables~$60–150K6–10 FTE across coordination + workflow + CDCDedicated rotations; quarterly quorum-loss drillsSeven figures if a shared cluster loses quorum at peak

The Principal framing: the platform's cost is dominated by people, not machines. Its ROI comes from not having the incident: if the org averages one double-effect incident per year at $100K plus a week of three engineers' time, a 3-person platform ($900K/year) is justified only when it also removes the per-team cost of building and operating coordination — typically 0.25–0.5 FTE per team that would otherwise hand-roll it. Across 20 teams that's 5–10 FTE saved.

The 3-Year Evolution Path#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionDoor TypeReversibility Cost
Lease TTL, renew intervalTwo-wayConfig change
Adding an outbox to a serviceTwo-wayDays; old dual-write path removed after
Public API idempotency-key semantics (scope, retention, error on body mismatch)One-wayClients depend on replay behavior; changing it breaks retries silently
Choosing 2PC across servicesOne-way-ishEvery participant couples to prepare/commit protocol; unwinding takes quarters
Workflow engine choice (durable history format)One-way-ishRunning workflows can last weeks; migration needs drain or replay
Sharing the Kubernetes control-plane etcd with application workloadsTwo-way to fix, one-way in the incidentCluster instability takes down scheduling and app coordination
Event schema published via outboxOne-way per fieldConsumers depend on it; evolve additively only

The Standard I'd Write#

RFC: Coordination and Effectively-Once Standard (v1)

Scope: All services that mutate state visible to customers, money, or other services.

MUST:

  1. Every mutating external and internal API accepts an idempotency key, stores it transactionally with the effect, and retains it ≥ 24h (or ≥ the longest client retry window, whichever is greater).
  2. Events derived from state changes MUST be published via transactional outbox (CDC or polling relay). Dual writes to a DB and a broker are prohibited.
  3. Single-owner roles MUST use the platform coordination library (etcd-backed lease). Every write from a lease holder MUST carry the fencing token, and the resource MUST reject stale tokens.
  4. Redis or cache-based locks MAY be used only where duplicate execution is harmless, and the design doc MUST say so.
  5. Cross-service transactions MUST use the workflow platform; each step MUST be idempotent and either compensatable or after a declared pivot.

SHOULD: Keep consensus off the per-request path; coordinate ownership, not operations. Alert on fence.rejected, outbox.publish_lag_ms, and saga.state_age_seconds.

Exceptions: Platform review; recorded with owner and expiry (max 2 quarters).

Success metrics: zero double-effect incidents per year; % of mutating endpoints with idempotency (target 100% of external, 90% internal); number of hand-rolled election/lock implementations trending to zero; p99 outbox lag < 2s.

What I'd Tell the VP#

Our systems run on many machines that must agree on who does what, and when they disagree we charge customers twice, lose orders, or send duplicate emails. Each team has been solving this on its own, and a few of those solutions have a flaw that only shows up under rare timing — which is how we had last year's billing incident. I'm proposing one supported way to do it: a small platform that makes retries safe, keeps our records and our event stream in sync, and handles multi-step transactions with automatic undo. It's roughly three engineers for a year, paid back by removing a class of costly incidents and by letting twenty teams stop maintaining their own versions. I need product teams to budget a few weeks each to migrate.

Principal Interview Signals#

SignalWhat It Sounds Like
Defect-class thinking"Double execution isn't one team's bug; it's the org's default until idempotency is mandatory on every mutating API."
Shared-fate awareness"Thirty teams on one etcd means one quorum loss is thirty incidents. I'd split by cell and give tenants quotas."
Deliberate non-standardization"The platform owns durability and replay; domain teams own what 'undo' means. Centralizing compensations would make us the bottleneck on every product change."
Pricing the paved road"Each hand-rolled coordination stack costs a team a quarter of an engineer to operate. Across twenty teams the platform pays for itself before counting incidents."
One-way door recognition"Idempotency-key semantics in the public API are forever. I'd get scope, retention and mismatch behavior right before launch."

Staff answers that L7 interviewers find insufficient:

  • "I'd add fencing tokens to this scheduler." — Correct locally; ignores the other fourteen hand-rolled elections in the org.
  • "We'll run Temporal for this saga." — Picks a tool for one flow without asking who operates it, who else needs it, or how running workflows migrate later.
  • "etcd is highly available with five nodes." — Doesn't ask what else shares that cluster or what happens to every tenant when it loses quorum.

In the Wild#

Google: Chubby and Sequencers#

Google's Chubby paper (OSDI 2006) describes a coarse-grained lock service built on Paxos, used for leader election and small metadata by systems like GFS and Bigtable. Crucially, the paper acknowledges that a lock holder can lose its lock while still acting — and introduces sequencers: opaque tokens describing the lock's state and generation that clients pass to servers, which check them before acting on requests.

Staff insight: The most influential lock service in the industry explicitly says the lock alone is not enough and pushes the check to the resource. Citing Chubby sequencers is the strongest way to justify fencing tokens in an interview.

Kubernetes: Leader Election via Lease Objects#

Kubernetes controllers (the controller-manager, scheduler, and most operators) run multiple replicas and elect a single active leader using Lease objects stored in etcd, renewed periodically, with defaults of a 15s lease duration, 10s renew deadline and 2s retry period. Standbys take over when the lease isn't renewed.

Staff insight: This is the mainstream, battle-tested lease pattern — and the numbers (15s/10s/2s) are a defensible default to quote. It also shows the tradeoff honestly: Kubernetes accepts ~15s of controller unavailability on failover in exchange for avoiding flapping.

Stripe: Idempotency Keys on Every Mutating Request#

Stripe's public API accepts an Idempotency-Key header on POST requests; the server saves the status code and body of the first request for a key and returns the same result for retries, and it rejects reuse of a key with different parameters. Keys can be pruned after 24 hours. Stripe's engineering blog has also written about implementing idempotency with database transactions and "atomic phases."

Staff insight: Stripe made effectively-once a property of the API contract, not of internal locking. That's the design-coordination-away move at the largest scale: clients retry freely, and correctness lives with the effect.


Practice Drill#

Prompt: "We run a job scheduler with three replicas. One is elected leader via a Redis lock and fires jobs, including 'send monthly invoice.' Last month some customers got two invoices and a few were charged twice. Fix it."

Staff Answer

Two invoices means two leaders acted at the same time — almost certainly a pause or Redis failover that let the lock be held twice. Switching lock vendors won't fix it; every lease can be outlived by its holder. I'd fix it at three layers. (1) Make the effect idempotent: the invoice and charge are keyed by (customer_id, billing_period) with a unique constraint in the ledger and an idempotency key passed to the payment processor. A second scheduler firing the job produces a no-op, not a charge — this alone closes the incident class. (2) Make leadership fenced: move the election to an etcd lease (15s TTL, renew every 5s, step down if renewal fails for 10s) and use the lease revision as a fencing token; the job_runs table accepts a state transition only if fence >= stored_fence. (3) Make job firing claim-based: a run is INSERT INTO job_runs (job_id, scheduled_for) ON CONFLICT DO NOTHING, so each scheduled instance can be claimed exactly once even with two active schedulers. Failover is ~15–20s; jobs scheduled in that window fire late, not twice — product signs off that invoices may be up to a minute late. Metrics: leader.count, fence.rejected, job_runs.duplicate_claims, billing.charges_per_customer_per_period (alert on > 1), GC pause p99. Owners: billing owns the ledger constraint; platform owns the election library; scheduler team owns the claim table.

Why this is L6:

  • Diagnoses the failure as lease-outliving-holder, not "Redis is bad," and knows ZooKeeper/etcd have the same property.
  • Puts correctness at the side effect (unique constraint, idempotency key, fencing) so the lock becomes an efficiency mechanism.
  • Quantifies failover and makes the late-not-twice tradeoff explicit with a product sign-off.

What L7 adds:

  • Asks which other services use Redis locks for correctness and proposes the org rule: cache-based locks only where duplicates are harmless.
  • Proposes idempotency on every mutating billing/payment API as a standard, with retention tied to the longest retry window.
  • Prices the incident class (refunds + fees + support + trust) against the cost of a shared coordination library, and names the platform team that will own it.

Staff Interview Application#

How to Introduce This Pattern#

"There are really two questions here: does this need a single owner, or does the effect just need to happen once? I'll make every mutating step idempotent first, since that removes most coordination. For the part that genuinely needs one owner, I'll use a lease from etcd and fence writes with its revision. Across services, it's an orchestrated saga with compensations, and events go out through an outbox — no dual writes."

When NOT to Use This Pattern#

  • Everything lives in one database: Use local ACID transactions and row locks. Distributed machinery inside one Postgres is pure risk.
  • Naturally partitioned work: If each key has one writer by construction (Kafka partition → consumer), you already have single ownership; don't add a lock on top.
  • Idempotent, commutative operations: Set-membership adds, max-timestamp wins, CRDT counters — concurrent execution converges without coordination.
  • Duplicates are harmless: Cache warmers, thumbnail regeneration — a best-effort lock or nothing at all.

Follow-Up Questions to Anticipate#

Interviewer AsksWhat They Are TestingHow to Respond
"Why not just use ZooKeeper for the lock?"Understanding that any lease can be outlived"ZooKeeper gives a linearizable grant, but a paused holder doesn't know its session expired. I'd still pass the zxid as a fence and have the resource reject older ones."
"Why not 2PC?"Blocking protocol awareness"Coordinator failure after prepare blocks every participant with locks held, and the payment processor can't participate in XA anyway. Saga with compensations, pivot at capture."
"What if a compensation fails?"Operational realism"Compensations are idempotent and retried with backoff; after 3 failures it pages the owning team, and the saga dashboard shows age per state."
"How long do you keep idempotency keys?"Retention reasoning"At least the longest client retry window — 24h is common. Partition by day, drop old partitions."
"How does the outbox avoid duplicates?"At-least-once honesty"It doesn't — publication is at-least-once. Consumers dedupe on event ID. What it prevents is lost or phantom events."
"How fast is failover?"Quantified tradeoffs"About TTL plus detection: 15–20s with a 15s lease. Shorter TTLs cause spurious failovers on GC pauses."

Evaluation Rubric#

DimensionSenior (L5)Staff (L6)Principal (L7)
FramingPicks a lock or 2PCSeparates single-owner from effect-once; designs coordination awayTreats coordination bugs as an org defect class
CorrectnessTTL-based locksFencing tokens, idempotency with the effect, outboxMandated via standard and middleware
Transactions2PC / "a transaction"Orchestrated saga, pivot step, compensationsWorkflow platform; domain teams own semantics
Failure"ZooKeeper handles it"Split-brain window, stuck sagas, CDC lag — with metricsShared-fate of consensus clusters; cells and quotas
CostNot discussedFailover time vs flappingPlatform headcount vs per-team cost and incident exposure

Strong Hire Signals

SignalWhat It Sounds Like
Fencing reflex"The resource checks the token; the lock is just an optimization."
Idempotency first"Make it safe to run twice, then worry about running it once."
Pivot awareness"Capture is the pivot; after it we only retry forward."
Dual-write detection"That's a dual write — I'd use an outbox."

Lean No-Hire Signals

SignalWhy It Misses the Bar
"Exactly-once delivery" without idempotent consumersConfuses transport with effect
Distributed lock on the per-request pathLatency and availability coupled to the coordinator
No story for a failed compensationSagas will get stuck in production

Common False Positives: Reciting Raft's election rules ≠ designing safe leadership. Naming Temporal ≠ designing compensations. Knowing Redlock's algorithm ≠ knowing when a lock is safe.


Capacity Planning Quick Reference#

Sizing Coordination#

failover_time          ≈ lease_ttl + detection + warmup          # 15s + 2s + warm-up
etcd_write_load        = leases × (1 / renew_interval) + elections + map updates
                         # 2,000 leases renewing every 5s = 400 writes/s — fine
idempotency_storage    = mutating_rps × 86,400 × retention_days × ~300 B
outbox_rows_per_sec    = events_per_tx × tx_per_sec              # prune published rows hourly
saga_concurrency       = saga_rate × avg_saga_duration_s

Key Numbers Worth Memorizing#

NumberContext
3 or 5Consensus cluster size (tolerates 1 or 2 failures)
~2–10 msetcd/ZooKeeper write latency within one region
Hundreds–low thousands/sComfortable consensus write rate; keep it off the request path
15s / 10s / 2sKubernetes lease duration / renew deadline / retry period defaults
1/3 TTLRenew interval; step down at 2/3 TTL without renewal
1–40 sObserved stop-the-world GC pauses on large heaps — longer than many lease TTLs
24 hCommon idempotency-key retention
100 ms–1 sTypical outbox/CDC publish lag
100 itemsDynamoDB transactional write limit per transaction
2+ round trips2PC cost, with locks held for the duration

Common Pitfalls Checklist#

  • Every lock is labeled "efficiency" or "correctness"; correctness locks have fencing at the resource
  • Mutating APIs accept idempotency keys stored in the same transaction as the effect
  • Events are published via outbox, never dual-written
  • Sagas declare a pivot step; pre-pivot steps have tested, idempotent compensations
  • Consensus is used for ownership and metadata, not per-request operations
  • Leaders step down before their lease expires if renewal fails
  • Stuck-saga age, fence rejections and outbox lag are alerted on
  • The consensus cluster's tenants, quotas and quorum-loss runbook are known
  1. Loading the index…