Walking Trees Concurrently with Bounded Fan-Out¶
Many async jobs walk a tree whose shape is only known as you go: a directory of an object store, an org chart behind an API, the pages of a site, nested comments. Each node's children come from one network call, so the work is naturally concurrent — and naturally unbounded, because the number of calls in flight grows with the tree's width. Measured on Python 3.14 with a tree of 3,906 nodes (branching factor 5, depth 5), each fetch taking 10 ms: a sequential walk took 39.6 s. Recursive gather at every level took 0.24 s by issuing up to 3,125 fetches at once — the whole bottom level — which no real API would accept. A semaphore around the fetch capped concurrency at 20 and took 2.16 s, but still created 3,907 tasks; a queue served by 20 workers took 2.13 s with 22 tasks and 0.2 MiB of traced memory instead of 5.5 MiB. Holding the semaphore while waiting for children deadlocked after 20 fetches. And when 1 in 50 fetches failed, a worker pool without per-item error handling lost all 20 workers and hung; with it, the walk finished and recorded 63 errors. This guide builds the bounded walker.
Prerequisites¶
- Python 3.11+ and an async client for the tree's API.
- Queues and workers, from Async Queue Management.
- The topic overview, Coroutine Design Patterns.
1. See why unbounded recursion is fast and unusable¶
The natural recursive walk fetches a node's children, then walks all of them concurrently:
async def walk(node) -> int:
children = await fetch_children(node)
return 1 + sum(await asyncio.gather(*(walk(c) for c in children)))
Measured: 0.24 s for all 3,906 nodes — about one round trip per level — with a peak of 3,125 fetches in flight and 3,132 tasks alive. That peak is the tree's widest level, so it grows with the data: a tree with a million leaves would open a million connections. Against a real service, this is a burst the server will throttle or refuse, and against your own process, a task and buffers per node. The sequential alternative, an explicit stack popped one node at a time, never had more than one fetch in flight and took 39.6 s. The goal is the middle: a fixed number of fetches in flight, whatever the tree looks like.
Verify: for a recursive walk, you know the peak number of concurrent calls it makes and that it is bounded.
2. Do not hold a semaphore while waiting for children¶
Adding a semaphore to the recursive walk is the obvious fix, and where it goes decides whether it works. Acquired around the whole node, it deadlocks:
async def walk(node) -> int:
async with SLOTS: # held while children are walked...
children = await fetch_children(node)
return 1 + sum(await asyncio.gather(*(walk(c) for c in children))) # ...and they need slots
Measured with 20 slots: 20 fetches completed, then nothing — every slot was held by a node waiting for children that were waiting for a slot. No exception, no timeout of its own; the walk was stopped by an outer five-second deadline. The rule is general: never hold a limited resource while waiting for work that needs the same resource. Acquire it only around the I/O:
async def walk(node) -> int:
async with SLOTS:
children = await fetch_children(node) # the slot covers the fetch only
return 1 + sum(await asyncio.gather(*(walk(c) for c in children)))
Measured: 3,906 nodes in 2.16 s with exactly 20 fetches in flight at the peak. Concurrency is bounded — but the number of tasks is not: every node still gets a task, and 3,907 of them were alive at once, most waiting on the semaphore.
Verify: no code path acquires a slot and then awaits work that needs another slot from the same pool.
3. Use a queue and a fixed set of workers¶
To bound tasks as well as calls, invert the structure: put nodes on a queue, and let a fixed number of workers take a node, fetch its children, and put them back on the queue. Queue.join() returns when every node put on the queue has been processed:
async def walk(root, workers: int = 20) -> int:
queue: asyncio.Queue = asyncio.Queue()
queue.put_nowait(root)
visited = 0
async def worker():
nonlocal visited
while True:
node = await queue.get()
try:
for child in await fetch_children(node):
queue.put_nowait(child)
visited += 1
finally:
queue.task_done()
tasks = [asyncio.create_task(worker()) for _ in range(workers)]
try:
await queue.join()
finally:
for t in tasks:
t.cancel()
await asyncio.gather(*tasks, return_exceptions=True)
return visited
Measured: 2.13 s, 20 fetches in flight, 22 tasks alive and 0.2 MiB of traced memory, against 5.5 MiB for the semaphore version — the queue holds plain node values, not suspended coroutines. The walk order is breadth-first-ish rather than depth-first; where order matters, use a LifoQueue for depth-first or a PriorityQueue keyed on depth. For graphs rather than trees, keep a seen set and add a node to the queue only the first time it appears, so cycles and shared children are fetched once.
Verify: the number of tasks during the walk stays at the worker count, and the visited count matches the tree's size.
4. Handle errors per node, or the walk hangs¶
In the worker pattern, an exception that escapes a worker ends that worker — and join() waits for nodes that no worker remains to process:
node = await queue.get()
try:
for child in await fetch_children(node): # raises ConnectionError for some nodes
queue.put_nowait(child)
finally:
queue.task_done() # counted done, but the worker has died
Measured with 1 in 50 nodes at depth 3 raising ConnectionError: workers died one by one, and after five seconds all 20 were gone, 181 nodes sat in the queue, and join() was still waiting. Catch per-node errors inside the worker, record them, and carry on:
try:
children = await fetch_children(node)
except (ConnectionError, TimeoutError) as exc:
errors.append((node, exc)) # the subtree below this node is skipped
children = []
Measured: the walk finished in 1.03 s, visited 2,016 nodes and recorded 63 errors — the failed nodes' subtrees were skipped, and the caller can retry them or report them. Errors that should stop the whole walk — authentication failures, a deleted root — should instead cancel it: raise them from the worker into a TaskGroup rather than catching them, as described in handling per-item errors in async pipelines.
Verify: a walk with injected per-node failures completes and reports them; a walk with a fatal error stops promptly.
5. Bound depth, size and time¶
A tree that comes from someone else's API can be deeper or wider than expected — or not a tree at all. Put limits on everything the walk can grow:
MAX_DEPTH, MAX_NODES = 12, 100_000
async def worker():
while True:
node, depth = await queue.get()
try:
if depth >= MAX_DEPTH or len(seen) >= MAX_NODES:
truncated.append(node)
continue
for child in await fetch_children(node):
if child.id not in seen:
seen.add(child.id)
queue.put_nowait((child, depth + 1))
finally:
queue.task_done()
async with asyncio.timeout(300): # the whole walk
await walk(root)
Each limit reports what it cut off rather than silently dropping it, so a truncated walk is visible. The overall timeout cancels the walker's join(); the finally in step 3 then cancels the workers, so a timed-out walk leaves no tasks behind. Choose the worker count from what the API allows, not from the tree: the walk's concurrency is the same 20 whether the tree has 4,000 nodes or 4 million, which is the property that makes it safe to point at data you have not seen. The same shape bounds a web crawler, as in limiting concurrency per host in an async crawler.
Verify: walks of deliberately deep, wide or cyclic test trees stop at their limits and report what they skipped.
Verification¶
A concurrent tree walk is bounded and robust when:
- In-flight calls are capped by workers or a semaphore that covers only the I/O.
- Tasks are bounded too, by using a queue and fixed workers for large trees.
- Per-node errors are caught and recorded, so workers never die and
join()always returns. - Depth, node count and total time are limited, with truncation reported.
Diagnostic Hook: export the walker's queue size and worker count alive. A queue that keeps growing means the tree is wider than the workers can keep up with — fine, if memory allows; a queue with items and no live workers means workers died on an unhandled error, and the walk will never finish on its own.
Pitfalls & edge cases¶
- Recursive gather on unknown data. Measured: 3,125 calls in flight at once.
- A semaphore held across recursion. Measured: deadlock after 20 fetches.
- A semaphore without a task bound. Measured: 3,907 tasks for a 20-call limit.
- Exceptions escaping workers. Measured: all 20 died and
join()hung.
Frequently Asked Questions¶
How do I limit concurrency when recursively crawling a tree in asyncio?
Use a queue with a fixed number of worker coroutines: 20 workers walked 3,906 nodes in 2.13 s with 20 calls in flight and 22 tasks. A semaphore around each fetch also bounds calls but created a task per node.
Why does my recursive asyncio walk with a semaphore hang?
The semaphore is held while waiting for children that need it too: with 20 slots, the walk stopped after 20 fetches. Acquire it only around the I/O call.
Why does queue.join() never return?
Usually because workers died on exceptions and the queue still has items. Catch per-item errors inside the worker; with that, a walk with 63 failing nodes finished in 1.03 s.
How fast is a concurrent tree walk compared with a sequential one?
For 3,906 nodes at 10 ms per fetch: 39.6 s sequentially, 2.13 s with 20 workers, 0.24 s with unbounded concurrency that no real API would accept.
Related¶
- Coroutine Design Patterns — up to the topic overview.
- Building an async event emitter — fan-out to listeners rather than children.
- Asyncio Fundamentals & Event Loop Architecture — the section overview.