跳到主要内容
极客日志极客日志面向AI+效率的开发者社区
首页博客我的书AI学习GitHub 精选镜像AI 生图工具UI配色美学关于
搜索内容 / 工具 / 仓库 / 镜像...⌘K搜索
注册
博客列表
Python

Python 异步编程实战:基于 async/await 的高并发实现

Python 异步编程的核心原理,包括协程工作机制与事件循环调度。详细讲解了 asyncio 核心组件(协程、任务、Future、事件循环)及原生 async/await 语法优势。通过高并发 HTTP 请求实战,演示了 aiohttp 使用、信号量控制并发及生产者 - 消费者模式。此外还涵盖异步上下文管理器、迭代器应用、性能优化原则、常见陷阱排查及最佳实践总结,帮助开发者掌握高并发 IO 处理技能。

字节跳动发布于 2026/3/29更新于 2026/9/869 浏览

一、异步编程的核心原理

1.1 什么是异步编程?

传统的同步编程中,代码按照顺序一行行执行,遇到 IO 操作(如网络请求、文件读写)时,程序会阻塞等待操作完成,导致 CPU 空闲浪费。而异步编程的核心思想是:遇到 IO 操作时自动切换,IO 操作完成后自动切回,在单线程内实现高并发。

1.2 协程的工作原理

协程是异步编程的基础单元,它的工作流程如下:

┌─────────────────────────────────────────────────────────┐
│ 事件循环 (Event Loop)
├─────────────────────────────────────────────────────────┤
│
│ ┌──────────┐ ┌──────────┐ ┌──────────┐
│ │ 任务 1 │ ──> │ 任务 2 │ ──> │ 任务 3 │
│ │ (协程 A) │ │ (协程 B) │ │ (协程 C) │
│ └────┬─────┘ └────┬─────┘ └────┬─────┘
│ │ │ │
│ ▼ await ▼ await ▼ await
│ ┌──────────┐ ┌──────────┐ ┌──────────┐
│ │ IO 操作 1 │ │ IO 操作 2 │ │ IO 操作 3 │
│ └──────────┘ └──────────┘ └──────────┘
│ │ │ │
│ └────────────────┼────────────────┘
│ ▼
│ IO 完成,回调通知
│
└─────────────────────────────────────────────────────────┘

关键机制:

  1. 遇到 IO 自动挂起:执行到 await 关键字时,协程主动让出 CPU
  2. 事件循环调度:事件循环负责管理所有协程,当某个协程挂起时,立即切换到下一个就绪协程
  3. IO 完成自动恢复:底层通过 select/epoll 等机制监听 IO 事件,完成后将对应协程放回就绪队列

二、Python 异步编程的演进

2.1 发展阶段对比
版本核心特性代表作缺点
Python 2.x无原生支持gevent(第三方)猴子补丁,魔法太重
Python 3.3yield fromasyncio 雏形生成器协程混淆
Python 3.5async/await原生协程-
Python 3.7+asyncio.run()极致简化-
2.2 原生协程的优势

Python 3.5 引入的 async/await 语法彻底改变了异步编程:

# 生成器协程(3.4 及以前)
@asyncio.coroutine
def hello():
    yield from asyncio.sleep(1)
    print('Hello')

# 原生协程(3.5+)
async def hello():
    await asyncio.sleep(1)
    print('Hello')
# 更清晰,无混淆

三、asyncio 核心组件详解

3.1 四大核心对象
import asyncio

# 1. 协程函数和协程对象
async def coro_func():
    return 42

coro = coro_func()  # 协程对象(此时未执行)

# 2. 任务(Task)- 对协程的封装
task = asyncio.create_task(coro())  # Python 3.7+ 推荐方式
# 或 task = asyncio.ensure_future(coro())

# 3. 事件循环 - 调度器
loop = asyncio.get_event_loop()  # 获取事件循环
loop.run_until_complete(coro())  # 运行直到完成

# Python 3.7+ 简单方式:
asyncio.run(coro())  # 自动创建/关闭事件循环
3.2 awaitable 对象的三种类型
# 1. 协程(Coroutine)
async def foo():
    return 123

# 2. 任务(Task)
async def main():
    task = asyncio.create_task(foo())
    result = await task

