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

如何用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())
解决方案

核心实现思路

  1. 保留asyncio.create_task()实现非阻塞消费,通过主动让出事件循环控制权确保POST请求立即发起。
  2. 补充异步任务的异常处理逻辑,避免任务静默失败。
  3. 规范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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 15:51:05