
本文介绍如何使用 asyncio.queue 构建一个可在运行时动态添加新任务的异步并行处理系统,适用于“边执行边发现新任务”的典型场景(如爬虫、图遍历、依赖解析等)。
本文介绍如何使用 asyncio.queue 构建一个可在运行时动态添加新任务的异步并行处理系统,适用于“边执行边发现新任务”的典型场景(如爬虫、图遍历、依赖解析等)。
在需要并发处理大量任务,且任务本身可能触发新任务生成的场景中(例如:处理一个网页时发现若干待抓取链接;解析一个模块时发现未导入的依赖),传统静态任务池(如 concurrent.futures.ProcessPoolExecutor 或 multiprocessing.Pool)难以满足需求——它们的任务队列通常在启动前固定,不支持运行时动态扩增。
asyncio.Queue 是 Python 异步生态中专为协程设计的线程/协程安全队列,天然支持“生产-消费”模型,并允许在任意协程中通过 put_nowait() 或 put() 动态入队新任务。配合 asyncio.create_task() 启动多个长期运行的工作协程,即可构建一个弹性、可扩展的异步任务调度器。
以下是一个完整、可运行的示例:
import asyncio
async def worker(queue: asyncio.Queue, worker_id: int):
"""工作协程:持续从队列获取任务并处理,支持动态派生新任务"""
while True:
try:
item = await queue.get()
print(f"[Worker-{worker_id}] Processing: {item}")
# 模拟业务逻辑:若任务含 'spawn',则生成两个新任务
if isinstance(item, str) and 'spawn' in item.lower():
print(f"[Worker-{worker_id}] → Generating new jobs: 'subtask_a', 'subtask_b'")
await queue.put('subtask_a')
await queue.put('subtask_b')
else:
# 模拟耗时操作(如网络请求、文件解析)
await asyncio.sleep(0.5)
print(f"[Worker-{worker_id}] ✅ Done: {item}")
queue.task_done()
except asyncio.CancelledError:
print(f"[Worker-{worker_id}] Shutting down...")
break
async def main():
# 创建异步队列
q = asyncio.Queue()
# 启动 3 个并发工作协程
workers = [
asyncio.create_task(worker(q, i))
for i in range(1, 4)
]
# 初始任务入队
initial_jobs = ['job_1', 'job_2', 'spawn_new_jobs', 'job_3']
for job in initial_jobs:
q.put_nowait(job)
# 等待所有已入队任务(含后续动态添加的)完成
await q.join()
# 取消所有工作协程
for w in workers:
w.cancel()
await asyncio.gather(*workers, return_exceptions=True)
if __name__ == '__main__':
asyncio.run(main())✅ 关键特性说明:
- queue.get() 是协程挂起点,自动实现任务分发与负载均衡;
- queue.put() / queue.put_nowait() 可在任意协程(包括正在执行的 worker 内)安全调用,实现队列动态增长;
- queue.task_done() 与 queue.join() 协同实现精确的完成状态跟踪,避免过早退出;
- 工作协程采用 while True + try/except CancelledError 结构,确保优雅终止。
⚠️ 注意事项:
- 不要混用 threading.Thread 或 multiprocessing.Process 与 asyncio.Queue——后者仅限协程环境;
- 若任务涉及 CPU 密集型计算,asyncio 并非最优选择,应考虑 concurrent.futures.ProcessPoolExecutor + 自定义线程安全队列(如 queue.Queue 配合 loop.run_in_executor);
- 生产环境中建议增加异常捕获、重试机制、最大队列长度限制(asyncio.Queue(maxsize=1000))及监控指标,防止内存无限增长。
该方案简洁、高效、符合 Python 异步编程范式,是实现“动态任务流”最自然的解法。

