# 3. Future(底层对象,通常不直接使用)
fut = asyncio.Future()
asyncio.ensure_future(set_after(fut, 1, 456))
result = await fut

四、实战:高并发 HTTP 请求

4.1 基础用法
import asyncio
import aiohttp
import time

async def fetch_one(session, url):
    """单个 HTTP 请求"""
    async with session.get(url) as response:
        # await 挂起直到响应返回
        return await response.text()

async def main_simple():
    """简单示例:请求单个 URL"""
    async with aiohttp.ClientSession() as session:
        html = await fetch_one(session, 'http://httpbin.org/get')
        print(f"响应长度:{len(html)}")

# Python 3.7+ 运行方式
asyncio.run(main_simple())
4.2 高并发批量请求

这是异步编程发挥威力的核心场景:

import asyncio
import aiohttp
import time
from typing import List, Dict

async def fetch_url(session: aiohttp.ClientSession, url: str) -> Dict:
    """单个 URL 请求,带错误处理"""
    start = time.time()
    try:
        async with session.get(url, timeout=10) as response:
            content = await response.text()
            return {
                'url': url,
                'status': response.status,
                'length': len(content),
                'time': time.time() - start,
                'success': True
            }
    except Exception as e:
        return {
            'url': url,
            'error': str(e),
            'time': time.time() - start,
            'success': False
        }

async def fetch_many(urls: List[str], max_concurrent: int = 10):
    """
    高并发请求多个 URL
    - 使用信号量控制并发数
    - 收集所有结果
    """
    # 控制并发量的信号量
    semaphore = asyncio.Semaphore(max_concurrent)

    async def bounded_fetch(url):
        async with semaphore:
            # 限制并发数
            return await fetch_url(session, url)

    async with aiohttp.ClientSession() as session:
        # 创建所有任务
        tasks = [bounded_fetch(url) for url in urls]
        # 并发执行并收集结果
        results = await asyncio.gather(*tasks, return_exceptions=True)
        return results

async def main():
    """性能对比演示"""
    # 测试 URL 列表(10 个不同请求)
    urls = [
        'http://httpbin.org/delay/1',  # 延迟 1 秒
        'http://httpbin.org/get',
        'http://httpbin.org/json',
        'http://httpbin.org/xml',
        'http://httpbin.org/robots.txt',
        'http://httpbin.org/anything',
        'http://httpbin.org/uuid',
        'http://httpbin.org/image',
        'http://httpbin.org/headers',
        'http://httpbin.org/ip'
    ] * 2  # 20 个请求

    # 1. 异步并发执行
    start = time.time()
    results = await fetch_many(urls, max_concurrent=5)
    async_time = time.time() - start

    # 2. 统计结果
    success_count = sum(1 for r in results if isinstance(r, dict) and r.get('success'))
    print(f"异步并发请求 {len(urls)} 个 URL:")
    print(f" 耗时:{async_time:.2f}秒")
    print(f" 成功:{success_count}/{len(results)}")

    # 3. 如果要对比同步版本(耗时通常是异步的 5-10 倍)
    # 请使用 requests 库顺序执行做对比测试

if __name__ == "__main__":
    asyncio.run(main())
4.3 高级模式:生产者 - 消费者
import asyncio
import aiohttp
from asyncio import Queue

async def producer(queue: Queue, urls: List[str]):
    """生产者:将 URL 放入队列"""
    for url in urls:
        await queue.put(url)
        print(f"生产者放入:{url}")
    # 发送结束信号
    for _ in range(3):  # 消费者数量
        await queue.put(None)

async def consumer(queue: Queue, session: aiohttp.ClientSession, name: str):
    """消费者:从队列取 URL 并请求"""
    while True:
        url = await queue.get()
        if url is None:
            queue.task_done()
            break
        print(f"消费者{name} 处理:{url}")
        try:
            async with session.get(url) as resp:
                text = await resp.text()
                print(f"消费者{name} 完成:{url}, 大小:{len(text)}")
        except Exception as e:
            print(f"消费者{name} 失败:{url}, 错误:{e}")
        finally:
            queue.task_done()

