跳至正文
来两杯美式
返回

Python 知识系列(六):异步编程

By 来两杯美式
发布于

1. 什么是异步编程

异步编程是一种非阻塞的编程范式,它允许程序在等待某些耗时操作(如网络请求、文件读写、数据库查询等I/O密集型任务)完成时,不必阻塞主线程的执行,从而可以切换去执行其他任务,极大地提高了程序的并发处理能力和整体性能。

1.1 异步编程与同步编程的核心区别

在Python中,异步编程主要通过协程 (Coroutine) 和事件循环 (Event Loop) 来实现。协程是轻量级的“线程”,可以在I/O操作时主动让出执行权;而事件循环则是调度中心,负责管理和调度所有协程的执行。

1.2 异步编程的优势

1.3 适用场景与局限性

主要应用场景

局限性

2. Python异步编程核心概念

2.1 事件循环 (Event Loop)

事件循环是异步编程的“心脏”,它是一个持续运行的循环,负责任务调度、监听I/O事件和执行回调。其底层通常利用操作系统的I/O多路复用机制(如Linux的epoll),使得它可以在等待多个I/O操作时保持非阻塞。

每个事件循环调用都是完全独立的,每个asyncio.run()有完整的生命周期,包括创建、关闭等。

import asyncio

async def main():
    print("Hello")
    await asyncio.sleep(1)  # 模拟I/O操作,此时会切换到其他任务
    print("World")

if __name__ == "__main__":
    # asyncio.run() 会自动创建和管理事件循环
    asyncio.run(main())

事件循环的工作原理(简化版)

  1. 维护一个可执行任务的队列。
  2. 循环从队列中取出任务并执行。
  3. 如果任务执行到await表达式(等待一个I/O操作),它会暂停执行并让出控制权。
  4. 事件循环会继续执行队列中的其他任务,同时监听之前暂停任务所等待的I/O事件。
  5. 当I/O事件完成时,事件循环会将对应的暂停任务重新放回可执行队列,等待下一次调度。

2.2 协程 (Coroutine)

协程是异步任务的基本单元,本质上是一个可以暂停和恢复执行的特殊函数。通过async def关键字定义。

import asyncio

async def my_coroutine():
    print("协程开始")
    await asyncio.sleep(1)  # 暂停点:释放CPU,允许事件循环执行其他任务
    print("协程结束")
    return "执行结果"

async def main():
    # 调用协程函数返回一个协程对象,此时代码并未执行
    coro = my_coroutine()
    print(f"协程对象类型: {type(coro)}")
    # 使用 await 关键字驱动协程执行,并等待其结果
    result = await coro
    print(f"协程返回值: {result}")

if __name__ == "__main__":
    asyncio.run(main())

执行流程

  1. 程序启动if __name__ == "__main__":条件满足,调用asyncio.run(main())
  2. 创建事件循环asyncio.run()创建新的事件循环
  3. 进入main协程:开始执行main()函数
  4. 创建协程对象coro = my_coroutine()(此时只是创建协程对象,并未执行)
  5. 打印协程类型:输出"协程对象类型: <class 'coroutine'>"
  6. await 触发执行await coro开始真正执行my_coroutine()
  7. 打印”协程开始”
  8. 遇到第一个awaitawait asyncio.sleep(1)(协程暂停,控制权返回事件循环)
  9. 休眠1秒:事件循环处理其他任务(本例中没有其他任务)
  10. 恢复执行:1秒后继续执行协程剩余部分
  11. 打印”协程结束”
  12. 返回结果:返回字符串"执行结果"
  13. 回到main()协程result变量获得返回值"执行结果"
  14. 打印返回值:输出"协程返回值: 执行结果"
  15. 事件循环关闭asyncio.run()自动关闭事件循环

关键点await是协程执行的”开关”,只有遇到await时协程才会真正开始执行并可能暂停。

2.4 任务 (Task)

任务 (asyncio.Task) 是对协程的进一步封装,它代表一个在事件循环中独立运行的并发单元。使用asyncio.create_task()创建任务后,该任务会被立即提交给事件循环,在不久的将来自动执行,无需手动await启动(但是结果需要使用await获取)。

import asyncio

async def task_func(name, delay):
    print(f"任务 {name} 开始,将执行 {delay} 秒")
    await asyncio.sleep(delay)
    print(f"任务 {name} 结束")
    return f"任务 {name} 完成"

async def main():
    # 创建任务,它们会“在后台”并发执行
    task1 = asyncio.create_task(task_func("A", 10))
    task2 = asyncio.create_task(task_func("B", 5))

    # 只有当执行到await task时,才会获取task的结果
    # 如果任务已经执行完了,就会立即返回结果,否则等待任务执行
    result1 = await task1
    result2 = await task2
    # 注意:这里会先等待task1完成(即使task2先完成,也要等task1结束后才能继续)
    print("所有任务完成,结果:", result1, result2)

