第 5 章 · 01 BaseGateway 网关契约与 on_xxx 回调


文档摘要

第 5 章 · 01 BaseGateway 网关契约与 onxxx 回调 本节摘要:本节精读 (273 行)——所有交易网关(CTP/IB/富途/币安)都必须继承的契约基类。VeighNa 之所以能同时对接全球 30+ 个交易所,核心就在于这份契约:它规定了网关"必须实现什么方法、必须推送哪些回调、必须遵守什么并发约束"。BaseGateway 本身不含任何网络代码,但它把"事件总线怎么投递"的细节封装成了 8 个 方法,让子类只需关心业务。读完本节,你看到一个普通 Python ABC 怎么靠"双话题推送"模式,把行情/订单流既广播给全局监听者、又精确路由到订阅了具体合约的监听者。 内容来源:原项目源码 (共 273 行),精读并套用体系化模板。

第 5 章 · 01 BaseGateway 网关契约与 on_xxx 回调

本节摘要:本节精读 vnpy/trader/gateway.py(273 行)——所有交易网关(CTP/IB/富途/币安)都必须继承的契约基类。VeighNa 之所以能同时对接全球 30+ 个交易所,核心就在于这份契约:它规定了网关"必须实现什么方法、必须推送哪些回调、必须遵守什么并发约束"。BaseGateway 本身不含任何网络代码,但它把"事件总线怎么投递"的细节封装成了 8 个 on_xxx 方法,让子类只需关心业务。读完本节,你看到一个普通 Python ABC 怎么靠"双话题推送"模式,把行情/订单流既广播给全局监听者、又精确路由到订阅了具体合约的监听者。

内容来源:原项目源码 vnpy/trader/gateway.py(共 273 行),精读并套用体系化模板。

学习目标

阅读完本节,你应当能够:

  1. 说清契约——BaseGateway 的 docstring 三大铁律(线程安全/非阻塞/自动重连)。
  2. 区分类属性(default_name/default_setting/exchanges)与实例属性,知道为什么没有 class_name
  3. 逐行读懂 on_event 通用入口和 8 个 on_xxx 回调的双话题推送模式。
  4. 解释为什么 EVENT_TICK + tick.vt_symbol 能拼成 "eTick.rb2401.SHFE"
  5. 区分 7 个抽象方法(必选)4 个默认方法(可选能力)

一、设计意图:三大铁律

gateway.py:33-70 给 BaseGateway 写了一份罕见的详细 docstring,直接规定了子类该怎么写。先看三段铁律(gateway.py:42-49):

42 * this class should be thread-safe: 43 * all methods should be thread-safe 44 * no mutable shared properties between objects. 46 * all methods should be non-blocked 48 * automatically reconnect if connection lost.

三条铁律翻译过来:

  1. 线程安全:所有方法必须线程安全;对象之间不共享可变属性。原因——CTP/IB 这类 C++ API 都在工作线程里回调,VeighNa 的 EventEngine 又是另一条线程,网关同时被多线程访问,稍微写点可变共享状态就死锁。
  2. 方法非阻塞:connect/subscribe/send_order/query_account 这些方法绝不能卡住调用者(主线程 GUI 或策略线程)。要等服务器响应怎么办?发完请求立刻 return,等回调来了再用 on_xxx 异步推上去。
  3. 自动重连:断线了网关自己重连,不要把异常抛给上层。VeighNa 是面向实盘的,断线恢复是基本生存能力。

docstring 还规定子类必须响应 6 个回调(gateway.py:55-61):on_tick/on_trade/on_order/on_position/on_account/on_contract,以及一个约束——传给 on_xxx 的 data 对象必须是不可变的(用 copy.copy 再传)。这是为了避免同一份对象在引擎里被多个处理器改坏。

💡 核心心法:这份 docstring 就是 BaseGateway 当作"接口规范"在用。VeighNa 没有用 interface/Protocol,而是用 ABC + 长文档+ @abstractmethod 三件套来定契约。子类作者(比如想接一个新交易所的志愿者)读完这一段,就知道什么必须做、什么不能做。这是开源项目降低协作成本的关键工程实践。

