Limiting active work with a semaphore does not necessarily limit task count: creating a task for every item still retains those tasks. A bounded queue plus a fixed worker set controls how much work is scheduled and buffered at once.
Before you start
Use Python 3.11 or later for TaskGroup. Know async functions, await and cooperative cancellation. This example processes integers using a simulated I/O operation. Production operations need their own deadlines and resource cleanup. It does not offload CPU-intensive Python code.
Step-by-step walkthrough
Step 1: Choose two independent bounds
The worker count limits active processing. Queue maxsize limits items waiting to start. When the queue is full, an awaited put pauses the producer. This is useful backpressure rather than unlimited accumulation, but the original input collection can still occupy memory if it was materialized earlier.
Step 2: Give every received item an outcome
Workers get an item, process it, and call task_done in finally. The bookkeeping must happen even when processing fails. Queue completion is distinct from business success: acknowledging a failed item’s queue ownership does not claim that its external operation succeeded.
Step 3: Scope shutdown and failures
After normal production, send one stop marker per worker. TaskGroup owns all producer and worker tasks. An unexpected worker failure propagates through that scope and cancels siblings; swallowing every exception would instead require an explicit per-item failure and retry policy.
Worked scenario
import asyncio
async def process(value):
await asyncio.sleep(0)
return value * 2
async def run(values, workers=3):
if workers < 1:
raise ValueError('workers must be positive')
queue = asyncio.Queue(maxsize=4)
results = []
async def produce():
for value in values:
await queue.put(value)
for _ in range(workers):
await queue.put(None)
async def consume():
while True:
value = await queue.get()
try:
if value is None:
return
results.append(await process(value))
finally:
queue.task_done()
async with asyncio.TaskGroup() as group:
group.create_task(produce())
for _ in range(workers):
group.create_task(consume())
return results
print(sorted(asyncio.run(run(range(8)))))
# [0, 2, 4, 6, 8, 10, 12, 14]There are three consumers and one producer, not eight processing tasks. The queue holds at most four waiting entries, including stop markers. None is reserved as termination, so this interface deliberately accepts integer work items rather than treating None as ordinary data. Sorting the displayed output avoids claiming completion order equals input order.
Common mistake
Returning a huge results list recreates unbounded memory on the output side even if the input queue is bounded. For large processing, emit results to a bounded downstream queue or durable sink. Similarly, eagerly loading every input before run begins is not made memory-efficient by this queue.
Verify the behavior
Track active processing and maximum queue size under delayed process calls. Assert activity never exceeds workers and buffering never exceeds maxsize. Test empty input, one worker and invalid worker count. Make process raise and confirm the task group reports failure rather than hanging a producer permanently against a full queue.
Interview exercise
How would you preserve input order and continue after individual failures?
Answer and reasoning
Attach an input index to each item and record an indexed result or failure. Decide whether ordering requires buffering unfinished earlier positions, which itself needs a bound. Catch only failures covered by the per-item policy; unexpected faults should still reach the owner. Continuing requires an explicit outcome for every item, not merely ignoring exceptions.
Continue learning
See async task ownership and Python interview questions. Review asyncio queues and TaskGroup.