返回文章列表

Python异步编程:asyncio实战指南

异步编程是现代Python开发中不可或缺的技能。当程序需要处理大量I/O操作(如网络请求、文件读写、数据库查询)时,异步编程能显著提升性能。本文将从同步与异步的基本概念讲起,系统讲解Python asyncio库的核心概念与实战应用。

一、同步与异步的概念

同步(Synchronous)代码按顺序逐行执行,一行代码执行完毕后才执行下一行。当遇到I/O操作(如网络请求)时,程序会阻塞等待,期间CPU空闲。异步(Asynchronous)代码在遇到I/O操作时不会阻塞,而是切换去执行其他任务,等I/O操作完成后再回来继续处理。

# 同步方式:串行执行,总时间 = 任务1 + 任务2 + 任务3
import time

def sync_fetch(name, delay):
    """模拟同步请求"""
    print(f"[同步] 开始获取 {name}...")
    time.sleep(delay)  # 阻塞等待
    print(f"[同步] {name} 获取完成(耗时{delay}秒)")
    return f"{name} 的数据"

start = time.time()
sync_fetch("页面A", 2)
sync_fetch("页面B", 1)
sync_fetch("页面C", 3)
print(f"[同步] 总耗时:{time.time() - start:.1f}秒")
# 总耗时:约 6 秒(2+1+3)

print("\n" + "="*50 + "\n")

# 异步方式:并发执行,总时间约等于最长的那个任务
import asyncio

async def async_fetch(name, delay):
    """模拟异步请求"""
    print(f"[异步] 开始获取 {name}...")
    await asyncio.sleep(delay)  # 非阻塞等待
    print(f"[异步] {name} 获取完成(耗时{delay}秒)")
    return f"{name} 的数据"

async def main():
    start = time.time()
    tasks = [
        async_fetch("页面A", 2),
        async_fetch("页面B", 1),
        async_fetch("页面C", 3)
    ]
    await asyncio.gather(*tasks)
    print(f"[异步] 总耗时:{time.time() - start:.1f}秒")
    # 总耗时:约 3 秒(最长的任务)

asyncio.run(main())

二、async/await语法基础

async def 定义一个协程函数,await 等待一个协程或可等待对象完成。协程是异步编程的基本单位。

import asyncio

# 定义协程函数
async def say_hello():
    """一个简单的协程"""
    print("Hello")
    await asyncio.sleep(1)  # 模拟异步操作
    print("World")

async def count_up(n):
    """从1数到n"""
    for i in range(1, n + 1):
        print(f"计数:{i}")
        await asyncio.sleep(0.5)

# 运行协程的三种方式
async def main():
    # 方式1:直接await
    await say_hello()

    # 方式2:创建任务并发执行
    task = asyncio.create_task(count_up(3))
    await task

    # 方式3:用asyncio.gather并发执行多个协程
    await asyncio.gather(
        say_hello(),
        count_up(2)
    )

asyncio.run(main())

三、事件循环(Event Loop)

事件循环是asyncio的核心调度器,负责执行协程、管理I/O事件和定时任务。asyncio.run() 会自动创建并管理事件循环。

import asyncio

async def worker(name, work_time):
    """模拟工作协程"""
    print(f"[{name}] 开始工作")
    await asyncio.sleep(work_time)
    print(f"[{name}] 工作完成(耗时{work_time}秒)")
    return f"{name} 完成了 {work_time} 秒的工作"

async def main():
    loop = asyncio.get_running_loop()

    # 查看当前时间
    print(f"当前时间:{loop.time():.2f}")

    # 创建多个任务
    tasks = [
        loop.create_task(worker("工人A", 2)),
        loop.create_task(worker("工人B", 1)),
        loop.create_task(worker("工人C", 3))
    ]

    # 等待所有任务完成
    results = await asyncio.gather(*tasks, return_exceptions=True)

    for i, result in enumerate(results):
        print(f"任务{i+1}结果:{result}")

