You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

异步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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.12 17:10:24