Python中asyncio模块结合TWS API使用问题及报错排查
TWS API结合asyncio的参数错误修复与事件循环整合方案
一、直接错误原因与修复
出现TypeError: subscribe() missing 1 required positional argument: 'contract'的原因是:
- 定义的
subscribe函数需要contract_id和contract两个参数,但启动线程时仅传入了(app,),参数数量不匹配。 - 另外,
subscribe是异步函数,不能直接用threading.Thread调用,异步函数必须在asyncio事件循环中执行。
如果只是临时修复参数问题,需要传入正确的contract_id和contract,但更合理的方案是放弃单独线程,将所有逻辑整合到同一事件循环中。
二、整合到同一事件循环的实现
TWS的EClient.run()是阻塞方法,我们可以用asyncio.to_thread(Python 3.9+)将其放到后台线程,同时在主事件循环中管理所有异步任务(包括订阅逻辑)。
重构后的完整代码
import asyncio import pandas as pd from ibapi.client import EClient from ibapi.wrapper import EWrapper from ibapi.common import decimalMaxString, floatMaxString master_df = pd.DataFrame() master_sub_req = 0 class TradingApp(EWrapper, EClient): def __init__(self): EClient.__init__(self, self) self.subscription_tasks = [] def position(self, account, contract, position, avgCost): super().position(account, contract, position, avgCost) print("Position.", "Account:", account, "Symbol:", contract.symbol, "SecType:", contract.secType, "Currency:", contract.currency, "Position:", decimalMaxString(position), "Avg cost:", floatMaxString(avgCost)) row = {'Contract ID': contract.conId, 'Contract Symbol': contract.symbol, 'Position': position, 'Avg cost':avgCost} global master_df master_df = pd.concat((master_df, pd.DataFrame([row]))) # 在主事件循环中创建订阅任务 loop = asyncio.get_running_loop() task = loop.create_task(self.subscribe(contract.conId, contract)) self.subscription_tasks.append(task) def positionEnd(self): print('position end') def pnlSingle(self, reqId, pos, dailyPnL, unrealizedPnL, realizedPnL, value): super().pnlSingle(reqId, pos, dailyPnL, unrealizedPnL, realizedPnL, value) print("Daily PnL Single. ReqId:", reqId, "Position:", decimalMaxString(pos), "DailyPnL:", floatMaxString(dailyPnL), "UnrealizedPnL:", floatMaxString(unrealizedPnL), "RealizedPnL:", floatMaxString(realizedPnL), "Value:", floatMaxString(value)) def tickPrice(self, reqId, tickType, price, attrib): super().tickPrice(reqId, tickType, price, attrib) print(attrib) print("TickPrice. TickerId:", reqId, "tickType:", tickType, "Price:", price) # 注意:原代码中subscription_requests_tick_price未定义,需补充或注释 # for key, value in subscription_requests_tick_price: # if master_df.loc[master_df['Contract Symbol'] == key]: # master_df.loc[master_df['Contract Symbol'] == key,'Tick Price Request ID'] = value # master_df.loc[master_df['Contract Symbol'] == key, "Market Price"] = price def error(self, reqId, errorCode, errorMsg): print(f"Error: {errorCode} - {errorMsg}") async def subscribe(self, contract_id, contract): global master_sub_req while True: master_sub_req += 1 self.reqPnLSingle(master_sub_req, "DU7058034", "", contract_id) await asyncio.sleep(2) master_sub_req += 1 self.reqMktData(reqId=master_sub_req, contract=contract, genericTickList="", snapshot=False, regulatorySnapshot=False, mktDataOptions=[]) await asyncio.sleep(2) async def main(): app = TradingApp() app.connect('127.0.0.1', 9225, clientId=4) # 将TWS的run()放入后台线程,不阻塞主事件循环 await asyncio.to_thread(app.run) # 请求持仓数据 app.reqPositions() # 等待所有订阅任务(长期运行可注释此句,或用信号监听实现优雅关闭) await asyncio.gather(*app.subscription_tasks) if __name__ == "__main__": try: asyncio.run(main()) except KeyboardInterrupt: print("程序已终止")
关键改动说明
- 移除独立线程:用
asyncio.to_thread托管app.run(),避免阻塞主事件循环,同时保证TWS的消息处理正常运行。 - 异步订阅整合:将
subscribe改为类的异步方法,在position回调中直接创建异步任务,无需手动启动线程,彻底解决参数传递错误问题。 - 任务管理:用类属性
subscription_tasks管理所有订阅任务,方便后续统一等待或取消。 - 修复错误方法:补全了被注释的
error方法,确保错误信息能正常输出。
三、额外提示
- 原代码中
subscription_requests_tick_price变量未定义,需要补充相关逻辑或注释掉对应代码块,否则会运行报错。 - TWS的回调函数运行在
app.run()所在的线程中,通过asyncio.get_running_loop()可以获取到主事件循环,确保异步任务能正确调度。 - 如需优雅关闭程序,可以监听
KeyboardInterrupt信号,主动调用app.disconnect()并取消所有订阅任务。
内容的提问来源于stack exchange,提问作者Hayden
相关产品推荐
相关产品推荐

