并发批处理模式
async def batch_process(urls: list[str]):
"""将一批任务并发执行,并收集结果"""
async def process(url: str) -> dict:
async with httpx.AsyncClient() as client:
resp = await client.get(url)
return {"url": url, "status": resp.status_code}
results = await asyncio.gather(*[process(url) for url in urls])
return results
限速器模式
class RateLimiter:
"""限制单位时间内的请求数"""
def __init__(self, max_rate: float, interval: float = 1.0):
self._sem = asyncio.Semaphore(max_rate)
self._interval = interval
async def acquire(self):
await self._sem.acquire()
def release():
self._sem.release()
loop = asyncio.get_running_loop()
loop.call_later(self._interval, release)
async def __aenter__(self):
await self.acquire()
async def __aexit__(self, *args):
pass
async def main():
rate_limiter = RateLimiter(max_rate=2, interval=1.0)
async def fetch(url: str) -> str:
async with rate_limiter:
await asyncio.sleep(0.3)
return f"{url} done"
results = await asyncio.gather(*[fetch(f"url{i}") for i in range(6)])
print(results) # 每秒最多完成 2 个请求
asyncio.run(main())
重试模式
async def retry(coro_factory, max_retries: int = 3, delay: float = 1.0):
"""为异步操作添加重试机制"""
for attempt in range(max_retries):
try:
return await coro_factory()
except Exception as e:
if attempt == max_retries - 1:
raise
print(f"第 {attempt + 1} 次失败,{delay} 秒后重试...")
await asyncio.sleep(delay)
async def main():
n = 0
async def unstable_request() -> str:
nonlocal n
n += 1
if n < 3:
raise ConnectionError(f"第 {n} 次请求失败")
return "成功响应"
result = await retry(unstable_request, max_retries=3, delay=0.5)
print(result) # 前 2 次失败,第 3 次成功
asyncio.run(main())
优雅关闭模式
import asyncio
import signal
class GracefulServer:
def __init__(self):
self._running = True
async def serve(self):
while self._running:
try:
await asyncio.sleep(1) # 模拟处理请求
print("正在处理请求...")
except asyncio.CancelledError:
print("收到取消信号,正在关闭...")
break
def shutdown(self):
print("开始优雅关闭...")
self._running = False
async def run(self):
loop = asyncio.get_running_loop()
stop = loop.create_future()
def signal_handler():
stop.set_result(None)
loop.add_signal_handler(signal.SIGINT, signal_handler) # Ctrl+C
loop.add_signal_handler(signal.SIGTERM, signal_handler) # 终止信号
task = asyncio.create_task(self.serve())
await stop # 等待关闭信号
self.shutdown()
task.cancel()
await task
async def main():
server = GracefulServer()
await server.run()
# 按 Ctrl+C 触发 SIGINT,程序会优雅退出而非直接崩溃
asyncio.run(main())