Python对接盈透TWS API多标的实时数据多接口整合问题求助
优化方案
核心改造点
- 移除冗余的reqId到标的代码的判断逻辑,直接用索引取值即可
- 在IBApi类中新增实例属性存储所有标的的最新行情数据,所有tick回调统一更新该存储结构
- 移除tickPrice中的
time.sleep(0.5),该行代码会阻塞IB的消息接收线程,导致行情延迟甚至丢失 - 把静态方法
myData改为实例方法,每次更新完单个字段后直接触发指标计算和预警校验
修改后完整代码
EWrapper类代码
import threading import pandas as pd from ibapi.wrapper import EWrapper from ibapi.client import EClient from ibapi.contract import Contract from ibapi.common import TickTypeEnum import time class IBApi(EWrapper, EClient): def __init__(self): self.syms = ['EN', 'DG', 'AI', 'ORA', 'RI', 'ENGI', 'AC', 'VIV', 'KER', 'CA', 'BN', 'WLN', 'OR', 'VIE', 'LR', 'ML', 'SGO', 'CAP', 'MC', 'ACA', 'ATO', 'UG', 'SU', 'HO', 'BNP', 'GLE', 'SAN', 'SW', 'AIR', 'TTE'] EClient.__init__(self, self) # 新增行情存储结构,key为标的代码,value为各字段最新值 self.market_data = {sym: {"HIGH": None, "LOW": None, "CLOSE": None, "VWAP": None} for sym in self.syms} # 可选:加线程锁保证数据安全,避免多线程更新冲突 self.lock = threading.Lock() # 接收实时数据 def tickString(self, reqId, tickType, value): super().tickString(reqId, tickType, value) try: # 一行完成reqId到标的代码映射,替换原来30个if判断 sym = self.syms[reqId] if tickType == 48: rtVolume = value.split(";") vwap = float(rtVolume[4]) self.update_market_data(sym, "VWAP", vwap) except Exception as e: print(e) def tickPrice(self, reqId, tickType, price, attrib): super().tickPrice(reqId, tickType, price, attrib) try: sym = self.syms[reqId] tick_type_str = TickTypeEnum.to_str(tickType) # 映射IB返回的tick类型到我们需要的字段 field_map = { "HIGH": "HIGH", "LOW": "LOW", "CLOSE": "CLOSE" # 有其他需要的字段可以在这里扩展映射 } if tick_type_str in field_map: self.update_market_data(sym, field_map[tick_type_str], price) except Exception as e: print(e) def update_market_data(self, sym, field, value): with self.lock: # 更新对应标的的对应字段 self.market_data[sym][field] = value # 打印更新后的数据,方便调试 print(f"{sym} {field} 更新为: {value}, 当前完整数据: {self.market_data[sym]}") # 这里直接调用预警逻辑 self.check_alert(sym) def check_alert(self, sym): # 先判断该标的所有需要的字段是否都已经拿到有效值 current_data = self.market_data[sym] if all(v is not None for v in current_data.values()): # 把当前数据转成DataFrame,用你之前单标的的逻辑计算指标即可 df = pd.DataFrame([current_data]) # 这里写你自己的预警条件,示例:当前价格低于VWAP时触发提醒 if df['CLOSE'].iloc[0] < df['VWAP'].iloc[0]: print(f"⚠️ 预警触发:{sym} 最新价 {df['CLOSE'].iloc[0]} 低于VWAP {df['VWAP'].iloc[0]}") def error(self, id, errorCode, errorMsg): print(f"错误 {errorCode}: {errorMsg}")
App类代码
class App: ib = None def __init__(self): self.ib = IBApi() self.ib.connect("127.0.0.1", 7496, 88) ib_thread = threading.Thread(target=self.run_loop, daemon=True) ib_thread.start() time.sleep(0.5) for req_num, sym in enumerate(self.ib.syms): self.marketData(req_num, self.symbolsForData(sym)) def symbolsForData(self, mySymbol, sec_type='STK', currency='EUR', exchange='SBF'): contract1 = Contract() contract1.symbol = mySymbol.upper() contract1.secType = sec_type contract1.currency = currency contract1.exchange = exchange return contract1 def marketData(self, req_num, contract1): self.ib.reqMktData(reqId=req_num, contract=contract1, genericTickList='233', snapshot=False, regulatorySnapshot=False, mktDataOptions=[]) def run_loop(self): self.ib.run() # 启动程序 if __name__ == "__main__": App()
扩展说明
如果需要存储历史行情计算时序指标,可以把market_data里每个字段的存储结构改成列表,每次更新的时候append新值,再定时截取最近N条数据计算均线之类的指标即可。
内容的提问来源于stack exchange,提问作者Mario Chacon
相关产品推荐
相关产品推荐

