Skip to content

Building an Async Event Emitter

An event emitter lets parts of an application react to something — an order placed, a file uploaded, a user signing in — without the code that raises the event knowing who listens. With async listeners, the emitter has to decide whether to wait for them, in what order they see events, and what happens when one is slow or fails. Measured on Python 3.14 with four listeners — fast (1 ms), slow (100 ms), flaky (fails on every tenth event) and jittery (alternating 20 ms and 1 ms) — and 100 events: awaiting listeners one by one stopped at the first event, when the flaky listener raised and the rest never ran. gather per event delivered everything but made emitting take 10.07 s, the slow listener's total. Starting a task per listener per event emitted in 7.7 ms but delivered the jittery listener's events out of order and left 10 "Task exception was never retrieved" reports. One queue and one worker per listener emitted in 0.3 ms, delivered every event in order, and captured all 10 errors. Bounding those queues at 10 events turned a slow listener into backpressure: the emit loop took 8.95 s but the backlog never exceeded 10. This guide builds the queue-based emitter.

Prerequisites

1. Decide what emit() promises

Before choosing an implementation, decide what the caller of emit() is told when it returns. The options measured here give four different answers:

class SequentialEmitter:                    # "every listener has finished, in order"
    async def emit(self, event):
        for listener in self.listeners:
            await listener(event)

class GatherEmitter:                         # "every listener has finished"
    async def emit(self, event):
        await asyncio.gather(*(l(event) for l in self.listeners), return_exceptions=True)

class FireAndForgetEmitter:                  # "listeners have been started"
    async def emit(self, event):
        for listener in self.listeners:
            asyncio.create_task(listener(event))

Measured with 100 events: the sequential emitter raised ValueError: flaky failed on 0 from the very first emit, after the fast and slow listeners had each handled one event and before the jittery listener saw any. The gather emitter delivered all events but took 10,071 ms, because each emit waited for the slow listener. The fire-and-forget emitter returned in 7.7 ms. "Every listener finished" couples the emitter to its slowest listener; "listeners started" decouples it, at the cost of everything in steps 2 and 3. For domain events in a service, the useful promise is usually a third one: "the event is queued for every listener", which is what a queue per listener gives.

Verify: the emitter's docstring states what has happened when emit() returns, and its behaviour matches.

100 events, four listeners (fast, slow, flaky, jittery) A grid of 4 rows by 5 columns. 100 events, four listeners (fast, slow, flaky, jittery) emitter emit loop time delivered (fast/slow/flaky) jittery in order errors await each listener in turn raised on event 0 1 / 1 / 0 - stopped everything gather per event 10,071 ms 100 / 100 / 90 yes discarded by return_exceptions create_task per listener 7.7 ms 100 / 100 / 90 no 10 'never retrieved' queue + worker per listener 0.3 ms 100 / 100 / 90 yes 10 captured Python 3.14; flaky fails on every tenth event.

2. Give each listener its own queue and worker

A queue per listener, drained by one worker task per listener, delivers each listener's events in the order they were emitted, isolates listeners from each other, and lets emit() return as soon as the event is queued:

class Emitter:
    def __init__(self, maxsize: int = 1000):
        self._queues: dict[str, asyncio.Queue] = {}
        self._workers: list[asyncio.Task] = []
        self.errors: list[tuple[str, Exception]] = []

    def subscribe(self, name: str, listener, maxsize: int = 1000) -> None:
        queue: asyncio.Queue = asyncio.Queue(maxsize)
        self._queues[name] = queue
        self._workers.append(asyncio.create_task(self._run(name, listener, queue), name=f"listener:{name}"))

    async def _run(self, name, listener, queue):
        while True:
            event = await queue.get()
            try:
                await listener(event)
            except Exception as exc:
                self.errors.append((name, exc))
                log.exception("listener %s failed on %r", name, event)
            finally:
                queue.task_done()

    async def emit(self, event) -> None:
        for queue in self._queues.values():
            await queue.put(event)

Measured: emitting 100 events took 0.3 ms; the fast, slow and flaky listeners received 100, 100 and 90 events; the jittery listener received its events in order even though its handling time alternated between 20 ms and 1 ms; and all 10 failures were captured with the listener's name. One worker per listener is what preserves order — a listener that needs concurrency should get several workers and give up ordering explicitly.

Verify: each listener receives events in emission order, and a failing listener neither stops others nor loses its own later events.

3. Never lose errors or tasks

The fire-and-forget design shows the two ways events disappear. Exceptions in tasks nobody awaits are reported only when the task is garbage-collected, as "Task exception was never retrieved" — measured, 10 such reports, one per failure, with no indication of which event caused them. And the event loop holds only weak references to tasks, so a task created without keeping a reference can be collected before it finishes. gather(..., return_exceptions=True) loses errors more quietly: it returns them as values, and an emitter that ignores the return value discards them.

