Actors & Supervision in asyncio¶
Shared mutable state is the hardest part of concurrent code, and asyncio only removes half of the problem. There are no data races between threads, but any await between reading state and writing it back lets another task interleave — and the result can be dramatic. Measured on Python 3.14, 10,000 concurrent increments with a single await inside the read-modify-write left a balance of 1. The actor pattern answers that by giving each piece of state to exactly one task, which receives messages through a mailbox and handles them one at a time. This section builds actors that are correct, observable and resilient, and measures each part: an actor handled the 10,000 increments correctly in 0.040 s against 0.083 s for a lock; an ask round trip cost 4.5 µs; a supervisor without backoff restarted a failing actor 19,806 times in 3 seconds; boolean connection flags produced a closed-but-connected zombie in 25% of fuzzed schedules where a transition table produced none; an unbounded mailbox at twice capacity reached a 4.6 s p99 wait; and sharding by key took throughput from 937 to 90,931 messages per second — until one hot key capped it at 4,906.
The parent section, Concurrent Execution & Worker Patterns, covers worker pools and queues; actors are what those become when each worker owns state rather than processing interchangeable jobs.
Architectural principles¶
- One owner per piece of state. Only the actor's task reads or writes it; everything else sends messages.
- Messages are immutable and typed. The protocol between actors is a set of small data classes, not shared objects.
- Every mailbox is bounded, with a deliberate policy for when it is full.
- Every ask has a timeout, and every actor fails its pending asks when it stops.
- Every actor has a supervisor that restarts it with backoff, caps restart intensity and decides what happens to the message that crashed it.
Execution model: one message at a time on one loop¶
An actor is a task in a loop around await mailbox.get(). Because the event loop runs one task at a time, and only the actor touches its state, every handler runs to completion — including across its own awaits — without any other code observing the state in between. That is the whole correctness argument, and it holds even when the handler awaits slow I/O. The cost is serialization: an actor's throughput is at most 1 / handler time. A handler that awaits a 1 ms write tops out near 1,000 messages per second, measured at 937.
The mailbox is an asyncio.Queue, and its mechanics determine the overheads: a put_nowait and get pair cost about 348 ns per message, or 2.87 million messages per second, so the mailbox is never the bottleneck — the handler is. Asks add a future per request and cost about 4.5 µs per round trip. Supervision lives one level up, as a coroutine that awaits the actor's run() and reacts when it raises; cancellation passes straight through it, so shutting down the supervisor shuts down the actor. When one actor is too slow, the model scales sideways: more actors, each owning a slice of the keys, all on the same loop.
Pattern catalogue¶
An actor with a mailbox¶
The basic unit: state, a bounded queue, and a loop that matches message types:
class AccountActor:
def __init__(self, maxsize: int = 1000) -> None:
self._balance = 0
self.mailbox: asyncio.Queue[Deposit | Withdraw | None] = asyncio.Queue(maxsize)
async def run(self) -> None:
while (msg := await self.mailbox.get()) is not None:
match msg:
case Deposit(amount):
b = self._balance
await audit_log(amount) # safe: nothing else touches _balance
self._balance = b + amount
case Withdraw(amount):
if amount <= self._balance:
self._balance -= amount
When to use it: state with rules spanning several operations, which needs a lifecycle and metrics. Trade-off: all operations on the state are serialized. See building an actor with an asyncio mailbox.
Request-reply with futures¶
A future in the message is the reply address. The actor must complete it — with a result, or with an exception if the actor dies:
async def ask(self, key: str, timeout: float = 1.0) -> int | None:
if self.closed:
raise ActorStopped(self.name)
msg = Get(key)
async with asyncio.timeout(timeout):
await self.mailbox.put(msg)
return await msg.reply
Without crash handling, three of four callers were still waiting a second after the actor died. See request-reply messaging between actors.
Supervision with backoff and intensity limits¶
The supervisor restarts the actor after a jittered, growing delay, and escalates when restarts exceed N in T seconds:
async def supervise(factory, max_restarts=5, window=10.0, base=0.05, cap=1.0) -> None:
history: list[float] = []
delay = base
while True:
try:
await factory()
return
except Exception as exc: # CancelledError passes through
now = time.monotonic()
history = [t for t in history if now - t < window] + [now]
if len(history) > max_restarts:
raise Escalate(f"{len(history)} crashes in {window}s") from exc
await asyncio.sleep(delay * random.uniform(0.5, 1.0))
delay = min(cap, delay * 2)
Measured: 19,806 restarts and a full CPU core in 3 s without the sleep; 8 with it; escalation on the sixth crash with a limit of 5 in 10 s. See supervising actors with restart strategies.
State machines for lifecycles¶
Connections and sessions keep one state variable, changed only through a transition table, with a re-check after every await:
async def connect(self) -> None:
if not self._transition("connect"):
return
sock = await open_socket()
if not self._transition("established"): # close() won the race
sock.close()
return
self.sock = sock
A flag-based connection ended closed-but-connected in 2,509 of 10,000 fuzzed schedules; the table-driven one in none. See modelling connection lifecycles as async state machines.
Bounded mailboxes and sharding¶
A mailbox's capacity times the handler time is the worst-case wait; the overload policy decides who pays. When the actor itself is too slow, shard it by a stable hash of the key:
def shard_for(self, key: str) -> Shard:
return self.shards[zlib.crc32(key.encode()) % len(self.shards)]
See bounding actor mailboxes under load and sharding state across actors by key.
Actors, locks, pools or a broker?¶
Actors are one of several ways to coordinate concurrent work, and they are not always the right one. A plain asyncio.Lock is simpler when the shared state is a single value with a short critical section and no lifecycle: a counter, a cached token. A worker pool is the better fit when jobs are interchangeable and carry their own data — any worker can process any job, and there is no state to own between jobs. An actor earns its keep when the state persists across messages and has rules: an account balance that must never go negative, a connection whose lifecycle must be consistent, a session whose messages must be applied in order.
The boundary with a message broker is about durability and distance. In-process actors lose their mailboxes when the process dies and cannot be reached from another process; that is acceptable for caches, connection managers and in-memory aggregation, and not for business events that must survive a crash. When messages must be durable or cross processes, the mailbox becomes a broker queue or stream, and the actor becomes a consumer — the same shape, with acknowledgement and redelivery semantics from Message Brokers & Event Streams replacing the in-memory queue. Many systems use both: a durable stream feeding sharded in-process actors that hold the hot state.
Resource boundaries¶
An actor system has more queues and tasks than it first appears, and each needs a limit:
- Mailbox capacity = tolerable latency ÷ handler time. 200 slots at about 1 ms per message gave a measured p99 wait of 220 ms.
- Ask timeouts shorter than the caller's own deadline, so a stalled actor surfaces as an error in the caller rather than a timeout further up.
- Shard count no higher than the backing store's useful concurrency: each shard is a concurrent client of whatever it writes to.
- Restart intensity — 5 restarts in 10 s is a reasonable default — so a persistent fault escalates within seconds instead of looping indefinitely.
- Redelivery attempts for a message that crashes its actor: three, then a dead-letter store.
Each limit converts a failure mode that would otherwise grow without bound — memory, latency, CPU, log volume — into a counter that can be alerted on.
Integrated production example¶
A sharded counter service: 16 shards routed by crc32, each a supervised actor with a bounded mailbox, asks that fail cleanly when a shard dies, and restarts with backoff and an intensity limit. In a test it processed 51,000 increments, survived an injected crash with one restart, and returned exact totals:
class CounterShard:
def __init__(self, name: str, capacity: int = 500) -> None:
self.name, self.closed = name, False
self.counts: dict[str, int] = {}
self.mailbox: asyncio.Queue[Incr | Read | None] = asyncio.Queue(maxsize=capacity)
async def run(self) -> None:
msg = None
try:
while (msg := await self.mailbox.get()) is not None:
match msg:
case Incr(key, by):
self.counts[key] = self.counts.get(key, 0) + by
case Read(key, reply) if not reply.done():
reply.set_result(self.counts.get(key, 0))
except BaseException as exc:
self.closed = True
err = ActorStopped(f"{self.name}: {exc!r}")
pending = [msg] + [self.mailbox.get_nowait() for _ in range(self.mailbox.qsize())]
for m in pending:
if isinstance(m, Read) and not m.reply.done():
m.reply.set_exception(err)
raise
class Counters:
def __init__(self, n_shards: int = 16) -> None:
self.shards = [CounterShard(f"shard-{i}") for i in range(n_shards)]
def _shard(self, key: str) -> CounterShard:
return self.shards[zlib.crc32(key.encode()) % len(self.shards)]
async def incr(self, key: str, by: int = 1) -> None:
shard = self._shard(key)
if shard.closed:
raise ActorStopped(shard.name)
async with asyncio.timeout(1.0):
await shard.mailbox.put(Incr(key, by))
async def read(self, key: str) -> int:
shard = self._shard(key)
if shard.closed:
raise ActorStopped(shard.name)
msg = Read(key)
async with asyncio.timeout(1.0):
await shard.mailbox.put(msg)
return await msg.reply
async def _supervise(self, i: int, max_restarts: int = 5, window: float = 10.0) -> None:
history: list[float] = []
delay = 0.05
while True:
shard = self.shards[i]
try:
await shard.run()
return # clean stop via sentinel
except Exception as exc:
now = time.monotonic()
history = [t for t in history if now - t < window] + [now]
if len(history) > max_restarts:
raise RuntimeError(f"{shard.name} keeps crashing") from exc
log.warning("%s crashed (%r); restarting", shard.name, exc)
await asyncio.sleep(delay * random.uniform(0.5, 1.0))
delay = min(1.0, delay * 2)
fresh = CounterShard(shard.name)
fresh.counts = await load_snapshot(shard.name) # durable state, not leftovers
self.shards[i] = fresh
async def run(self) -> None:
async with asyncio.TaskGroup() as tg: # escalation cancels the rest
for i in range(len(self.shards)):
tg.create_task(self._supervise(i), name=f"supervisor-{i}")
The router and supervisor are the same object here because the shard set is fixed; clients call incr and read and never see a mailbox. A shard's crash fails its pending reads with ActorStopped, marks it closed so new calls fail fast during the backoff, and is replaced by a fresh instance loaded from durable state. An escalation from any supervisor raises out of the TaskGroup, cancelling the other shards — the one-for-all behaviour you want at the top of a tree, where the next step is the process exiting and being restarted. In the test run, the restored instance kept its in-memory counts, which is why the totals matched exactly; in production that line loads from storage.
Diagnostic hook callout¶
Per actor, export four series: mailbox depth, mailbox wait (time from put to get), handler duration, and restarts. Alert on:
- Mailbox wait p99 above the latency budget of the requests that send to the actor — the actor is saturated; shard it or batch.
- Depth pinned at capacity together with a rising drop or rejection count — sustained overload is being shed.
- Restart rate pinned at one per backoff cap — a dependency is down and the supervisor is waiting it out.
- Any escalation — the supervisor gave up; something above it must act.
- Ask timeouts across all callers of one actor — it stalled or died without failing its pending asks.
Per-shard message counts complete the picture: a flat distribution means sharding works, and one shard far above the others names a hot key.
Failure modes¶
| Failure mode | Root cause | Detection | Fix |
|---|---|---|---|
| Lost updates | await between read and write of shared state |
Concurrent test checks final state | Actor owns the state, or a lock |
| Callers hang forever | Actor crashed without failing pending futures | Ask timeouts across all callers | Fail in-flight and queued asks; closed flag |
| Restart storm, CPU at 100% | Supervisor restarts without delay | Restart counter rate; CPU | Exponential backoff with jitter |
| Stuck on one message | Poison message redelivered forever | Same message ID in every crash | Limited attempts, then dead-letter |
| Zombie connection | Flags set after an await without re-check |
Fuzz test invariants | One state + transition table |
| Seconds of latency, growing memory | Unbounded mailbox at overload | Mailbox wait and depth | Bounded mailbox with a policy |
| Throughput flat as shards are added | Hot key | Per-shard message counts | Batch in the hot shard; split the key |
| Keys move between shards | hash() randomized per process |
Different routing per process | Stable hash such as crc32 |
Frequently Asked Questions¶
What is the actor model in asyncio?
A task that exclusively owns some state and handles messages from an asyncio.Queue one at a time. Because no other code touches the state, awaits inside a handler cannot cause races; in testing it handled 10,000 contended updates correctly in 0.040 s.
When should I use an actor instead of an asyncio.Lock?
When the state has rules across several operations, needs a lifecycle (start, stop, restart) and metrics, or needs request-reply semantics. For a single counter with a short critical section a lock is simpler.
How do I supervise an asyncio task like Erlang does?
Wrap it in a coroutine that restarts it after an exponential, jittered delay, caps restarts at N within T seconds and raises past that limit so a parent can act. Without the delay, an always-failing task restarted 19,806 times in 3 seconds.
How do asyncio actors scale beyond one task?
By sharding: several actors, with each key routed by a stable hash to one of them so per-key order is kept. With 1 ms of I/O per message, 256 shards processed 90,931 messages per second against 937 for one actor.
Is there an actor library for asyncio?
Several exist, but the core is small enough to write with asyncio.Queue, futures and TaskGroup, as this section shows; the parts libraries add are mainly remoting across processes and supervision trees with configuration.
Related¶
- Building an actor with an asyncio mailbox — state, mailbox and lifecycle.
- Request-reply messaging between actors — asks that always answer.
- Supervising actors with restart strategies — backoff, intensity and poison messages.
- Modelling connection lifecycles as async state machines — no zombies.
- Bounding actor mailboxes under load — overload policies.
- Sharding state across actors by key — scaling sideways.
- Worker Pool Implementations — the stateless counterpart.
- Concurrent Execution & Worker Patterns — the parent section.