二、三个类属性:网关的"身份证"

gateway.py:72-79:

72 # Default name for the gateway. 73 default_name: str = "" 74 75 # Fields required in setting dict for connect function. 76 default_setting: dict[str, str | int | float | bool] = {} 77 78 # Exchanges supported in the gateway. 79 exchanges: list[Exchange] = []

三个类属性(不是实例属性),所有该网关的实例共享同一份:

属性 类型 用途
default_name str 网关名,如 "CTP" / "BINANCE";后续作为 gateway_name 默认值,也用于 UI 菜单显示
default_setting dict 连接配置表单 schema,UI 据此自动生成连接对话框(字段名 → 类型 → 默认值)
exchanges list[Exchange] 该网关支持的交易所枚举列表,订阅/下单前会校验 vt_symbol 的 exchange 是否在内

举两个具体例子帮助理解:

  • CTP 网关的 default_setting 大致是 {"用户名": "", "密码": "", "经纪商代码": "", "交易服务器": "", "行情服务器": "", "产品名称": "", "授权码": "", ...},VeighNa 的 UI 拿到这个 dict,自动画出对应的输入框——这就是为什么不同网关的连接对话框长得不一样,但代码逻辑统一。
  • 币安网关的 exchanges[Exchange.BINANCE, Exchange.BINANCE_FUTURES, ...],你订阅一个不在表里的交易所,网关会拒收。

⚠️ 重要说明:BaseGateway 没有 class_name 属性。网关类名靠 type(gateway).__name__ 获取;gateway_name(实例属性,见下节)才是用户层面的标识,允许同一类网关起多个实例(比如同时连两个 CTP 账户,分别叫 "CTP""CTP_2")。如果你看到老教程里写 gateway.class_name,那是错的——核心包从没这个字段。

三、init:只持有两样东西

gateway.py:81-84:

81 def __init__(self, event_engine: EventEngine, gateway_name: str) -> None: 82 """""" 83 self.event_engine: EventEngine = event_engine 84 self.gateway_name: str = gateway_name

整个 __init__ 只有两个赋值:

  • event_engine:事件总线引用(第 2 章详讲),网关所有回调都靠它向外推。
  • gateway_name:实例名,通常等于 default_name,但用户可以重命名。

注意:这里没有创建任何线程、没有打开任何 socket、没有起任何连接。BaseGateway 是个"纯回调外壳",连接逻辑全部在子类的 connect 里。__init__ 这么轻,是为了让 MainEngine 可以先实例化所有网关占位(UI 显示菜单),用户真正点"连接"时才触发昂贵的网络握手。

四、on_event:通用投递入口

所有 on_xxx 的底层都走这一个方法。gateway.py:86-91:

86 def on_event(self, type: str, data: object = None) -> None: 87 """ 88 General event push. 89 """ 90 event: Event = Event(type, data) 91 self.event_engine.put(event)

三步:把 (type, data) 包成 Event,然后 put 进事件引擎的队列。这就是网关和事件总线唯一的接触点——子类只要调 on_event,事件就进了第 2 章那套 Queue+Thread 分发体系。

为什么不直接 event_engine.put(Event(...))?多一层封装有两个好处:一是统一入口便于以后加 hook(比如日志/统计),二是子类不需要 import Event 类,只调 self.on_event 即可。

五、on_tick:双话题推送模式

这是全篇最关键的设计。gateway.py:93-99:

93 def on_tick(self, tick: TickData) -> None: 94 """ 95 Tick event push. 96 Tick event of a specific vt_symbol is also pushed. 97 """ 98 self.on_event(EVENT_TICK, tick) 99 self.on_event(EVENT_TICK + tick.vt_symbol, tick)

每个 tick 推两次——同一个 tick 对象,投到两个不同的话题:

  1. 全局话题:type = EVENT_TICK = "eTick."(注意末尾带点)。
  2. 专属话题:type = EVENT_TICK + tick.vt_symbol,例如 "eTick." + "rb2401.SHFE" = "eTick.rb2401.SHFE"

