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.
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.
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.
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.
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_lockfails. 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.
Related¶
- Implementing a Redis lock with fencing tokens — the lease, the release script and the fence.
- Using Postgres advisory locks from asyncio — session versus transaction locks and the pooling trap.
- Leader election for asyncio workers — one owner per duty, with measured failover.
- Implementing per-key async locks — the in-process counterpart.
- Concurrent Execution & Worker Patterns — the parent section.