如何用asyncio和aiohttp立即发送异步POST并异步处理响应
问题描述
基于asyncio和aiohttp实现生产者-消费者模型:生产者持续向队列推送数值,消费者处理队列数据时,若遇到偶数则发送HTTP POST请求到指定服务器。现有两种实现方案存在缺陷:
- 方案1:直接
await session.post(),功能正常但会阻塞消费者协程,必须等待POST请求响应返回后才能继续处理下一个队列元素(生产者不受影响)。 - 方案2:通过
asyncio.create_task()创建独立任务处理POST请求,不会阻塞消费者,但无法保证POST请求被立即发起——任务仅被调度到事件循环,执行时机由事件循环调度策略决定。
需求目标:立即发起POST请求,同时异步处理服务器响应验证,且不阻塞消费者继续处理队列中的后续数据。
原始代码如下:
import asyncio import aiohttp import random async def my_send_post(session, url, headers, params): resp = await session.post(url=url, headers=headers, json=params) resp = await resp.json() print(resp) return async def producer(queue): x = 0 while True: print(f'producing {x}') await asyncio.sleep(1) await queue.put(x) x += 1 async def consumer(queue): session = aiohttp.ClientSession() url = 'https://api.testurl.com/' params = {} headers = request_headers(params) # 外部库生成请求头 while True: item = await queue.get() if item %2 == 0: # METHOD 1 resp = await session.post(url=url, headers=headers, json=params) resp = await resp.json() print(resp) # METHOD 2 asyncio.create_task(my_send_post(session=session, url=url, headers=headers, params=params)) else: print(f"consuming odd number{item}") async def main(): queue = asyncio.Queue() await asyncio.gather(producer(queue), consumer(queue)) if __name__ == "__main__": asyncio.run(main())
解决方案
核心实现思路
- 保留
asyncio.create_task()实现非阻塞消费,通过主动让出事件循环控制权确保POST请求立即发起。 - 补充异步任务的异常处理逻辑,避免任务静默失败。
- 规范
aiohttp.ClientSession的生命周期管理,防止连接资源泄漏。
优化后的代码
import asyncio import aiohttp import random async def my_send_post(session, url, headers, params, item): try: # 发起POST请求(await会立即触发请求发送) resp = await session.post(url=url, headers=headers, json=params) resp.raise_for_status() # 捕获4xx/5xx等HTTP错误状态码 resp_data = await resp.json() print(f"Item {item} POST响应: {resp_data}") except Exception as e: print(f"Item {item} POST请求失败: {str(e)}") async def producer(queue): x = 0 while True: print(f'producing {x}') await asyncio.sleep(1) await queue.put(x) x += 1 async def consumer(queue): # 使用async with自动管理ClientSession的创建与关闭 async with aiohttp.ClientSession() as session: url = 'https://api.testurl.com/' params = {} headers = request_headers(params) # 假设该函数返回合法请求头 while True: item = await queue.get() if item % 2 == 0: # 创建异步任务处理POST请求 asyncio.create_task( my_send_post(session=session, url=url, headers=headers, params=params, item=item) ) # 主动让出事件循环控制权,让新任务立即执行(发起POST请求) await asyncio.sleep(0) else: print(f"consuming odd number {item}") # 标记队列任务处理完成(用于队列的join()操作,可选) queue.task_done() async def main(): queue = asyncio.Queue() await asyncio.gather(producer(queue), consumer(queue)) if __name__ == "__main__": asyncio.run(main())
关键细节说明
- 立即发起请求:
await asyncio.sleep(0)会让当前消费者协程主动放弃事件循环的执行权,事件循环会优先调度刚创建的POST任务,从而立即触发请求发送。这是协作式异步模型中强制任务调度的常用技巧。 - 非阻塞消费:消费者创建异步任务后无需等待响应,直接继续处理下一个队列元素,完全不影响消费流程的效率。
- 异常处理:在POST请求处理函数中添加完整的异常捕获,涵盖网络错误、HTTP状态码错误、响应解析错误等场景,避免异步任务静默失败导致问题难以排查。
- 资源管理:通过
async with管理ClientSession,确保会话在消费者退出时自动关闭,避免TCP连接泄漏。
若需要控制并发请求数量(防止短时间内发起过多请求触发限流),可以在my_send_post中引入asyncio.Semaphore进行并发控制:
# 在consumer中初始化信号量 semaphore = asyncio.Semaphore(5) # 限制最多同时发起5个POST请求 # 在my_send_post中使用 async with semaphore: resp = await session.post(url=url, headers=headers, json=params) # ...后续逻辑
内容的提问来源于stack exchange,提问作者Airplaaaaaane
相关产品推荐
相关产品推荐

