第 2 章 · 02 register/put/process 三大核心机制


文档摘要

第 2 章 · 02 register/put/process 三大核心机制 本节摘要:本节继续精读 ,聚焦事件总线的三大核心机制: (投递事件)、 (分发事件)、 / (订阅/退订)。你会看到 VeighNa 的发布订阅设计有一个独特之处——它同时支持"专属处理器"(按 type 订阅,绝大多数场景)和"通用处理器"(监听全量事件,用于日志/监控)。这种双轨制让"全局观察者"组件不必为每个 type 重复注册一次。读完本节,你完整掌握了 EventEngine 的工作原理,为理解后续所有"谁在 put、谁在 register"打下基础。 内容来源:原项目源码 ,精读并套用体系化模板。 学习目标 阅读完本节,你应当能够: 说清 方法的线程安全性和它的唯一性(投递事件的唯一入口)。

第 2 章 · 02 register/put/process 三大核心机制

本节摘要:本节继续精读 engine.py,聚焦事件总线的三大核心机制:put(投递事件)、_process(分发事件)、register/unregister(订阅/退订)。你会看到 VeighNa 的发布订阅设计有一个独特之处——它同时支持"专属处理器"(按 type 订阅,绝大多数场景)和"通用处理器"(监听全量事件,用于日志/监控)。这种双轨制让"全局观察者"组件不必为每个 type 重复注册一次。读完本节,你完整掌握了 EventEngine 的工作原理,为理解后续所有"谁在 put、谁在 register"打下基础。

内容来源:原项目源码 vnpy/event/engine.py,精读并套用体系化模板。

学习目标

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

  1. 说清 put 方法的线程安全性和它的唯一性(投递事件的唯一入口)。
  2. 理解 _process两段式分发(专属 + 通用)。
  3. 区分 register(按 type)与 register_general(全局)的适用场景。
  4. 解释 defaultdictregister 里的作用。
  5. 读懂 EVENT_TIMER 心跳在实际策略里的典型用法。

一、put:投递事件的唯一入口

engine.py:105-109:

105 def put(self, event: Event) -> None: 109 self._queue.put(event)

就一行。所有外部代码(网关、策略、定时器线程本身)都通过 put 与引擎通信——这是投递事件的唯一入口

谁会调用 put?

  • 网关:CtpGateway 收到行情 → on_tickevent_engine.put(Event("eTick.", tick))
  • 定时器线程:_run_timer 每秒 put(Event("eTimer"))
  • 策略引擎:CtaEngine 在每个 bar 触发 → 通过事件通知 UI。
  • MainEngine:write_logput(Event("eLog.", log))

queue.Queue 内部有锁,所以 put 是线程安全的——多个线程同时 put 不会出错,事件按到达顺序排队。

💡 核心心法:put 的极简设计体现了"单一职责"——投递者只管往队列塞,不关心谁来处理、什么时候处理。这种解耦让"生产者"和"消费者"完全独立演化。

二、_process:两段式分发

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 影响其它订阅者。

三、register / unregister:专属处理器

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) # 清空后移除键,避免字典膨胀

几个细节:

  1. defaultdict 的妙用(116 行):self._handlers[type] 即使 type 不存在也会自动建空 list,不会 KeyError。这就是 __init__ 里用 defaultdict(list) 的收益——register 时不用先判断 if type not in _handlers: _handlers[type] = []

  2. 去重(117 行):同一个 handler 对同一 type 只注册一次(用 in 判断)。避免重复注册导致同一事件被同一函数处理多次。

  3. 清空后移除键(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)

四、register_general / unregister_general:通用处理器

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。逻辑更简单——加进去/移除,去重即可。

典型场景:

  • 日志网关:把所有事件落盘,用于事后复盘。
  • 监控统计:统计每秒事件数、各 type 占比。
  • GUI 全量事件流:开发调试时展示所有流转的事件。

五、专属 vs 通用的对照

维度 专属处理器(register) 通用处理器(register_general)
存储 _handlers defaultdict,按 type 分桶 _general_handlers 扁平 list
触发条件 事件 type 匹配 任何事件都触发
适用 绝大多数业务订阅(OMS/策略/UI) 日志/监控/调试等"全局观察者"
性能开销 按 type O(1) 查找 每事件 O(n) 遍历(n 通常很小)

实际项目中,99% 的订阅用 register(按需订阅特定 type),只有少数"我要看所有事件"的组件用 register_general。

六、EVENT_TIMER 心跳的实际用途

第 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 事件驱动的全貌。

本节要点回顾

  1. put 是投递事件的唯一入口,基于 queue.Queue 线程安全。
  2. _process 两段式:先按 type 调专属处理器,再调通用处理器。
  3. register(按 type)用于绝大多数业务订阅,register_general(全局)用于日志/监控。
  4. defaultdict 让 register 不必预判键是否存在,自动建空 list。
  5. 去重 + 清空移除键:register 防重复,unregister 清空后 pop 键保内存卫生。
  6. EVENT_TIMER 心跳:策略引擎/算法引擎/风控订阅它做周期任务。
  7. 完整流转:网关 put → 队列 → 主循环 → _process 两段分发 → 订阅者处理,全解耦。

第 2 章结束。你已经掌握了 VeighNa 的心脏。下一章(第 3 章)往上看一层——Event 的 data 载荷到底是什么?答案是 dataclass 数据建模(object.py + constant.py)。


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