第 5 章 · 01 BaseGateway 网关契约与 onxxx 回调 本节摘要:本节精读 (273 行)——所有交易网关(CTP/IB/富途/币安)都必须继承的契约基类。VeighNa 之所以能同时对接全球 30+ 个交易所,核心就在于这份契约:它规定了网关"必须实现什么方法、必须推送哪些回调、必须遵守什么并发约束"。BaseGateway 本身不含任何网络代码,但它把"事件总线怎么投递"的细节封装成了 8 个 方法,让子类只需关心业务。读完本节,你看到一个普通 Python ABC 怎么靠"双话题推送"模式,把行情/订单流既广播给全局监听者、又精确路由到订阅了具体合约的监听者。 内容来源:原项目源码 (共 273 行),精读并套用体系化模板。
本节摘要:本节精读
vnpy/trader/gateway.py(273 行)——所有交易网关(CTP/IB/富途/币安)都必须继承的契约基类。VeighNa 之所以能同时对接全球 30+ 个交易所,核心就在于这份契约:它规定了网关"必须实现什么方法、必须推送哪些回调、必须遵守什么并发约束"。BaseGateway 本身不含任何网络代码,但它把"事件总线怎么投递"的细节封装成了 8 个on_xxx方法,让子类只需关心业务。读完本节,你看到一个普通 Python ABC 怎么靠"双话题推送"模式,把行情/订单流既广播给全局监听者、又精确路由到订阅了具体合约的监听者。
内容来源:原项目源码
vnpy/trader/gateway.py(共 273 行),精读并套用体系化模板。
阅读完本节,你应当能够:
on_xxx 回调的双话题推送模式。EVENT_TICK + tick.vt_symbol 能拼成 "eTick.rb2401.SHFE"。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.
三条铁律翻译过来:
connect/subscribe/send_order/query_account 这些方法绝不能卡住调用者(主线程 GUI 或策略线程)。要等服务器响应怎么办?发完请求立刻 return,等回调来了再用 on_xxx 异步推上去。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 是否在内 |
举两个具体例子帮助理解:
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,那是错的——核心包从没这个字段。
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_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 即可。
这是全篇最关键的设计。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 对象,投到两个不同的话题:
type = EVENT_TICK = "eTick."(注意末尾带点)。type = EVENT_TICK + tick.vt_symbol,例如 "eTick." + "rb2401.SHFE" = "eTick.rb2401.SHFE"。为什么这样设计?回顾第 2 章的 register(type, handler):
engine.register(EVENT_TICK, handler)。engine.register("eTick.rb2401.SHFE", handler)。一个事件两种订阅粒度,靠 type 字符串前缀分发。EVENT_TICK 末尾的 "." 是关键的命名空间分隔符,保证 "eTick." 不会误匹配 "eTick.rb2401.SHFE" 的订阅者——后者完整字符串只在专属话题里出现。
💡 核心心法:这是"主题订阅 + 通配订阅"的轻量实现。VeighNa 不引入 MQTT/Redis 那套完整的 pub-sub 中间件,只靠 Python 字符串拼接 + defaultdict,就把"按合约路由"做出来了。代价是每个 tick 在队列里放两次——但 Queue.put 极快(纳秒级),换来的是订阅逻辑的极简。
剩下几个回调全部是 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。
@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 | 查询持仓 |
重点看 connect 和 send_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 个方法没有 @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
它们对应"非通用能力":
default_setting,主要给 UI 用。⚠️ 重要说明:这 4 个方法的默认实现不是简单的
pass——send_quote返回空字符串,query_history返回空列表,BaseDatafeed风格的"失败静默"。这样上层调用者不用 isinstance 判断,直接bars = gateway.query_history(req),没实现就拿到空结果。这是 Python"鸭子类型 + 默认实现"的典型用法,与第 2 节要讲的 BaseDatafeed 思路一致。
default_name/default_setting(UI 自动生成连接表单)/exchanges;没有 class_name,实例名是 gateway_name。EVENT_TICK = "eTick." 末尾带点是命名空间分隔符。EVENT_LOG = "eLog" 不带点。下一节,我们离开网关,看 VeighNa 的另一个抽象层级——BaseApp 应用契约(cta_strategy/spread_trading 等业务模块的基类),以及和网关"用户手动装配"完全不同的另一套发现机制:数据库 BaseDatabase 和数据服务 BaseDatafeed 靠配置自动 import_module 发现。