PrepZone Logo
PrepZone

Distributed Locks and Fencing Tokens

Why naive Redis locks fail and how monotonic fencing tokens prevent stale leader writes.

Why this matters

  • Locks are easy to get wrong and hard to debug — paused GC, clock skew, and network delays create subtle double-execution.
  • Martin Kleppmann's fencing token article is interview gold — storage must validate tokens, not just the lock service.
  • ShardPay uses locks only for batch reconciliation jobs, not per-transfer — idempotency and single-leader writes cover the hot path.
  • Redlock and TTL-only patterns fail without fencing — know why before proposing locks in design reviews.
Old leaderToken 41 — rejected
StorageChecks token > last write
New leaderToken 42 — accepted
Monotonic token rejects writes from a leader that lost the lease.

Locks that actually work

ShardPay's end-of-day reconciliation job acquires a distributed lock before scanning unsettled transfers. The lock has a 30-second lease renewed every 10 seconds while the worker is healthy. The object store that receives reconciliation reports checks a monotonic fencing token on every write — a paused worker whose lease expired cannot overwrite a newer run's output.

Per-transfer locking was rejected: leader serialization plus idempotency keys achieve exclusivity without Redis in the payment path.

Lock primitives

  • Lease — Lock with TTL; holder must renew before expiry. ShardPay's reconciliation worker calls renewLease() on a heartbeat thread; missed renewal releases lock for standby.
  • Fencing token — Monotonic number issued with lock grant; storage rejects writes with token < highest seen. ShardPay's S3 writer stores fenceToken in object metadata; stale workers get 412 Precondition Failed.
  • Linearizable lock service — Lock grant must be consistent (etcd, ZooKeeper, Consul with CP semantics). ShardPay uses etcd transactional lock: create /locks/reconcile if not exists with lease.
  • Redlock critique — Multiple independent Redis masters without shared fencing — clock drift and delayed processes can double-acquire. ShardPay does not use Redlock for financial exclusivity; etcd + fencing for batch jobs only.
  • Prefer alternatives — Single Raft leader, idempotency keys, compare-and-set, database SELECT FOR UPDATE. ShardPay transfers use leader append + transferId dedup — no lock.

Safe lock pattern

Walkthrough: Reconciliation with fencing

  1. Worker W1 acquires etcd lock /locks/reconcile; receives fenceToken=42, lease 30s.
  2. W1 scans ledger for unsettled transfers; GC pauses worker for 45 seconds — lease expires.
  3. Standby W2 acquires lock; receives fenceToken=43.
  4. W1 resumes, completes scan, attempts PUT report.xml with token 42.
  5. Object store rejects: maxFenceSeen=43 — W1's write discarded.
  6. W2 writes report with token 43 — correct single reconciliation output.
Java
// ShardPay etcd lock with fencing token
public LockHandle acquireLock(String path) {
    long token = tokenGenerator.incrementAndGet();
    Lease lease = etcdClient.grant(30);
    boolean acquired = etcdClient.txn()
        .If(new Cmp(path, Cmp.Op.CREATE, CmpTarget.createRevision(0)))
        .Then(Op.put(path, String.valueOf(token), lease))
        .commit();
    if (!acquired) throw new LockNotAcquiredException();
    return new LockHandle(path, token, lease);
}

// Storage layer — reject stale fence
public void writeReport(String key, byte[] data, long fenceToken) {
    if (fenceToken < store.getMaxFence(key)) {
        throw new StaleFenceException(fenceToken);
    }
    store.put(key, data, fenceToken);
}

When locks are and are not appropriate

Decision guide

  • Batch singleton jobs — One reconciliation, one settlement file generator. ShardPay locks here — duplicate runs waste compute; fencing prevents corrupt output.
  • Per-request payment processing — Never lock per transfer at scale. ShardPay uses partition leader + idempotency — 50k TPS through Redis locks is a bottleneck and failure mode.
  • Leader election vs lock — Long-lived leader is implicit lock on write path. ShardPay's ledger shard leader holds write exclusivity by protocol — not a separate lock service.
  • Database advisory locks — pg_advisory_lock for migrations only. ShardPay schema migrations take advisory lock per shard; not used in request path.
  • Lock renewal failure — Treat as fatal: stop work, do not assume lock held. ShardPay worker exits process on failed renew so orchestrator restarts clean standby.

Example: TTL lock without fencing incident

Early ShardPay nightly job used Redis SET reconcile-lock NX EX 60 without fencing. A slow node lost lock, second node started, both uploaded CSV to finance SFTP — duplicate settlement totals. Finance caught $2M double-count before wire. Migration to etcd lease + S3 fence tokens eliminated class of bug.

Quick recall

Everything you need if you only revisit this box.

  1. TTL locks alone are unsafe — storage must reject stale fencing tokens.
  2. ShardPay uses locks only for batch jobs; transfers use leader + idempotency.
  3. Prefer leader election, CAS, or idempotency over distributed locks on hot paths.

Test yourself

Answer these before moving on — recall is what makes it stick.