# 运行
asyncio.run(main())

四、asyncio.gather并发执行

asyncio.gather 是最常用的并发执行工具,它可以同时启动多个协程并等待它们全部完成。

import asyncio
import random

async def download_image(url, name):
    """模拟下载图片"""
    delay = random.uniform(0.5, 2.0)
    print(f"[下载] 开始下载 {name} ({url})...")
    await asyncio.sleep(delay)
    print(f"[下载] {name} 下载完成({delay:.2f}秒)")
    return {"name": name, "url": url, "time": delay}

async def main():
    urls = [
        ("https://example.com/img1.jpg", "图片1"),
        ("https://example.com/img2.jpg", "图片2"),
        ("https://example.com/img3.jpg", "图片3"),
        ("https://example.com/img4.jpg", "图片4"),
        ("https://example.com/img5.jpg", "图片5"),
    ]

    print("=== 开始批量下载 ===")
    start = asyncio.get_event_loop().time()

    # gather: 并发执行所有下载任务
    tasks = [download_image(url, name) for url, name in urls]
    results = await asyncio.gather(*tasks)

    elapsed = asyncio.get_event_loop().time() - start

    print("\n=== 下载结果 ===")
    total = 0
    for r in results:
        total += r["time"]
        print(f"  {r['name']}: {r['time']:.2f}秒")
    print(f"\n并发总耗时:{elapsed:.2f}秒")
    print(f"串行总耗时:{total:.2f}秒")
    print(f"节省时间:{total - elapsed:.2f}秒({(1-elapsed/total)*100:.0f}%)")

asyncio.run(main())

五、asyncio.wait与超时控制

asyncio.wait 提供更灵活的等待策略,可以设置超时时间,也可以按完成顺序处理结果。

import asyncio

async def delayed_task(name, delay, success=True):
    """模拟可能失败的任务"""
    print(f"[{name}] 开始(预计{delay}秒)")
    try:
        await asyncio.sleep(delay)
        if not success:
            raise ValueError(f"{name} 模拟失败")
        print(f"[{name}] 成功完成")
        return f"{name} 的结果"
    except Exception as e:
        print(f"[{name}] 失败:{e}")
        raise

async def main():
    # asyncio.wait 的使用
    tasks = {
        asyncio.create_task(delayed_task("任务A", 1.0, True)),
        asyncio.create_task(delayed_task("任务B", 2.0, True)),
        asyncio.create_task(delayed_task("任务C", 0.5, False)),
    }

    # 策略1:等待全部完成(不管成功或失败)
    print("=== 等待全部完成 ===")
    done, pending = await asyncio.wait(
        tasks, return_when=asyncio.ALL_COMPLETED
    )
    for t in done:
        if t.exception():
            print(f"  失败: {t.exception()}")
        else:
            print(f"  成功: {t.result()}")

    print()

    # 策略2:等待第一个完成
    tasks2 = {
        asyncio.create_task(delayed_task("快速任务", 0.3, True)),
        asyncio.create_task(delayed_task("慢速任务", 3.0, True)),
    }
    print("=== 等待第一个完成 ===")
    done, pending = await asyncio.wait(
        tasks2, return_when=asyncio.FIRST_COMPLETED
    )
    for t in done:
        print(f"  第一个完成: {t.result()}")
    for t in pending:
        t.cancel()  # 取消剩余任务
    print("  已取消剩余任务")

    print()

    # asyncio.wait_for:超时控制
    print("=== 超时控制 ===")
    try:
        result = await asyncio.wait_for(
            delayed_task("限时任务", 5.0, True),
            timeout=2.0
        )
        print(f"  结果: {result}")
    except asyncio.TimeoutError:
        print("  任务超时了!")

asyncio.run(main())

六、异步上下文管理器

异步上下文管理器使用 async with 语法,常用于管理需要异步获取和释放的资源(如网络连接)。

