Skip to content

Latest commit

 

History

History
182 lines (127 loc) · 8.58 KB

File metadata and controls

182 lines (127 loc) · 8.58 KB

第 5 阶段:asyncio 异步并发教程

0. 学习目标

这一阶段的目标是:能理解 Python 异步代码,能写并发 HTTP 调用、并发 Agent 评估脚本。

学完后,你应该能:

  • 解释 async def / await / coroutine / event loop 的关系。
  • asyncio.run 作为同步世界进入异步世界的唯一入口。
  • asyncio.gather 并发跑一批任务。
  • asyncio.Semaphore 限制最大并发。
  • asyncio.wait_for / asyncio.timeout 给任务加超时。
  • 做异步异常隔离:一条任务失败不影响其它任务。

这是你迁移成本最低的一段——asyncio 的并发模型和 CompletableFuture 高度对应。但有一个关键心智差异必须先建立(见第 2 节)。

1. Java 对照速查

Java Python asyncio 说明
CompletableFuture<T> coroutine(async def 的返回)/ Task 一个"将来会有结果"的异步操作
future.get() await coro 等结果;asyncio 里 await 不阻塞线程,而是让出 loop
CompletableFuture.allOf(...) asyncio.gather(*coros) 并发等一批完成
响应式 WebClient httpx.AsyncClient 非阻塞 HTTP
Executor / 线程池 event loop 单线程跑协程,靠 await 切换
线程池大小 / 背压 / 限流 asyncio.Semaphore 控制同时在跑的任务数
future.get(timeout, MS) asyncio.wait_for(coro, timeout) 超时控制
try { f.get() } catch 协程里 try/except await 异步异常处理

2. 一个必须先建立的心智差异

asyncio 是 单线程 + 协作式 并发,不是多线程:

  • 没有抢占。只有当某个协程 await 一个 IO 时,event loop 才会切去跑别的协程。
  • 所以它适合 IO 密集(网络、磁盘),不适合 CPU 密集(纯计算占着 loop 不让出,并发就退化成串行)。
  • 协程里绝对不能调阻塞函数time.sleep(1)requests.get(...) 会卡死整个 loop(所有协程一起卡)。必须用 await 版本:asyncio.sleephttpx.AsyncClient

Java 直觉里"开个线程跑就行"在这里不成立。记住:阻塞调用 = 毒药。一个同步阻塞调用能让你所有"并发"瞬间归零。

3. async / await / event loop

import asyncio


async def say(msg: str) -> str:        # async def 定义协程函数
    await asyncio.sleep(0.1)            # await 一个异步 IO,期间让出 loop
    return msg.upper()


async def main() -> None:
    result = await say("hi")           # 在协程里用 await 取结果
    print(result)


asyncio.run(main())                    # 同步世界 -> 异步世界的唯一入口
  • async def 的函数一调用并不执行,只返回一个 coroutine 对象(像没 join 的 future)。必须 await 它或丢给 loop 才真正跑。
  • asyncio.run(coro) 创建 event loop、跑协程、关 loop。整个程序通常只有一处。

4. 任务 1:gather 并发

async def map_concurrent(items, worker):
    return await asyncio.gather(*(worker(item) for item in items))
  • gather(*coros) 同时启动所有协程,全部完成后返回结果列表。
  • 结果顺序 == 输入顺序(不是完成顺序)——和 CompletableFuture.allOf 后按原序 join 一致。
  • 默认任一协程抛异常,gather 立即抛出(其余被取消)。要"不炸整批",见任务 4。
  • ⚠️ gather 不限并发:1000 个 item 就同时发 1000 个请求。要限流见任务 2。

5. 任务 2:Semaphore 限流

async def map_concurrent_bounded(items, worker, *, concurrency):
    sem = asyncio.Semaphore(concurrency)

    async def guarded(item):
        async with sem:                # 超过 concurrency 的协程在这排队
            return await worker(item)

    return await asyncio.gather(*(guarded(item) for item in items))

asyncio.Semaphore(n) 是异步信号量:最多 n 个协程能同时进入 async with sem,其余在此 await 排队。这就是背压 / 限流,等价于固定大小线程池 —— 避免一次性把下游 API 打爆。

参考测试用一个"记录最大并发数"的 worker,断言 concurrency=3 时运行期间峰值恰好是 3,真实证明限流生效(见 tests/test_stage5_asyncio_tasks.py)。

6. 任务加超时

