Skip to content

Distributed Locks & Coordination in asyncio

An asyncio.Lock coordinates tasks on one event loop. The moment a service runs as several processes or replicas — which is almost immediately in production — anything that must happen once needs coordination outside the process: one worker generating the nightly invoice run, one replica owning a websocket fan-out, one consumer applying migrations, one task refreshing a shared cache. Distributed locks are the usual tool, and they are subtler than local ones for a single reason: a lock holder can stop without releasing. A process pauses for garbage collection, a network partition hides it, a container is killed. Every distributed lock therefore expires, and every expiring lock admits the case where the old holder wakes up and keeps working after someone else has taken over.

This section treats that case as the design centre. Measured against Redis 7.4 and PostgreSQL 17: a worker whose 300 ms lease expired while it was paused found its lock taken by another worker, and a storage layer that checked fencing tokens rejected the stale worker's write (token 1 against the current 2); a leader-election loop failed over between 674 ms and 1,011 ms after the leader died with a 1,000 ms lease, depending on when the dead leader had last renewed; and session-level advisory locks behaved differently across Postgres drivers — asyncpg's pool released them when a connection was returned, while psycopg's pool kept them held, leaking the lock. The parent section, Concurrent Execution & Worker Patterns, covers the in-process primitives these build on.

Scope of this section:

  • Redis locks with lease expiry, safe release and fencing tokens.
  • PostgreSQL advisory locks from asyncpg and psycopg, session and transaction scoped.
  • Leader election and failover timing for asyncio workers.
  • Renewing leases from a heartbeat task, and stopping work when renewal fails.
  • Running scheduled jobs exactly once across replicas.

Architectural principles

  • Every distributed lock is a lease. It expires. Design for the holder that outlives its lease, not just the one that releases cleanly.
  • Fence the resource, not just the lock. A monotonically increasing token issued with each acquisition, checked by the resource being protected, is the only way to reject a stale holder's writes.
  • Prefer the database you already write to. If the protected resource is a Postgres table, a Postgres advisory lock or row lock gives you atomicity with the write itself; a separate Redis lock does not.
  • Renew from a separate task, and stop work when renewal fails. A lease that cannot be renewed means you may no longer be the holder; continuing is the bug.
  • Failover time equals lease time. Short leases fail over fast and are sensitive to pauses; long leases are robust and slow to recover. Choose deliberately.
Where coordination lives in a multi-replica service 4 stacked layers. Where coordination lives in a multi-replica service asyncio.Lock, Semaphore tasks on one event loop only Postgres advisory / row locks atomic with your writes Redis leases (SET NX PX) fast, cross-service, expiring fencing token check at the resource rejects stale holders Each layer covers what the one above it cannot; only fencing protects against a holder that outlived its lease.

Execution model: leases, pauses and asyncio

In a single process, the lock holder and the lock live and die together. Across processes they do not, and asyncio adds its own pauses to the usual causes. An event loop blocked by a synchronous call, a long garbage-collection cycle, a CPU-bound handler, or simply a loop overloaded with ready tasks can delay a lease-renewal task by hundreds of milliseconds. If that delay exceeds the remaining lease, another process acquires the lock while this one still believes it holds it — and the code that was waiting on the loop resumes and writes.

That is why the measured scenario matters: worker A acquired a 300 ms lease and then paused for 400 ms; worker B acquired the expired lock and wrote with fencing token 2; A woke and tried to write with token 1, and the storage layer, which remembered the highest token it had seen, rejected it. Without the fence, both writes would have landed and the later one — A's stale write — would have won. A second safety property was also measured: A's attempt to release the lock afterwards did not delete B's lock, because release compared the stored token before deleting.

Coordination primitives also interact with the event loop's own concurrency. A lease-renewal task competes for the loop like any other task, so the more work the loop does, the less predictable renewal timing becomes. Keeping renewal intervals at a fraction of the lease — a third is common — leaves margin for loop lag, which you should measure with the techniques in measuring event loop lag in production.

