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能带来显著的性能提升。