Leader Election for asyncio Workers¶
Some duties must have exactly one owner across a fleet: the scheduler that enqueues periodic jobs, the consumer of a partition that cannot be processed in parallel, the process that holds an upstream connection allowing only one client. Running the duty on every replica duplicates it; hard-coding replica 0 leaves no owner when replica 0 dies. Leader election solves both: replicas campaign for a lease, the winner leads and renews, and when it stops renewing — crash, partition, deploy — another replica takes over. Tested against Redis 7.4 with a 1,000 ms lease, a follower took over 674 ms after the leader process was killed with SIGKILL (five runs, 673–675 ms), and 1,011 ms in a harness where the leader had just renewed — failover is about one lease, by construction. This guide builds the election loop, the stepping-down logic that keeps it safe, and the measurements that tell you whether it is healthy.
Prerequisites¶
- Python 3.11+,
redis.asyncioand Redis. - The underlying lock, from implementing a Redis lock with fencing tokens.
- Lease renewal, from renewing lock leases with a heartbeat task.
1. Campaign in a loop¶
Every replica runs the same loop: try to acquire the leadership key; if it wins, lead while renewing; if not, wait and try again. The loop never ends while the replica lives:
import asyncio
import logging
import uuid
log = logging.getLogger("election")
class Elector:
def __init__(self, r, duty: str, ttl_ms: int = 5000, replica: str = "") -> None:
self.r, self.key, self.ttl = r, f"leader:{duty}", ttl_ms
self.replica = replica or str(uuid.uuid4())[:8]
self.token: str | None = None
self.fence: int | None = None
self.is_leader = False
async def _try_acquire(self) -> bool:
token = f"{self.replica}:{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")
return True
return False
async def campaign(self, lead) -> None:
while True:
try:
if await self._try_acquire():
await self._lead(lead) # returns when leadership ends
except Exception:
log.exception("election error")
await asyncio.sleep(self.ttl / 1000 / 3)
The token embeds the replica name, so GET leader:scheduler tells an operator who leads. Followers poll at a third of the lease, which bounds how long after expiry someone notices: failover ≈ remaining lease + up to one poll interval.
Verify: start three replicas; exactly one reports leadership at any time.
2. Lead only while renewal succeeds¶
Leading means running the duty and renewing the lease concurrently, with any renewal failure ending the duty immediately:
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")
async def _lead(self, lead) -> None:
self.is_leader = True
log.info("%s became leader (fence %d)", self.replica, self.fence)
async def renew():
while True:
await asyncio.sleep(self.ttl / 1000 / 3)
if not await self.r.eval(EXTEND, 1, self.key, self.token, self.ttl):
log.error("%s lost leadership", self.replica)
return
try:
renewer = asyncio.create_task(renew())
duty = asyncio.create_task(lead(self.fence))
await asyncio.wait({renewer, duty}, return_when=asyncio.FIRST_COMPLETED)
for t in (renewer, duty):
t.cancel()
await asyncio.gather(renewer, duty, return_exceptions=True)
finally:
self.is_leader = False
await self.r.eval(RELEASE, 1, self.key, self.token) # hand over quickly if still ours
The duty receives the fence and must attach it to its writes, so a former leader that resumes after a pause cannot overwrite the new leader's work. Releasing in finally matters for planned shutdowns: a leader that steps down on SIGTERM releases immediately, so the next leader takes over within one poll interval instead of waiting out the lease.
Verify: stop the leader gracefully; a new leader appears within one poll interval. Kill it with SIGKILL; a new leader appears within one lease.
3. Measure failover and choose the lease¶
Failover after an ungraceful death is bounded by the lease plus the followers' poll interval:
import sys
async def measure_failover(r) -> float:
loop = asyncio.get_running_loop()
leader = await asyncio.create_subprocess_exec(sys.executable, "elector_n1.py",
stdout=asyncio.subprocess.PIPE)
await leader.stdout.readline() # n1 prints once it leads
took_over = asyncio.Event()
when: dict[str, float] = {}
async def lead(fence):
when["t"] = loop.time()
took_over.set()
await asyncio.Event().wait()
follower = asyncio.create_task(Elector(r, "demo", 1000, "n2").campaign(lead))
await asyncio.sleep(1.0)
killed = loop.time()
leader.kill() # SIGKILL: no release, no cleanup
await leader.wait()
await asyncio.wait_for(took_over.wait(), 5)
follower.cancel()
return when["t"] - killed
Killing a separate process matters: cancelling an in-process elector runs its finally and releases the key, which measures a graceful handover (4 ms in testing), not a crash.
Measured with a 1,000 ms lease: 673–675 ms across five SIGKILL runs, because the dead leader had renewed about a third of a lease before dying; the worst case is a full lease plus one follower poll interval. The lease trade-off from renewing lock leases with a heartbeat task applies directly: shorter leases fail over faster and are lost during event loop stalls; longer ones ride out stalls and leave the duty unowned longer after a crash. For a periodic scheduler, a 10–30 s gap is usually fine; for a stream consumer, measure the lag you can tolerate.
Verify: run the failover measurement in your environment; it should be close to the lease, and never much more.
4. Watch for split brain and flapping¶
Two failure patterns show up in metrics before they show up as incidents:
- Split brain — two replicas both acting as leader. With leases and fencing it is brief and harmless, but each occurrence means a leader outlived its lease: a stall longer than the margin. Count fence rejections and overlapping
is_leaderreports. - Flapping — leadership changing hands repeatedly at stable load. The lease is too short for normal loop lag, or the leader is overloaded so its renewal task runs late.
async def report(elector: Elector, metrics, every: float = 5.0) -> None:
while True:
metrics.gauge("is_leader", 1 if elector.is_leader else 0,
duty=elector.key, replica=elector.replica)
metrics.gauge("leader_fence", elector.fence or 0, duty=elector.key)
await asyncio.sleep(every)
Summing is_leader across replicas should be exactly 1, except momentarily during failover. Leadership changes per hour should be close to the number of deploys and crashes. If a leader is overloaded by its duty, move the duty's heavy work off the leader's event loop so its renewal keeps time, as discussed in how the GIL affects async services.
Verify: the sum of is_leader across replicas stays at 1 over a day, and leadership changes match deploy events.
5. Elect per partition to spread load¶
One leader per duty puts all of that duty's work on one replica. When the duty divides naturally — tenants, shards, partitions — elect a leader per partition, and each replica ends up leading some of them:
async def lead_partitions(r, partitions: list[str], replica: str, handle) -> None:
electors = [Elector(r, f"partition:{p}", ttl_ms=10_000, replica=replica) for p in partitions]
async with asyncio.TaskGroup() as tg:
for e, p in zip(electors, partitions):
tg.create_task(e.campaign(lambda fence, p=p: handle(p, fence)))
With every replica campaigning for every partition, the first replica to start tends to win them all. Add random jitter before each campaign attempt, or a simple cap — skip campaigning when already leading more than partitions / replicas — so leadership spreads. This is the same balancing problem consumer groups solve for Kafka partitions, covered in consuming Kafka topics with aiokafka; if your data already lives in a system with consumer groups, use them instead of building elections.
Verify: with N replicas and P partitions, each replica leads roughly P / N partitions, and a replica's death redistributes its partitions within one lease.
Verification¶
Leader election is healthy when:
- Exactly one replica leads each duty, except briefly during failover.
- The duty stops as soon as renewal fails, and every duty write carries the fence.
- Failover time is close to the lease after an ungraceful death, and one poll interval after a graceful one.
- Leadership changes match deploys and crashes, not normal operation.
Diagnostic Hook: alert when the sum of is_leader for a duty is 0 for longer than the lease plus one poll interval (nobody leads), or above 1 for longer than a few seconds (split brain). Both are cheap to compute from the per-replica gauge and catch the two ways elections fail in production.
Pitfalls & edge cases¶
- Duty work on the leader's event loop starving renewal. An overloaded leader loses its lease to its own load.
- No release on graceful shutdown. Every deploy waits out a full lease with no leader.
- Assuming the leader is unique. Fence writes; leases only make overlap brief.
- Everyone campaigning for every partition at once. The first replica wins everything; add jitter or caps.
Frequently Asked Questions¶
How do I elect a leader among asyncio workers?
Each replica repeatedly tries to set a Redis key with NX and an expiry; the winner leads and renews the lease from a heartbeat task, stepping down as soon as a renewal fails. Followers keep trying at a fraction of the lease.
How long does leader failover take?
About one lease after an ungraceful death: between 674 ms and 1,011 ms with a 1,000 ms lease in testing, depending on when the dead leader last renewed. A leader that releases the key on graceful shutdown hands over within one poll interval.
Can two replicas be leader at the same time?
Briefly, if a leader pauses longer than its lease and resumes before noticing. Fence every write the duty makes so the stale leader's writes are rejected.
Should I use leader election or Kafka consumer groups?
If the duty is consuming a partitioned stream, use the broker's consumer groups. Use leader election for duties with no such owner, such as a scheduler or a singleton connection.
Related¶
- Distributed Locks & Coordination — up to the topic overview.
- Running singleton jobs across replicas — when a per-slot marker is enough instead of a leader.
- Concurrent Execution & Worker Patterns — the section overview.