async def producer_consumer_demo():
    """生产者 - 消费者模式示例"""
    urls = ['http://httpbin.org/get'] * 20
    queue = Queue(maxsize=5)  # 缓冲队列

    async with aiohttp.ClientSession() as session:
        # 启动消费者
        consumers = [
            asyncio.create_task(consumer(queue, session, f"{i}"))
            for i in range(3)
        ]
        # 启动生产者
        producer_task = asyncio.create_task(producer(queue, urls))
        # 等待所有任务完成
        await asyncio.gather(producer_task, *consumers)
        await queue.join()  # 等待队列清空

asyncio.run(producer_consumer_demo())

五、异步上下文管理器与异步迭代器

5.1 异步上下文管理器(async with)
class AsyncResource:
    """模拟需要异步初始化和清理的资源"""
    async def __aenter__(self):
        print("正在获取资源...")
        await asyncio.sleep(1)  # 模拟异步初始化
        print("资源已获取")
        return self

    async def __aexit__(self, exc_type, exc_val, exc_tb):
        print("正在释放资源...")
        await asyncio.sleep(0.5)  # 模拟异步清理
        print("资源已释放")

    async def work(self):
        return "资源使用中"

async def use_async_context():
    async with AsyncResource() as resource:
        result = await resource.work()
        print(result)

# 实际应用:数据库连接
class DatabaseConnection:
    async def __aenter__(self):
        self.conn = await create_db_connection()
        return self.conn

    async def __aexit__(self, *args):
        await self.conn.close()

async def query_db():
    async with DatabaseConnection() as conn:
        return await conn.execute("SELECT * FROM users")
5.2 异步迭代器(async for)
import asyncio

class AsyncRange:
    """异步范围迭代器"""
    def __init__(self, start, end, delay=0.1):
        self.start = start
        self.end = end
        self.delay = delay
        self.current = start

    def __aiter__(self):
        return self

    async def __anext__(self):
        if self.current >= self.end:
            raise StopAsyncIteration
        await asyncio.sleep(self.delay)  # 模拟异步操作
        self.current += 1
        return self.current - 1

async def main():
    async for num in AsyncRange(1, 5):
        print(f"异步生成:{num}")

asyncio.run(main())

# 实际应用:分页 API 请求
class PaginatedAPI:
    async def __aiter__(self):
        return self

    async def __anext__(self):
        page = await self.fetch_page()
        if not page:
            raise StopAsyncIteration
        return page

    async def fetch_page(self):
        # 实际 API 请求逻辑
        pass

async def fetch_all_pages():
    async for page in PaginatedAPI():
        await process_page(page)

六、性能优化与最佳实践

6.1 七项核心原则
# 1. 正确创建任务
async def good():
    task = asyncio.create_task(coro())  # 立即调度
    await task

# 错误:协程不会并发执行
async def bad():
    await coro()  # 等价于同步调用
    await coro2()

# 2. 使用 gather/wait 并发收集
async def fetch_all():
    results = await asyncio.gather(
        fetch_url(url1),
        fetch_url(url2),
        return_exceptions=True  # 防止单个失败影响全部
    )

# 3. 控制并发数量(信号量)
sem = asyncio.Semaphore(10)
async def bounded_fetch(url):
    async with sem:
        return await fetch_url(url)

# 4. 设置超时
async def fetch_with_timeout(url):
    try:
        return await asyncio.wait_for(fetch_url(url), timeout=5.0)
    except asyncio.TimeoutError:
        return None

# 5. 使用异步库,不用阻塞调用
# aiohttp | requests
# aiomysql | pymysql
# asyncpg | psycopg2

# 6. 避免跨线程使用事件循环
loop = asyncio.new_event_loop()
asyncio.set_event_loop(loop)
# 然后在子线程执行 loop.run_forever()

# 7. 正确处理取消
async def cancellable():
    try:
        await asyncio.sleep(10)
    except asyncio.CancelledError:
        print("任务被取消")
        await cleanup()
        raise  # 重新抛出
6.2 性能对比测试框架
import asyncio
import aiohttp
import requests
import time
from functools import wraps

def timeit(func):
    """性能计时装饰器"""
    @wraps(func)
    async def async_wrapper(*args, **kwargs):
        start = time.perf_counter()
        result = await func(*args, **kwargs)
        cost = time.perf_counter() - start
        print(f"{func.__name__} 耗时:{cost:.3f}秒")
        return result
    return async_wrapper

