如何控制事件循环实现代码并发,避免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
相关产品推荐
相关产品推荐