if __name__ == "__main__":
    asyncio.run(main())

# 输出
# 任务 A 开始,将执行 10 秒
# 任务 B 开始,将执行 5 秒
# 任务 B 结束
# 任务 A 结束
# 所有任务完成,结果: 任务 A 完成 任务 B 完成

补充Future对象是一个代表异步操作最终结果的底层抽象。它是一个容器,表示“未来某个时刻会完成的操作”。任务 (Task) 是Future的一个具体子类。通常我们直接使用Task,较少直接操作Future

3. asyncio模块详解与核心API

3.1 async/await 语法

3.2 并发任务管理(选学)

asyncio.gather() - 批量并发与结果聚合

gather()用于同时运行多个可等待对象,并按输入顺序将它们的结果聚合到一个列表中。

import asyncio

async def task(i):
    await asyncio.sleep(i)
    return f"任务 {i} 完成"

async def main():
    # 并发执行3个任务
    results = await asyncio.gather(
        task(1),
        task(2),
        task(0.5)
    )
    # 结果按输入顺序返回,而不是完成顺序
    print(results)  # ['任务 1 完成', '任务 2 完成', '任务 0.5 完成']

if __name__ == "__main__":
    asyncio.run(main())

asyncio.wait() - 更灵活的完成条件控制

wait()提供了更细粒度的控制,可以指定等待条件(如等待第一个任务完成、等待所有任务完成或出现第一个异常)。

import asyncio

async def task(i):
    await asyncio.sleep(i)
    return i

async def main():
    tasks = [
        asyncio.create_task(task(1)),
        asyncio.create_task(task(2)),
        asyncio.create_task(task(0.5))
    ]

    # 等待第一个任务完成
    done, pending = await asyncio.wait(tasks, return_when=asyncio.FIRST_COMPLETED)
    print(f"已完成的任务: {[t.result() for t in done]}")
    print(f"未完成的任务数: {len(pending)}")

    # 取消未完成的任务
    for t in pending:
        t.cancel()

if __name__ == "__main__":
    asyncio.run(main())

4. with关键字

在实际编程中,系统资源(如文件、数据库连接、网络套接字等)在使用后必须正确关闭。但是这种写法存在明显缺陷:

使用with关键字,可实现:

使用示例

# 优先使用 with 管理资源
# 对于文件、网络连接、锁等资源,总是优先考虑使用 with 语句
# 良好的实践 - 资源管理清晰明确
with open('data.txt', 'r') as file:
    # 只包含与资源相关的操作
    data = file.read()

5. 高级特性与性能优化

5.1 并发限制 (Concurrency Limiting)

在爬虫或API客户端等场景中,同时发起过多请求可能会压垮服务器或耗尽本地资源。使用asyncio.Semaphore可以有效控制并发任务的数量。

import asyncio
import aiohttp

async def controlled_fetch(semaphore, session, url):
    async with semaphore:
        # 在这个代码块内,同时运行的任务数不会超过Semaphore的限制
        print(f"开始获取 {url}")
        await asyncio.sleep(1) # 模拟网络延迟
        print(f"完成获取 {url}")
        return f"Result from {url}"

async def main():
    # 创建一个信号量,最多只允许3个任务同时执行
    semaphore = asyncio.Semaphore(3)
    urls = [f"http://example.com/{i}" for i in range(10)]

    async with aiohttp.ClientSession() as session:
        tasks = [controlled_fetch(semaphore, session, url) for url in urls]
        await asyncio.gather(*tasks)

if __name__ == "__main__":
    asyncio.run(main())

5.2 优雅的异常处理

使用asyncio.gather时,默认情况下一个任务失败会导致整个gather调用抛出异常。通过设置return_exceptions=True,可以将异常作为结果返回,从而不中断其他任务。

import asyncio

async def successful_task():
    await asyncio.sleep(1)
    return "成功"

async def failing_task():
    await asyncio.sleep(0.5)
    raise ValueError("任务失败")

async def main():
    results = await asyncio.gather(
        successful_task(),
        failing_task(),
        return_exceptions=True
    )

    for result in results:
        if isinstance(result, Exception):
            print(f"捕获到一个异常: {result}")
        else:
            print(f"获取到成功结果: {result}")

if __name__ == "__main__":
    asyncio.run(main())

5.3 其他最佳实践

6. Python Web服务并发模型与架构选型

6.1 痛点分析

不同于Java中的Tomcat、Jetty等多线程Web服务器,Python的服务器是单进程、单线程的同步Web服务器。例如当一个请求在处理过程中执行数据库查询或其它耗时操作时,该请求所在的进程/线程会被完全阻塞,无法处理任何其他请求直至上一个请求结束。

