异步POST请求成功/失败计数实现及优化方案咨询
异步POST请求优化方案:重试、轮询与指标统计
需求概述
需要实现异步POST请求的成功(状态码=200)与失败(状态码≠200)计数逻辑,同时优化现有异步调用方案,支持重试、状态轮询,并输出请求指标。当前用asyncio.run_in_executor把同步requests调用包装成异步,作为协程新手想拿到更优的实现方案。
现有代码问题
当前方案本质是把同步IO塞进线程池模拟异步,没发挥异步IO的并发优势;缺少重试、状态轮询逻辑;指标统计只覆盖了部分异常场景,没完整统计非200状态码的失败情况。
优化方案要点
- 用
aiohttp替代requests:实现真正的异步HTTP请求,避免线程池开销,提升并发效率 - 加重试机制:针对网络错误、5xx状态码等可重试场景,配置最大重试次数与间隔
- 状态轮询逻辑:如果接口需要异步任务状态查询,请求成功后定时轮询直到拿到目标状态或超时
- 完善指标统计:区分成功(200)、非200失败、异常失败三类场景分别计数
完整代码实现
1. 异步HTTP客户端封装
import asyncio import aiohttp from datetime import datetime import logging from typing import Dict, Optional, List logger = logging.getLogger(__name__) class AsyncAPIClient: def __init__(self, service_uri: str, timeout: int = 10, max_retries: int = 3, retry_delay: float = 1.0): self.service_uri = service_uri.rstrip("/") self.timeout = aiohttp.ClientTimeout(total=timeout) self.max_retries = max_retries self.retry_delay = retry_delay # 指标统计容器 self.metrics = { "success": 0, "failure_non_200": 0, "failure_exception": 0 } async def _request(self, api_path: str, method: str = "GET", **kwargs) -> Optional[aiohttp.ClientResponse]: """内部异步请求方法,包含重试逻辑""" url = f"{self.service_uri}{api_path}" headers = kwargs.pop("headers", {}) headers.setdefault("Content-Type", "application/json") for attempt in range(self.max_retries + 1): try: async with aiohttp.ClientSession(timeout=self.timeout) as session: start_time = datetime.now() async with session.request(method, url, headers=headers, **kwargs) as response: duration = (datetime.now() - start_time).total_seconds() * 1000 logger.debug(f"'{method}'请求 {url} 耗时 {round(duration, 3)} ms,状态码: {response.status}") # 更新指标统计 if response.status == 200: self.metrics["success"] += 1 else: self.metrics["failure_non_200"] += 1 logger.warning(f"请求 {url} 失败,状态码: {response.status}") # 非5xx状态码直接返回,不重试 if response.status < 500: return response # 5xx状态码触发重试 if attempt < self.max_retries: logger.info(f"请求 {url} 遇到5xx错误,第 {attempt+1} 次重试,间隔 {self.retry_delay}s") await asyncio.sleep(self.retry_delay) except aiohttp.ClientError as e: self.metrics["failure_exception"] += 1 logger.error(f"请求 {url} 异常: {str(e)}") if attempt < self.max_retries: logger.info(f"第 {attempt+1} 次重试,间隔 {self.retry_delay}s") await asyncio.sleep(self.retry_delay) else: return None logger.error(f"请求 {url} 达到最大重试次数 {self.max_retries},失败") return None async def publish_actual(self, event_name: str, custom_payload: Dict = {}, event_message_params: List = []): """异步POST请求方法,可根据实际需求填充请求体""" json_data = {} # 这里根据业务逻辑组装请求体,比如: # json_data["event_name"] = event_name # json_data.update(custom_payload) path = "/some/path" response = await self._request(path, "POST", json=json_data) if response: # 按需处理响应数据,比如转JSON # return await response.json() return await response.text() return None async def poll_status(self, status_path: str, target_status: str, poll_interval: float = 2.0, timeout: float = 30.0) -> bool: """通用状态轮询方法""" start_time = datetime.now() while (datetime.now() - start_time).total_seconds() < timeout: response = await self._request(status_path, "GET") if response and response.status == 200: data = await response.json() if data.get("status") == target_status: logger.info(f"轮询到目标状态: {target_status}") return True logger.info(f"未轮询到目标状态,{poll_interval}s后重试") await asyncio.sleep(poll_interval) logger.error(f"状态轮询超时 {timeout}s") return False def get_metrics(self) -> Dict: """获取当前请求指标统计""" return self.metrics.copy()
2. 使用示例
async def main(): # 初始化客户端,配置服务地址、超时、重试参数 client = AsyncAPIClient("https://your-service-url.com", timeout=10, max_retries=3) # 发送异步POST请求 await client.publish_actual("test_event", {"key": "value"}) # 如果需要轮询任务状态,取消注释下面一行 # await client.poll_status("/status/123", "completed") # 输出当前指标统计 print("请求指标统计:", client.get_metrics()) if __name__ == "__main__": asyncio.run(main())
关键说明
- 纯异步实现:用
aiohttp替代requests,不需要线程池,真正发挥异步IO的并发优势 - 可配置重试:针对网络异常和5xx错误自动重试,重试次数和间隔可自定义
- 完整指标统计:三类失败场景分别计数,方便监控和排查问题
- 灵活轮询:封装通用轮询逻辑,支持自定义目标状态、间隔和超时时间
- 清晰日志:不同场景的日志分级输出,便于定位问题
内容的提问来源于stack exchange,提问作者harsh solanki
相关产品推荐
相关产品推荐