# wait_for:到点抛 asyncio.TimeoutError,并取消被包裹的协程
value = await asyncio.wait_for(worker(item), timeout=5)

# Python 3.11+ 更现代的写法:超时上下文管理器
async with asyncio.timeout(5):
    value = await worker(item)

Java 对照:future.get(5, TimeUnit.SECONDS)。区别是 asyncio 超时会主动取消那个协程(向它抛 CancelledError),不是单纯不等了。

7. 任务 3+4:并发评估 + 超时 + 异常隔离

把限流、超时、异常隔离合到一个 run_eval 里(见参考脚本):

async def run_one(item):
    async with sem:                                  # 限并发
        try:
            value = await asyncio.wait_for(worker(item), timeout)  # 超时
            return TaskResult(id=key(item), ok=True, value=value)
        except asyncio.TimeoutError:
            return TaskResult(id=key(item), ok=False, error="timeout")
        except Exception as exc:                      # 异常隔离:本条挂了,其它继续
            return TaskResult(id=key(item), ok=False, error=str(exc))

results = await asyncio.gather(*(run_one(it) for it in items))

为什么不用 gather(..., return_exceptions=True) 那样异常对象会混进结果列表,还得在外面 zip 回 item 判断类型。把 try/except 放进 run_one,输出直接就是干净的结构化结果(沿用第 2/3 阶段 {id, ok, value, error} 风格),拿到就能写盘 / 统计。

两种异常隔离手段都要会:return_exceptions=True(第 3 阶段 fetch_all 用过)和"每条内部 try/except"(本阶段)。后者更适合要结构化输出的评估场景。

落盘成 JSONL(每行一个 JSON 对象,便于流式追加和 pandas 读取,第 8 阶段会用):

def write_jsonl(path, results):
    lines = [json.dumps(asdict(r), ensure_ascii=False) for r in results]
    Path(path).write_text("\n".join(lines) + "\n", encoding="utf-8")

8. Agent 场景里什么时候需要异步

  • 并发调用多个工具 / 多个 embedding
  • 并发跑评估集(100 条问题同时打模型)
  • 并发抓网页做 RAG 语料
  • FastAPI 的异步接口(第 6 阶段)
  • LangGraph / AutoGen 里的异步节点和异步消息流

9. 常见坑(Java 程序员高频踩)

  • 在协程里调阻塞函数requeststime.sleep、重 CPU 计算)→ 卡死整个 loop。用 await 版本,重活丢 asyncio.to_thread(fn, ...)
  • 忘了 awaitworker(x) 只创建协程不执行,还会报 "coroutine was never awaited"。
  • 不限并发直接 gather 几千个 → 打爆对端或本地端口耗尽。上 Semaphore。
  • 裸 gather 一个失败炸全场:要么 return_exceptions=True,要么每条 try/except。
  • 混用同步异步:异步函数只能在协程里 await 调;顶层用 asyncio.run 进入一次就好。

10. 可运行验证

跑参考脚本("假 LLM" 并发评估 11 条、写 JSONL,无需网络):

.\.venv\Scripts\python.exe scripts\stage5_asyncio_tasks.py --output data\stage5_eval.jsonl

运行测试(全部离线,含并发上限的真实验证):

.\.venv\Scripts\python.exe -m pytest tests\test_stage5_asyncio_tasks.py

11. 本阶段掌握标准

  • 能解释 async / await / coroutine / event loop,以及"单线程协作式"意味着什么。
  • 能用 asyncio.run + asyncio.gather 并发跑任务并保持结果顺序。
  • 能用 asyncio.Semaphore 限制最大并发,并能说清为什么要限。
  • 能用 wait_for / asyncio.timeout 加超时。
  • 能做异步异常隔离,输出结构化结果。
  • 能识别并避免"协程里阻塞调用"这个致命坑。

12. 练习任务(自己动手)

参考脚本已实现 roadmap 第 8 节 4 个任务,可在其上扩展:

  1. run_eval 加进度日志:每完成 N 条打一行(用一个共享计数器 + asyncio.Lock)。
  2. fetch_all 加重试:失败的 URL 退避后再试 K 次(结合第 3 阶段的 RetryPolicy)。
  3. asyncio.as_completed 改写,让结果"谁先完成先处理",对比和 gather 的差别。
  4. 把一段 CPU 密集函数用 asyncio.to_thread 包起来并发跑,观察它不再卡 loop。
  5. 故意在协程里写一个 time.sleep(1),测量总耗时,亲眼看到"并发退化成串行"。