为什么同步服务会被阻塞?这涉及到Python的全局解释器锁(GIL)。

GIL 的存在:确保在任何一个时刻,单个Python进程内只有一个线程在执行Python字节码。

GIL的设计初衷是为了简化CPython的内存管理。通过GIL,CPython的内存管理机制变成了线程安全的,避免了复杂的锁机制来防止多个线程同时修改Python对象,从而使得大量的C语言扩展库可以轻松集成。

重要事实:Python标准库和绝大多数第三方库在执行阻塞式I/O操作时,会主动释放GIL!

这意味着,当线程A发起一个数据库查询时,它会释放GIL并进入等待状态。此时,线程B可以立即获得GIL,开始处理另一个请求。当线程A的数据库查询返回结果后,它会重新排队等待获取GIL以继续执行。

6.2 Python Web服务并发模型演进

基于对GIL的理解,我们可以清晰地梳理出Python Web服务的并发模型演进路径。

模型一:单进程单线程同步模型 (The “Blocking” Model)

模型二:多进程 / 多线程同步模型 (The “Worker” Model - WSGI)

为了解决单线程阻塞问题,社区制定了WSGI (Web Server Gateway Interface) 规范。它是一种标准接口,解耦了Web服务器与Web应用框架(如Flask, Django)。

像Gunicorn、uWSGI这样的生产级WSGI服务器,通过创建多个工作单元(Worker)来实现并发。

工作方式(以Gunicorn为例)

  1. Gunicorn启动一个主进程(Master Process)。
  2. 主进程会创建(fork)多个子进程,即工作进程(Worker Processes)。
  3. 每个Worker都是一个独立的Python进程,拥有自己的内存空间和自己的GIL。
  4. 当请求到来时,主进程会将其分发给一个空闲的Worker。

并发原理

模型三:单线程异步模型 (The “Non-Blocking” Model - ASGI)

随着异步编程在Python中的兴起(以asyncio库为核心),诞生了新一代的规范ASGI (Asynchronous Server Gateway Interface)。

核心思想: 异步模型不再依赖多线程/多进程来处理并发,而是在单个线程内通过事件循环(Event Loop)和协程(Coroutine)来实现。

与多线程的对比

代表框架与服务器

代码示例 (FastAPI)

# aio_app.py
import asyncio
from fastapi import FastAPI
from datetime import datetime

app = FastAPI()

@app.get("/")
async def handle_request():
    request_time = datetime.now().strftime('%H:%M:%S')
    print(f"[{request_time}] Request received. Starting 5-second async task...")

    # 使用 asyncio.sleep(5) 模拟一个非阻塞的I/O操作
    await asyncio.sleep(5)

    completion_time = datetime.now().strftime('%H:%M:%S')
    print(f"[{completion_time}] Task finished. Sending response.")
    return {"message": f"Processed at {completion_time}"}
**部署示例**:
    # 安装 FastAPI 和 Uvicorn
# pip install fastapi "uvicorn[standard]"

# 启动服务
uvicorn aio_app:app --host 0.0.0.0 --port 8000

6.3 总结与架构选型建议

特性/模型单进程同步 (Flask Dev)多进程同步 (Gunicorn + Flask)单线程异步 (Uvicorn + FastAPI)
并发机制无(阻塞)多进程并行 / 多线程并发单线程事件循环 + 协程
GIL影响-每个进程有独立的GIL,通过进程实现并行单线程内,GIL不是并发瓶颈
资源消耗极低高(进程开销大)非常低(协程开销小)
理论并发数1取决于Worker数量(几十到几百)非常高(上千到数万)
适用场景开发、调试CPU密集型、传统I/O密集型应用、已有大量同步库的系统高并发I/O密集型应用、流式处理、WebSockets
编程复杂度简单简单(传统同步代码)中等(需要async/await和异步生态)

架构选型指南

  1. 传统同步应用

    • 适用场景:已有成熟的同步代码库、CPU密集型任务、简单的CRUD应用
    • 推荐:Flask/Django + Gunicorn/uWSGI
    • 配置建议:Worker数量 = CPU核心数 × 2 + 1
  2. 高并发I/O密集型应用

    • 适用场景:API网关、实时通信、WebSocket服务、爬虫系统
    • 推荐:FastAPI/Starlette + Uvicorn
    • 优势:内存占用小、并发能力高、响应速度快
  3. 混合架构

    • 适用场景:部分模块需要高性能异步处理,部分模块依赖同步库
    • 方案:FastAPI(异步主框架)+ 线程池处理同步阻塞操作

7. 异步编程实战应用

7.1 网络请求 (aiohttp)

aiohttp是一个流行的异步HTTP客户端/服务器库。

并发爬取网页示例

import asyncio
import aiohttp

