在Python的世界里,并发编程一直是个绕不开的话题。多线程受制于GIL,多进程开销又大,而异步编程凭借事件驱动模型,用单线程就能支撑起高并发的IO密集型任务。这篇文章带你从事件循环的底层机制讲起,一步步搞懂异步任务队列的设计思路,并动手实现一个完整的生产者消费者模型。

一、事件驱动到底是怎么回事
要理解异步任务队列,先得弄明白事件驱动这个核心概念。传统同步代码是"你等我、我等他"的串行模式,一个函数调用必须等返回结果才能继续往下走。而事件驱动把控制权交给了事件循环,程序注册好"某个事件发生后要执行什么回调",然后事件循环不断轮询,谁就绪了就调度谁。
在Python的asyncio中,这个调度中枢就是事件循环。它维护着几个关键结构:就绪任务列表、定时器堆、以及IO多路复用器。每当循环运行一轮,它会先把到期的定时器任务拎出来执行,再通过select或epoll等待IO事件,谁的文件描述符可读写就把谁加入就绪队列。整个过程单线程完成,没有线程切换开销,也没有锁竞争问题。
协程是事件循环里被调度的基本单位。当协程执行到await处,它会主动让出控制权,把当前状态(局部变量、执行位置)挂起,事件循环趁机去执行其他就绪的任务。这种协作式调度决定了异步代码的一个铁律:协程内部不能有阻塞调用。你如果在协程里调用了time.sleep(10),整个事件循环都会被卡住十秒,所有任务全部停摆。
import asyncio
import time
# 错误示范:阻塞调用会冻结整个事件循环
async def bad_worker():
time.sleep(2) # 所有协程都被卡住
# 正确做法:使用异步sleep
async def good_worker(name):
print(f"{name} 开始工作")
await asyncio.sleep(2) # 让出控制权,其他任务继续跑
print(f"{name} 完成")
return f"{name} 的结果"
async def main():
results = await asyncio.gather(
good_worker("任务A"),
good_worker("任务B"),
good_worker("任务C")
)
print(results)
asyncio.run(main())上面这段代码中,三个任务的总耗时约等于单个任务的耗时,这就是并发带来的收益。如果把await asyncio.sleep(2)换成time.sleep(2),总耗时就会变成六秒左右,性能差异一目了然。
二、asyncio.Queue:异步任务队列的核心组件
有了事件循环的基础,接下来看任务队列。asyncio.Queue是标准库提供的线程安全(准确说是事件循环内安全)的队列,它与普通queue.Queue最大的区别在于:它的put和get都是协程方法,队列满时生产者会挂起,队列空时消费者会挂起,全程不阻塞事件循环。
一个典型的生产者消费者模型包含三类角色:生产者协程负责往队列里塞任务,消费者协程从队列里取任务并处理,主协程负责协调它们的生命周期。任务队列在这里起到了削峰填谷的作用——当生产速度远超消费速度时,队列缓存住多余任务,避免系统被压垮。
下面是一个完整的实现,包含优雅关闭逻辑,值得仔细阅读注释部分:
import asyncio
import random
async def producer(queue: asyncio.Queue, producer_id: int):
"""生产者:模拟不断产生任务"""
for i in range(5):
item = f"生产者{producer_id}-任务{i}"
await queue.put(item) # 队列满(maxsize)时会自动挂起
print(f"投递: {item}")
await asyncio.sleep(random.uniform(0.1, 0.3))
return None
async def consumer(queue: asyncio.Queue, consumer_id: int):
"""消费者:从队列取任务并处理"""
while True:
item = await queue.get() # 队列空时会挂起等待
try:
print(f"消费者{consumer_id} 正在处理: {item}")
await asyncio.sleep(random.uniform(0.2, 0.5)) # 模拟IO耗时
except Exception as e:
print(f"处理失败: {e}")
finally:
queue.task_done() # 必须调用,join依赖它判断结束
async def main():
queue = asyncio.Queue(maxsize=10)
# 启动3个生产者、2个消费者
producers = [asyncio.create_task(producer(queue, i)) for i in range(3)]
consumers = [asyncio.create_task(consumer(queue, j)) for j in range(2)]
# 等所有生产者完成任务投递
await asyncio.gather(*producers)
# 等队列中所有任务被处理完
await queue.join()
# 取消所有消费者(它们在死循环中等待)
for c in consumers:
c.cancel()
await asyncio.gather(*consumers, return_exceptions=True)
print("全部任务处理完毕")
asyncio.run(main())这段代码有几个容易踩坑的点。第一,task_done()一定要放在finally块里,否则任务抛异常后queue.join()会永远等待。第二,取消消费者任务后要再用gather收集一次,并传return_exceptions=True,否则CancelledError会让主流程报错退出。第三,maxsize参数不是必须的,但设置了它就能实现背压机制,防止队列无限膨胀吃光内存。
三、单进程队列不够用?分布式方案怎么选
asyncio.Queue的局限很明显:它只存在于单个进程的内存里,进程一挂任务就丢了,也无法跨机器分发任务。当业务规模上来后,你需要引入消息中间件,Celery和RQ是最常见的两个选择。
Celery是Python生态里最成熟的分布式任务队列,支持RabbitMQ、Redis等多种broker,功能非常全面:任务重试、定时调度、结果存储、任务链、优先级队列一应俱全。代价是配置复杂、概念较多,学习曲线陡峭。RQ则走极简路线,只依赖Redis,API简单直观,几分钟就能上手,但功能相对单薄,没有复杂的任务编排能力。
| 方案 | 部署复杂度 | 适用场景 | 可靠性 |
|---|---|---|---|
| asyncio.Queue | 零依赖 | 单进程内IO并发 | 进程挂则任务丢失 |
| RQ | 依赖Redis | 中小项目简单异步任务 | Redis持久化可恢复 |
| Celery | 依赖Redis/RabbitMQ | 复杂分布式任务编排 | 完善的重试与确认机制 |
还有一种折中思路:用FastAPI或aiohttp搭一个异步Web服务,把任务接口暴露出去,再结合asyncio.Queue在服务内部做流量整形。这种方案适合任务产生方和消费方都在同一服务内的场景,省去了外部中间件的运维成本。
四、超时控制与异常处理实战
生产环境的任务队列不能只有正常流程,超时和异常处理才是体现工程功底的地方。设想一个爬虫任务队列,某个页面的请求卡死了,如果不做超时控制,一个消费者协程就永久占用了。
asyncio提供了多层超时机制。任务级用asyncio.wait_for,它会在超时后取消内部任务并抛出TimeoutError;批任务级用asyncio.wait,可以设置整体超时并区分成功与失败的任务。此外Python 3.11之后还引入了asyncio.timeout上下文管理器,写法更加简洁。
import asyncio
async def fetch_with_retry(queue, max_retries=3):
while True:
item = await queue.get()
for attempt in range(1, max_retries + 1):
try:
# 单任务最多执行3秒,超时抛TimeoutError
result = await asyncio.wait_for(
handle(item), timeout=3
)
print(f"成功: {item}, 结果: {result}")
break
except asyncio.TimeoutError:
print(f"第{attempt}次尝试超时: {item}")
if attempt == max_retries:
print(f"任务放弃,写入死信队列: {item}")
except Exception as e:
print(f"未知异常,直接放弃: {e}")
break
finally:
queue.task_done()
async def handle(item):
await asyncio.sleep(5) # 模拟一个卡死的任务
return "done"上面展示了重试加超时的组合拳。超过最大重试次数后,任务应写入死信队列留待人工排查,而不是悄悄丢弃——线上任务凭空消失,排查起来是最痛苦的。另外建议给每个任务加上唯一ID,处理时记录结构化日志,这样任务流转的全链路都能追踪。
总结一下选型思路:单进程内的高并发IO,直接用asyncio.Queue配合多个消费者协程就够了;需要跨进程或跨机器分发,Redis加RQ能快速落地;任务链复杂、要定时调度和精细的失败控制,就上Celery。先把事件循环的机制理解透,再去用这些框架,你会发现它们不过是把今天讲的这些模式封装得更完善而已。
Python异步任务队列事件驱动asyncio修改时间:2026-09-16 20:30:52