results = await asyncio.gather(*(l(event) for l in self.listeners), return_exceptions=True)
for listener, result in zip(self.listeners, results):
    if isinstance(result, Exception):
        log.error("listener %s failed", listener.__name__, exc_info=result)

The queue-based emitter avoids both problems by construction: workers are long-lived and referenced by the emitter, and errors are caught where the listener runs. If you do start tasks per event, keep them in a set and attach a done-callback that logs exceptions, as described in debugging unawaited coroutines in large codebases.

Verify: a listener failure appears in logs with the listener's name and the event, and no "never retrieved" messages appear.

4. Choose what a slow listener does to the emitter

With a queue per listener, a slow listener builds a backlog. The queue's size decides whether that backlog grows or pushes back on the emitter:

emitter.subscribe("audit", write_audit_log, maxsize=10_000)    # absorb bursts
emitter.subscribe("search", reindex, maxsize=10)               # push back when behind

Measured with the 100 ms listener and 100 events: with maxsize=1000, the emit loop took 0.1 ms, the backlog peaked at 100 events, and delivery finished 10.04 s later; with maxsize=10, the emit loop took 8,951 ms — it waited for space — the backlog never exceeded 10, and delivery finished at 10.06 s. A large queue keeps the producer fast and puts the cost in memory and delivery delay; a small one bounds memory and slows the producer to the slowest listener. A third option is to drop events for a listener that falls behind, using put_nowait and counting QueueFull, for listeners where only recent events matter, such as live dashboards.

Verify: each listener's queue size reflects a decision — absorb, push back, or drop — and backlog per listener is observable.

A 100 ms listener and 100 events A grid of 2 rows by 4 columns. A 100 ms listener and 100 events listener queue size emit loop peak backlog all delivered after 1,000 0.1 ms 100 events 10.04 s 10 8,951 ms 10 events 10.06 s The slow listener sets delivery time either way; the queue decides who waits.

5. Drain on shutdown and test delivery

Events queued but not delivered when the process stops are lost. Give the emitter a shutdown that drains each queue within a deadline, then cancels the workers:

async def aclose(self, timeout: float = 5.0) -> None:
    try:
        async with asyncio.timeout(timeout):
            for queue in self._queues.values():
                await queue.join()                     # everything queued has been handled
    except TimeoutError:
        log.warning("emitter shutdown: %d events undelivered",
                    sum(q.qsize() for q in self._queues.values()))
    finally:
        for worker in self._workers:
            worker.cancel()
        await asyncio.gather(*self._workers, return_exceptions=True)

In tests, await emitter.aclose() (or a drain() that only joins the queues) after emitting makes assertions deterministic: every queued event has been handled when it returns. For events that must survive a crash — payments, emails — an in-process emitter is the wrong tool; write them to a durable queue or an outbox table in the same transaction as the change that caused them, as in Background Jobs & Task Queues.

Verify: shutdown delivers or reports every queued event, and tests drain before asserting.

Which emitter fits? A decision on What must emit() guarantee with 4 outcomes. Which emitter fits? What must emit() guarantee? all listeners finished gather + inspect results slowest listener sets emit time queued, ordered, isolated queue + worker per listener 0.3 ms, in order latest only, never block put_nowait, drop and count for dashboards survives a crash durable queue / outbox not in-process Fire-and-forget tasks are never the answer: order and errors are lost.

Verification

An async event emitter is sound when:

  • What emit() promises is stated and matches the implementation.
  • Each listener has its own queue and worker, preserving order and isolating failures.
  • Errors are caught and logged per listener, and no tasks are left unreferenced.
  • Queue sizes reflect a backpressure decision, and shutdown drains within a deadline.

Diagnostic Hook: export per-listener queue depth and handling time. One listener's depth climbing while the others stay at zero identifies the slow consumer immediately — and shows whether it is about to push back on every producer, if its queue is bounded.

Pitfalls & edge cases

  • Awaiting listeners in turn. Measured: the first failure stopped delivery to everyone.
  • gather per event. Measured: emitting 100 events took 10 s.
  • Tasks per event. Measured: out-of-order delivery and 10 unretrieved exceptions.
  • Unbounded backlogs. A slow listener grows memory without limit.

Frequently Asked Questions

How do I build an event emitter with async listeners in Python?

Give each listener an asyncio.Queue and a worker task that awaits the listener for each event, catching its exceptions; emit() puts the event on every queue. In testing it emitted 100 events in 0.3 ms with order preserved.

Should emit() await all listeners?

Only if callers need that guarantee: with gather, emitting 100 events took 10 s because of one 100 ms listener, and awaiting in turn stopped at the first exception.

Why are my listener exceptions only reported as 'Task exception was never retrieved'?

The listeners ran in tasks that nobody awaited. Catch exceptions where the listener runs, or keep the tasks and add a done-callback that logs them.

How do I stop a slow listener from slowing the emitter?

Give it a large queue to absorb the backlog, or drop events for it when its queue is full. A queue of 10 made the emitter wait 8.95 s for a 100 ms listener.