Building a Countdown Latch in asyncio¶
A countdown latch opens once a fixed number of events have happened: N shards have loaded, N replicas have acknowledged, N warm-up requests have completed. Any number of tasks can wait on it, and the tasks that count down do not wait at all. Java has CountDownLatch; asyncio does not, but it has Barrier (3.11+), which looks similar and behaves differently — a barrier makes the arriving tasks wait for each other, and its wait() must be called by exactly the participating parties. A latch built on asyncio.Event takes a dozen lines: in a test with three shards finishing at 10, 20 and 30 ms, all four waiters were released together at 30 ms, the moment the last shard counted down. This guide builds it, adds failure propagation, and shows when a TaskGroup makes it unnecessary.
Prerequisites¶
- Python 3.11+, stdlib only.
- Event semantics, from signalling with asyncio.Event set() and clear().
- Barrier, from using asyncio.Barrier to start tasks together.
1. Build the latch on a one-shot Event¶
import asyncio
class CountDownLatch:
def __init__(self, count: int) -> None:
if count < 0:
raise ValueError("count must be >= 0")
self._count = count
self._done = asyncio.Event()
if count == 0:
self._done.set()
@property
def count(self) -> int:
return self._count
def count_down(self) -> None:
if self._count > 0:
self._count -= 1
if self._count == 0:
self._done.set()
async def wait(self) -> None:
await self._done.wait()
count_down() is synchronous, so it can be called from callbacks and from code that must not suspend. On a single event loop the decrement and the check cannot interleave with another task, so no lock is needed. Extra calls after zero are ignored rather than going negative, which matches Java's semantics and avoids a class of off-by-one crashes when a shard reports twice.
The underlying Event is only ever set, never cleared, so it is a pure latch: waiters arriving after it opened return immediately.
Verify: with count 3, two count_down() calls leave waiters blocked and the third releases all of them at once.
2. Know when it beats Barrier and TaskGroup¶
The three tools answer different questions:
- Latch: "wait until N events have happened" — the counters and the waiters are different tasks, and counters never block.
- Barrier: "wait until N tasks have all reached this point" — the same N tasks are counters and waiters, and the barrier can cycle for repeated phases. Verified: three tasks arriving at 0, 10 and 20 ms each got their arrival index (0, 1, 2) and all left at 20 ms.
- TaskGroup: "wait until these N tasks have finished" — when you own the tasks, exiting the
async withblock is the latch.
If the N events are the completion of N tasks you created, use a TaskGroup and skip the latch. A latch earns its place when the events are not task completions — N acknowledgements arriving on a socket, N workers becoming ready while continuing to run, N files appearing — or when the tasks that count down are owned elsewhere.
Verify: for each latch in your code, check whether the counting events are simply task completions; if so, a TaskGroup is clearer.
3. Propagate failure instead of waiting forever¶
A plain latch has no notion of failure: if one shard crashes before counting down, waiters block forever. Add an explicit failure path that releases waiters with an error:
class FailableLatch(CountDownLatch):
def __init__(self, count: int) -> None:
super().__init__(count)
self._error: BaseException | None = None
def fail(self, exc: BaseException) -> None:
if not self._done.is_set():
self._error = exc
self._done.set()
async def wait(self) -> None:
await self._done.wait()
if self._error is not None:
raise self._error
async def load_shard(i: int, latch: FailableLatch) -> None:
try:
await fetch_shard(i)
except Exception as exc:
latch.fail(RuntimeError(f"shard {i} failed: {exc!r}"))
raise
latch.count_down()
Every waiter now either proceeds because all shards loaded, or fails promptly with the first error. Raising the same exception object in several waiters is safe in Python, though each raise appends to its traceback; wrap it in a fresh exception per waiter if clean tracebacks matter.
Verify: make one shard raise; every waiter receives the error within one loop iteration.
4. Put a deadline on the wait¶
Even with failure propagation, a counter that hangs — a shard whose network call never returns — keeps waiters blocked. Bound the wait, and make the timeout tell you which events are missing:
class TrackingLatch(FailableLatch):
def __init__(self, names: list[str]) -> None:
super().__init__(len(names))
self.pending = set(names)
def count_down_for(self, name: str) -> None:
if name in self.pending:
self.pending.discard(name)
self.count_down()
async def wait_all_shards(latch: TrackingLatch, timeout: float) -> None:
try:
async with asyncio.timeout(timeout):
await latch.wait()
except TimeoutError:
raise TimeoutError(f"still waiting for: {sorted(latch.pending)}") from None
Naming the counters turns "startup timed out" into "startup timed out waiting for shard-7", which is the difference between a five-minute and a fifty-minute incident. Naming also makes duplicate reports from the same counter harmless.
Verify: stall one shard; the timeout error lists exactly that shard.
5. Use it for readiness gating¶
The most common production use is startup: a service should report ready only when its background components are running, but those components are long-lived tasks that never complete, so a TaskGroup exit cannot be the signal:
async def main() -> None:
ready = TrackingLatch(["consumer", "cache-warmer", "http"])
async with asyncio.TaskGroup() as tg:
tg.create_task(run_consumer(on_ready=lambda: ready.count_down_for("consumer")))
tg.create_task(warm_cache(on_ready=lambda: ready.count_down_for("cache-warmer")))
tg.create_task(run_http(on_ready=lambda: ready.count_down_for("http")))
await wait_all_shards(ready, timeout=60)
health.mark_ready()
Each component counts down once when it is ready to serve and then keeps running. The latch and the TaskGroup complement each other: the latch says "everything is up", the TaskGroup says "if anything dies, everything stops". The probe side is covered in implementing health and readiness probes for asyncio.
Verify: delay one component's readiness; the health endpoint stays not-ready until it reports, and a startup timeout names it.
Verification¶
The latch is correct when:
- Waiters are released exactly when the count reaches zero, together.
- Extra count-downs are harmless and the count never goes negative.
- A failing counter releases waiters with its error instead of hanging them.
- Timed-out waits name the counters that never reported.
Diagnostic Hook: log each count-down with the counter's name and the elapsed time since the latch was created, and export the time-to-open as a metric. Over several deploys this shows which component dominates startup; a time-to-open that drifts upward is an early warning that some dependency is getting slower to come up.
Pitfalls & edge cases¶
- Using
asyncio.Barrieras a latch. Itswait()blocks the callers that count; non-participants callingwait()change the party count semantics. - Counting down from another thread.
count_downtouches anasyncio.Event; from a thread useloop.call_soon_threadsafe(latch.count_down). - Count mismatches. A latch for 3 with 4 counters opens early; one with 2 counters never opens. Named counters catch both.
- Reusing a latch. It is one-shot by design; create a new one per phase or use a Barrier for cyclic phases.
Frequently Asked Questions¶
Does asyncio have a CountDownLatch?
No, but one is a few lines on top of asyncio.Event: keep a counter, decrement it in a synchronous count_down method, and set the event when it reaches zero. Waiters await the event.
What is the difference between a latch and asyncio.Barrier?
A latch lets any tasks wait until N events have happened, and the code that signals those events never blocks. A Barrier makes N participating tasks wait for each other at the same point, and can be reused for repeated phases.
Should I use a latch or a TaskGroup?
If the events you are counting are the completion of tasks you own, a TaskGroup already waits for them and handles failures. Use a latch when the events are not task completions, such as long-running components becoming ready.
How do I stop a latch wait hanging forever?
Add a failure path that releases waiters with an exception, wrap waits in asyncio.timeout, and track counters by name so the timeout error says which ones never reported.
Related¶
- Synchronization Primitives — up to the topic overview.
- Structured concurrency with asyncio.TaskGroup — the alternative when you own the tasks.
- Asyncio Fundamentals & Event Loop Architecture — the section overview.