第 2 章 · 02 register/put/process 三大核心机制 本节摘要:本节继续精读 ,聚焦事件总线的三大核心机制: (投递事件)、 (分发事件)、 / (订阅/退订)。你会看到 VeighNa 的发布订阅设计有一个独特之处——它同时支持"专属处理器"(按 type 订阅,绝大多数场景)和"通用处理器"(监听全量事件,用于日志/监控)。这种双轨制让"全局观察者"组件不必为每个 type 重复注册一次。读完本节,你完整掌握了 EventEngine 的工作原理,为理解后续所有"谁在 put、谁在 register"打下基础。 内容来源:原项目源码 ,精读并套用体系化模板。 学习目标 阅读完本节,你应当能够: 说清 方法的线程安全性和它的唯一性(投递事件的唯一入口)。
本节摘要:本节继续精读
engine.py,聚焦事件总线的三大核心机制:put(投递事件)、_process(分发事件)、register/unregister(订阅/退订)。你会看到 VeighNa 的发布订阅设计有一个独特之处——它同时支持"专属处理器"(按 type 订阅,绝大多数场景)和"通用处理器"(监听全量事件,用于日志/监控)。这种双轨制让"全局观察者"组件不必为每个 type 重复注册一次。读完本节,你完整掌握了 EventEngine 的工作原理,为理解后续所有"谁在 put、谁在 register"打下基础。
内容来源:原项目源码
vnpy/event/engine.py,精读并套用体系化模板。
阅读完本节,你应当能够:
put 方法的线程安全性和它的唯一性(投递事件的唯一入口)。_process 的两段式分发(专属 + 通用)。register(按 type)与 register_general(全局)的适用场景。defaultdict 在 register 里的作用。engine.py:105-109:
105 def put(self, event: Event) -> None: 109 self._queue.put(event)
就一行。所有外部代码(网关、策略、定时器线程本身)都通过 put 与引擎通信——这是投递事件的唯一入口。
谁会调用 put?
on_tick → event_engine.put(Event("eTick.", tick))。_run_timer 每秒 put(Event("eTimer"))。write_log → put(Event("eLog.", log))。queue.Queue 内部有锁,所以 put 是线程安全的——多个线程同时 put 不会出错,事件按到达顺序排队。
💡 核心心法:put 的极简设计体现了"单一职责"——投递者只管往队列塞,不关心谁来处理、什么时候处理。这种解耦让"生产者"和"消费者"完全独立演化。
engine.py:66-78:
66 def _process(self, event: Event) -> None: 74 if event.type in self._handlers: 75 [handler(event) for handler in self._handlers[event.type]] 76 77 if self._general_handlers: 78 [handler(event) for handler in self._general_handlers]
_run 主循环拿到事件后调 _process。它做两段分发:
第一段(74-75 行):专属处理器。检查 _handlers 字典里有没有这个 type 的处理器列表,有就逐个调用。比如 type 是 "eTick.",就只调订阅了 "eTick." 的那些 handler。
第二段(77-78 行):通用处理器。检查 _general_handlers 列表是否非空,非空就对所有事件逐个调用。通用处理器不区分 type,拿到所有事件。
💡 核心心法:双轨制的设计动机是"全局观察者"。日志、监控、UI 全量事件流这类组件需要感知所有事件——如果只有专属处理器,它们得为每个 type 都注册一次(几十个 type 几十次 register)。有了 general_handlers,注册一次就能收到全量事件,代码极简。
注意第 75/78 行用列表推导式 [handler(event) for ...] 而非 for 循环——这是 VeighNa 的风格选择,仅为副作用执行(不收集返回值),功能等价于:
for handler in self._handlers[event.type]: handler(event)
⚠️ 重要说明:处理器在主循环线程内同步调用,所以单个 handler 抛异常会中断后续分发。实战中 handler 应自行 try/except,避免一个 bug 影响其它订阅者。
engine.py:111-130:
111 def register(self, type: str, handler: HandlerType) -> None: 116 handler_list: list = self._handlers[type] 117 if handler not in handler_list: 118 handler_list.append(handler) 120 def unregister(self, type: str, handler: HandlerType) -> None: 124 handler_list: list = self._handlers[type] 126 if handler in handler_list: 127 handler_list.remove(handler) 129 if not handler_list: 130 self._handlers.pop(type) # 清空后移除键,避免字典膨胀
几个细节:
defaultdict 的妙用(116 行):self._handlers[type] 即使 type 不存在也会自动建空 list,不会 KeyError。这就是 __init__ 里用 defaultdict(list) 的收益——register 时不用先判断 if type not in _handlers: _handlers[type] = []。
去重(117 行):同一个 handler 对同一 type 只注册一次(用 in 判断)。避免重复注册导致同一事件被同一函数处理多次。
清空后移除键(129-130 行):unregister 后若 list 空了,主动 pop 掉键,保持 _handlers 紧凑。这是内存卫生——长时间运行的进程里,如果一直往 _handlers 加 type 从不删,字典会无限膨胀。
典型用法(摘自后续章节会讲到的 OmsEngine):
# 订阅行情事件:每个 tick 来都回调 process_tick_event event_engine.register(EVENT_TICK, self.process_tick_event) event_engine.register(EVENT_ORDER, self.process_order_event)
engine.py:132-145:
132 def register_general(self, handler: HandlerType) -> None: 137 if handler not in self._general_handlers: 138 self._general_handlers.append(handler) 140 def unregister_general(self, handler: HandlerType) -> None: 144 if handler in self._general_handlers: 145 self._general_handlers.remove(handler)
通用处理器没有 type 维度,是一个扁平 list。逻辑更简单——加进去/移除,去重即可。
典型场景:
| 维度 | 专属处理器(register) | 通用处理器(register_general) |
|---|---|---|
| 存储 | _handlers defaultdict,按 type 分桶 |
_general_handlers 扁平 list |
| 触发条件 | 事件 type 匹配 | 任何事件都触发 |
| 适用 | 绝大多数业务订阅(OMS/策略/UI) | 日志/监控/调试等"全局观察者" |
| 性能开销 | 按 type O(1) 查找 | 每事件 O(n) 遍历(n 通常很小) |
实际项目中,99% 的订阅用 register(按需订阅特定 type),只有少数"我要看所有事件"的组件用 register_general。
第 01 节讲过定时器线程每秒 put(Event("eTimer"))。那么谁会订阅这个 eTimer 事件?典型用途:
class SomeEngine(BaseEngine): def __init__(self, ...): self.event_engine.register(EVENT_TIMER, self.process_timer_event) def process_timer_event(self, event: Event) -> None: # 每秒触发一次 self.check_pending_orders() # 检查挂单是否超时 self.update_statistics() # 刷新统计指标
CTA 策略引擎、算法交易引擎(做 TWAP/VWAP 拆单)、风控引擎(定时检查持仓)等,都会订阅 EVENT_TIMER 做周期性任务。
⚠️ 重要说明:
EVENT_TIMER的 type 字符串是"eTimer",没有按 vt_symbol 细分(因为心跳是全局的,不针对具体合约)。这与业务事件不同——业务事件 type 往往带合约后缀,如"eTick.rb2401.SHFE",第 5 章讲网关on_tick时会看到。
假设 CTP 收到一个 tick,完整流转:
1. CtpGateway.on_tick(tick) │ 构造 Event("eTick.rb2401.SHFE", tick) ▼ 2. event_engine.put(event) │ event 进 Queue ▼ 3. _run 主循环 queue.get() 拿到 event ▼ 4. _process(event): ├─ 第一段:_handlers["eTick.rb2401.SHFE"] + _handlers["eTick."] │ ├─ OmsEngine.process_tick_event → 更新 ticks 缓存 │ ├─ CtaEngine.process_tick_event → 喂给 BarGenerator │ └─ UI 行情面板 → 刷新表格 └─ 第二段:_general_handlers └─ 日志/监控组件 → 记录"收到 tick"
整条链路:网关 put → 队列 → 主循环 get → _process 两段分发 → 各订阅者处理。没有任何一方直接调用另一方,全靠事件 type 解耦。这就是 VeighNa 事件驱动的全貌。
queue.Queue 线程安全。第 2 章结束。你已经掌握了 VeighNa 的心脏。下一章(第 3 章)往上看一层——Event 的 data 载荷到底是什么?答案是 dataclass 数据建模(object.py + constant.py)。