An actor is a unit of computation that owns private state and a mailbox. It processes one message at a time, and the only way to affect it is to send it a message. That rule removes shared mutable state from application code, and with it most locks and data races. The concepts are introduced in the actor model overview. This article is about the machinery underneath: what a runtime has to build so that millions of actors can share a handful of threads, and the operational decisions that make actor systems either robust or quietly overloaded.
It walks through the parts of a runtime, builds a small but working one in Python asyncio with request-reply and supervision, sizes a real workload with a worked example, and ends with distribution, failure modes and a checklist. The ideas carry over to Erlang/OTP, Akka, Orleans and similar systems; where details differ between them, the article says so.
Anatomy of an actor runtime
A runtime has five parts. An actor reference is an address you can send to; it hides where the actor lives and whether it has been restarted. A mailbox holds messages until the actor processes them. A behaviour is the code and private state that handles one message at a time. A dispatcher decides which actor runs next on which thread. A supervisor decides what happens when an actor fails.
The data flow for one message is short. The sender enqueues into the target's mailbox. If the actor was idle, the enqueue also schedules it on the dispatcher's ready queue. A thread picks it up, runs its behaviour for a batch of messages, then either reschedules the actor if more messages are waiting or marks it idle. The key invariant is that an actor is scheduled on at most one thread at a time, which is what makes its state safe without locks.
Mailboxes: bounded, unbounded and what happens when full
Mailboxes are usually multi-producer, single-consumer queues: any number of senders, one actor draining. Lock-free linked queues are common because enqueue is on every sender's hot path. Erlang process mailboxes and Akka's default mailbox are unbounded, which is simple and never blocks a sender, but means a slow actor can accumulate messages until the process runs out of memory.
A bounded mailbox turns that slow failure into an explicit choice. When the mailbox is full, the runtime can:
- Block or suspend the sender until there is room. This propagates backpressure, but a sender that is itself an actor stalls its own mailbox, and cycles of actors waiting on each other can deadlock.
- Reject the message and tell the sender, which can retry later, shed load, or return an error to its client.
- Drop the newest or oldest message. Acceptable for telemetry or position updates where only the latest value matters.
Some runtimes also offer priority or control mailboxes, so that a stop or a configuration change is processed before a backlog of ordinary work. Whichever you choose, measure mailbox depth: it is the actor equivalent of queue length, and its growth is the earliest sign of overload. Flow control between stages without mailboxes is covered in Reactive Streams backpressure.
Dispatching: fairness on shared threads
Actors are cheap because they are not threads. A dispatcher multiplexes many actors onto a pool of threads, often one per core. Two settings shape its behaviour. The first is how many messages an actor may process before giving its thread back. Akka calls this throughput, and its reference configuration sets it to 5. A high value improves cache locality and reduces scheduling overhead; a low value improves fairness, so one busy actor cannot delay thousands of others. The second is the pool itself: its size and whether it steals work between threads, as described in work stealing.
The rule that matters most in practice: never block a dispatcher thread. An actor that makes a synchronous database call holds one of a few threads for the whole call. Eight such actors on an eight-thread pool stop every other actor in the system. Either use asynchronous I/O and send the result back to the actor as a message, or run blocking actors on a separate, dedicated pool sized for the blocking work. Erlang avoids most of this because its scheduler preempts processes after a budget of work, but long-running native code can still stall a scheduler there.
A minimal runtime in Python asyncio
This runtime fits in about 50 lines and runs on Python 3.11 or later. Each actor is an asyncio task that drains a bounded asyncio.Queue. tell is fire-and-forget and raises asyncio.QueueFull when the mailbox is full; ask waits for room, then waits for a reply with a timeout. The supervisor restarts the behaviour with fresh state after a crash, and gives up when there are more than intensity restarts within period seconds.
import asyncio
import time
STOP = object()
def settle(reply, result=None, exc=None):
if reply is not None and not reply.done(): # caller may have timed out
reply.set_exception(exc) if exc else reply.set_result(result)
class Ref:
def __init__(self, mailbox):
self.mailbox = mailbox
def tell(self, msg):
self.mailbox.put_nowait((msg, None)) # raises asyncio.QueueFull
async def ask(self, msg, timeout=1.0):
reply = asyncio.get_running_loop().create_future()
await self.mailbox.put((msg, reply)) # waits while full
return await asyncio.wait_for(reply, timeout)
async def supervise(factory, mailbox, intensity=3, period=5.0):
restarts = []
while True:
behaviour = factory() # fresh state on every start
try:
while True:
msg, reply = await mailbox.get()
if msg is STOP:
return
try:
settle(reply, behaviour.handle(msg))
except Exception as exc:
settle(reply, exc=exc)
raise
except Exception:
now = time.monotonic()
restarts = [t for t in restarts if now - t < period] + [now]
if len(restarts) > intensity:
raise # escalate to our parent
class Account:
def __init__(self):
self.balance = 0
def handle(self, msg):
kind, amount = msg
if kind == "withdraw" and amount > self.balance:
raise ValueError("insufficient funds")
self.balance += amount if kind == "deposit" else -amount
return self.balance
async def main():
box = asyncio.Queue(maxsize=100)
acct = Ref(box)
async with asyncio.TaskGroup() as tg:
tg.create_task(supervise(Account, box))
print(await acct.ask(("deposit", 50))) # 50
try:
await acct.ask(("withdraw", 80))
except ValueError as e:
print("rejected:", e) # actor crashed and restarted
print(await acct.ask(("deposit", 5))) # 5, not 55
acct.tell(STOP)
asyncio.run(main())The last line of output is the lesson. The withdrawal raised, the supervisor restarted the actor, and the restart wiped the balance, exactly as designed. Restarting with clean state is what makes "let it crash" safe for bugs and corrupted state, but it means two things for real systems. Expected business outcomes such as insufficient funds should be ordinary replies, not exceptions. And state that must survive a restart must be persisted, for example as an event log the actor replays when it starts, or loaded from a database in its constructor.
The mailbox lives outside the behaviour, so the ref stays valid across restarts and messages queued during a crash are not lost. A message that crashes the actor is lost, though, and if it crashes every time, the restart limit is what stops an endless loop.
Supervision strategies and restart intensity
Erlang/OTP defined the vocabulary most runtimes use. A supervisor's strategy decides which children restart when one fails: one-for-one restarts only the failed child; one-for-all restarts all children, for groups that share assumptions; rest-for-one restarts the failed child and every child started after it, for pipelines where later stages depend on earlier ones.
Restart intensity limits the damage of a persistent fault. If more than intensity restarts happen within period seconds, the supervisor stops its children and itself fails, passing the problem to its own supervisor. OTP's defaults are an intensity of 1 and a period of 5 seconds. Tune them per subtree: a child that crashes on a bad message deserves a few restarts, while one that cannot reach its database should escalate quickly so a higher level can do something different, such as backing off or failing over.
Design the tree around failure domains. Put actors that fail independently under one-for-one supervisors, isolate risky integrations in their own subtrees, and keep the root small so that a full restart is rare and fast.
Request-reply, ordering and delivery guarantees
Most runtimes guarantee less than people assume. Delivery is usually at-most-once: a message to a crashed actor, a full mailbox or a remote node behind a partition can be lost without an error. Ordering is usually guaranteed only per sender-receiver pair: if A sends m1 then m2 to B, B sees m1 before m2, but messages from A and C to B can interleave in any order. Erlang and Akka both give this pairwise guarantee.
That leads to three rules. Every ask needs a timeout, and a timeout means "unknown", not "failed": the actor may have processed the request. Requests that change state should carry an idempotency key so retries are safe. And when you need at-least-once delivery, build it explicitly with acknowledgements, retries and deduplication, or put a durable log between actors.
Request-reply between two actors that each wait on the other is the actor version of a deadlock. Prefer one-way messages, with the reply arriving as just another message, and keep blocking ask calls at the edges of the system. Channels, which give some of the same isolation without addresses, are compared in channels.
Worked example: one actor per device
A fleet of 200,000 sensors each reports every 10 seconds, giving 20,000 messages per second. Each device actor validates a reading, updates a rolling window and occasionally emits an alert. Measured handling time is 0.15 ms per message.
- CPU. 20,000 messages per second at 0.15 ms each is 3 core-seconds per second. On an 8-core node that is about 38% utilisation, leaving headroom for bursts and garbage collection.
- Memory. At about 2 KB of state and mailbox overhead per actor, 200,000 actors take roughly 400 MB, which is reasonable for one node and a clear signal of when to shard.
- Mailbox bound. Each device sends one message per 10 seconds, so a healthy mailbox holds 0 or 1 messages. A bound of 16 tolerates a burst after a network reconnect; beyond that, keep only the newest reading.
- Slow path. Writing alerts to a database is blocking I/O, so device actors send alerts to a small pool of writer actors on a dedicated dispatcher, whose mailbox depth is the metric to alert on.
Device actors are supervised one-for-one with a modest restart limit. A device that sends malformed data crashes and restarts its own actor without affecting the other 199,999.
Distribution and virtual actors
Because senders hold refs rather than objects, an actor can live on another node without changing calling code. This is location transparency, and it hides real differences: remote sends can be lost, are slower by orders of magnitude, and need serialisable messages. Cluster sharding assigns actors to nodes by a key, such as device ID, and moves them when nodes join or leave.
Orleans introduced virtual actors, called grains, which always exist logically. The runtime activates a grain on some node the first time it receives a message and deactivates it when idle, so callers never create or destroy them. This removes lifecycle code, but it does not remove distributed-systems problems. During membership changes, check your runtime's documented guarantees about duplicate activations, and protect persisted state with optimistic concurrency such as version checks.
Failure modes and trade-offs
- Unbounded mailbox growth. A slow consumer accumulates messages until memory runs out, with latency rising long before the crash.
- Blocked dispatcher. Synchronous I/O inside actors starves every other actor that shares the pool.
- Restart loops. A poison message or a missing dependency crashes an actor repeatedly; without an intensity limit it consumes CPU and floods logs.
- Lost state on restart. Treating actor memory as durable loses data the first time a supervisor does its job.
- Ask chains. Long synchronous chains of
askcalls add latency at every hop and can deadlock. - Hot actors. One actor per key serialises all work for that key; a popular key becomes a bottleneck that adding nodes cannot fix.
Actors trade shared-state bugs for message-flow bugs. They are a good fit for many independent stateful entities, such as sessions, devices, game objects or accounts. They are a poor fit for bulk data-parallel computation, where a parallel loop over arrays is simpler and faster, and for problems that need atomic changes across many entities, which turn into distributed transactions between actors.
What to do next
- Run the asyncio runtime above, then change it so that insufficient funds is returned as a reply instead of raised, and confirm the balance survives.
- List the actors in your system and decide, for each, whether its state must survive a restart; persist the ones that must.
- Set a mailbox bound and an overflow policy for every actor type, and export mailbox depth as a metric.
- Find every blocking call inside actor code and move it to async I/O or a dedicated dispatcher.
- Draw your supervision tree, choose a strategy and restart intensity for each supervisor, and test what happens when a child crashes repeatedly.
- Add timeouts to every request-reply call and idempotency keys to every state-changing request.
- Load-test with realistic key distributions to find hot actors before production does.