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

如何控制事件循环实现代码并发,避免TradingStream运行阻塞程序?

Alpaca Trading Stream 并发执行解决方案

问题根源

trading_stream.run()是同步阻塞方法,会独占事件循环,导致后续代码无法执行。直接包装成asyncio Task报错,是因为该方法内部已自行启动事件循环,无需手动嵌套。

实现并发的正确方式

通过asyncio显式管理事件循环,将Trading Stream的异步启动逻辑与其他任务并行调度:

import asyncio
from alpaca.trading.stream import TradingStream

trading_stream = TradingStream('api-key', 'secret-key', paper=True)

async def update_handler(data):
    print(data)

trading_stream.subscribe_trade_updates(update_handler)

# 示例:自定义并发任务
async def concurrent_task():
    while True:
        print("并发任务运行中...")
        await asyncio.sleep(3)

async def main():
    # 启动Trading Stream异步任务
    stream_task = asyncio.create_task(trading_stream._run_async())
    # 启动自定义并发任务
    custom_task = asyncio.create_task(concurrent_task())
    
    # 保持事件循环运行,直到所有任务结束
    await asyncio.gather(stream_task, custom_task)

if __name__ == "__main__":
    asyncio.run(main())

核心要点

  • 替换run()为_run_async():后者是Trading Stream的异步入口,可被asyncio事件循环调度,不会阻塞主线程。
  • 用asyncio.create_task()创建并行任务:事件循环会自动在多个协程间切换,保证handler打印数据的同时,其他任务正常执行。
  • asyncio.gather()用于等待所有任务:避免主协程提前退出,维持所有并发任务的持续运行。

常见误区规避

  • 禁止直接调用trading_stream.run():该方法会启动独立的事件循环并阻塞,导致其他代码无法执行。
  • 所有并发逻辑必须封装为async协程:非协程代码无法被asyncio事件循环调度,无法实现真正的并发。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 19:24:26