A paused holder outlives its lease 3 lanes over time. A paused holder outlives its lease worker A acquire, token 1 paused (GC, blocked loop) write with 1: rejected lock in Redis held by A, 300 ms expired held by B worker B acquire, token 2 write with 2: accepted time (not to scale) → Measured against Redis 7.4: the fence rejected A's stale write and A's release left B's lock intact.

Pattern catalogue

Redis lease with safe release and a fencing token

When to use: coordination across services or where the protected resource is not your database.

import uuid

RELEASE = "if redis.call('get',KEYS[1])==ARGV[1] then return redis.call('del',KEYS[1]) else return 0 end"


class RedisLock:
    def __init__(self, r, name: str, ttl_ms: int) -> None:
        self.r, self.key, self.ttl = r, f"lock:{name}", ttl_ms
        self.token: str | None = None
        self.fence: int | None = None

    async def acquire(self) -> bool:
        token = str(uuid.uuid4())
        if await self.r.set(self.key, token, nx=True, px=self.ttl):
            self.token = token
            self.fence = await self.r.incr(self.key + ":fence")   # monotonic per lock
            return True
        return False

    async def release(self) -> bool:
        return bool(await self.r.eval(RELEASE, 1, self.key, self.token))

Trade-off: Redis is fast and simple, but a single Redis node is a single point of failure and failover can lose a lock. The fence makes that survivable. Full treatment in implementing a Redis lock with fencing tokens.

Postgres advisory lock, transaction scoped

When to use: the protected work writes to the same Postgres database.

async def run_exclusive(pool, key: int, work) -> bool:
    async with pool.acquire() as conn, conn.transaction():
        if not await conn.fetchval("select pg_try_advisory_xact_lock($1)", key):
            return False                                # someone else holds it
        await work(conn)                                # released at commit or rollback
        return True

Trade-off: transaction-scoped locks release automatically at commit or rollback — no leaks — but the transaction stays open for the duration of the work. Session-scoped locks avoid that and introduce a pooling hazard: verified, asyncpg's pool unlocks them on release, psycopg's pool does not. See using Postgres advisory locks from asyncio.

Leader election with a renewed lease

When to use: one replica must own a long-running responsibility — a scheduler, a fan-out, a consumer that cannot be parallelised.

async def campaign(lock: RedisLock, lead, interval: float) -> None:
    while True:
        if await lock.acquire():
            try:
                await lead_while_renewing(lock, lead, interval)
            finally:
                await lock.release()
        await asyncio.sleep(interval)

Trade-off: failover time is roughly one lease: measured between 674 ms and 1,011 ms for a 1,000 ms lease. Detail in leader election for asyncio workers.

Heartbeat renewal that cancels the work

When to use: any lock held for longer than its lease.

async def hold_while(lock: RedisLock, work, renew_every: float) -> None:
    async with asyncio.TaskGroup() as tg:
        job = tg.create_task(work())

        async def renew():
            while not job.done():
                await asyncio.sleep(renew_every)
                if not await lock.extend():
                    job.cancel()                       # we may no longer be the holder
                    return
        tg.create_task(renew())

Trade-off: renewal keeps long work safe, but only if losing the lease actually stops the work. See renewing lock leases with a heartbeat task.

Run a scheduled job once per slot

When to use: cron-style jobs in a service with several replicas.

async def run_once_per_slot(r, name: str, slot: str, job) -> bool:
    if not await r.set(f"ran:{name}:{slot}", "1", nx=True, ex=86_400):
        return False                                    # another replica took this slot
    await job()
    return True

Trade-off: a "did it run" marker per schedule slot is simpler than a lock and survives crashes in a defined way (the slot is skipped, not repeated). See running singleton jobs across replicas.

