基于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
相关产品推荐
相关产品推荐

