> ## Content Index
> Fetch the complete content index at: https://blog.vercanti.com/llms.txt
> Use this file to discover other available public pages before exploring further.

# asyncio 异步编程完全指南
- URL: https://blog.vercanti.com/asyncio-yi-bu-bian-cheng-wan-quan-zhi-nan/
- Published: 2026-08-28T14:34:32.000Z
- Updated: 2026-08-28T14:56:44.000Z
- Description: asyncio 适合 I/O 密集型 任务（网络请求、文件读写、数据库查询）。CPU 密集型任务用 multiprocessing 或 concurrent.futures.ProcessPoolExecutor。 Python 3.7+ 推荐的入口，创建事件循环、运行协程、关闭循环： gather 参数说明： return_when 可选值： 异步代码中调用同步阻塞函数（如读文件、调用同步库）会阻塞整个事件循环。应用 run_in_executor 放到线程池/进程池执行： TaskGroup 在任意子任务失败时会自动取消其余任务，并将所有异常汇总为
- Author: yellowdog
- Tags: Python, 基础

> 官方文档：<https://docs.python.org/zh-cn/3/library/asyncio.html>  
> 最后更新：2026-03-29

---

## 1\. 核心概念

### 事件循环、协程、任务

| 概念               | 说明                                 |
| ---------------- | ---------------------------------- |
| 事件循环（Event Loop） | asyncio 的调度核心，负责运行协程、处理 I/O 事件     |
| 协程（Coroutine）    | 用 async def 定义的函数，调用后返回协程对象，不会立即执行 |
| 任务（Task）         | 对协程的封装，提交给事件循环并发运行                 |
| Future           | 低层原语，表示一个异步操作的最终结果                 |
| await            | 挂起当前协程，将控制权交还事件循环，等待目标完成后继续        |

### 同步 vs 异步 I/O 模型

```
同步 I/O：
  请求 A → 等待 → 结果 A → 请求 B → 等待 → 结果 B   （串行，等待期间 CPU 空闲）

异步 I/O：
  请求 A → 挂起 → 请求 B → 挂起 → 结果 A 到达 → 继续 A → 结果 B 到达 → 继续 B
  （并发，等待期间处理其他任务）

```

asyncio 适合 **I/O 密集型** 任务（网络请求、文件读写、数据库查询）。CPU 密集型任务用 `multiprocessing` 或 `concurrent.futures.ProcessPoolExecutor`。

---

## 2\. 基础语法

### async def 和 await

```python
import asyncio

async def fetch_data(name: str, delay: float) -> str:
    print(f"{name} 开始")
    await asyncio.sleep(delay)  # 模拟 I/O 等待，不阻塞事件循环
    print(f"{name} 完成")
    return f"{name} 的结果"

# 运行协程
asyncio.run(fetch_data("任务A", 1.0))

```

### asyncio.run()

Python 3.7+ 推荐的入口，创建事件循环、运行协程、关闭循环：

```python
async def main():
    result = await fetch_data("任务A", 1.0)
    print(result)

asyncio.run(main())  # 程序入口，只调用一次

```

---

## 3\. 并发运行多个协程

### asyncio.gather() — 并发等待所有结果

```python
import asyncio

async def task(name: str, delay: float) -> str:
    await asyncio.sleep(delay)
    return f"{name} done"

async def main():
    # 并发运行，总耗时约等于最长的那个（2s），而非串行的 1+2+3=6s
    results = await asyncio.gather(
        task("A", 1),
        task("B", 2),
        task("C", 3),
    )
    print(results)  # ['A done', 'B done', 'C done']（按提交顺序返回）

asyncio.run(main())

```

`gather` 参数说明：

| 参数                   | 类型                 | 默认值   | 说明                   |
| -------------------- | ------------------ | ----- | -------------------- |
| \*coros\_or\_futures | coroutine / Future | 必填    | 要并发运行的协程或 Future     |
| return\_exceptions   | bool               | False | True 时异常作为结果返回而非向上抛出 |

```python
# return_exceptions=True：某个任务失败不影响其他任务
results = await asyncio.gather(
    task("A", 1),
    broken_task(),    # 会抛出异常
    task("C", 1),
    return_exceptions=True,
)
# results = ['A done', SomeException(...), 'C done']
for r in results:
    if isinstance(r, Exception):
        print(f"任务失败：{r}")

```

### asyncio.create\_task() — 立即启动任务