Coordination tools compared A grid of 5 rows by 4 columns. Coordination tools compared tool best for expiry main hazard Redis lease + fence cross-service exclusion TTL stale holder without a fence Postgres xact advisory lock work on the same database commit or rollback long open transactions Postgres session advisory lock long work, no open txn session end pool may keep it held leader election one owner of a duty lease failover takes a full lease run marker per slot cron jobs on replicas marker TTL crash skips the slot Pick by where the protected resource lives and how long the work takes.

Choosing a coordination backend

The tool choice follows from two questions: where does the protected resource live, and what happens if two holders act at once?

The resource is a table in your Postgres database. Use the database. A transaction-scoped advisory lock, or simply SELECT … FOR UPDATE on the rows involved, makes the lock and the write atomic: if the holder dies, the transaction rolls back and the lock is gone with it, and there is no window in which a stale holder can write, because its writes die with its transaction. This is strictly stronger than any external lock, and it costs nothing extra to operate. The limit is duration: the transaction stays open while the work runs, so long work should take a session lock on a dedicated connection, or restructure into short transactions with a claim column, as in building a durable job queue on Postgres with asyncio.

The resource is elsewhere — an external API, a file store, another service. An external lock cannot be atomic with the effect, so it can only make concurrent holders rare. Make them harmless too: fence writes where the resource supports a conditional write (an object store's If-Match, a version column, an idempotency key), and design the work so that a duplicate run is safe. Redis is the common choice here because it is already present in most stacks and its operations are fast enough to renew leases frequently.

Correctness depends on the lock alone. If two holders acting at once would cause real damage and the resource cannot be fenced, a single Redis node is not enough: its failover can lose a lock. Use a consensus-backed store (etcd, ZooKeeper, Consul) for the lock, or — usually better — change the design so the resource itself enforces exclusivity.

Which backend should hold this lock? A decision on Where does the protected resource live with 3 outcomes. Which backend should hold this lock? Where does the protected resource live? my Postgres database advisory or row lock atomic with the write elsewhere, can be fenced Redis lease + fence duplicates made harmless elsewhere, cannot be fenced consensus store or redesign lock alone must be right The closer the lock is to the data, the fewer failure cases it has.

Resource boundaries

Resource Constraint Guidance
Lease TTL failover time vs pause tolerance 3–10× the renewal interval; measure loop lag
Renewal interval must beat worst-case loop lag a third of the TTL is common
Redis round trips one per acquire, release, renew fine for coarse locks; not per request
Postgres connections session locks pin a connection dedicate one, outside the request pool
Open transactions xact locks hold one open keep the locked work short
Fence storage resource must remember max token a column, a row, or a key per resource

The connection row deserves emphasis: a session-level advisory lock is attached to one database connection, so the process must hold that connection for as long as it holds the lock. Taking it from the request pool reduces the pool's capacity, and returning it to a pool that does not reset session state leaks the lock — which psycopg's pool did in testing. Use a dedicated connection for long-held session locks.

Integrated production example

A coordinator that elects a leader through Redis, renews the lease from a heartbeat, fences every write, and stops leading the moment renewal fails:

import asyncio
import logging
import uuid

import redis.asyncio as aioredis

log = logging.getLogger("coord")

EXTEND = ("if redis.call('get',KEYS[1])==ARGV[1] then "
          "return redis.call('pexpire',KEYS[1],ARGV[2]) else return 0 end")
RELEASE = "if redis.call('get',KEYS[1])==ARGV[1] then return redis.call('del',KEYS[1]) else return 0 end"


class Coordinator:
    def __init__(self, r: aioredis.Redis, name: str, ttl_ms: int = 3000) -> None:
        self.r, self.key, self.ttl = r, f"leader:{name}", ttl_ms
        self.token: str | None = None
        self.fence: int | None = None

    async def _acquire(self) -> bool:
        token = str(uuid.uuid4())
        if not await self.r.set(self.key, token, nx=True, px=self.ttl):
            return False
        self.token, self.fence = token, await self.r.incr(self.key + ":fence")
        log.info("became leader with fence %d", self.fence)
        return True

    async def _renew_until_lost(self) -> None:
        while True:
            await asyncio.sleep(self.ttl / 3000)
            if not await self.r.eval(EXTEND, 1, self.key, self.token, self.ttl):
                log.error("lost leadership (fence %d)", self.fence)
                return                                       # returning ends the TaskGroup below

    async def run(self, lead) -> None:
        while True:
            try:
                if await self._acquire():
                    try:
                        async with asyncio.TaskGroup() as tg:
                            renew = tg.create_task(self._renew_until_lost())
                            work = tg.create_task(lead(self.fence))
                            done, _ = await asyncio.wait({renew, work},
                                                         return_when=asyncio.FIRST_COMPLETED)
                            for t in (renew, work):
                                t.cancel()                   # leadership ends both
                    finally:
                        await self.r.eval(RELEASE, 1, self.key, self.token)
                await asyncio.sleep(self.ttl / 3000)
            except (aioredis.ConnectionError, aioredis.TimeoutError) as exc:
                log.warning("coordination store unavailable: %r", exc)
                await asyncio.sleep(1)


async def lead(fence: int) -> None:
    while True:
        await apply_scheduled_work(fence=fence)              # every write carries the fence
        await asyncio.sleep(5)

Diagnostic Hook: export which replica is leader (a gauge per replica, 1 or 0), the current fence value, leadership changes per hour, and renewal latency. More than one replica reporting leadership at the same time is the incident this whole design exists to tolerate — it should be rare and brief, and the fence keeps it harmless; frequent leadership changes at stable load mean the lease is too short for the loop lag you actually have.

Diagnostic hook callout

Three signals cover distributed coordination in production:

  • Renewal latency versus lease. Export the time each renewal takes and the gap between renewals. Alert when the gap exceeds half the lease: the holder is one hiccup away from losing it.
  • Fence rejections. Count writes rejected for a stale token. Each is a holder that outlived its lease; a rising count means pauses are longer than the lease allows.
  • Lock wait and acquisition failures. For locks guarding work, record how long acquisitions wait and how often try_lock fails. A lock that is contended most of the time is serialising work that may not need it.

Failure modes

Failure mode Root cause Detection Fix
Two holders write lease expired during a pause fence rejections, duplicate effects fence every write; renew from a heartbeat
Another holder's lock deleted release without checking the token lock vanishes while holder works compare-and-delete script
Lock never released holder crashed, no TTL work stops, key without expiry always set a TTL
Advisory lock leaked session lock on a pooled connection pg_locks shows it with no owner xact locks, or a dedicated connection
Slow failover lease too long leader gap after crashes shorter lease with renewal
Flapping leadership lease shorter than loop lag spikes frequent leader changes longer lease; fix loop lag
Cron job ran on every replica no coordination per slot duplicate job output run marker or lock per schedule slot

Frequently Asked Questions

How do I make sure only one asyncio worker runs a task across several processes?

Use a distributed lock that expires: a Redis key set with NX and an expiry, or a Postgres advisory lock. Renew it while working, stop work if renewal fails, and fence writes with a token so a holder that outlived its lease cannot overwrite newer work.

What is a fencing token?

A number that increases with every lock acquisition and is attached to each write. The protected resource remembers the highest token it has accepted and rejects writes with a lower one, which stops a paused former holder from writing after someone else took over.

Should I use Redis or Postgres for distributed locks?

Postgres advisory locks when the protected work writes to the same Postgres database, because the lock and the write can share a transaction. Redis for cross-service coordination or when the resource is not in your database.

Why did my Postgres advisory lock stay held after my code finished?

It was a session-level lock on a pooled connection that was returned without unlocking. In testing, asyncpg's pool released such locks on return but psycopg's pool did not. Use transaction-scoped locks or a dedicated connection.

How long should a lock lease be?

Long enough to survive your worst realistic event loop pause with renewal at a third of the lease, and short enough that failover after a crash is acceptable, since failover takes about one full lease.