【python】协程 (Coroutine) 详解
·
协程 (Coroutine) 详解
文章目录
协程是 Python 异步编程的核心概念,它是一种比线程更轻量级的并发执行单元。
一.什么是协程?
基本概念
协程是可以暂停执行并在之后恢复的函数,它允许在单个线程中实现并发。与线程不同,协程的调度由程序控制,而不是操作系统。
关键特性:
- 可以暂停和恢复执行
- 保持局部状态
- 通过协作而不是抢占来实现多任务
- 比线程更轻量级(一个线程可以运行成千上万个协程)
二.如何定义协程
1. 使用 async def 定义协程函数
import asyncio
# 定义一个简单的协程
async def simple_coroutine():
print("Hello from coroutine!")
return "Done"
# 协程可以包含 await 表达式
async def coroutine_with_await():
print("开始执行")
await asyncio.sleep(1) # 模拟异步操作
print("1秒后继续执行")
return "完成"
2. 协程的层次结构
async def level1():
print("Level 1 开始")
result = await level2() # 等待另一个协程
print(f"Level 1 收到: {result}")
return "Level 1 完成"
async def level2():
print("Level 2 开始")
await asyncio.sleep(0.5)
print("Level 2 结束")
return "Level 2 的结果"
async def level3():
print("Level 3 开始")
await asyncio.sleep(0.2)
print("Level 3 结束")
return "Level 3 的结果"
三.如何运行协程
1. 使用 asyncio.run()(推荐方式)
import asyncio
async def main():
print("开始主协程")
result = await simple_coroutine()
print(f"结果: {result}")
# 运行协程的入口点
asyncio.run(main())
2. 使用事件循环(传统方式)
import asyncio
async def my_coroutine():
print("协程执行中")
await asyncio.sleep(1)
return "结果"
# 传统方式 - 手动管理事件循环
loop = asyncio.get_event_loop()
try:
result = loop.run_until_complete(my_coroutine())
print(f"得到结果: {result}")
finally:
loop.close()
四.协程的并发执行
1. 使用 asyncio.gather() 并发运行多个协程
import asyncio
import time
async def task(name, duration):
print(f"{time.strftime('%X')}: {name} 开始")
await asyncio.sleep(duration)
print(f"{time.strftime('%X')}: {name} 完成")
return f"{name} 的结果"
async def main():
# 并发执行多个协程
results = await asyncio.gather(
task("任务A", 2),
task("任务B", 1),
task("任务C", 3)
)
print(f"所有任务完成: {results}")
asyncio.run(main())
输出:
14:30:00: 任务A 开始
14:30:00: 任务B 开始
14:30:00: 任务C 开始
14:30:01: 任务B 完成
14:30:02: 任务A 完成
14:30:03: 任务C 完成
所有任务完成: ['任务A 的结果', '任务B 的结果', '任务C 的结果']
2. 使用 asyncio.create_task() 创建任务
import asyncio
async def background_task(name, seconds):
print(f"{name} 开始运行")
for i in range(seconds):
await asyncio.sleep(1)
print(f"{name} 运行了 {i+1} 秒")
return f"{name} 完成"
async def main():
print("主程序开始")
# 创建任务(立即开始执行,但不等待)
task1 = asyncio.create_task(background_task("后台任务1", 3))
task2 = asyncio.create_task(background_task("后台任务2", 2))
# 主程序可以继续做其他事情
print("任务已启动,主程序继续执行")
await asyncio.sleep(0.5)
print("主程序做了一些其他工作")
# 等待任务完成
result1 = await task1
result2 = await task2
print(f"任务结果: {result1}, {result2}")
asyncio.run(main())
五.协程的高级用法
1. 协程与生成器的关系
import asyncio
import types
# 协程本质上是一种特殊的生成器
async def async_coroutine():
print("步骤 1")
await asyncio.sleep(1)
print("步骤 2")
return "完成"
# 检查协程类型
coro = async_coroutine()
print(f"类型: {type(coro)}") # <class 'coroutine'>
print(f"是协程: {asyncio.iscoroutine(coro)}") # True
print(f"是协程函数: {asyncio.iscoroutinefunction(async_coroutine)}") # True
2. 协程的状态管理
import asyncio
import inspect
async def stateful_coroutine():
print("协程开始")
print(f"当前状态: {inspect.getcoroutinestate(stateful_coroutine)}")
await asyncio.sleep(1)
print("第一阶段完成")
await asyncio.sleep(1)
print("第二阶段完成")
return "最终结果"
async def monitor_coroutine():
coro = stateful_coroutine()
# 协程的不同状态
print(f"初始状态: {coro.cr_running}") # False - 未运行
# 开始执行
task = asyncio.create_task(coro)
await asyncio.sleep(0.1)
print(f"运行中: {not task.done()}") # True - 运行中
# 等待完成
result = await task
print(f"已完成: {task.done()}") # True
print(f"结果: {result}")
asyncio.run(monitor_coroutine())
3. 协程的异常处理
import asyncio
async def risky_coroutine():
print("开始有风险的操作")
await asyncio.sleep(1)
raise ValueError("出错了!")
return "正常结果"
async def safe_coroutine():
print("开始安全操作")
await asyncio.sleep(0.5)
return "安全结果"
async def main():
try:
# 单个协程的异常处理
result = await risky_coroutine()
except ValueError as e:
print(f"捕获到异常: {e}")
# 多个协程的异常处理
try:
results = await asyncio.gather(
safe_coroutine(),
risky_coroutine(),
return_exceptions=True # 将异常作为结果返回而不是抛出
)
print(f"所有结果: {results}")
except Exception as e:
print(f" gather 异常: {e}")
asyncio.run(main())
六.协程的实际应用模式
1. 生产者-消费者模式
import asyncio
import random
async def producer(queue, name, count):
"""生产者协程"""
for i in range(count):
item = f"{name}-产品{i}"
await asyncio.sleep(random.uniform(0.1, 0.3))
await queue.put(item)
print(f"📦 生产: {item}")
await queue.put(None) # 结束信号
async def consumer(queue, name):
"""消费者协程"""
while True:
item = await queue.get()
if item is None:
queue.put(None) # 传递给其他消费者
break
print(f"🛒 {name} 消费: {item}")
await asyncio.sleep(random.uniform(0.2, 0.4))
queue.task_done()
async def main():
queue = asyncio.Queue(maxsize=5)
# 创建生产者和消费者任务
producers = [
asyncio.create_task(producer(queue, "工厂A", 3)),
asyncio.create_task(producer(queue, "工厂B", 2))
]
consumers = [
asyncio.create_task(consumer(queue, "消费者1")),
asyncio.create_task(consumer(queue, "消费者2"))
]
# 等待所有生产者完成
await asyncio.gather(*producers)
# 等待队列清空
await queue.join()
# 取消消费者
for c in consumers:
c.cancel()
asyncio.run(main())
2. 协程池模式
import asyncio
from asyncio import Semaphore
async def worker(semaphore, name, task_id):
"""工作协程"""
async with semaphore: # 限制并发数量
print(f"👷 {name} 开始任务 {task_id}")
await asyncio.sleep(1) # 模拟工作
print(f"✅ {name} 完成任务 {task_id}")
return f"任务{task_id}结果"
async def limited_concurrency():
"""限制并发数量的协程池"""
semaphore = Semaphore(3) # 最多同时运行3个协程
tasks = []
for i in range(10):
task = asyncio.create_task(
worker(semaphore, f"Worker{i%3}", i)
)
tasks.append(task)
results = await asyncio.gather(*tasks)
print(f"所有任务完成: {len(results)} 个结果")
asyncio.run(limited_concurrency())
协程与函数的区别
| 特性 | 普通函数 | 协程 |
|---|---|---|
| 定义方式 | def function(): | async def coroutine(): |
| 调用方式 | function() | await coroutine() |
| 执行 | 立即执行到结束 | 可以暂停和恢复 |
| 返回值 | return value | return value |
| 暂停机制 | 无 | await expression |
| 并发能力 | 无 | 可以在单个线程中并发执行多个协程 |
总结
协程的核心要点:
- 定义:使用
async def定义协程函数 - 调用:使用
await调用其他协程 - 运行:通过
asyncio.run()或事件循环运行 - 并发:使用
asyncio.gather()或asyncio.create_task()实现并发 - 优势:轻量级、高效率、适合I/O密集型任务
适用场景:
- 网络请求
- 文件I/O操作
- 数据库查询
- Web服务器
- 实时数据处理
不适用场景:
- CPU密集型计算
- 需要真正并行执行的场景
协程通过协作式多任务提供了比线程更高效的并发解决方案,是现代Python异步编程的基石。
更多推荐


所有评论(0)