新闻详情

新闻详情

首页 / 资讯中心 / 详情

asyncio.Queue 有界队列实战:背压、task_done、join 与异常收尾

发布时间:2026/10/1 7:09:35来源:尧图网络
asyncio.Queue 有界队列实战:背压、task_done、join 与异常收尾
把并发数设成20内存就一定稳定吗如果先为几十万个输入创建Task再让它们等信号量正在执行的数量受限等待任务却已经留在内存里。asyncio.Queue(maxsizeN)可以控制排队项数但前提是生产者逐项await put()消费者数量固定上下游也遵守容量约定。本篇用无网络的模拟流水线验证满队列等待、空队列仍未完成、正常收尾及异常传播。分清三种数量qsize()只统计还在队列里的项目worker 已取走、尚未处理完的项目不在其中。并发数是当前同时执行的工作量未完成计数则由put()增加、task_done()减少。例如队列容量3worker有2个每个worker一次只持有1项单个生产者准备1项后等待入队。在这个受限模型里相关项目可能分布为“排队3处理中2生产者手里1”而不是总共只有3项。若输入每项大小不同这也不是一个固定字节数上限。有界队列满时await put()挂起当前生产者让它暂时别继续取新输入。这个反馈就是背压。它让积压有上限不会凭空让消费者变快也不是每秒固定请求数的限速器。一个完整、可退出的最小流水线本次实际运行macOS arm64、CPython 3.13.132026-09-30。示例要求Python 3.11及以上因为用了TaskGroup没有使用3.13才增加的Queue.shutdown。这里只以异步短暂等待模拟I/O不代表真实爬虫或模型接口性能。保存为queue_lab.py运行python3 queue_lab.pyimportasyncioasyncdefboundaries():qasyncio.Queue(maxsize2)awaitq.put(A)awaitq.put(B)enteredasyncio.Event()asyncdefput_third():entered.set()awaitq.put(C)putterasyncio.create_task(put_third())awaitentered.wait()assertnotputter.done()andq.qsize()2print(full queue: third put is waiting)assertawaitq.get()Aq.task_done()awaitputterassertawaitq.get()Bq.task_done()assertawaitq.get()Cassertq.empty()waiterasyncio.create_task(q.join())awaitasyncio.sleep(0)assertnotwaiter.done()print(empty queue: join still waits for task_done)q.task_done()awaitwaiterprint(acknowledged: join finished)asyncdefpipeline(fail_atNone):qasyncio.Queue(maxsize3)stopobject()stats{active:0,peak_active:0,peak_queued:0,success:0}asyncdefworker():whileTrue:itemawaitq.get()try:ifitemisstop:returnstats[active]1stats[peak_active]max(stats[peak_active],stats[active])try:awaitasyncio.sleep(0.001)# Simulated I/O, no network requests.ifitemfail_at:raiseRuntimeError(synthetic failure)stats[success]1finally:stats[active]-1finally:q.task_done()# Accounting; does not mean business success.asyncwithasyncio.TaskGroup()asgroup:for_inrange(2):group.create_task(worker())foriteminrange(20):# Lazy source, no list of 20 tasks.awaitq.put(item)stats[peak_queued]max(stats[peak_queued],q.qsize())for_inrange(2):awaitq.put(stop)awaitq.join()returnstatsasyncdefmain():awaitboundaries()statsawaitpipeline()assertstats[success]20andstats[active]0assertstats[peak_active]2andstats[peak_queued]3print(normal:,stats)try:awaitpipeline(fail_at4)except*RuntimeError:print(failure: TaskGroup raised, no success report)else:raiseAssertionError(failure was swallowed)if__name____main__:asyncio.run(main())本机实际输出full queue: third put is waiting empty queue: join still waits for task_done acknowledged: join finished normal: {active: 0, peak_active: 2, peak_queued: 3, success: 20} failure: TaskGroup raised, no success report为什么空队列仍然等不到 join实验取出了最后一个项目C但还没有调用它对应的task_done()。这时q.empty()已经为真q.join()仍在等待。取出、处理与确认结束是不同动作。示例每次成功get()后进入try/finally确保相应计数被结算一次结束哨兵也是入过队的项目因此也会结算。task_done()没有成功或失败参数。这里故意把业务成功数量单独计入success而不是根据join返回推断全部成功。如果实际工作失败不能仅在finally结算后吞掉异常并报“全部完成”。本示例发生合成错误时TaskGroup取消其他成员等待它们退出再把异常交给上层打印的是失败分支没有正常成功统计。取消不会撤销已经完成的写入恢复任务仍需要独立的业务状态和幂等设计。让容量约定贯穿整条链路生产者使用惰性的range(20)逐项入队并没有先构造20个Task。真实项目可逐页读数据或逐批读取文件避免在排队之前就把全量响应或所有URL加载进内存。也不要对每次q.put(item)再包一层无界create_task否则等待会转移到队列外。worker数量控制处理并发队列容量控制等待项数每秒请求预算、目标站点规则、连接池上限与单项超时需要另外设置。如果目标接口允许的持续吞吐低于生产速度任何有限队列最终都会填满系统必须选择等待、拒绝、落盘或降采样等明确策略。输出端同样会积压。示例只保留固定大小的计数真实代码如果把全部结果追加到列表输入队列再小输出列表仍会持续增长。大响应应考虑流式读取、大小限制和及时写入持久化存储。两个容易漏掉的故障一是动态爬虫“所有消费者都回填同一个满队列”。如果每个worker都在等put()却没有worker继续get()就可能相互等待。不能因为有Queue就保证任何拓扑都不会死锁。需要重新设计任务发现与调度边界例如让专门调度器管理待抓取集合并给发现结果的传递设置明确容量和溢出策略。二是把失败后的join()当作兜底清理。TaskGroup失败会中止这一批队列里可能还有未取走的项目不应退出异常处理后再无条件等待同一个join。示例让生产、消费者和等待过程处于同一个TaskGroup作用域失败直接传播调用方据此决定如何持久化与重启。这份内存队列本身不耐进程崩溃不保证恰好一次执行也不跨线程安全。这里用它验证单事件循环内的背压和生命周期需要跨进程、持久化与可靠交付时要采用相应的任务存储和确认协议。怎样选择容量先观察单项大小、平均与尾部处理耗时、允许等待时间和上下游速度再确定worker数与缓冲量。容量加大只是允许更多等待不是吞吐优化的证据。监控至少应区分排队量、在途量、最老任务等待时间、成功数和失败数。正常运行结果里的peak_queued3与peak_active2是本次合成实验观察20项全部成功也是本地结果。它们证明示例遵守自己的数量约束不能推导线上内存占用、吞吐或收益。参考Python asyncio.Queue 文档、TaskGroup 官方文档。
网站建设高端定制企业官网
RELATED

