,

python常用异步编程模式

并发批处理模式 限速器模式 重试模式 优雅关闭模式

并发批处理模式

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())

发表回复

您的邮箱地址不会被公开。 必填项已用 * 标注

About the Author

每个人都有自己得时区,在自己得时区里,一切都是准时的。

BlockSpare — News, Magazine and Blog Addons for (Gutenberg) Block Editor