Python
并发编程与 asyncio
依据 I/O 与 CPU 特征选择线程、进程或 asyncio,并处理取消、超时和共享状态。
发布于 2026年7月23日
并发编程与 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 密集任务速度。
十一、练习与自测
- 用本地 HTTP 服务模拟不同延迟,并限制并发数。
- 让一个 TaskGroup 子任务失败,观察其他任务取消和异常组。
- 分别用顺序、线程池和 asyncio 测量同一组阻塞请求。
自测:
- I/O 密集与 CPU 密集任务分别适合哪些工具?
- TaskGroup 如何约束子任务生命周期?
to_thread解决什么问题,又不解决什么问题?
十二、官方资料
上一篇:HTTP、JSON 与 API 客户端 | 下一篇:测试、调试与代码质量