浏览知识库目录

Python

并发编程与 asyncio

依据 I/O 与 CPU 特征选择线程、进程或 asyncio,并处理取消、超时和共享状态。

并发编程与 asyncio

并发不是让代码“自动更快”。选择模型前,先判断任务是在等待网络或磁盘,还是持续占用 CPU;再确定共享状态、取消、超时和错误传播方式。


一、学习目标

  • 区分并发、并行、I/O 密集与 CPU 密集
  • 使用 concurrent.futures 管理线程和进程
  • 使用 asyncio.TaskGroup 建立结构化并发
  • 正确处理阻塞调用、取消与异常
  • 为任务同步选择合适模型

二、先测量再并发

适合并发的常见场景:

  • 同时等待多个 HTTP 请求;
  • 读取多个独立文件;
  • 执行互不依赖的计算任务。

不适合盲目并发:

  • 数据量很小;
  • 顺序操作已经足够快;
  • 多个任务竞争同一数据库写锁;
  • 业务步骤有严格先后顺序。

并发会增加调度、错误组合和测试成本。


三、线程池处理阻塞 I/O

from concurrent.futures import ThreadPoolExecutor, as_completed


def sync_all(tasks: list[Task]) -> list[SyncResult]:
    results: list[SyncResult] = []
    with ThreadPoolExecutor(max_workers=8) as executor:
        futures = {
            executor.submit(push_task, API_URL, task): task
            for task in tasks
        }
        for future in as_completed(futures):
            task = futures[future]
            try:
                results.append(future.result())
            except RuntimeError as exc:
                print(f"#{task.id} 同步失败:{exc}")
    return results

线程适合现有阻塞 I/O API。限制工作线程数量,避免给远端服务、文件描述符或数据库造成压力。


四、进程池处理 CPU 密集任务

from concurrent.futures import ProcessPoolExecutor


def checksum(path: str) -> str:
    ...


with ProcessPoolExecutor() as executor:
    checksums = list(executor.map(checksum, paths))

进程可以真正并行执行 CPU 密集的 Python 工作,但参数和返回值必须可序列化,启动与进程通信也有成本。

入口代码需要保护:

if __name__ == "__main__":
    main()

不同平台的进程启动方式不同,不要依赖只在某个系统偶然成立的全局状态。


五、asyncio 基础

import asyncio


async def main() -> None:
    await asyncio.sleep(0.1)
    print("完成")


asyncio.run(main())

async def 调用返回协程对象;只有被 await、创建任务或交给事件循环后才执行。

协程应在等待 I/O 时主动让出控制权。长时间 CPU 循环会阻塞整个事件循环。


六、TaskGroup 结构化并发

import asyncio


async def push_tasks(
    url: str,
    tasks: list[Task],
    token: str | None = None,
) -> list[SyncResult]:
    async with asyncio.TaskGroup() as group:
        jobs = [
            group.create_task(
                asyncio.to_thread(push_task, url, task, token)
            )
            for task in tasks
        ]
    return [job.result() for job in jobs]

TaskGroup 离开上下文前会等待所有子任务。如果一个任务失败,其余任务会被取消,异常以异常组形式传播。这让子任务生命周期不会偷偷超出父操作。

asyncio.to_thread 用线程运行阻塞函数,适合逐步接入现有同步库;它不会把 CPU 密集代码变成高效并行。


七、超时和取消

async def sync_with_timeout(tasks: list[Task]) -> list[SyncResult]:
    async with asyncio.timeout(30):
        return await push_tasks(API_URL, tasks)

取消是协程的正常控制流程。清理资源后应继续传播 CancelledError,不要把它当作普通业务错误吞掉。

async def worker() -> None:
    try:
        await do_work()
    finally:
        await close_resources()

八、共享状态

最安全的共享状态是没有共享。让并发任务返回结果,由父任务统一合并:

results = await push_tasks(url, tasks)
accepted_ids = {
    result["task_id"]
    for result in results
    if result["accepted"]
}

若必须共享,应使用相应模型的锁或队列,并把临界区保持很短。线程锁不能直接替代 asyncio 锁。


九、Python 3.14 的并发背景

Python 3.14 中,自由线程构建已成为受支持但仍可选的构建模式;标准发行仍可能使用传统 GIL 构建。3.14 也在标准库提供多解释器能力。

应用不应仅因新能力出现就移除锁或假设第三方扩展全部兼容。确认所用解释器构建、依赖兼容性和性能测试结果后再采用。


十、常见错误

  • 在协程中直接调用长时间阻塞 I/O。
  • 创建任务后不保存引用、不等待结果。
  • 无限并发请求远端服务。
  • 多线程无锁修改共享列表或连接。
  • 认为增加线程一定提升 CPU 密集任务速度。

十一、练习与自测

  1. 用本地 HTTP 服务模拟不同延迟,并限制并发数。
  2. 让一个 TaskGroup 子任务失败,观察其他任务取消和异常组。
  3. 分别用顺序、线程池和 asyncio 测量同一组阻塞请求。

自测:

  • I/O 密集与 CPU 密集任务分别适合哪些工具?
  • TaskGroup 如何约束子任务生命周期?
  • to_thread 解决什么问题,又不解决什么问题?

十二、官方资料

上一篇:HTTP、JSON 与 API 客户端 | 下一篇:测试、调试与代码质量