异步编程 (asyncio) 工作原理详解

协程 (Coroutine) 详解


协程是 Python 异步编程的核心概念,它是一种比线程更轻量级的并发执行单元。

一.什么是协程?

基本概念

协程是可以暂停执行并在之后恢复的函数,它允许在单个线程中实现并发。与线程不同,协程的调度由程序控制,而不是操作系统。

关键特性:

  • 可以暂停和恢复执行
  • 保持局部状态
  • 通过协作而不是抢占来实现多任务
  • 比线程更轻量级(一个线程可以运行成千上万个协程)
就绪
已完成
事件循环开始
从任务队列取一个任务
任务状态
执行任务
从队列中移除
执行到await表达式
挂起当前任务
并将控制权交还给事件循环
事件循环继续下一个任务
还有任务?
事件循环结束

二.如何定义协程

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 valuereturn value
暂停机制await expression
并发能力可以在单个线程中并发执行多个协程

总结

协程的核心要点:

  1. 定义:使用 async def 定义协程函数
  2. 调用:使用 await 调用其他协程
  3. 运行:通过 asyncio.run() 或事件循环运行
  4. 并发:使用 asyncio.gather()asyncio.create_task() 实现并发
  5. 优势:轻量级、高效率、适合I/O密集型任务

适用场景:

  • 网络请求
  • 文件I/O操作
  • 数据库查询
  • Web服务器
  • 实时数据处理

不适用场景:

  • CPU密集型计算
  • 需要真正并行执行的场景

协程通过协作式多任务提供了比线程更高效的并发解决方案,是现代Python异步编程的基石。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