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

基于Asyncio与aiohttp的PUT请求:错误处理与请求管控问题

问题解决:aiohttp + Asyncio 批量PUT请求的错误处理、流程中止与请求管控

一、原代码的核心问题

你当前的代码存在一个关键逻辑错误:在循环中每次添加任务后立即执行await asyncio.gather(*tasks),这会导致请求串行执行,完全失去异步IO的性能优势。正确的做法是先批量创建所有任务,再统一等待执行结果,或者用实时遍历的方式处理完成的任务。

二、错误处理方案

1. 请求级异常捕获

在send_payload函数中捕获网络错误、超时、HTTP服务端错误等异常,返回统一格式的结果,避免单个请求失败导致整个任务集崩溃:

async def send_payload(session, url, semaphore):
    payload = '{"name":"John", "age":30, "car":null}'        
    token = {'apikey':'abcd'}
    async with semaphore:
        try:
            async with session.put(url, data=payload, headers=token) as resp:
                resp_text = await resp.text()
                # 返回(状态码, 响应内容, 异常信息)的统一结构
                return (resp.status, resp_text, None)
        except (aiohttp.ClientError, asyncio.TimeoutError) as e:
            # 捕获网络类异常:连接失败、超时、SSL错误等
            return (None, None, str(e))

2. 全局异常兜底

如果希望单个请求失败不阻断其他请求,可在asyncio.gather中添加return_exceptions=True参数,此时异常会被作为结果返回而非抛出。但更推荐在任务内部捕获异常,便于统一处理结果格式。

三、中止请求处理流程

若需要在特定条件(如错误数达到阈值、遇到致命错误)下立即停止所有未完成请求,可通过取消任务实现:

实时监控+中途中止(推荐)

使用asyncio.as_completed遍历任务,实时处理完成的请求,一旦触发中止条件就取消剩余任务:

async def main():
    timeout = aiohttp.ClientTimeout(total=3)
    semaphore = asyncio.Semaphore(10)  # 控制并发数
    abort_threshold = 10  # 错误数达10则中止

    async with aiohttp.ClientSession(connector=aiohttp.TCPConnector(ssl=False), timeout=timeout) as session:
        count = 0
        errcount = 0
        tasks = []

        # 批量创建所有任务
        for _ in range(1, 151):
            url = "https://myapi.com/test"
            task = asyncio.create_task(send_payload(session, url, semaphore))
            tasks.append(task)

        # 实时处理完成的任务
        for task in asyncio.as_completed(tasks):
            count += 1
            try:
                status, resp_text, exc = await task
                if exc:
                    print(f"请求{count}失败: {exc}")
                    errcount += 1
                elif status != 200:
                    print(f"请求{count}错误: 状态码{status}")
                    errcount += 1
                else:
                    print(f"请求{count}成功: {status} {resp_text}\n")

                # 触发中止条件:取消所有未完成任务
                if errcount >= abort_threshold:
                    print(f"错误数达阈值{abort_threshold},中止剩余请求")
                    for t in tasks:
                        if not t.done():
                            t.cancel()
                    break
            except asyncio.CancelledError:
                print(f"请求{count}被取消")
                continue

四、管控每个PUT请求

1. 并发数控制

用asyncio.Semaphore限制同时发起的请求数量,避免对API造成压力或触发限流。比如Semaphore(10)表示同时最多10个请求在执行,可根据API的限流规则调整数值。

2. 单个请求超时

可在全局会话设置超时,也可为单个请求单独指定超时:

# 单个请求单独设置5秒超时
async with session.put(url, data=payload, headers=token, timeout=aiohttp.ClientTimeout(total=5)) as resp:

3. 临时错误重试

对5xx服务端错误、临时网络波动等情况进行重试,提升请求成功率:

async def send_payload(session, url, semaphore, max_retries=3):
    payload = '{"name":"John", "age":30, "car":null}'        
    token = {'apikey':'abcd'}
    retry_count = 0
    while retry_count < max_retries:
        async with semaphore:
            try:
                async with session.put(url, data=payload, headers=token) as resp:
                    resp_text = await resp.text()
                    # 5xx状态码重试
                    if 500 <= resp.status < 600:
                        retry_count += 1
                        await asyncio.sleep(1)  # 重试前等待1秒
                        continue
                    return (resp.status, resp_text, None)
            except (aiohttp.ClientError, asyncio.TimeoutError) as e:
                retry_count += 1
                await asyncio.sleep(1)
                continue
    # 重试失败返回最终结果
    return (None, None, f"重试{max_retries}次后失败")

五、完整优化代码

import aiohttp
import asyncio
import time

start_time = time.time()

async def send_payload(session, url, semaphore, max_retries=3):
    payload = '{"name":"John", "age":30, "car":null}'        
    token = {'apikey':'abcd'}
    retry_count = 0
    while retry_count < max_retries:
        async with semaphore:
            try:
                async with session.put(url, data=payload, headers=token) as resp:
                    resp_text = await resp.text()
                    if 500 <= resp.status < 600:
                        retry_count += 1
                        await asyncio.sleep(1)
                        continue
                    return (resp.status, resp_text, None)
            except (aiohttp.ClientError, asyncio.TimeoutError) as e:
                retry_count += 1
                await asyncio.sleep(1)
                continue
    return (None, None, f"重试{max_retries}次失败: 未知错误")

async def main():
    timeout = aiohttp.ClientTimeout(total=3)
    semaphore = asyncio.Semaphore(10)
    abort_threshold = 10
    
    async with aiohttp.ClientSession(connector=aiohttp.TCPConnector(ssl=False), timeout=timeout) as session:
        count = 0
        errcount = 0
        tasks = []
        
        for _ in range(1, 151):
            url = "https://myapi.com/test"
            task = asyncio.create_task(send_payload(session, url, semaphore))
            tasks.append(task)
        
        for task in asyncio.as_completed(tasks):
            count += 1
            try:
                status, resp_text, exc = await task
                if exc:
                    print(f"请求{count}失败: {exc}")
                    errcount += 1
                elif status != 200:
                    print(f"请求{count}错误: 状态码{status}, 响应: {resp_text}")
                    errcount += 1
                else:
                    print(f"请求{count}成功: 状态码{status}, 响应: {resp_text}\n")
                
                if errcount >= abort_threshold:
                    print(f"错误数达{abort_threshold},中止剩余请求")
                    for t in tasks:
                        if not t.done():
                            t.cancel()
                    break
            except asyncio.CancelledError:
                print(f"请求{count}被取消")
                continue
            except Exception as e:
                print(f"请求{count}未捕获异常: {str(e)}")
                errcount += 1

if __name__ == "__main__":
    asyncio.run(main())
    print(f"处理完成: {count}条记录, 错误{errcount}条, 耗时{time.time() - start_time:.2f}秒")

关键说明

  • 并发控制:通过Semaphore避免请求过载,根据API限流规则调整并发数;
  • 错误处理:任务内部捕获异常+重试,确保单个请求失败不影响全局;
  • 流程中止:通过task.cancel()取消未完成任务,结合as_completed实现实时监控;
  • 性能优化:修复原代码串行执行的问题,保证异步IO的高效性。

内容的提问来源于stack exchange,提问作者BigPig89

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 23:20:45