System Design
T6

Distributed locks and fencing

A lock that expires keeps a crashed client from blocking everyone. The same expiry lets a paused client wake up and write after another client took the lock.

Not startedSaved in this browser only.
  1. 1A lock with an expiry is a lease. A paused client can still act after its lease ends.
  2. 2Only the storage can stop a stale write: it rejects a token lower than one it has seen.
  3. 3A Redis lock is fine for efficiency, to avoid duplicate work. Do not trust it for correctness.
  4. 4For correctness, lock in the same transaction as the data, or fence every write.
T6
    A

    The paused client

    expiry, then a stale write
    0 ms100 ms200 ms300 ms400 ms500 msClient 1Lockin RedisClient 2Storagethe documentpaused: stop-the-world GCtoken 1, TTL 300 msfreetoken 2holds the lockv0by c2by c1stale write✕ At 500 ms both clients act as the holder. c2's write is lost.
    • A garbage collector, swap, a VM migration or a slow network can stop a client for longer than the TTL.
    • The client gets no signal that its lock expired.
    • Times match the recorded run in panel C.

    A lock with a TTL cannot stop a client that paused past the TTL. That client still believes it holds the lock.

    B

    Which lock, for what

    efficiency or correctness
    lockefficiencycorrectness
    Redis SET NX PX, one nodeApprovedNot approved
    Redlock, 5 Redis nodesCosts 5 nodesNot approved
    Any lock plus a fenced writeApprovedDurable tokens
    etcd or ZooKeeper leaseApprovedWith fencing
    Postgres lease row with fenceApprovedStorage checks
    Advisory lock, work in the same transactionApprovedApproved
    A mutex in one processNot approvedNot approved

    Efficiency: a failure runs some work twice. Correctness: a failure corrupts data.

    I ask what a double run costs. If it corrupts data, the lock needs fencing or a transaction.

    C

    See it: fencing off, then on

    recorded from Redis and Postgres
    0 ms100 ms200 ms300 ms400 ms500 ms600 msClient 1LockClient 2Storagetoken 1, expires at 300 msv0

    The document now

    body
    v0
    fence
    not stored

    With fencing, the stale write fails at the storage, and the client learns that it lost the lock.

    D

    The fenced path

    click a step; its path lights up
    Callerjob or requestWorker serviceN copies, any can pauseRedislock and fence counterPostgresstorage, checks the token

    Step 1: Acquire

    • Run the acquire script: set the lock only if it is absent.
    • In the same step, add 1 to the fence counter. That number is the token.
    • The lock expires after its TTL, so a crashed holder frees it.

    If it fails

    The lock is busy: retry after a random delay, or answer "busy". Redis is down: no new work starts.

    I take a token with the lock, send it with every write, and let the storage reject old tokens.

    E

    Capabilities used

    what each tool gives you
    toolcapabilitywhat it gives this designalso used for
    RedisSET with NX and PXCreate the lock only if it is absent, with an expiry, in one command.Idempotency keys, dedupe
    RedisServer-side scripts (EVAL)Compare the token, then delete or extend, as one step.Rate limits, holds
    RedisINCRA counter that only grows: the fencing token.IDs, counters
    RedisKey expiry (TTL)A crashed holder frees the lock with no cleanup job.Sessions, caches
    RedisCluster hash tags: lock:{report}The lock and its counter share one slot, so one script uses both.Any multi-key script
    RedisAsynchronous replicationLimit A failover can lose a lock and repeat a token.
    PostgresTransaction-level advisory lockA named lock with no table. Commit or rollback ends it.Singleton cron jobs, migrations
    PostgresSession-level advisory lockEnds when the connection ends. A transaction-mode pooler breaks it.Leader election on one database
    PostgresINSERT ... ON CONFLICT DO UPDATE ... WHERETake an expired lease and raise its fence in one statement.Upserts, claims
    PostgresConditional UPDATE: fence <= tokenThe storage rejects a stale write. The row count tells the writer.Versioned writes
    etcdLeases with keepalive; a revision on every writeThe lock key ends with the lease. Its revision only grows, so it is a fencing token.Leader election, config
    ZooKeeperEphemeral sequential znodes; zxidWaiters queue in order, and a lost session deletes the lock.Group membership, leaders

    Redis gives me an atomic set-if-absent with expiry and a counter. Postgres gives me locks that end with a transaction, and a conditional write that checks the token.

    F

    Lock and release in Redis

    pseudo code
    acquire, release, extendpseudo code
    acquire(name, ttl):                    // ONE script, atomic1
      IF lock:{name} exists: RETURN busy
      token = INCR fence:{name}            // only grows2
      SET lock:{name} = token, expire after ttl3
      RETURN token
    
    release(name, token):                  // ONE script
      IF GET lock:{name} == token: DEL4 lock:{name}
    
    extend(name, token, ttl):              // ONE script
      IF GET lock:{name} == token: reset expiry to ttl5
    1. 1No other command runs between the check and the set. Two clients cannot both see the lock as free.
    2. 2Each holder gets a higher token than every earlier holder.
    3. 3Set in the same step, so a crash cannot leave the lock forever.
    4. 4Delete only your own lock. After expiry it may belong to another client.
    5. 5Extend before the TTL ends. A false reply means the lock is already gone.
    Tested source Redis script: acquire with a token · Redis script: release · Redis script: extend
    Redis script: acquire with a tokenlua
    -- KEYS[1] = lock key, KEYS[2] = fence counter. ARGV[1] = TTL in ms.
    -- Take the lock only if it is free. The value is a new fencing token, one higher than the last.
    if redis.call('EXISTS', KEYS[1]) == 1 then
      return false
    end
    local token = redis.call('INCR', KEYS[2])
    redis.call('SET', KEYS[1], token, 'PX', ARGV[1])
    return token
    Redis script: releaselua
    -- KEYS[1] = lock key, ARGV[1] = the caller's token.
    -- Delete the lock only if the caller still holds it. After expiry it may belong to another client.
    if redis.call('GET', KEYS[1]) == ARGV[1] then
      return redis.call('DEL', KEYS[1])
    end
    return 0
    Redis script: extendlua
    -- KEYS[1] = lock key, ARGV[1] = the caller's token, ARGV[2] = new TTL in ms.
    -- Extend the lock only if the caller still holds it.
    if redis.call('GET', KEYS[1]) == ARGV[1] then
      return redis.call('PEXPIRE', KEYS[1], ARGV[2])
    end
    return 0
    • Without a token, use SET key random-value NX PX ttl, and compare that value on release.
    • The hash tag puts the lock and its counter in one cluster slot.

    I acquire with a script that sets the lock only if it is free and returns a new token. I release only if the lock still holds my token.

    G

    Fencing at the storage

    figure, then pseudo code
    Client 2token 2Client 1token 1, latedocumentsbody = by c2fence = 2✓ 2 ≥ 2: accepted✕ 1 < 2: rejectedThe check runs inside the write, so no pause can split it.
    a fenced writepseudo code
    write(doc, body, token):
      rows = set body = body, fence = token
             WHERE name = doc AND fence <= token1
      IF rows == 0: RETURN rejected        // a newer holder wrote2
      RETURN accepted
    1. 1The check and the write are one statement. The same holder can write again with its token.
    2. 2Stop the work. Retrying with the same token fails the same way.
    Tested source SQL: the documents table and the fenced write
    SQL: the documents table and the fenced writesql
    -- The resource the lock protects. fence is the highest token that has written.
    CREATE TABLE documents (
      name  text   PRIMARY KEY,
      body  text   NOT NULL,
      fence bigint NOT NULL DEFAULT 0
    );
    
    -- A lock row: a lease with an owner, a fencing token and an expiry.
    CREATE TABLE leases (
      name       text        PRIMARY KEY,
      owner      text        NOT NULL,
      fence      bigint      NOT NULL,
      expires_at timestamptz NOT NULL
    );
    
    -- Accept the write only if no holder with a newer token has written.
    UPDATE documents SET body = $2, fence = $3
    WHERE name = $1 AND fence <= $3;

    The storage stores the highest token that wrote, and it refuses any write with a lower one.

    H

    Redlock and its critique

    5 nodes, one clock jump
    Client 1A, B, C: 3 of 5Client 2C, D, E: 3 of 5Redis Aheld by c1Redis Bheld by c1Redis Cc1, then c2Redis Dheld by c2Redis Eheld by c2clock jumps +11 skey expires earlyTTL 10 s. Both clients hold a majority, so both act as the holder.
    • A client takes the lock on a majority of 5 independent Redis nodes, within the TTL.
    • Each node expires the key by its own clock.
    • A test in the lab replays both failures on 5 simulated nodes and gets two holders each time.
    Redlock assumeswhat breaks itresult
    Clocks run at about the same rateAn NTP step or a manual change on one nodeThat node expires the key early. Two clients can each hold 3 of 5.
    Process pauses are shorter than the TTLA GC pause, swap or VM migrationThe lock expires on every node. A second client takes all 5.
    Network delay is shortA write packet held in a queueThe write arrives after the lock moved on.
    Nodes keep keys after a restartA restart with no persisted dataThe node forgets the lock. A second majority forms.

    Redlock improves availability over one node. It does not fix the paused client, because it issues no token the storage can check.

    Redlock needs bounded pauses, delays and clock drift. When one assumption fails it gives two holders, and it has no token to fence them.

    I

    Leases on a consensus store

    etcd, ZooKeeper
    capabilitywhat it givesstatus
    Writes agreed by a quorum (Raft, ZAB)A failover does not lose a lock.Approved
    Lease or session with keepaliveThe lock ends when the holder stops renewing.Approved
    Revision (etcd), zxid (ZooKeeper)A number that only grows across the cluster: a durable fencing token.Approved
    Watches on the keyWaiters get notified. No polling.Approved
    The lease aloneA paused holder still acts after expiry.Not approved
    • Every lock change is a quorum write, so lock rates are lower than on Redis.
    • Run it as 3 or 5 nodes. It is a cluster to operate, not a library.

    On etcd I use a lease for the lock and the key's revision as the fencing token. The storage must still check it.

    J

    Locks in Postgres

    pseudo code
    advisory lock and lease rowpseudo code
    with_lock(name, work):                // work in the same database
      BEGIN
        IF NOT try_advisory_xact_lock1(name): ROLLBACK; RETURN busy
        work()                              // same transaction
      COMMIT                                // ends the lock too2
    
    lease(name, owner, ttl):                // work outside the database
      insert (name, owner, fence 1, expires now + ttl)
      ON CONFLICT: take it only IF expires_at < now3,
                   and set fence = fence + 14
      RETURN fence, or busy if a live lease exists
    1. 1Returns false at once if another transaction holds the lock. It never waits.
    2. 2Commit or rollback releases the lock. A dead client rolls back, so it cannot write later.
    3. 3One Postgres clock decides expiry. Client clocks do not matter.
    4. 4The row survives release, so every new holder gets a higher token.
    Tested source SQL: lease row and advisory lock
    SQL: lease row and advisory locksql
    -- Take the lease if it is free or expired. Each new holder gets the next fencing token.
    INSERT INTO leases (name, owner, fence, expires_at)
    VALUES ($1, $2, 1, now() + make_interval(secs => $3))
    ON CONFLICT (name) DO UPDATE
    SET owner = EXCLUDED.owner,
        fence = leases.fence + 1,
        expires_at = EXCLUDED.expires_at
    WHERE leases.expires_at < now()
    RETURNING fence;
    
    -- Give the lease back early. Keep the row, so the next token is still higher.
    UPDATE leases SET expires_at = '-infinity'
    WHERE name = $1 AND fence = $2;
    
    -- A lock that lives as long as this transaction. Commit or rollback releases it.
    SELECT pg_try_advisory_xact_lock(hashtext($1));

    If the work writes to the same database, I take an advisory lock in that transaction. Then the lock and the work end together.

    K

    Failure cases

    what breaks, and what keeps the data right
    eventresultwhy the data stays rightsaved by
    A client pauses past the TTLAnother client takes the lock.The storage rejects the old token.Fencing
    A write is delayed in the networkIt arrives after the lock moved on.Same check: its token is lower than the last one that wrote.Fencing
    Redis fails over before it replicatesThe new primary has no lock, and the counter can repeat a token.Take tokens from a durable store: a lease row or an etcd revision.Lease row
    The holder crashesThe lock stays until its TTL.The TTL frees it. The next holder starts the work again.TTL
    A slow client releases after expiryThe key holds another token.The release script deletes nothing.Token check
    A job runs longer than the TTLThe lock is lost during the job.Extend before the TTL ends, and stop if the extend fails. Fence anyway.Extend
    One Redlock node's clock jumpsTwo clients hold a majority.Redlock cannot detect it. Only fencing at the storage can.Fencing
    An advisory lock holder disconnectsPostgres rolls back its transaction.The lock and the work end together.Transaction
    L

    Scale ladder

    start simple; climb only on a signal
    Each step adds one component1Advisory lock2+ lease row, fence3+ Redis lock4+ Redis Clustermore load →
    Lock and release pairs a second, 32 clients1k10k100kPostgres advisory lock: 36,150 pairs per secondPostgres advisory lock36,150Postgres lease row: 26,770 pairs per secondPostgres lease row26,770Redis SET NX PX: 23,144 pairs per secondRedis SET NX PX23,144Redis, fenced script: 19,367 pairs per secondRedis, fenced script19,367Redis, one hot lock: 1,728 pairs per secondRedis, one hot lock1,728pairs per second, log scale
    stepaddit handlesmove up when you see
    1A Postgres advisory lock in the same transaction as the work.About 36,000 lock and release pairs a second. No fencing needed.The work is outside Postgres: an API call, a file, another store.
    2A lease row with a fencing token. The storage checks the token.About 27,000 pairs a second, with durable tokens.Lock traffic competes with the primary's own writes.
    3A Redis lock with a token. Keep the fenced write.About 23,000 pairs a second on one Redis, off the primary.One Redis is busy with locks.
    4Redis Cluster: locks spread by name over many nodes.Each node adds its own rate. Use etcd instead when the lock picks a leader.Top of the ladder.

    On this laptop Postgres used every core and Redis ran commands on one thread, so Postgres was faster here. Move locks to Redis to unload the primary, not for speed.

    I start with a Postgres advisory lock. I add Redis locks only when lock traffic should leave the primary, and I keep the fence.

    M

    Drill

    predict, then reveal

    0 of 8 known

    1. Client 1 pauses for 500 ms with a 300 ms lock. Client 2 takes the lock and writes. What happens when client 1 resumes and writes?

    2. Why can client 1 not check "do I still hold the lock?" just before it writes?

    3. Why must the release compare the token before it deletes the key?

    4. What does Redlock assume, and what breaks it?

    5. Is a Redis INCR a safe fencing token?

    6. When do you need no fencing at all?

    7. What is the difference between a lock for efficiency and a lock for correctness?

    8. The storage is a third-party API that cannot check tokens. What do you do?

    N

    Numbers to say

    measured in the lab
    Redis
    About 23,000 lock and release pairs a second; 19,000 with a fencing token.
    hot lock
    About 54,000 tries a second on one lock, and 1,700 of them win.
    Postgres
    About 36,000 advisory lock pairs a second; 27,000 lease row pairs.
    TTL
    Longer than the longest pause or network delay you expect. Fence anyway: no TTL covers every pause.
    Redlock
    5 nodes; a majority is 3.

    Redis 8 and Postgres 16 on an 8-core laptop, 32 clients, 3 s per case, shared with other workloads. Use these as orders of magnitude.