@timeit
async def async_benchmark():
    """异步版本性能测试"""
    urls = ['http://httpbin.org/get'] * 20
    async with aiohttp.ClientSession() as session:
        tasks = [fetch_url(session, url) for url in urls]
        return await asyncio.gather(*tasks)

def sync_benchmark():
    """同步版本性能测试(对比用)"""
    urls = ['http://httpbin.org/get'] * 20
    results = []
    start = time.perf_counter()
    for url in urls:
        results.append(requests.get(url).text)
    print(f"sync_benchmark 耗时:{time.perf_counter() - start:.3f}秒")
    return results

# 运行对比测试
asyncio.run(async_benchmark())
sync_benchmark()

七、常见陷阱与解决方案

7.1 常见错误排查
# 1. 忘记 await
async def mistake1():
    coro()  # 协程未执行,警告:coroutine never awaited

# 正确
async def correct1():
    await coro()

# 2. 在同步函数中创建事件循环多次
def mistake2():
    asyncio.run(main())  # OK
    asyncio.run(main())  # 事件循环已关闭

# 正确:一个程序只有一个入口
if __name__ == "__main__":
    asyncio.run(main())

# 3. 阻塞事件循环
async def mistake3():
    time.sleep(1)  # 阻塞所有协程!
    await asyncio.sleep(0)

# 正确
async def correct3():
    await asyncio.sleep(1)

# 4. 任务创建后不 await
async def mistake4():
    asyncio.create_task(work())  # 任务可能未完成就结束
    # 程序退出,任务被取消

# 正确
async def correct4():
    task = asyncio.create_task(work())
    await task  # 或 await asyncio.gather(task)
7.2 调试技巧
import asyncio
import logging

# 启用调试日志
logging.basicConfig(level=logging.DEBUG)

# 开启 asyncio 调试模式
asyncio.run(main(), debug=True)

# 或手动设置
loop = asyncio.get_event_loop()
loop.set_debug(True)

# 查看未等待的任务
pending = asyncio.all_tasks(loop)
for task in pending:
    print(f"未完成任务:{task}")

八、总结与展望

8.1 异步编程的核心要点
  1. 原理理解:协程遇 IO 自动切换,事件循环统一调度
  2. 语法掌握:async/await、async with、async for、asyncio.gather()
  3. 库选择:使用 aiohttp、asyncpg 等原生异步库
  4. 模式应用:信号量限流、生产者 - 消费者、超时控制
  5. 性能意识:避免阻塞调用,合理设置并发数
8.2 适用场景
场景推荐度原因
网络爬虫五星大量 IO 等待,异步收益明显
Web 应用五星FastAPI、Sanic 等框架原生支持
数据库访问四星连接池 + 异步驱动,吞吐量提升
CPU 密集型两星多进程更合适
简单脚本三星视 IO 密集程度而定
·
8.3 未来演进

Python 3.11+ 引入了更高效的 asyncio 实现,性能进一步提升。异步编程已成为 Python 生态中处理高并发 IO 任务的标准方案,掌握它是现代 Python 开发者必备的核心技能。

