Sharding State Across asyncio Actors by Key¶
One actor processes one message at a time, so an actor whose handler awaits I/O can do at most 1 / handler time messages per second — about a thousand for a 1 ms Redis write. Sharding splits the state by key across several actors, each with its own mailbox, so different keys proceed in parallel while each key is still handled by exactly one actor, in order. Measured on Python 3.14 with 10,000 messages over 1,000 keys and 1 ms of awaited I/O per message: one actor processed 937 messages per second; 4 shards 3,615; 16 shards 12,674; 64 shards 38,336; 256 shards 90,931 — and per-key order was preserved in every configuration. With a skewed (Zipf 1.1) key distribution, where the hottest key carried 18% of all messages, throughput stopped at 4,906 messages per second however many shards were added: the shard holding the hot key had to process its messages one at a time. This guide shards actors and handles the keys that refuse to spread.
Prerequisites¶
- Python 3.11+; standard library only.
- An actor with a bounded mailbox, from bounding actor mailboxes under load.
- Per-key ordering in pools, from processing queue items in order per key.
1. Route each key to one shard with a stable hash¶
Each shard is an ordinary actor. A router picks the shard from the key with a hash that is stable across processes and restarts:
import asyncio
import zlib
from collections import defaultdict
class Shard:
def __init__(self) -> None:
self.mailbox: asyncio.Queue[tuple[str, object] | None] = asyncio.Queue(maxsize=1000)
self.state: dict[str, list] = defaultdict(list)
async def run(self) -> None:
while (msg := await self.mailbox.get()) is not None:
key, value = msg
await write_through(key, value) # ~1 ms of I/O
self.state[key].append(value)
class Router:
def __init__(self, n: int) -> None:
self.shards = [Shard() for _ in range(n)]
def shard_for(self, key: str) -> Shard:
return self.shards[zlib.crc32(key.encode()) % len(self.shards)]
async def tell(self, key: str, value: object) -> None:
await self.shard_for(key).mailbox.put((key, value))
zlib.crc32 is used instead of the built-in hash() because string hashing is randomized per process (PYTHONHASHSEED), so hash(key) % n sends the same key to different shards in different processes — harmless within one process, wrong as soon as two processes or a restart must agree. Run in three separate processes, hash('user:42') % 16 gave shards 11, 1 and 8; zlib.crc32(b'user:42') % 16 gave 6 every time. Measured: with uniformly distributed keys, throughput rose from 937 messages per second with one shard to 90,931 with 256, and every key's values were recorded in the order they were sent.
Verify: for each key, the sequence of values the shard processed equals the sequence sent.
2. Choose the shard count from the handler's wait¶
Shards help in proportion to how much of the handler is waiting rather than computing. With 1 ms of awaited I/O, 16 shards gave 13.5× the single-actor throughput and 64 gave 41×. Beyond that the gains flattened — 256 shards gave 97× rather than 256× — as the event loop itself, the router and the per-message overhead began to show:
# Rough sizing: enough shards to cover the target rate at the handler's latency
target_rate = 20_000 # messages/s
handler_seconds = 0.001 # measured p50 of the awaited I/O
shards = math.ceil(target_rate * handler_seconds * 1.5) # 30, with 50% headroom
The shard count is also an upper bound on concurrent I/O against the backing store, so it doubles as a concurrency limit: 256 shards means up to 256 concurrent Redis writes. If the store cannot take that, the right number of shards is what the store can absorb, not what the event loop can drive — the same reasoning as sizing async connection pools for throughput. Handlers that are CPU-bound gain nothing from more shards on one event loop; they need processes.
Verify: throughput at the chosen shard count meets the target with the backing store's latency still at its baseline.
3. Find and handle hot keys¶
Sharding spreads keys, not load. When one key carries a large share of the traffic, its shard is the bottleneck:
from collections import Counter
counts = Counter(keys)
hottest, n = counts.most_common(1)[0]
print(hottest, n / len(keys)) # measured: k0 carried 18% of messages
Measured with a Zipf 1.1 distribution: 2,455 messages per second at 4 shards, 4,292 at 16, 4,665 at 64 and 4,906 at 256. The busiest shard handled 19–20% of all messages from 64 shards upwards — essentially the hot key alone — and since that shard processes one message per millisecond, it alone took about 1.9 seconds for the run, setting the total. Three remedies, in order of preference: make the hot key's handling cheaper (batch its messages, as in building an actor with an asyncio mailbox); split the hot key's state if its operations commute (counters can be sub-sharded and summed); or accept the ceiling and protect everyone else by giving the hot key its own shard so it does not slow the keys that happen to hash next to it.
Verify: export per-shard message counts; the busiest shard's share should be close to 1 / shards, and a share far above it names a hot key.
4. Batch inside the hot shard¶
The cheapest fix for a hot key is to stop paying the I/O cost per message. A shard that drains its mailbox and applies everything it found in one write turns a hot key's burst into a single round trip:
async def run(self) -> None:
while (first := await self.mailbox.get()) is not None:
batch = [first]
while len(batch) < 200 and not self.mailbox.empty():
nxt = self.mailbox.get_nowait()
if nxt is None:
await self._flush(batch)
return
batch.append(nxt)
await self._flush(batch)
async def _flush(self, batch: list[tuple[str, object]]) -> None:
by_key: dict[str, list] = defaultdict(list)
for key, value in batch:
by_key[key].append(value) # order within each key is kept
await write_many(by_key) # one round trip
for key, values in by_key.items():
self.state[key].extend(values)
Grouping by key inside the batch keeps per-key order, because the batch was taken from the shard's mailbox in arrival order and each key's values are appended in that order. Under light load batches stay at size one and latency is unchanged; under a hot-key burst they grow toward the cap and the shard's throughput grows with them.
Verify: under a skewed load test, the hot shard's batch sizes rise above one and its mailbox wait falls.
5. Keep routing stable when the shard count changes¶
crc32(key) % n remaps most keys when n changes, which is fine for in-process shards rebuilt at startup and wrong for anything that persists per-shard state. When shards own durable state or caches, use a fixed number of virtual slots mapped onto however many shard actors exist:
SLOTS = 1024 # fixed forever
def slot_for(key: str) -> int:
return zlib.crc32(key.encode()) % SLOTS
class SlotRouter:
def __init__(self, n_shards: int) -> None:
self.shards = [Shard() for _ in range(n_shards)]
self.owner = [i % n_shards for i in range(SLOTS)] # slot -> shard
def shard_for(self, key: str) -> Shard:
return self.shards[self.owner[slot_for(key)]]
A key's slot never changes; rebalancing moves whole slots between shards, and only the keys in moved slots need their state transferred. This is the scheme Redis Cluster uses with 16,384 slots, and it is the same idea as consistent hashing with less machinery. For in-process actors whose state is rebuilt on start, plain modulo is enough.
Verify: changing the shard count changes the owner of moved slots only, and every key in an unmoved slot stays on its shard.
Verification¶
Sharded actors scale correctly when:
- Keys are routed with a stable hash (
crc32, nothash()), and each key always reaches the same shard. - Per-key order is preserved, checked by comparing processed sequences with sent ones.
- The shard count covers the target rate at the measured handler latency, within what the backing store can absorb.
- Per-shard load is exported, and hot keys are handled by batching, splitting or isolation.
Diagnostic Hook: chart messages processed per shard. A flat profile means sharding is working; one shard far above the rest names a hot key, and that shard's mailbox wait is the latency every message with that key is paying.
Pitfalls & edge cases¶
hash(key) % n. String hashes are randomized per process; routing differs between processes.- Hot keys. Measured: a key with 18% of traffic capped throughput at 4,906 msg/s.
- More shards than the store can take. Each shard is a concurrent client of the backing store.
- Modulo routing with durable state. Changing
nremaps most keys.
Frequently Asked Questions¶
How do I scale an asyncio actor beyond one task?
Shard it: run several actors and route each message by a stable hash of its key, so each key is always handled by the same actor in order. With 1 ms of I/O per message, throughput rose from 937 to 90,931 messages per second at 256 shards.
Does sharding preserve message order?
Per key, yes: all messages for a key go to one shard's mailbox and are processed in arrival order. There is no ordering across keys on different shards.
Why doesn't adding shards increase throughput?
Usually a hot key: with Zipf-distributed keys, one key carried 18% of messages and throughput stopped near 4,900 messages per second regardless of shard count. Batch within that shard or split the key's state.
Why not use hash() to pick a shard?
Python randomizes string hashing per process, so the same key maps to different shards in different processes or after a restart; use a stable hash such as zlib.crc32.
Related¶
- Actors & Supervision — up to the topic overview.
- Building an actor with an asyncio mailbox — the unit being sharded.
- Concurrent Execution & Worker Patterns — the section overview.