整体说明
- 协程的本质是函数 (具体来说是定义时被
async修饰的函数) asyncio是 Python 内置的 异步 I/O 框架 ,核心用于编写高效的并发代码,尤其适合处理 I/O 密集型任务(如网络请求、文件读写、数据库操作等)asyncio基于 协程(coroutine) 实现,通过事件循环(Event Loop)调度任务,避免了多线程的上下文切换开销,效率更高asyncio在较高的 Python 版本中才可以使用,建议在 Python 3.7 以上使用- 常见应用场景包括
- 网络请求:并发调用 API(如结合
aiohttp库) - 文件操作:异步读写文件(如结合
aiofiles库) - 数据库操作:异步操作数据库(如
asyncpg用于 PostgreSQL,motor用于 MongoDB) - WebSocket 服务:实现高并发的实时通信(如
websockets库) - 定时任务:通过
asyncio.create_task()+ 循环实现简单定时任务
- 网络请求:并发调用 API(如结合
- 仅适用于 I/O 密集型任务:
asyncio是单线程的,CPU 密集型任务会阻塞事件循环,需结合loop.run_in_executor()提交到线程池/进程池- 注:
run_in_executor(executor, func, *args)将同步函数func放到一个线程池(或进程池)中执行,返回一个 可等待的 Future- 如果第一个参数是 None,表示使用默认的线程池(ThreadPoolExecutor)
- 注:
await只能在协程中使用:await关键字不能在普通函数中使用,必须在async def定义的协程中- 事件循环是单线程的:协程的并发是“协作式”的,需通过
await主动交出执行权,否则会独占事件循环
asyncio 相关核心概念
协程(Coroutine)
- 协程是可暂停、可恢复的函数 ,用
async def定义(语法),是asyncio的核心执行单元 - 协程函数调用后不会立即执行 ,而是返回一个协程对象(coroutine object),需通过事件循环调度才能运行
事件循环(Event Loop)
asyncio的“大脑”,负责调度所有协程任务:- 管理任务的暂停/恢复、监听 I/O 事件、分发任务执行权
- 通常通过
asyncio.run()自动创建和管理事件循环(推荐用法)- 注:这行代码会直到整个时间循环运行完成才最终返回
等待对象(Awaitable)
- 可被
await关键字修饰的对象- 包括:协程对象、Task、Future
await会暂停当前协程,等待目标对象完成后再恢复,期间事件循环可调度其他协程执行(实现并发)
Task(任务)
- 对协程的封装,将协程注册到事件循环中,使其可被调度执行
- 通过
asyncio.create_task()创建,会自动加入事件循环并运行
Future
- 表示异步操作的“未来结果”,是低层级的对象(通常无需手动创建,Task 继承自 Future)
asyncio 基础用法
定义和运行协程
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15import asyncio
# 定义协程函数(async def 关键字)
async def hello(name):
print(f"Hello, {name}! (开始)")
# 模拟 I/O 等待(必须用 await 修饰可等待对象)
# 注:await 后面必须跟一个“可等待对象”(比如还可以是另一个协程函数)
await asyncio.sleep(1) # 这里的含义是暂停 1 秒,期间事件循环可执行其他任务,即让当前的协程暂停执行 1 秒钟,同时把控制权交还给事件循环
print(f"Hello, {name}! (结束)")
# 运行协程(Python 3.7+ 推荐用 asyncio.run())
asyncio.run(hello("asyncio")) # asyncio.run 启动事件循环并以 hello 协程为入口执行协程(等待执行完成才最终返回)
# Hello, asyncio! (开始)
## 【等待】... 1s
# Hello, asyncio! (结束)- 补充理解:
await asyncio.sleep(1)- 当协程执行到
await asyncio.sleep(1)时,它会向事件循环注册一个延迟任务,然后自身进入挂起(suspended)状态 - 此时,事件循环会立刻接管控制权,去执行其他准备好的协程
- 1 秒钟后,事件循环会重新唤醒这个协程继续往下执行
- 当协程执行到
- 补充理解:
await 对接协程函数
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30import asyncio
# 内部协程:返回一个数字
async def get_number():
await asyncio.sleep(1)
return 42
# 外部协程:使用 return await
async def wrapper_coroutine():
# # 下面 await get_number():显式地等待子协程完成,如果子协程内部抛出异常,异常会在这里被捕获或处理
return await get_number() # await 后接协程函数,等待 get_number() 执行完成,return 负责返回对应执行结果
# # 等价写法1:
# x = await get_number() # 获取协程执行结果
# return x
# # 等价写法2:
# x = get_number() # 直接调用协程时,只会返回协程对象,不会执行内部代码
# x = await x # 等待协程执行完成
# return x
# # 错误写法:
# return get_number() # 直接调用协程时,只会返回协程对象,不会执行内部代码
async def main():
# result 的值就是 42
result = await wrapper_coroutine() # await 后接协程函数
print(result)
asyncio.run(main())特别说明:
- await 只能在 协程中被使用,普通函数中不能使用 await
- 协程调用的三种方式总结:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27import asyncio
async def my_coroutine():
print("Hello from coroutine!")
return 42
# 调用方式一:使用 asyncio.run() 作为程序入口调用
# 下面会 打印 "Hello from coroutine!"
asyncio.run(my_coroutine()) # 启动事件循环并以 my_coroutine 协程为入口 执行协程
# 调用方式二:在另一个协程中使用 await 调用,然后外面的协程被其他调用方式调用
async def main():
# 下面这行会打印 打印 "Hello from coroutine!"
result = await my_coroutine() # 这里会真正执行并等待结果
# # 等价于 下面两行:
# x = my_coroutine() # 调用后返回一个 <coroutine object>,此时函数体内的代码完全没有执行,必须手动使用 await 关键字去等待它,它才会开始运行
# result = await x # 这里会真正执行并等待结果
print(result) # 输出: 42
asyncio.run(main()) # 启动事件循环并以 my_coroutine 协程为入口 执行协程
# 调用方式三:使用 asyncio.create_task() 并发调度
async def main():
task = asyncio.create_task(my_coroutine()) # 提交给事件循环并发执行,注意是立即后台巡行,调用后返回一个 <Task pending> 对象
# print(task) # <Task pending name='Task-8' coro=<my_coroutine() running at /Users/sanye/Workspace/IdeaProjects/Python/llm_demo/temp/temp.py:3>>
result = await task
print(result)
asyncio.run(main()) # 启动事件循环并以 my_coroutine 协程为入口 执行协程
并发执行多个协程
通过
asyncio.gather()或asyncio.create_task()实现并发(多个任务同时执行,总耗时接近最长任务的耗时)方式 1:
asyncio.gather()(批量等待多个协程)1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25import asyncio
async def task1():
await asyncio.sleep(2)
return "Task 1 完成"
async def task2():
await asyncio.sleep(1)
return "Task 2 完成"
async def main():
# 并发执行 task1 和 task2,等待所有完成后返回结果(按传入顺序)
result1, result2 = await asyncio.gather(task1(), task2())
# # 等价与下面的两行
# results = asyncio.gather(task1(), task2())
# result1, result2 = results
print(result1)
print(result2)
# 总耗时大约 2 秒(而非 2+1=3 秒)
asyncio.run(main()) # 启动事件循环并以 main 协程为入口 执行协程
# Task 1 完成
# Task 2 完成方式 2:
asyncio.create_task()(手动创建任务)1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19import asyncio
async def task(name, delay):
await asyncio.sleep(delay)
print(f"Task {name} 完成(延迟 {delay} 秒)")
async def main():
# 创建任务并自动加入事件循环
task1 = asyncio.create_task(task("A", 2)) # 立即开始执行
task2 = asyncio.create_task(task("B", 1)) # 立即开始执行
# 等待任务完成(可单独等待,也可一起等待)
await task1
await task2
asyncio.run(main()) # 启动事件循环并以 main 协程为入口 执行协程
# Task B 完成(延迟 1 秒)
# Task A 完成(延迟 2 秒)
附录:处理异常
- 协程中的异常需通过
try/except捕获,或在gather()中通过return_exceptions=True收集异常1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26import asyncio
async def task1():
await asyncio.sleep(2)
return "Task 1 完成"
async def faulty_task():
await asyncio.sleep(1)
raise ValueError("任务执行失败!") # 模拟抛出异常
async def main():
# 方式1:捕获单个协程的异常
try:
await faulty_task()
except ValueError as e:
print(f"捕获异常:{e}") # 捕获异常:任务执行失败!
# 方式2:批量捕获多个协程的异常(return_exceptions=True)
results = await asyncio.gather(
task1(), # 正常任务,输出 "Task 1 完成"
faulty_task(), # 异常任务,抛出异常 ValueError("任务执行失败!")
return_exceptions=True # 不终止,返回异常对象
)
print(results) # 输出:["Task 1 完成", ValueError("任务执行失败!")]
asyncio.run(main())
协程的进阶特性
超时控制(asyncio.wait_for())
- 限制协程的执行时间,超时则抛出
TimeoutError1
2
3
4
5
6
7
8
9
10
11
12
13import asyncio
async def long_task():
await asyncio.sleep(3) # 模拟耗时 3 秒的任务
async def main():
try:
# 限制任务 2 秒内完成,超时则取消任务并抛出异常
result = await asyncio.wait_for(long_task(), timeout=2)
except asyncio.TimeoutError:
print("任务超时被取消!")
asyncio.run(main())
任务取消(Task.cancel())
- 手动取消正在执行的任务,被取消的任务会抛出
CancelledError1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18import asyncio
async def endless_task():
try:
while True:
print("任务运行中...")
await asyncio.sleep(1)
except asyncio.CancelledError:
print("任务被取消!")
raise # 可选:重新抛出,让调用方知道任务被取消
async def main():
task = asyncio.create_task(endless_task())
await asyncio.sleep(2) # 运行 2 秒后取消
task.cancel()
await task # 必须等待任务处理取消逻辑
asyncio.run(main())
异步上下文管理器(async with)
- 用于异步资源的获取和释放(如异步数据库连接、异步文件),需实现
__aenter__和__aexit__方法1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21import asyncio
class AsyncResource:
async def __aenter__(self):
print("获取异步资源")
await asyncio.sleep(0.5)
return self
async def __aexit__(self, exc_type, exc_val, exc_tb):
print("释放异步资源")
await asyncio.sleep(0.5)
async def main():
async with AsyncResource() as res:
print("使用异步资源")
asyncio.run(main())
# 获取异步资源
# 使用异步资源
# 释放异步资源
异步迭代器(async for)
- 用于迭代异步生成的数据(如异步流、分页接口),需实现
__aiter__和__anext__方法1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26import asyncio
class AsyncIterator:
def __init__(self, limit):
self.limit = limit
self.count = 0
def __aiter__(self):
return self
async def __anext__(self):
if self.count >= self.limit:
raise StopAsyncIteration
self.count += 1
await asyncio.sleep(0.5) # 模拟异步获取数据
return self.count
async def main():
async for num in AsyncIterator(3):
print(f"迭代得到:{num}")
asyncio.run(main())
# 迭代得到:1
# 迭代得到:2
# 迭代得到:3
附录:协程函数可以被继承和重写
- 在 Python 中,协程函数(使用
async def定义的函数)本质上仍然是类的方法- 它们遵循 Python 所有的面向对象规则,因此可以被继承、重写(Override),甚至可以通过 super() 调用父类的协程方法
- 相关示例:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30import asyncio
class BaseAPI:
async def fetch_data(self, url):
print(f"[BaseAPI] 正在从 {url} 获取数据...")
await asyncio.sleep(1) # 模拟网络请求
return {"status": 200, "data": "基础数据"}
class AdvancedAPI(BaseAPI):
# 1. 重写父类的协程方法
async def fetch_data(self, url):
print(f"[AdvancedAPI] 拦截请求: {url}")
# 2. 通过 super() 调用父类的协程方法(注意必须加 await)
result = await super().fetch_data(url)
# 增加额外的处理逻辑
result["enhanced"] = True
return result
async def main():
api = AdvancedAPI()
data = await api.fetch_data("https://example.com")
print("最终结果:", data)
asyncio.run(main())
# [AdvancedAPI] 拦截请求: https://example.com
# [BaseAPI] 正在从 https://example.com 获取数据...
# 最终结果: {'status': 200, 'data': '基础数据', 'enhanced': True}
借助协程能实现 async with 语句
async with是 Python 异步编程中的异步上下文管理器语法 ,基于协程实现- 可以把它理解为同步代码中
with语句的异步版本,专门用于在协程环境中优雅、安全地管理那些需要异步初始化和清理的资源(如异步网络连接、数据库会话、异步文件读写等)
- 可以把它理解为同步代码中
async with的存在完全依赖于协程机制,只能在async def定义的协程函数内部使用
with 对比 async with
- 同步
with:底层依赖对象的__enter__和__exit__方法- 这两个方法执行时是同步的,会阻塞当前线程
- 异步
async with:底层依赖对象的__aenter__和__aexit__方法- 这两个方法必须使用
async def定义 (即它们本身就是协程函数) - 当执行
async with时,解释器会自动await这两个方法,从而在等待资源获取或释放时挂起当前协程,将控制权交还给事件循环,实现非阻塞并发
- 这两个方法必须使用
async with 核心工作机制
- 当写下
async with obj as x:时,Python 解释器在底层会做以下事情:- 1)调用
obj.__aenter__()并await它,获取资源对象并赋值给x - 2)执行
with代码块内部的逻辑 - 3)无论代码块是正常结束还是抛出异常,都会自动调用
obj.__aexit__(exc_type, exc_val, exc_tb)并await它,完成资源的清理
- 1)调用
使用示例:实际项目中的常见应用
- 在实际开发中,不需要每次都手写
__aenter__,而是直接使用第三方异步库提供的对象:1
2
3
4
5
6
7
8import aiohttp
async def fetch_data(url):
# aiohttp.ClientSession 是一个异步上下文管理器
async with aiohttp.ClientSession() as session: # 这个类非常常用
# session.get() 返回的响应对象也是一个异步上下文管理器
async with session.get(url) as response:
return await response.text()
自定义异步上下文管理器示例
- 自定义异步上下文管理器示例代码
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19import asyncio
class AsyncDatabaseSession:
async def __aenter__(self):
print("正在异步建立数据库连接...")
await asyncio.sleep(1) # 模拟异步网络IO,不阻塞事件循环
return self
async def __aexit__(self, exc_type, exc_val, exc_tb):
print("正在异步关闭数据库连接...")
await asyncio.sleep(0.5) # 模拟异步清理操作
async def main():
# 只能在协程中使用 async with
async with AsyncDatabaseSession() as session:
print("正在安全地使用数据库会话...")
await asyncio.sleep(1)
asyncio.run(main())
asyncio.Semaphore() 的介绍和使用
asyncio.Semaphore是 Python 异步编程中用于控制并发任务数量的核心同步原语- 理解:就像一个”限流器”,能确保同时执行的协程数量不超过预设的上限(常用于 控制网络请求并发数等)
标准用法
在异步编程中,
asyncio.Semaphore的标准用法是配合async with管理上下文:1
2
3
4
5
6
7
8sem = asyncio.Semaphore(max_conccurency)
async def worker():
# 统计耗时位置1:统计排队耗时 + 执行耗时
async with sem: # 获取令牌(若满则阻塞等待)
# 执行受限的 I/O 密集型操作(如网络请求、数据库读写)
# # 统计耗时位置2:仅统计到执行耗时
await do_io()
# 离开上下文自动释放令牌- 最佳实践:将需要消耗物理资源(网络连接、文件句柄)的代码 放在
async with内部,将本地计算/日志统计 放在外部 - 若要
- 最佳实践:将需要消耗物理资源(网络连接、文件句柄)的代码 放在
信号量(Semaphore)与锁(Lock)的区别
- Semaphore 用于 限流(控制并发任务数)
- 防止同时发起过多 请求,打爆 API 或数据库连接池
- Lock 用于 互斥(保护共享数据)
- 防止多个协程同时修改指定数据
高阶通用模式示例
- 通用模板:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22# 1. 初始化限流器与共享锁
sem = asyncio.Semaphore(concurrency)
lock = asyncio.Lock()
results = []
async def worker(data):
# 2. 外部记录开始(用于计算排队延迟)
start = time.time()
# 3. 获取信号量,执行受限操作
async with sem:
result = await expensive_io_task(data) # 仅这里受并发限制
# 4. 退出信号量立即计算总耗时(含排队)
result['total_latency'] = time.time() - start
# 5. 使用 Lock 保护共享状态写入(不占用信号量)
async with lock:
results.append(result)
update_summary(result)
await write_to_file(result) # 异步写文件(或用同步写但持有锁)
pbar.update(1) # 更新进度条
附录:使用 asyncio.create_subprocess_exec 创建子进程
官方说明文档:3.14.6 Documentation » Python 标准库 » 网络和进程间通信 » asyncio — 异步 I/O » 子进程
asyncio.create_subprocess_exec是 Python 用于异步创建和管理子进程的高层 API- 可以在不阻塞事件循环的前提下执行外部程序,并与其输入/输出流进行交互
asyncio.create_subprocess_exec函数签名与参数1
2async asyncio.create_subprocess_exec(program, *args, stdin=None, stdout=None,
stderr=None, limit=None, **kwds)program(必需): 要执行的程序路径(字符串)*args:程序的命令行参数 ,作为可变参数传入,可以一次输出多个(可按照顺序随意列出)stdin,stdout,stderr:标准输入、输出、错误流的重定向方式,常用取值有:asyncio.subprocess.PIPE: 创建一个管道,此时通过Process对象的stdin/stdout/stderr进行读写asyncio.subprocess.DEVNULL: 重定向到os.devnull- 也可以是已有的文件描述符或文件对象
limit:为stdout和stderr的StreamReader包装器设置缓冲区大小上限(字节数)**kwds:其他关键字参数(如env,cwd等),具体可参考loop.subprocess_exec()文档- 该协程返回一个
asyncio.subprocess.Process实例
返回值说明
asyncio.subprocess.Process实例可以进行如下操作:- 等待进程结束 :使用
await proc.wait() - 获取输出 :通过
proc.stdout.read()等方法读取 - 发送输入 :使用
proc.stdin.write() - 终止进程 :调用
proc.terminate()或proc.kill() - 获取进程的
returncode、pid等信息
- 等待进程结束 :使用
使用注意
- 这是一个协程,必须使用
await调用 - 如果
Process对象在子进程仍在运行时被垃圾回收 ,子进程会被杀掉 create_subprocess_exec直接执行程序 ,不经过 shell 解析,因此没有 shell 注入风险 ,安全性更高- 可以配合
asyncio.gather等工具并行执行和管理多个子进程
asyncio.create_subprocess_exec 使用示例
- 异步执行
ls -l -a命令并读取其输出:1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22import asyncio
import asyncio.subprocess
async def run_ls():
# 创建子进程,捕获 stdout 和 stderr
proc = await asyncio.create_subprocess_exec(
'ls', '-l', '-a',
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE
)
# 等待进程结束并获取输出
stdout, stderr = await proc.communicate()
print(f"进程退出码: {proc.returncode}")
if stdout:
print(f"[标准输出]\n{stdout.decode()}")
if stderr:
print(f"[标准错误]\n{stderr.decode()}")
# 运行异步函数
asyncio.run(run_ls())