目录

  1. 一、异步编程的核心原理
  2. 1.1 什么是异步编程?
  3. 1.2 协程的工作原理
  4. 二、Python 异步编程的演进
  5. 2.1 发展阶段对比
  6. 2.2 原生协程的优势
  7. 生成器协程(3.4 及以前)
  8. 原生协程(3.5+)
  9. 更清晰,无混淆
  10. 三、asyncio 核心组件详解
  11. 3.1 四大核心对象
  12. 1. 协程函数和协程对象
  13. 2. 任务(Task)- 对协程的封装
  14. 或 task = asyncio.ensure_future(coro())
  15. 3. 事件循环 - 调度器
  16. Python 3.7+ 简单方式:
  17. 3.2 awaitable 对象的三种类型
  18. 1. 协程(Coroutine)
  19. 2. 任务(Task)
  20. 3. Future(底层对象,通常不直接使用)
  21. 四、实战:高并发 HTTP 请求
  22. 4.1 基础用法
  23. Python 3.7+ 运行方式
  24. 4.2 高并发批量请求
  25. 4.3 高级模式:生产者 - 消费者
  26. 五、异步上下文管理器与异步迭代器
  27. 5.1 异步上下文管理器(async with)
  28. 实际应用:数据库连接
  29. 5.2 异步迭代器(async for)
  30. 实际应用:分页 API 请求
  31. 六、性能优化与最佳实践
  32. 6.1 七项核心原则
  33. 1. 正确创建任务
  34. 错误:协程不会并发执行
  35. 2. 使用 gather/wait 并发收集
  36. 3. 控制并发数量(信号量)
  37. 4. 设置超时
  38. 5. 使用异步库,不用阻塞调用
  39. aiohttp | requests
  40. aiomysql | pymysql
  41. asyncpg | psycopg2
  42. 6. 避免跨线程使用事件循环
  43. 然后在子线程执行 loop.run_forever()
  44. 7. 正确处理取消
  45. 6.2 性能对比测试框架
  46. 运行对比测试
  47. 七、常见陷阱与解决方案
  48. 7.1 常见错误排查
  49. 1. 忘记 await
  50. 正确
  51. 2. 在同步函数中创建事件循环多次
  52. 正确:一个程序只有一个入口
  53. 3. 阻塞事件循环
  54. 正确
  55. 4. 任务创建后不 await
  56. 正确
  57. 7.2 调试技巧
  58. 启用调试日志
  59. 开启 asyncio 调试模式
  60. 或手动设置
  61. 查看未等待的任务
  62. 八、总结与展望
  63. 8.1 异步编程的核心要点
  64. 8.2 适用场景
  65. 8.3 未来演进

更多推荐文章

查看全部
  • OpenClaw 安装与飞书接入实操
  • ZLibrary 反爬机制深度解析:JS 混淆、签名与频率限制绕过
  • 基于迁移学习的个性化 AKI 预测模型开发与验证
  • GSD 元提示系统:深度拆解解决 AI 编程上下文遗忘问题
  • 大模型降低 AIGC 率指令策略与实战指南
  • 大模型会取代程序员吗?哪些岗位风险最高
  • 基于 C++11 手写 Promise 实现
  • OpenClaw对接飞书机器人高频踩坑实战指南:从插件安装到回调配对全解析
  • LangChain 构建智能 AI 客服系统实战
  • QGroundControl 跨平台安装指南:Windows macOS Linux Android 部署详解
  • 基于 ECharts 与 Three.js 的碳排放可视化大屏实现
  • Cute_Animal_For_Kids_Qwen_Image 儿童专属 AI 绘画工具实战
  • GitHub 启用双因素身份验证(2FA)配置指南:TOTP.app 动态验证码设置
  • Windows 环境下 OpenClaw 接入飞书机器人配置指南
  • Python 异步爬虫结合 K8S 弹性伸缩构建高并发采集引擎
  • 机器人多备用电池与主电池不断电切换管理模块原理及应用
  • 在 Cursor 中配置并使用 MCP 服务实战指南
  • Online Softmax 算法原理与 Flash Attention 应用解析
  • Python 驱动的 Web 与 App 端自动化测试实践
  • LLM-AWQ多模态基准:10个INT4量化模型在20项任务上的全面评估

相关免费在线工具

  • curl 转代码

    解析常见 curl 参数并生成 fetch、axios、PHP curl 或 Python requests 示例代码。 在线工具,curl 转代码在线工具,online

  • Base64 字符串编码/解码

    将字符串编码和解码为其 Base64 格式表示形式即可。 在线工具,Base64 字符串编码/解码在线工具,online

  • Base64 文件转换器

    将字符串、文件或图像转换为其 Base64 表示形式。 在线工具,Base64 文件转换器在线工具,online

  • Markdown转HTML

    将 Markdown(GFM)转为 HTML 片段,浏览器内 marked 解析;与 HTML转Markdown 互为补充。 在线工具,Markdown转HTML在线工具,online

  • HTML转Markdown

    将 HTML 片段转为 GitHub Flavored Markdown,支持标题、列表、链接、代码块与表格等;浏览器内处理,可链接预填。 在线工具,HTML转Markdown在线工具,online

  • JSON 压缩

    通过删除不必要的空白来缩小和压缩JSON。 在线工具,JSON 压缩在线工具,online