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

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("程序已终止")

关键改动说明

  1. 移除独立线程:用asyncio.to_thread托管app.run(),避免阻塞主事件循环,同时保证TWS的消息处理正常运行。
  2. 异步订阅整合:将subscribe改为类的异步方法,在position回调中直接创建异步任务,无需手动启动线程,彻底解决参数传递错误问题。
  3. 任务管理:用类属性subscription_tasks管理所有订阅任务,方便后续统一等待或取消。
  4. 修复错误方法:补全了被注释的error方法,确保错误信息能正常输出。

三、额外提示

  • 原代码中subscription_requests_tick_price变量未定义,需要补充相关逻辑或注释掉对应代码块,否则会运行报错。
  • TWS的回调函数运行在app.run()所在的线程中,通过asyncio.get_running_loop()可以获取到主事件循环,确保异步任务能正确调度。
  • 如需优雅关闭程序,可以监听KeyboardInterrupt信号,主动调用app.disconnect()并取消所有订阅任务。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 21:06:07