为什么这样设计?回顾第 2 章的 register(type, handler):

  • 想监听所有 tick(比如 K 线图表聚合所有合约):engine.register(EVENT_TICK, handler)
  • 只想监听 rb2401(比如策略只跑这一个合约):engine.register("eTick.rb2401.SHFE", handler)

一个事件两种订阅粒度,靠 type 字符串前缀分发。EVENT_TICK 末尾的 "." 是关键的命名空间分隔符,保证 "eTick." 不会误匹配 "eTick.rb2401.SHFE" 的订阅者——后者完整字符串只在专属话题里出现。

💡 核心心法:这是"主题订阅 + 通配订阅"的轻量实现。VeighNa 不引入 MQTT/Redis 那套完整的 pub-sub 中间件,只靠 Python 字符串拼接 + defaultdict,就把"按合约路由"做出来了。代价是每个 tick 在队列里放两次——但 Queue.put 极快(纳秒级),换来的是订阅逻辑的极简。

六、其他 on_xxx:同构套路

剩下几个回调全部是 on_tick 的同构变体,只是专属话题的 key 不同:

回调 全局话题 专属话题 key 行号
on_tick EVENT_TICK tick.vt_symbol 93-99
on_trade EVENT_TRADE trade.vt_symbol 101-107
on_order EVENT_ORDER order.vt_orderid 109-115
on_position EVENT_POSITION position.vt_symbol 117-123
on_account EVENT_ACCOUNT account.vt_accountid 125-131
on_quote EVENT_QUOTE quote.vt_symbol 133-139

注意 key 的选择有讲究:

  • 行情/持仓/报价vt_symbol(合约维度)——策略只关心"自己订阅的合约"。
  • 订单vt_orderid——下单后能精确跟踪"这一笔"的状态变化。
  • 账户vt_accountid——多账户场景下分账户监听。

最后两个回调只推全局话题,没有专属话题:

141 def on_log(self, log: LogData) -> None: 145 self.on_event(EVENT_LOG, log) 147 def on_contract(self, contract: ContractData) -> None: 151 self.on_event(EVENT_CONTRACT, contract)

原因——log 和 contract 没有"按合约订阅"的需求。日志是给人看的(全局显示在 UI 日志面板);合约是连接时一次性广播的元数据(订阅列表),没必要做合约级路由。注意 EVENT_LOG = "eLog" 末尾不带点(其他都带),也佐证了"它只有全局话题"——没有拼接专属话题的需要。

辅助方法 write_log(gateway.py:153-158)是 on_log 的语法糖,把字符串包成 LogData 再调 on_log:

153 def write_log(self, msg: str) -> None: 157 log: LogData = LogData(msg=msg, gateway_name=self.gateway_name) 158 self.on_log(log)

子类网关里写 self.write_log("登录成功") 即可,不用手动构造 LogData。

七、7 个抽象方法:必选能力

@abstractmethod 装饰的方法子类必须实现,否则实例化时报 TypeError。一共 7 个:

方法 行号 返回 作用
connect(setting) 160-180 None 连接服务器 + 登录 + 首次查询合约/账户/持仓/订单/成交
close() 182-187 None 关闭连接
subscribe(req) 189-194 None 订阅 tick 行情(SubscribeRequest)
send_order(req) 196-212 str 发新单,返回 vt_orderid
cancel_order(req) 214-221 None 撤单(CancelRequest)
query_account() 248-253 None 查询账户资金
query_position() 255-260 None 查询持仓

重点看 connectsend_order 的 docstring,它们写明了子类的实现步骤。

connect 的规定(gateway.py:164-174):

165 to implement this method, you must: 167 * log connected if all necessary connection is established 168 * do the following query and response corresponding on_xxxx and write_log 169 * contracts : on_contract 170 * account asset : on_account 171 * account holding: on_position 172 * orders of account: on_order 173 * trades of account: on_trade 174 * if any of query above is failed, write log.