相关资讯

更多精彩内容,欢迎继续阅读

较早相关资讯

最新相关资讯

FRI 与 KG-TOWER 二次开发教程(20):收官——FRI/KG-TOWER 二次开发水力学核算工具包(完整项目) 2026/10/1 11:06:33

FRI 与 KG-TOWER 二次开发教程(20):收官——FRI/KG-TOWER 二次开发水力学核算工具包(完整项目)

FRI 与 KG-TOWER 二次开发教程(20):收官——FRI/KG-TOWER 二次开发水力学核算工具包(完整项目)版本与事实声明 项目名 fri_kt,版本 0.1.0;环境:Python 3.8,numpy 2.2.6、…

阅读更多 →
AI+CAD落地实战:从DXF/DWG解析到FreeCAD批量改图的工程避坑指南 2026/10/1 11:06:33

AI+CAD落地实战:从DXF/DWG解析到FreeCAD批量改图的工程避坑指南

1. 为什么“AI CAD”看起来很美,落地却处处碰壁过去两年,我参与过三个跟“AI 辅助 CAD”相关的内部项目,也帮朋友的公司做过几次技术选型评估。一个非常明显的感受是:Demo 满天飞,工程走不通。你在网上能看到大量“上…

阅读更多 →
B3616队列模板题全解:从FIFO到单调队列与消息队列 2026/10/1 11:06:32

B3616队列模板题全解:从FIFO到单调队列与消息队列

刷过洛谷“模板”系列的人,大概率都跟这道 B3616 打过照面。它挂着“【模板】队列”的名头,看起来就是一道入门的不能再入门的裸题,但很多新手恰恰就是在这里翻了车——不是不会队列,而是不会“正确地模拟队列”。这道题表面上在考…

阅读更多 →
JDK 8u131安装与生产环境适配实战指南 2026/10/1 11:06:26

JDK 8u131安装与生产环境适配实战指南

1. 为什么现在还要讲 JDK 8u131?这不是“古董”吗?JDK 8u131 这个版本,乍一看确实像考古现场——它发布于2017年4月,距今已超七年。但如果你正在维护一套运行在金融核心系统、电力调度平台、大型国企ERP或老版本Spring Boot 1.x微…

阅读更多 →
GAMMA 2020在Ubuntu 20.04安装全指南:从解压到License配置 2026/10/1 11:06:26

GAMMA 2020在Ubuntu 20.04安装全指南:从解压到License配置

用了大半天时间,总算在实验室那台Ubuntu 20.04工作站的坑坑洼洼里把GAMMA 2020版装利索了。这期间朋友问得最多的就是:GAMMA到底怎么装?和网上搜出来的gamma校正是一回事吗?这里先把最关键的结论放在开头:如果你做的是…

阅读更多 →
Redis三件套实战:redis-cli、hiredis与redis-benchmark完全指南 2026/10/1 11:06:12

Redis三件套实战:redis-cli、hiredis与redis-benchmark完全指南

1. 开篇:Redis三件套,到底指的是哪三样如果你跟Redis打交道超过一周,迟早会在各种文档、招聘要求、生产事故复盘里撞见这三个名字:redis-cli、hiredis 和 redis-benchmark。它们偶尔被混为一谈,但实际上分工完全不同。…

阅读更多 →

今日资讯

本周资讯

本月资讯

看完文章仍有疑问?

联系尧图顾问,获取一对一建站咨询

立即免费咨询 📞 400-888-8888
📞 ✉