Skip to content

Queue、Producer / Consumer 与 Race

1. 为什么这组知识要一起学

掌握 Agent Runtime 的 I/O 并发底座,理解 Structured Concurrency、Cancellation、Backpressure 与并发模型选择。

这几个知识点处在同一条工程链上。如果只记单个名词,很容易在真实系统里把责任放错层:例如让模型管理程序事实、让数据库 transaction 承担外部 API 原子性,或把一个 provider SDK 的行为误认为 Agent 的通用规律。

2. Mental Model

asyncio 通过 cooperative scheduling 提高 I/O 并发;Structured Concurrency 让 child task 的生命周期归属于明确父作用域。

text
Coroutine → Task/Event Loop → await I/O → other tasks progress → completion/cancellation/ExceptionGroup

3. 核心机制

Queue

bounded asyncio.Queue 用 put 的阻塞实现 backpressure;worker pool 消费任务,shutdown 时要停止接收并 drain/cancel。

Producer / Consumer

bounded asyncio.Queue 用 put 的阻塞实现 backpressure;worker pool 消费任务,shutdown 时要停止接收并 drain/cancel。

Race Condition

单线程 async 仍可能因 await 交错产生逻辑 race;共享状态跨 await 读改写时要用 lock 或更好的 ownership/message passing。

4. 最小实现 / 伪代码

下面代码只表达边界和生命周期,不要求照抄到项目中:

python
sem = asyncio.Semaphore(8)
async def run_one(call):
    async with sem:
        async with asyncio.timeout(call.timeout):
            return await call.execute()

真正实现时应把 I/O、状态持久化、错误翻译和策略注入拆成可测试组件,而不是把示例扩成一个巨型函数。

5. 在 Shadow Harness 中怎么落地

ParallelToolExecutor 用 TaskGroup/Queue/Semaphore;Runtime policy 决定 partial failure,而不是盲从库默认。

建议为本章涉及的行为留下明确的 domain object、interface 和 trace event;只要一个关键行为只能通过读日志猜测,就说明 Runtime contract 仍不够清晰。

6. Production Engineering 检查项

  • blocking I/O 不进 event loop
  • 保存 task 生命周期
  • 正确传播 CancelledError
  • bounded queue 实现 backpressure
  • CPU-bound 转 thread/process

7. Failure Modes

7.1 async def 就以为并发

当出现「async def 就以为并发」时,查看 event loop 是否被 blocking I/O 占用、task 是否 orphan、cancellation 是否被吞、queue/semaphore 是否有界。 这类问题通常需要修改 contract、policy、state 或 adapter,而不是只追加 Prompt。

7.2 无限 create_task

当出现「无限 create_task」时,查看 event loop 是否被 blocking I/O 占用、task 是否 orphan、cancellation 是否被吞、queue/semaphore 是否有界。 这类问题通常需要修改 contract、policy、state 或 adapter,而不是只追加 Prompt。

7.3 吞 cancellation

当出现「吞 cancellation」时,查看 event loop 是否被 blocking I/O 占用、task 是否 orphan、cancellation 是否被吞、queue/semaphore 是否有界。 这类问题通常需要修改 contract、policy、state 或 adapter,而不是只追加 Prompt。

7.4 gather 异常语义与业务预期不符

当出现「gather 异常语义与业务预期不符」时,查看 event loop 是否被 blocking I/O 占用、task 是否 orphan、cancellation 是否被吞、queue/semaphore 是否有界。 这类问题通常需要修改 contract、policy、state 或 adapter,而不是只追加 Prompt。

8. Trade-offs

TaskGroup fail-fast 语义强;gather 更灵活地收集独立结果;Queue/worker pool 适合持续大量任务。

设计记录最好明确:当前约束是什么、备选方案有哪些、为什么现在选这个、未来什么条件出现时需要重构。 这样 ADR 才能随着模型和基础设施变化被重新审视。

9. Experiment / Evaluation

20 个 fake I/O tool 对比 sequential/concurrency 2/5/20,记录 wall time、最大 in-flight、失败/取消传播。

实验应固定数据集、版本和环境,至少记录 success、latency、token/cost、attempt/step 数以及失败类型;涉及随机模型时需要重复运行而不是只看一次结果。

10. 常见问题

基础:Queue 最容易被误解的点是什么?

bounded asyncio.Queue 用 put 的阻塞实现 backpressure;worker pool 消费任务,shutdown 时要停止接收并 drain/cancel。

机制:这些能力在一次 Run 的哪个生命周期阶段生效?

沿着 Coroutine → Task/Event Loop → await I/O → other tasks progress → completion/cancellation/ExceptionGroup 找位置,并明确它的输入、输出、持久化事实和失败传播。

工程:如果这一层失败,应该由谁恢复?

先区分 transient failure、invalid input、permission、semantic failure 与 irreversible side effect。恢复策略属于拥有该状态与副作用的 Runtime/adapter,而不是交给模型自由决定。

设计:规模扩大 10 倍后,哪个假设最先失效?

优先检查 context/token、并发/连接池、catalog 大小、持久化吞吐、trace 体积、身份与租户隔离。不要默认“多加机器”能解决语义和一致性问题。

11. Sources

12. 本文结论

掌握本文的标准不是能背定义,而是能画出数据流、写出最小 contract、解释失败恢复,并用测试或 benchmark 证明设计没有只停留在概念层。

AI Engineering · Agent · Harness · Backend Systems