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

如何在asyncio循环中并行运行函数处理WebSocket数据包?

解决方案:Asyncio环境下并行处理WebSocket商品条目

你的核心需求是在asyncio驱动的WebSocket监听逻辑中,并行处理每个商品条目,避免串行执行导致错失抢购机会。针对你遇到的问题,给出两种可行方案:

方案1:将同步阻塞函数包装为异步任务(适用于无法修改现有同步代码)

如果你的add_to_cart是基于同步爬虫(比如requests库)的阻塞逻辑,可通过asyncio.to_thread(Python3.9+)或run_in_executor将其包装为异步任务,再通过asyncio.gather并发执行。

修正后的完整代码

import asyncio

# 假设全局定义的价格阈值
median_price = 4

def add_to_cart(item):
    # 此处为你的同步爬虫逻辑,比如发起HTTP请求添加购物车
    print(f"已将 {item['item']} 加入购物车")

def check_listing(item):
    # 修正原代码的索引错误,直接传入单个商品条目而非列表+索引
    if item['price'] <= median_price:
        add_to_cart(item)

@client.listen("saleFeed")
async def on_sale_feed(data):
    if data['eventType'] == 'listed':
        print("收到数据")
        # 生成所有异步任务
        tasks = [asyncio.to_thread(check_listing, item) for item in data['sales']]
        # 并发执行所有任务
        await asyncio.gather(*tasks)

兼容Python3.8及以下版本的写法

如果你的Python版本低于3.9,改用run_in_executor:

@client.listen("saleFeed")
async def on_sale_feed(data):
    if data['eventType'] == 'listed':
        print("收到数据")
        loop = asyncio.get_running_loop()
        tasks = [loop.run_in_executor(None, check_listing, item) for item in data['sales']]
        await asyncio.gather(*tasks)

方案2:将爬虫逻辑改为异步实现(推荐,效率更高)

如果可以重构add_to_cart的爬虫逻辑,使用异步HTTP库(如aiohttp)替代同步库,直接用asyncio原生异步机制处理,无需线程池开销。

示例代码

import asyncio
import aiohttp

median_price = 4

async def async_add_to_cart(item):
    # 用aiohttp发起异步请求
    async with aiohttp.ClientSession() as session:
        # 替换为实际的添加购物车API地址和参数
        await session.post(
            "https://example.com/api/add-cart",
            json={"item_id": item['item'], "price": item['price']}
        )
        print(f"已将 {item['item']} 加入购物车")

async def async_check_listing(item):
    if item['price'] <= median_price:
        await async_add_to_cart(item)

@client.listen("saleFeed")
async def on_sale_feed(data):
    if data['eventType'] == 'listed':
        print("收到数据")
        tasks = [async_check_listing(item) for item in data['sales']]
        await asyncio.gather(*tasks)

常见问题说明

你之前尝试run_in_executor未成功,大概率是以下原因:

  • 未正确获取当前asyncio事件循环
  • 未用await asyncio.gather()等待所有任务完成
  • 原代码存在索引错误(data[i]应改为data[position],更优的方式是直接传入单个商品条目)

asyncio环境下不推荐用多进程,因为多进程会带来进程间通信的额外开销,而你的场景是IO密集型(网页请求),线程池或异步IO完全可以满足需求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 03:43:39