
本文介绍一种基于 asyncio.queue 的轻量、高效、线程安全的异步任务队列方案,支持在任务执行过程中动态追加新任务,适用于爬虫发现新 url、图遍历、依赖解析等场景。
本文介绍一种基于 asyncio.queue 的轻量、高效、线程安全的异步任务队列方案,支持在任务执行过程中动态追加新任务,适用于爬虫发现新 url、图遍历、依赖解析等场景。
在实际开发中,我们常遇到一类“边处理、边发现”的任务场景:初始任务集只是起点,每个任务执行时可能生成新的待处理项(如解析网页时发现新链接、遍历目录时发现子目录、处理消息时触发下游任务)。此时,传统静态任务池(如 concurrent.futures.ProcessPoolExecutor 或 multiprocessing.Pool)无法满足需求——它们的输入队列在提交后即冻结,不支持运行时动态扩容。
asyncio.Queue 正是解决该问题的理想工具。它天然支持协程间安全通信、异步阻塞获取、动态增删,且无需加锁即可保证线程/协程安全。以下是一个完整、健壮的实现示例:
import asyncio
import logging
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
async def worker(queue: asyncio.Queue, worker_id: int):
"""工作协程:持续从队列取任务并执行,支持运行时动态入队"""
while True:
try:
item = await queue.get()
logger.info(f"[Worker-{worker_id}] Processing: {item}")
# 模拟业务逻辑:若为特定标记任务,则动态添加新任务
if isinstance(item, str) and item.startswith("spawn:"):
new_jobs = [f"dynamic-{worker_id}-{i}" for i in range(2)]
logger.info(f"[Worker-{worker_id}] Spawning: {new_jobs}")
for job in new_jobs:
await queue.put(job) # 使用 await put() 更稳妥(自动等待队列空间)
else:
# 模拟耗时处理
await asyncio.sleep(0.5)
logger.info(f"[Worker-{worker_id}] Done: {item}")
except asyncio.CancelledError:
logger.info(f"[Worker-{worker_id}] Shutting down.")
break
finally:
queue.task_done()
async def main():
# 创建任务队列(可设 maxsize 实现背压控制)
queue = asyncio.Queue(maxsize=100)
# 启动 3 个并发工作协程
workers = [
asyncio.create_task(worker(queue, i))
for i in range(3)
]
# 初始任务(含触发动态扩展的任务)
initial_jobs = ["job-1", "job-2", "spawn:trigger", "job-3"]
for job in initial_jobs:
await queue.put(job)
# 等待所有已入队任务完成(包括动态添加的)
await queue.join()
# 取消所有工作协程
for w in workers:
w.cancel()
await asyncio.gather(*workers, return_exceptions=True)
if __name__ == "__main__":
asyncio.run(main())
✅ 关键优势说明:
- queue.get() 是异步阻塞调用,空队列时自动挂起,无忙等待;
- queue.put_nowait() / await queue.put() 均线程安全,后者在队列满时自动等待;
- queue.task_done() 与 queue.join() 配合,精准跟踪任务完成状态,避免过早退出;
- 工作协程使用 while True + try/except CancelledError,确保优雅终止;
- 支持设置 maxsize 实现背压(backpressure),防止内存无限增长。
⚠️ 注意事项:
- 不要在同步函数中直接调用 queue.put_nowait() 后立即 await queue.join() —— 必须在 async 上下文中使用;
- 若任务需 CPU 密集型计算,应改用 loop.run_in_executor() 封装,避免阻塞事件循环;
- 动态添加任务时建议做去重或幂等校验(如用 set 缓存已入队 ID),防止无限递归;
- 生产环境建议增加异常捕获与重试机制(如 try/except Exception 包裹任务体)。
该方案兼顾简洁性与工程鲁棒性,无需引入复杂框架(如 Celery、Dask),即可构建高并发、自适应、可监控的动态任务系统。











