如何在持续运行的Python异步循环中并行执行代码且不阻塞主循环
实现方案
核心使用asyncio.create_task()将交易函数作为异步任务提交到事件循环调度,不会阻塞当前websocket行情接收主逻辑。
注意:run_a_trade_function需要先定义为异步协程函数,如果是同步耗时逻辑需要单独处理,避免阻塞事件循环
修改后完整代码
import asyncio from binance import AsyncClient, BinanceSocketManager # 异步交易函数 async def run_a_trade_function(res): # 你的交易逻辑写在此处 try: print(f"触发交易逻辑,对应行情流:{res.get('stream')}") # 示例:模拟交易逻辑异步耗时 await asyncio.sleep(0.3) except Exception as e: print(f"交易任务执行异常: {e}") async def main(): # 替换为你自己的api配置和订阅流 api_key = "你的币安API_KEY" api_secret = "你的币安API_SECRET" multi = ["btcusdt@kline_1m", "ethusdt@trade"] client = await AsyncClient.create(api_key, api_secret) bm = BinanceSocketManager(client) ms = bm.multiplex_socket(multi) async with ms as tscm: while True: res = await tscm.recv() if res: # 提交任务到事件循环,无需等待执行完成,主循环直接继续接收下一条行情 asyncio.create_task(run_a_trade_function(res)) if __name__ == "__main__": loop = asyncio.get_event_loop() loop.run_until_complete(main())
其他注意事项
- 如果你现有的交易逻辑是同步函数,不要直接放到异步协程里运行,会阻塞整个事件循环。改用
asyncio.to_thread()包裹调度同步逻辑:# 原有同步交易函数 def sync_trade_logic(res): # 同步计算、请求等逻辑 print("同步交易逻辑执行") async def run_a_trade_function(res): # 把同步函数丢到单独线程运行,不阻塞事件循环 await asyncio.to_thread(sync_trade_logic, res) - 高并发场景下建议对提交的任务做数量限制,避免任务堆积占用过多内存,可以用
asyncio.Semaphore做并发数控制。 - 交易函数内部一定要加异常捕获,避免单个交易任务抛出未处理的异常导致整个程序退出。
内容的提问来源于stack exchange,提问作者user17546065
相关产品推荐
相关产品推荐