```python
async def main():
    # create_task 立即将协程提交给事件循环，不需要等待
    task_a = asyncio.create_task(task("A", 2))
    task_b = asyncio.create_task(task("B", 1))

    # 此时 A 和 B 已经在并发运行
    result_a = await task_a
    result_b = await task_b
    print(result_a, result_b)

```

### asyncio.TaskGroup — Python 3.11+（推荐）

```python
async def main():
    async with asyncio.TaskGroup() as tg:
        task_a = tg.create_task(task("A", 1))
        task_b = tg.create_task(task("B", 2))
    # 退出 with 块时自动等待所有任务完成
    # 任意一个任务失败，其余任务会被取消，并抛出 ExceptionGroup
    print(task_a.result(), task_b.result())

```

### asyncio.wait() — 更细粒度的控制

```python
import asyncio

async def main():
    tasks = [asyncio.create_task(task(f"T{i}", i)) for i in range(1, 4)]

    # 等待第一个完成就返回
    done, pending = await asyncio.wait(tasks, return_when=asyncio.FIRST_COMPLETED)

    for t in done:
        print(t.result())

    # 取消未完成的任务
    for t in pending:
        t.cancel()

```

`return_when` 可选值：

| 值                | 说明            |
| ---------------- | ------------- |
| FIRST\_COMPLETED | 第一个任务完成时返回    |
| FIRST\_EXCEPTION | 第一个任务抛出异常时返回  |
| ALL\_COMPLETED   | 所有任务完成时返回（默认） |

---

## 4\. 超时控制

### asyncio.timeout() — Python 3.11+（推荐）

```python
async def main():
    try:
        async with asyncio.timeout(3.0):  # 3 秒内未完成则取消
            result = await slow_operation()
    except TimeoutError:
        print("操作超时")

```

### asyncio.wait\_for() — 兼容旧版本

```python
async def main():
    try:
        result = await asyncio.wait_for(slow_operation(), timeout=3.0)
    except asyncio.TimeoutError:
        print("操作超时")

```

---

## 5\. 同步原语

### Lock — 互斥锁

```python
import asyncio

lock = asyncio.Lock()
shared_resource = []

async def write(value: int):
    async with lock:  # 同一时刻只允许一个协程进入
        shared_resource.append(value)
        await asyncio.sleep(0.01)  # 模拟写入耗时

```

### Semaphore — 控制并发数

```python
import asyncio
import httpx

sem = asyncio.Semaphore(10)  # 最多同时 10 个并发请求

async def fetch(client: httpx.AsyncClient, url: str) -> str:
    async with sem:
        response = await client.get(url)
        return response.text

async def main(urls: list[str]):
    async with httpx.AsyncClient() as client:
        tasks = [fetch(client, url) for url in urls]
        results = await asyncio.gather(*tasks)
    return results

```

### Event — 事件通知

```python
import asyncio

event = asyncio.Event()

async def producer():
    await asyncio.sleep(2)
    print("数据准备好了")
    event.set()  # 通知等待者

async def consumer():
    await event.wait()  # 阻塞直到 event 被 set
    print("开始消费数据")

async def main():
    await asyncio.gather(producer(), consumer())

```

### Queue — 生产者消费者

```python
import asyncio

async def producer(queue: asyncio.Queue):
    for i in range(5):
        await queue.put(i)
        print(f"生产 {i}")
        await asyncio.sleep(0.5)
    await queue.put(None)  # 发送结束信号

async def consumer(queue: asyncio.Queue):
    while True:
        item = await queue.get()
        if item is None:
            break
        print(f"消费 {item}")
        queue.task_done()  # 标记任务完成

async def main():
    queue = asyncio.Queue(maxsize=3)  # 最多缓存 3 个
    await asyncio.gather(producer(queue), consumer(queue))

```

---

## 6\. 在线程/进程中运行阻塞代码

异步代码中调用同步阻塞函数（如读文件、调用同步库）会阻塞整个事件循环。应用 `run_in_executor` 放到线程池/进程池执行：

```python
import asyncio
import time
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor

def blocking_io(n: int) -> int:
    time.sleep(1)  # 阻塞操作
    return n * 2

def cpu_bound(n: int) -> int:
    return sum(range(n))  # CPU 密集

async def main():
    loop = asyncio.get_event_loop()

    # 线程池（适合 I/O 阻塞）
    with ThreadPoolExecutor(max_workers=4) as pool:
        result = await loop.run_in_executor(pool, blocking_io, 10)
        print(result)  # 20

    # 进程池（适合 CPU 密集）
    with ProcessPoolExecutor() as pool:
        result = await loop.run_in_executor(pool, cpu_bound, 10_000_000)
        print(result)

```