import asyncio

class AsyncTimer:
    """异步计时上下文管理器"""

    def __init__(self, name):
        self.name = name
        self.elapsed = 0

    async def __aenter__(self):
        print(f"[计时] {self.name} 开始...")
        self._start = asyncio.get_event_loop().time()
        return self

    async def __aexit__(self, exc_type, exc_val, exc_tb):
        self.elapsed = asyncio.get_event_loop().time() - self._start
        print(f"[计时] {self.name} 耗时:{self.elapsed:.2f}秒")
        return False  # 不吞掉异常

async def process_data(name, delay):
    """模拟数据处理"""
    async with AsyncTimer(f"处理 {name}"):
        await asyncio.sleep(delay)
        return f"{name} 处理完成"

# 自定义异步资源管理器
class AsyncConnection:
    """模拟异步数据库连接"""

    def __init__(self, host, port):
        self.host = host
        self.port = port
        self.connected = False

    async def connect(self):
        """异步连接"""
        await asyncio.sleep(0.1)  # 模拟连接耗时
        self.connected = True
        print(f"[连接] 已连接到 {self.host}:{self.port}")

    async def close(self):
        """异步关闭"""
        await asyncio.sleep(0.05)
        self.connected = False
        print(f"[连接] 已断开 {self.host}:{self.port}")

    async def __aenter__(self):
        await self.connect()
        return self

    async def __aexit__(self, exc_type, exc_val, exc_tb):
        await self.close()
        return False

    async def query(self, sql):
        """异步查询"""
        if not self.connected:
            raise RuntimeError("未连接")
        await asyncio.sleep(0.1)  # 模拟查询耗时
        return f"查询结果:{sql}"

async def main():
    await process_data("用户数据", 1.5)

    print()

    # 使用异步连接
    async with AsyncConnection("localhost", 3306) as conn:
        result = await conn.query("SELECT * FROM users")
        print(result)

asyncio.run(main())

七、实战案例:异步批量网页抓取

以下是一个使用aiohttp实现的异步批量网页抓取示例,展示如何在实际项目中运用asyncio。

import asyncio
import time

# 注意:运行此代码需要安装 aiohttp
# pip install aiohttp

async def fetch_page(session, url, name):
    """异步获取单个网页"""
    import aiohttp
    start = time.time()
    try:
        async with session.get(url, timeout=aiohttp.ClientTimeout(total=10)) as resp:
            status = resp.status
            text = await resp.text()
            elapsed = time.time() - start
            print(f"[成功] {name} - {url} (状态:{status}, {len(text)}字符, {elapsed:.2f}s)")
            return {
                "name": name,
                "url": url,
                "status": status,
                "size": len(text),
                "time": elapsed
            }
    except asyncio.TimeoutError:
        print(f"[超时] {name} - {url}")
        return {"name": name, "url": url, "status": "timeout", "error": "超时"}
    except Exception as e:
        elapsed = time.time() - start
        print(f"[失败] {name} - {url}: {e} ({elapsed:.2f}s)")
        return {"name": name, "url": url, "status": "error", "error": str(e)}

async def batch_fetch(urls, max_concurrent=5):
    """批量异步抓取,控制并发数"""
    import aiohttp

    semaphore = asyncio.Semaphore(max_concurrent)  # 并发限制

    async def limited_fetch(session, url, name):
        async with semaphore:  # 控制并发数
            return await fetch_page(session, url, name)

    connector = aiohttp.TCPConnector(limit=max_concurrent, ssl=False)
    async with aiohttp.ClientSession(connector=connector) as session:
        tasks = [
            limited_fetch(session, url, name)
            for name, url in urls
        ]
        results = await asyncio.gather(*tasks, return_exceptions=True)

    # 统计结果
    success = sum(1 for r in results if isinstance(r, dict) and r.get("status") == 200)
    failed = len(results) - success
    total_time = sum(r.get("time", 0) for r in results if isinstance(r, dict))

    print(f"\n{'='*50}")
    print(f"  总计:{len(results)}个 | 成功:{success} | 失败:{failed}")
    print(f"  累计请求时间:{total_time:.2f}秒(串行)")
    print(f"{'='*50}")

    return results