async def fetch(session, url):
    """异步获取单个URL内容"""
    try:
        async with session.get(url, timeout=10) as response:
            response.raise_for_status()  # 如果状态码不是2xx,则抛出异常
            return {
                "url": url,
                "status": response.status,
                "length": len(await response.text())
            }
    except Exception as e:
        return {"url": url, "error": str(e)}

async def main():
    urls = [
        "https://api.github.com",
        "http://httpbin.org/get",
        "https://invalid-url-for-testing.com"
    ]

    # 推荐在async with块中创建Session,它会自动管理连接池
    async with aiohttp.ClientSession() as session:
        tasks = [fetch(session, url) for url in urls]
        results = await asyncio.gather(*tasks)

    for res in results:
        if "error" in res:
            print(f"{res['url']}: 请求失败 - {res['error']}")
        else:
            print(f"{res['url']} | 状态: {res['status']} | 内容长度: {res['length']}")

if __name__ == "__main__":
    asyncio.run(main())

7.2 异步文件操作 (aiofiles)

标准的文件操作是阻塞的。aiofiles库提供了异步的文件读写接口。

import asyncio
import aiofiles

async def write_file(filename, content):
    async with aiofiles.open(filename, 'w', encoding='utf-8') as f:
        await f.write(content)

async def read_file(filename):
    async with aiofiles.open(filename, 'r', encoding='utf-8') as f:
        return await f.read()

async def main():
    files_to_write = {
        "file1.txt": "这是第一个异步文件。",
        "file2.txt": "这是第二个异步文件。"
    }

    # 并发写入文件
    write_tasks = [write_file(name, content) for name, content in files_to_write.items()]
    await asyncio.gather(*write_tasks)
    print("文件写入完成。")

    # 并发读取文件
    read_tasks = [read_file(name) for name in files_to_write.keys()]
    contents = await asyncio.gather(*read_tasks)

    for name, content in zip(files_to_write.keys(), contents):
        print(f"读取 {name}: {content}")

if __name__ == "__main__":
    asyncio.run(main())

7.3 异步数据库操作 (asyncpg, aiomysql)

使用异步数据库驱动可以避免在数据库查询时阻塞整个应用。

PostgreSQL (asyncpg) 示例

import asyncio
import asyncpg

async def run_query():
    # 推荐使用连接池进行管理
    conn = await asyncpg.connect(user='user', password='password', database='test', host='127.0.0.1')
    values = await conn.fetch('SELECT * FROM users LIMIT 5')
    print(values)
    await conn.close()

# if __name__ == "__main__":
#     asyncio.run(run_query())

使用连接池的推荐方式

import asyncio
import asyncpg

async def run_with_pool():
    # 创建连接池
    pool = await asyncpg.create_pool(
        user='user',
        password='password',
        database='test',
        host='127.0.0.1',
        min_size=5,
        max_size=20
    )

    async with pool.acquire() as conn:
        values = await conn.fetch('SELECT * FROM users LIMIT 5')
        print(values)

    await pool.close()

7.4 在异步代码中调用同步阻塞库

有时必须使用一个不支持异步的同步库(如requests)。直接调用会阻塞事件循环。正确的做法是使用loop.run_in_executor()将其放入一个独立的线程池中执行。

import asyncio
import requests  # 一个同步阻塞库

def sync_blocking_get(url):
    """这是一个耗时的同步函数"""
    print(f"开始同步请求 {url}...")
    response = requests.get(url)
    print(f"完成同步请求 {url}。")
    return response.status_code

async def main():
    loop = asyncio.get_running_loop()
    urls = ["https://www.google.com", "https://www.github.com"]

    # 使用 run_in_executor 将同步函数提交到线程池执行
    # 第一个参数为 None 时使用默认的线程池
    tasks = [loop.run_in_executor(None, sync_blocking_get, url) for url in urls]
    results = await asyncio.gather(*tasks)

    print("所有同步任务在异步环境中完成,结果:", results)

if __name__ == "__main__":
    asyncio.run(main())

分享这篇文章:
通过邮件分享这篇文章✓ 链接已复制
所属专题
Python
第 6 / 8 篇
查看系列全部文章
  1. 01.Python 知识系列(一):极简核心语法
  2. 02.Python 知识系列(二):面向对象编程
  3. 03.Python 知识系列(三):从 0 开始理解函数
  4. 04.Python 知识系列(四):装饰器详解
  5. 05.Python 知识系列(五):模块和包
  6. 06.Python 知识系列(六):异步编程
  7. 07.Python 知识系列(七):虚拟环境与依赖管理
  8. 08.Python 知识系列(八):企业级开发场景最佳实践

上一篇
Python 知识系列(七):虚拟环境与依赖管理
下一篇
Python 知识系列(五):模块和包