### asyncio.to\_thread() — Python 3.9+（更简洁）

```python
import asyncio

async def main():
    result = await asyncio.to_thread(blocking_io, 10)
    print(result)

```

---

## 7\. 异步上下文管理器与迭代器

### 异步上下文管理器

```python
class AsyncDBConnection:
    async def __aenter__(self):
        self.conn = await connect_db()
        return self.conn

    async def __aexit__(self, *args):
        await self.conn.close()

async def main():
    async with AsyncDBConnection() as conn:
        await conn.execute("SELECT 1")

```

### 异步迭代器

```python
class AsyncRange:
    def __init__(self, stop: int):
        self.stop = stop
        self.current = 0

    def __aiter__(self):
        return self

    async def __anext__(self):
        if self.current >= self.stop:
            raise StopAsyncIteration
        await asyncio.sleep(0.1)
        value = self.current
        self.current += 1
        return value

async def main():
    async for i in AsyncRange(5):
        print(i)

```

---

## 8\. 常用代码段

### 限速并发请求

```python
import asyncio
import httpx
from typing import Any

async def fetch_all(urls: list[str], concurrency: int = 10) -> list[Any]:
    sem = asyncio.Semaphore(concurrency)

    async def fetch(client: httpx.AsyncClient, url: str):
        async with sem:
            r = await client.get(url, timeout=10)
            r.raise_for_status()
            return r.json()

    async with httpx.AsyncClient() as client:
        return await asyncio.gather(*[fetch(client, u) for u in urls])

```

### 带重试的异步任务

```python
import asyncio

async def with_retry(coro_fn, *args, retries: int = 3, delay: float = 1.0):
    for attempt in range(retries):
        try:
            return await coro_fn(*args)
        except Exception as e:
            if attempt == retries - 1:
                raise
            await asyncio.sleep(delay * (2 ** attempt))  # 指数退避

```

### 异步定时任务

```python
import asyncio

async def periodic(interval: float):
    while True:
        await do_work()
        await asyncio.sleep(interval)

async def main():
    task = asyncio.create_task(periodic(60))
    # 运行主程序...
    await asyncio.sleep(3600)
    task.cancel()

```

### 超时 + 重试组合

```python
import asyncio

async def fetch_with_timeout_retry(url: str, timeout: float = 5.0, retries: int = 3):
    for attempt in range(retries):
        try:
            async with asyncio.timeout(timeout):
                return await fetch(url)
        except (TimeoutError, Exception) as e:
            if attempt == retries - 1:
                raise
            await asyncio.sleep(2 ** attempt)

```

---

## 9\. 最佳实践

### 只在顶层调用 asyncio.run()

```python
# 正确：顶层入口调用一次
asyncio.run(main())

# 错误：在协程内部调用 asyncio.run()（会嵌套事件循环）
async def bad():
    asyncio.run(some_coro())  # 报错

```

### 不要在协程中调用阻塞函数

```python
import time

# 错误：直接调用阻塞函数，会卡住整个事件循环
async def bad():
    time.sleep(1)       # 阻塞！
    requests.get(url)   # 阻塞！

# 正确
async def good():
    await asyncio.to_thread(time.sleep, 1)
    await httpx.AsyncClient().get(url)  # 使用异步 HTTP 库

```

### 用 TaskGroup 代替裸 gather（Python 3.11+）

`TaskGroup` 在任意子任务失败时会自动取消其余任务，并将所有异常汇总为 `ExceptionGroup`，比 `gather` 行为更安全可预期。

### 合理设置并发上限

不要无限制地 `gather` 大量任务，始终通过 `Semaphore` 限制并发数，避免连接池耗尽或目标服务被压垮。

---

## 10\. 踩坑与注意事项

### 协程对象未被 await 会静默丢弃

```python
async def main():
    fetch_data("A", 1)  # 没有 await，协程对象被创建但从未执行
    # Python 会警告：RuntimeWarning: coroutine 'fetch_data' was never awaited

```

### asyncio.sleep(0) 主动让出控制权

在长时间 CPU 运算中，插入 `await asyncio.sleep(0)` 可以让事件循环有机会处理其他事件：

```python
async def heavy_work(items):
    for i, item in enumerate(items):
        process(item)
        if i % 100 == 0:
            await asyncio.sleep(0)  # 每处理 100 个让出一次

```