也就是说,connect 不只是建立 socket,还包括登录后的"首次全量同步"——拉合约列表、账户资金、持仓、当日委托、当日成交,每一种都用对应的 on_xxx 推上去。失败的话不能抛异常,要 write_log 告诉用户。

send_order 的规定(gateway.py:201-211):

202 * create an OrderData from req using OrderRequest.create_order_data 203 * assign a unique(gateway instance scope) id to OrderData.orderid 204 * send request to server 205 * if request is sent, OrderData.status should be set to Status.SUBMITTING 206 * if request is failed to sent, OrderData.status should be set to Status.REJECTED 207 * response on_order: 208 * return vt_orderid 210 :return str vt_orderid for created OrderData

五步:① 用 req.create_order_data() 造 OrderData;② 分配网关实例范围内唯一的 orderid;③ 发服务器(成功置 SUBMITTING,失败置 REJECTED);④ 调 on_order 把状态推给上层;⑤ 返回 vt_orderid。注意"返回值"和"on_order 回调"是同一笔订单的两种传递方式——返回值给同步调用者(下单的瞬间拿到 id),on_order 给异步订阅者(后续状态变化)。

八、4 个默认实现:可选能力

剩下 4 个方法没有 @abstractmethod,父类给了空/默认实现,子类按需重写:

223 def send_quote(self, req: QuoteRequest) -> str: 238 return "" 240 def cancel_quote(self, req: CancelRequest) -> None: 246 return 262 def query_history(self, req: HistoryRequest) -> list[BarData]: 266 return [] 268 def get_default_setting(self) -> dict[str, str | int | float | bool]: 272 return self.default_setting

它们对应"非通用能力":

  • send_quote / cancel_quote:做市报价(双边挂单),只对支持做市的交易所(期权/期货做市)有意义,A 股股票网关用不到。
  • query_history:网关内置的历史数据下载——少数交易所(如币安)直接提供 REST 接口拉 K 线,CTP 没有这个能力。
  • get_default_setting:返回类属性 default_setting,主要给 UI 用。

⚠️ 重要说明:这 4 个方法的默认实现不是简单的 pass——send_quote 返回空字符串,query_history 返回空列表,BaseDatafeed 风格的"失败静默"。这样上层调用者不用 isinstance 判断,直接 bars = gateway.query_history(req),没实现就拿到空结果。这是 Python"鸭子类型 + 默认实现"的典型用法,与第 2 节要讲的 BaseDatafeed 思路一致。

本节要点回顾

  1. 三大铁律:线程安全(无可变共享状态)、方法非阻塞、自动重连——写在 docstring 里作为契约。
  2. 三个类属性:default_name/default_setting(UI 自动生成连接表单)/exchanges;没有 class_name,实例名是 gateway_name
  3. on_event 是唯一投递入口:包 Event → put 进队列;所有 on_xxx 都走它。
  4. 双话题推送:tick/trade/order/position/account/quote 既推全局又推专属话题;key 各有讲究(合约/订单/账户维度)。EVENT_TICK = "eTick." 末尾带点是命名空间分隔符。
  5. on_log/on_contract 只推全局:无按合约订阅需求;EVENT_LOG = "eLog" 不带点。
  6. 7 个抽象方法:connect/close/subscribe/send_order/cancel_order/query_account/query_position 必须实现;send_order 返回 vt_orderid 同时调 on_order。
  7. 4 个默认实现:send_quote/cancel_quote/query_history/get_default_setting 是"可选能力",空返回静默失败。

下一节,我们离开网关,看 VeighNa 的另一个抽象层级——BaseApp 应用契约(cta_strategy/spread_trading 等业务模块的基类),以及和网关"用户手动装配"完全不同的另一套发现机制:数据库 BaseDatabase 和数据服务 BaseDatafeed 靠配置自动 import_module 发现


发布者: 作者: 灏天文库 转发
评论区 (0)
U