1. 什么是异步编程
异步编程是一种非阻塞的编程范式,它允许程序在等待某些耗时操作(如网络请求、文件读写、数据库查询等I/O密集型任务)完成时,不必阻塞主线程的执行,从而可以切换去执行其他任务,极大地提高了程序的并发处理能力和整体性能。
1.1 异步编程与同步编程的核心区别
- 同步编程:代码按顺序逐行执行。当遇到一个耗时的I/O操作时,程序会阻塞等待,直到该操作完成后才能继续执行下一行代码。在处理大量并发I/O时,这种方式效率低下,因为CPU大部分时间都处于空闲等待状态。
- 异步编程:当遇到耗时的I/O操作时,程序会“挂起”当前任务,并立即切换到其他已就绪的任务继续执行。当之前的I/O操作完成后,事件循环会得到通知,并在适当的时候恢复执行被挂起的任务。
在Python中,异步编程主要通过协程 (Coroutine) 和事件循环 (Event Loop) 来实现。协程是轻量级的“线程”,可以在I/O操作时主动让出执行权;而事件循环则是调度中心,负责管理和调度所有协程的执行。
1.2 异步编程的优势
- 高并发处理能力:在单线程内实现高并发,避免了多线程模式下昂贵的上下文切换开销和线程安全问题。
- 高性能:针对I/O密集型任务,性能提升显著,通常是同步方式的数倍到数十倍。
- 资源效率:相比多线程,协程占用的内存和CPU资源更少。
1.3 适用场景与局限性
主要应用场景:
- 网络编程:如构建高并发的Web服务器、API网关、网络爬虫等。
- I/O密集型任务:处理大量文件读写、数据库查询等操作。
- 实时数据处理:如处理WebSocket消息流、传感器数据、日志流等。
- 高并发服务:如实时通信系统(聊天室)、消息队列消费者等。
局限性:
- 不适用于计算密集型任务:异步编程在单线程内切换协程。如果任务需要大量CPU计算(如图像处理、复杂算法),它会长时间占用CPU,导致事件循环无法调度其他任务,反而会降低整体性能。这类任务更适合使用多进程。
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())
事件循环的工作原理(简化版):
- 维护一个可执行任务的队列。
- 循环从队列中取出任务并执行。
- 如果任务执行到
await表达式(等待一个I/O操作),它会暂停执行并让出控制权。 - 事件循环会继续执行队列中的其他任务,同时监听之前暂停任务所等待的I/O事件。
- 当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())
执行流程:
- 程序启动:
if __name__ == "__main__":条件满足,调用asyncio.run(main()) - 创建事件循环:
asyncio.run()创建新的事件循环 - 进入main协程:开始执行
main()函数 - 创建协程对象:
coro = my_coroutine()(此时只是创建协程对象,并未执行) - 打印协程类型:输出
"协程对象类型: <class 'coroutine'>" - await 触发执行:
await coro开始真正执行my_coroutine() - 打印”协程开始”
- 遇到第一个await:
await asyncio.sleep(1)(协程暂停,控制权返回事件循环) - 休眠1秒:事件循环处理其他任务(本例中没有其他任务)
- 恢复执行:1秒后继续执行协程剩余部分
- 打印”协程结束”
- 返回结果:返回字符串
"执行结果" - 回到main()协程:
result变量获得返回值"执行结果" - 打印返回值:输出
"协程返回值: 执行结果" - 事件循环关闭:
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 语法
async def:用于定义一个协程函数。await:用于暂停当前协程的执行,等待一个可等待对象(如另一个协程、任务或Future)完成。它只能在async def函数内部使用。
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关键字
在实际编程中,系统资源(如文件、数据库连接、网络套接字等)在使用后必须正确关闭。但是这种写法存在明显缺陷:
- 容易忘记关闭资源
- 代码冗长,需要多行代码完成简单操作
- 异常处理复杂,需要手动编写try-finally块
使用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 其他最佳实践
- 批量操作:尽量合并多个小的I/O操作为一个大的批量操作,以减少网络或磁盘I/O的次数。
- 连接复用:对于数据库和HTTP请求,应使用连接池 (
asyncpg.create_pool,aiohttp.ClientSession) 来复用连接,避免频繁创建和销毁连接的开销。 - 异步缓存:对于频繁访问且不常变化的数据,使用异步缓存(如
aioredis)可以显著降低延迟。
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为例):
- Gunicorn启动一个主进程(Master Process)。
- 主进程会创建(fork)多个子进程,即工作进程(Worker Processes)。
- 每个Worker都是一个独立的Python进程,拥有自己的内存空间和自己的GIL。
- 当请求到来时,主进程会将其分发给一个空闲的Worker。
并发原理:
- 多进程并发:这是最核心的并发方式。由于每个Worker是独立的进程,它们可以在多核CPU上并行执行。如果一个Worker因为处理请求而被阻塞(无论是CPU计算还是I/O等待),完全不影响其他Worker接收和处理新的请求。
- 多线程并发(可选):Gunicorn还允许为每个Worker配置多个线程。在这种模式下,一个请求的I/O操作会释放其所在线程的GIL,让同一Worker内的其他线程能够处理其他请求。这在I/O密集型应用中可以进一步提升单个Worker的并发处理能力。
模型三:单线程异步模型 (The “Non-Blocking” Model - ASGI)
随着异步编程在Python中的兴起(以asyncio库为核心),诞生了新一代的规范ASGI (Asynchronous Server Gateway Interface)。
核心思想: 异步模型不再依赖多线程/多进程来处理并发,而是在单个线程内通过事件循环(Event Loop)和协程(Coroutine)来实现。
- 协程:可以看作是一种轻量级的、可被程序自身调度挂起和恢复的“线程”。
- 事件循环:一个中央调度器,负责管理和执行所有协程。当一个协程遇到I/O操作(例如
await some_db_query())时,它不会阻塞整个线程,而是会主动让出控制权给事件循环。事件循环会立即切换到另一个准备就绪的协程去执行。当之前的I/O操作完成后,事件循环会得到通知,并在合适的时机恢复执行原来的协程。
与多线程的对比:
- 资源开销:协程的创建和切换成本远低于操作系统线程,因此单线程异步模型可以用极低的内存开销支持极高的并发连接数(C10K/C100K问题)。
- 编程范式:需要使用
async和await关键字,整个技术栈(Web框架、数据库驱动、HTTP客户端等)都必须是异步兼容的。
代表框架与服务器:
- Web框架:FastAPI, Starlette, Sanic
- ASGI服务器:Uvicorn, Hypercorn, Daphne
代码示例 (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和异步生态) |
架构选型指南:
-
传统同步应用:
- 适用场景:已有成熟的同步代码库、CPU密集型任务、简单的CRUD应用
- 推荐:Flask/Django + Gunicorn/uWSGI
- 配置建议:Worker数量 = CPU核心数 × 2 + 1
-
高并发I/O密集型应用:
- 适用场景:API网关、实时通信、WebSocket服务、爬虫系统
- 推荐:FastAPI/Starlette + Uvicorn
- 优势:内存占用小、并发能力高、响应速度快
-
混合架构:
- 适用场景:部分模块需要高性能异步处理,部分模块依赖同步库
- 方案: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())