基于alpaca-py的交易Bot中WebSocket协程未触发问题求助
使用alpaca-py开发交易Bot,目标是买入股票后通过WebSocket获取订单成交状态,再卖出对应份额。未集成WebSocket前买卖功能正常,但集成后WebSocket协程无法触发,trade_update_handler从未被调用。
已尝试的方案:
- 调整协程中await语句位置
- 修改方法位置并移除部分async标识
- 查阅Alpaca文档(但alpaca-py 2023年发布,多数文档过时)
- 阅读TradingStream源码确认实现无误
- 修改asyncio.gather调用方式,结果一致
- 添加日志确认
trade_update_handler未被调用 - 使用
run()替代_run_forever(),导致Webhook报错
使用Django BaseCommand运行Bot,认为Django不是问题根源。相关代码如下:
class TradeMaker(): def __init__(self, **kwargs): self.paper_bool = kwargs.get('paper_bool', True) self.random_bool = kwargs.get('random', True) self.symbol_or_symbols = kwargs.get('symbol_or_symbols', 'AAPL') self.amount = kwargs.get('amount', 40000) self.seconds_between = kwargs.get('seconds_between', 4) self.log = kwargs.get('log') self.trading_client, self.trading_stream, self.account = self.open_client() self.trade_update_info = None self.order_filled = False self.shares_bought = 0 self.current_symbol = None def open_client(self): trading_client = TradingClient(ALPACA_ID, ALPACA_KEY, paper=self.paper_bool) trading_stream = TradingStream(ALPACA_ID, ALPACA_KEY, paper=self.paper_bool) try: account = trading_client.get_account() except Exception as e: logger.error(f"Exception in login: {e}") return trading_client, trading_stream, account async def trade_update_handler(self, data): logger.info('Trade Update called') print("Trade Update:", data) if data.event == TradeEvent.FILL: if data.order.side == OrderSide.BUY: self.order_filled = True self.shares_bought = data.order.filled_qty self.current_symbol = data.order.symbol async def run_stream(self): logger.info('Subscribing to trade updates') self.trading_stream.subscribe_trade_updates(self.trade_update_handler) logger.info('Preparing stream') await self.trading_stream._run_forever() async def stop_stream(self): logger.info('Stopping stream') trading_stream = TradingStream(ALPACA_ID, ALPACA_KEY, paper=self.paper_bool) await trading_stream.stop() def get_symbol(self): if self.random_bool: symbol = random.choice(self.symbol_or_symbols) return symbol else: symbol = self.symbol_or_symbols return symbol def buy(self): symbol = self.get_symbol() market_order_data = MarketOrderRequest( symbol=symbol, qty=1, side=OrderSide.BUY, time_in_force=TimeInForce.DAY ) try: market_order_buy = self.trading_client.submit_order( order_data=market_order_data ) except Exception as e: logger.error(f"Failed to buy {symbol}: {e}") return None return symbol, market_order_buy def sell(self, symbol): symbol = symbol shares = self.shares_bought market_order_data = MarketOrderRequest( symbol=symbol, qty=250, side=OrderSide.SELL, time_in_force=TimeInForce.DAY ) try: market_order_sell = self.trading_client.submit_order( order_data=market_order_data ) except Exception as e: logger.error(f"Failed to sell {symbol}: {e}") return None return market_order_sell async def make_trades(self): market_close = datetime.datetime.now().replace(hour=14, minute=0, second=0, microsecond=0) while datetime.datetime.now() < market_close: seconds = self.seconds_between try: symbol, market_order_buy = self.buy() print(f"Bought {symbol}: {market_order_buy}") except Exception as e: logger.error(f"Failed to buy during trade: {e}") return None while not self.order_filled: logger.info('Waiting for order status update') await asyncio.sleep(1) sleep(seconds) try: market_order_sell = self.sell(symbol=symbol) print(f"Sold {self.current_symbol}: {market_order_sell}") except Exception as e: logger.error(f"Failed to sell during trade: {e}") return None self.order_filled = False self.shares_bought = 0 sleep(seconds) print('Market closed, shutting down.') class Command(BaseCommand): help = """This bot trades the target stock. If you want it to choose randomly, pass it a list and set the variable random=True """ model = None def add_arguments(self, parser): parser.add_argument( '--paper', type=bool, help='Set false to live trade.', default=True ) parser.add_argument( '--folder', type=str, help='source folder for files', default='' ) parser.add_argument( '--symbol', type=str, help='target symbol, or list of symbols', default='AAPL' ) parser.add_argument( '--random', type=bool, help="Set to true if passing a list of symbols to choose randomly from.", default=False ) parser.add_argument( '--tradevalue', type=int, help="The amount the bot should trade. e.g. $40000", default=40000 ) parser.add_argument( '--seconds', type=int, help="The number of seconds the bot should wait between each trade.", default=4 ) def handle(self, **options): paper_bool = options['paper'] random_bool = options['random'] symbol_or_symbols = options['symbol'] amount = options['tradevalue'] seconds_between = options['seconds'] log = options['folder'] tm = TradeMaker( paper_bool = paper_bool, random = random_bool, symbol_or_symbols = symbol_or_symbols, amount = amount, seconds_between = seconds_between, log = log ) loop = asyncio.get_event_loop() try: loop.run_until_complete(asyncio.gather( tm.run_stream(), tm.make_trades() )) except KeyboardInterrupt: tm.stop_stream() print("Stopped with Interrupt") finally: tm.stop_stream() loop.close()
运行命令后终端输出(敏感信息已脱敏):
python manage.py trade_maker_v5 2023-03-30 11:51:48,342 - INFO - Subscribing to trade updates 2023-03-30 11:51:48,342 - INFO - Preparing stream Bought AAPL: id=UUID('foo') client_order_id='bar' created_at=datetime.datetime(2023, 3, 30, 18, 51, 49, 995853, tzinfo=datetime.timezone.utc) updated_at=datetime.datetime(2023, 3, 30, 18, 51, 49, 995921, tzinfo=datetime.timezone.utc) submitted_at=datetime.datetime(2023, 3, 30, 18, 51, 49, 994623, tzinfo=datetime.timezone.utc) filled_at=None expired_at=None canceled_at=None failed_at=None replaced_at=None replaced_by=None replaces=None asset_id=UUID('foo') symbol='AAPL' asset_class=<AssetClass.US_EQUITY: 'us_equity'> notional=None qty='1' filled_qty='0' filled_avg_price=None order_class=<OrderClass.SIMPLE: 'simple'> order_type=<OrderType.MARKET: 'market'> type=<OrderType.MARKET: 'market'> side=<OrderSide.BUY: 'buy'> time_in_force=<TimeInForce.DAY: 'day'> limit_price=None stop_price=None status=<OrderStatus.PENDING_NEW: 'pending_new'> extended_hours=False legs=None trail_percent=None trail_price=None hwm=None 2023-03-30 11:51:48,480 - INFO - Waiting for order status update 2023-03-30 11:51:49,493 - INFO - Waiting for order status update 2023-03-30 11:51:50,500 - INFO - Waiting for order status update
但单独运行以下代码时,能正常获取Bot交易的所有数据:
from alpaca.trading.stream import TradingStream trading_stream = TradingStream(ALPACA_ID, ALPACA_KEY, paper=True) async def update_handler(data): print(data) trading_stream.subscribe_trade_updates(update_handler) trading_stream.run()
疑问:为何单独运行WebSocket代码正常,集成到协程中却无法触发?
1. 阻塞sleep占用事件循环
make_trades方法中使用了标准库的sleep(seconds),这是阻塞调用,会独占整个事件循环,导致WebSocket协程无法获得执行时间。
修复:将所有sleep()替换为await asyncio.sleep(seconds),确保导入asyncio。
2. stop_stream方法操作错误实例
当前stop_stream中重新创建了新的TradingStream对象,而非使用初始化时的self.trading_stream,导致无法正确停止已运行的连接。
修复:
async def stop_stream(self): logger.info('Stopping stream') await self.trading_stream.stop()
3. 使用私有方法导致异常
_run_forever是TradingStream的内部私有方法,官方推荐使用公开的run()方法。之前使用run()报错是因为事件循环被阻塞sleep占用,修复sleep后即可正常使用。
修复修改run_stream方法:
async def run_stream(self): logger.info('Subscribing to trade updates') self.trading_stream.subscribe_trade_updates(self.trade_update_handler) logger.info('Preparing stream') await self.trading_stream.run()
4. 可选优化:线程安全锁
多协程共享self.order_filled等状态变量时,建议使用asyncio.Lock避免竞态条件:
def __init__(self, **kwargs): # 其他初始化代码 self.lock = asyncio.Lock() async def trade_update_handler(self, data): logger.info('Trade Update called') print("Trade Update:", data) if data.event == TradeEvent.FILL: if data.order.side == OrderSide.BUY: async with self.lock: self.order_filled = True self.shares_bought = data.order.filled_qty self.current_symbol = data.order.symbol async def make_trades(self): # 其他代码 while True: logger.info('Waiting for order status update') await asyncio.sleep(1) async with self.lock: if self.order_filled: break # 其他代码
验证
修改后重新运行Bot,WebSocket协程会正常接收订单更新,trade_update_handler会被触发,order_filled状态会被正确更新,买卖流程可正常衔接。
内容的提问来源于stack exchange,提问作者Lewis F