### 不能在非异步上下文中直接 await

```python
# 普通函数中无法使用 await
def sync_func():
    result = await some_coro()  # SyntaxError

# 如果必须在同步代码中调用协程
result = asyncio.run(some_coro())      # 如果当前没有运行中的事件循环
result = asyncio.get_event_loop().run_until_complete(some_coro())  # 旧写法

```

### CancelledError 不能被吞掉

```python
# 错误：捕获了 CancelledError 却没有重新抛出，导致任务无法被取消
async def bad():
    try:
        await asyncio.sleep(10)
    except Exception:  # 意外捕获了 CancelledError
        pass

# 正确：单独处理 CancelledError
async def good():
    try:
        await asyncio.sleep(10)
    except asyncio.CancelledError:
        # 做清理工作
        raise  # 必须重新抛出
    except Exception as e:
        handle(e)

```

---

## 最佳实践

**`asyncio.gather` vs `asyncio.TaskGroup`**：Python 3.11+ 优先用 `TaskGroup`，其中任一任务失败时会取消其余任务，比 `gather` 的默认行为（继续运行其他任务）更安全：

```python
async with asyncio.TaskGroup() as tg:
    t1 = tg.create_task(fetch(url1))
    t2 = tg.create_task(fetch(url2))
# 任一失败则两个都取消

```

**控制并发量用 `asyncio.Semaphore`**：批量请求时不设上限会耗尽连接池或触发限流，用 `Semaphore` 限制同时进行的协程数：

```python
sem = asyncio.Semaphore(10)
async def limited_fetch(url):
    async with sem:
        return await client.get(url)

```

**避免在协程中混用 `asyncio.sleep(0)` 作为让步**：`await asyncio.sleep(0)` 让出控制权，但过于频繁会降低吞吐量，真正需要让步的场景是长计算循环，而非 IO 等待（IO 本身已让步）。

**用 `asyncio.wait_for` 设置单个协程超时**：比在客户端设置全局超时更精确：

```python
try:
    result = await asyncio.wait_for(coro(), timeout=5.0)
except asyncio.TimeoutError:
    handle_timeout()

```

**`loop.run_in_executor` 桥接同步阻塞代码**：无异步版本的库（`boto3`、同步 DB 驱动）通过线程池运行，不阻塞事件循环：

```python
result = await loop.run_in_executor(None, sync_blocking_call, arg1)

```

---

## 常见陷阱

### 陷阱：`asyncio.create_task` 的任务被 GC 回收

**现象：** 创建的后台任务有时无征兆地停止，日志中可能出现 `Task was destroyed but it is pending!`。  
**原因：** `create_task` 返回的 `Task` 对象若没有被引用，可能被垃圾回收器回收，导致任务取消。  
**解决：** 将任务保存到集合中，任务完成后移除：

```python
background_tasks = set()
task = asyncio.create_task(coro())
background_tasks.add(task)
task.add_done_callback(background_tasks.discard)

```

### 陷阱：在同步上下文中调用 `asyncio.run` 嵌套事件循环

**现象：** 在已有事件循环（如 Jupyter、FastAPI 路由）中调用 `asyncio.run(coro())` 报 `This event loop is already running`。  
**原因：** `asyncio.run` 创建新事件循环，但当前线程已有运行中的循环，不允许嵌套。  
**解决：** 在已有循环的上下文中用 `await coro()` 直接调用；Jupyter 中可用 `nest_asyncio.apply()`。

### 陷阱：`shield` 不能阻止外部取消

**现象：** 用 `asyncio.shield(coro())` 期望保护协程不被取消，但协程仍然被取消了。  
**原因：** `shield` 只保护内部协程不因外层 Future 取消而停止，但若持有 shield Future 的任务本身被取消，shield 也会被取消。  
**解决：** 真正不可中断的清理操作用独立 Task + `shield`，并用 `try/finally` 确保执行完毕。

---

## 参见

- [装饰器与函数高级](https://blog.vercanti.com/python-zhuang-shi-qi-yu-han-shu-gao-ji-yong-fa/)
- [FastAPI完全指南](https://blog.vercanti.com/fastapi-wan-quan-zhi-nan/)
- [httpx完全指南](https://blog.vercanti.com/httpx-wan-quan-zhi-nan/)
- [并发编程完全指南](https://blog.vercanti.com/bing-fa-bian-cheng-wan-quan-zhi-nan/)