async def main():
    # 测试URL列表(使用公开的测试网站)
    test_urls = [
        ("HTTPBin-GET", "https://httpbin.org/get"),
        ("HTTPBin-IP", "https://httpbin.org/ip"),
        ("HTTPBin-Headers", "https://httpbin.org/headers"),
        ("HTTPBin-Delay1", "https://httpbin.org/delay/1"),
        ("HTTPBin-Delay2", "https://httpbin.org/delay/2"),
        ("HTTPBin-UUID", "https://httpbin.org/uuid"),
        ("JSONPlaceholder", "https://jsonplaceholder.typicode.com/posts/1"),
        ("JSONPlaceholder-Comments", "https://jsonplaceholder.typicode.com/comments?postId=1"),
    ]

    print("=== 开始异步批量抓取 ===")
    start = time.time()

    results = await batch_fetch(test_urls, max_concurrent=5)

    elapsed = time.time() - start
    print(f"\n实际并发耗时:{elapsed:.2f}秒")

asyncio.run(main())

八、异步队列与生产者-消费者模式

import asyncio
import random

async def producer(queue, name, count):
    """生产者:向队列中放入数据"""
    for i in range(count):
        item = f"{name}-商品{i+1}"
        await asyncio.sleep(random.uniform(0.1, 0.5))  # 模拟生产耗时
        await queue.put(item)
        print(f"[生产者 {name}] 生产了 {item} (队列大小: {queue.qsize()})")

async def consumer(queue, name):
    """消费者:从队列中取出数据处理"""
    while True:
        item = await queue.get()
        print(f"[消费者 {name}] 开始处理 {item}")
        await asyncio.sleep(random.uniform(0.2, 0.8))  # 模拟处理耗时
        print(f"[消费者 {name}] 处理完成 {item}")
        queue.task_done()

async def main():
    queue = asyncio.Queue(maxsize=10)

    # 启动2个生产者和3个消费者
    producers = [
        asyncio.create_task(producer(queue, "工厂A", 5)),
        asyncio.create_task(producer(queue, "工厂B", 5)),
    ]

    consumers = [
        asyncio.create_task(consumer(queue, "工人1")),
        asyncio.create_task(consumer(queue, "工人2")),
        asyncio.create_task(consumer(queue, "工人3")),
    ]

    # 等待所有生产者完成
    await asyncio.gather(*producers)
    print("\n所有生产者已完成,等待消费者处理完剩余商品...")

    # 等待队列中的所有任务完成
    await queue.join()

    # 取消消费者(它们是无限循环)
    for c in consumers:
        c.cancel()

    print("所有任务处理完毕!")

asyncio.run(main())

通过以上内容,你应该对Python异步编程有了系统的理解。asyncio的核心概念包括协程、事件循环、并发执行(gather/wait)、异步上下文管理器和异步队列。在实际开发中,当需要处理大量I/O密集型任务时(如网络爬虫、API调用、数据库查询),asyncio能带来显著的性能提升。

动手挑战

学到这里,不妨动手试一试以下练习,巩固你的理解:

  1. 基础练习:回顾本文核心概念,用自己的话总结关键知识点。
  2. 进阶实践:将文中的示例代码运行一遍,尝试修改参数观察变化。
  3. 拓展思考:想一想这个技术/方法还能应用在哪些场景中?

小贴士:遇到问题时,先独立思考,再查阅资料,最后请教他人——这是成长最快的学习方式。

赞赏支持

本文更新于 2026-08-22,环境 Node.js 20 / 现代浏览器