NautilusTrader 自定义数据(Custom Data)全解析:Python/Rust 双模式注册、Parquet 持久化与运行时路由实战
2026/9/12 11:47:05 网站建设 项目流程

NautilusTrader 自定义数据(Custom Data)全解析:Python/Rust 双模式注册、Parquet 持久化与运行时路由实战

【免费下载链接】nautilus_traderProduction-grade Rust-native trading engine with deterministic event-driven architecture项目地址: https://gitcode.com/GitHub_Trending/na/nautilus_trader

导读

NautilusTrader 的事件驱动引擎以内置的高频数据模型(QuoteTick、TradeTick、Bar 等)为核心,但真实交易场景往往需要携带额外的自定义信息流——从交易所原始快照、舆情情绪分数,到自研信号的中间结果。本文基于仓库中的 custom_data.md 概念文档,深入讲解 NautilusTrader 的自定义数据体系:如何在纯 Python 下定义数据类而不写一行 Rust,如何在同二进制 Rust 下用过程宏产出原生编解码,以及这两条路径如何统一收敛到同一个 PyO3CustomData包装器、走完注册、序列化、Parquet/Feather 持久化、消息总线路由与策略订阅的全流程。读完本文,你将掌握自定义数据的端到端接入方法、注册表与 Arrow C FFI 桥的底层原理,以及 ParquetDataCatalog 动态类型注册的持久化机制。

一、设计目标:为什么需要一套"自定义数据"架构

内置数据类型的 schema 与编解码器在 Rust 二进制中静态已知;而用户自定义数据在编译期不可预知。文档明确列出该架构要满足的五项需求,这也决定了后续每一处设计取舍:

  1. 纯 Python 可定义:用户无需编写 Rust 代码即可定义自定义数据;
  2. Rust 原生路径:Rust 侧定义的自定义数据使用原生 Rust JSON 与 Arrow 处理器;
  3. 统一边界包装器:在 PyO3 边界只保留一个面向用户的CustomData包装器;
  4. 动态注册持久化ParquetDataCatalog通过动态类型注册(而非硬编码 schema)支持持久化;
  5. 完整路由能力:自定义数据可走与内置数据完全相同的数据引擎、actor、策略订阅流程。

从源码结构看,这套需求被拆解为三部分实现:进程级注册表(registry.rs)、统一包装器CustomData与 trait(custom.rs)、持久化编排层(custom.rs)。

二、高层模型:两种编写模式的统一

自定义数据支持两种编写形式,它们最终都汇入同一个外层CustomData包装器与同一个DataType身份模型:

模式编写形式注册路径编码/解码路径包装器后端
纯 Python带 JSON 与 Arrow 方法的类register_custom_data_class(...)Python 回调 + Arrow C FFIPythonCustomDataWrapper
同二进制 Rust#[custom_data]#[custom_data(pyo3)]类型ensure_custom_data_registered::<T>()+ 提取器注册原生 Rust原生 Rust payload

关键点在于:无论底层 payload 是 Python 对象包装还是原生 Rust 值,用户面对的 API 始终是同一个CustomData。这避免了 Python 用户接触 FFI 细节,也让 Rust 用户获得零 GIL 开销的原生路径。

三、端到端数据流:从定义到查询

文档用一张时序图描述了完整的生命周期,其核心链路如下(省略参与者标注,保留关键调用序列):

这条链路的两个关键观察点:一是注册发生在写入/查询之前,且是进程级、按type_name键控的;二是写入与查询对称——两侧都通过注册表按type_name解析处理器,因此只要注册保持一致,读写即可互相匹配。

四、核心组件逐个拆解

4.1 Registry 模块:进程级 JSON/Arrow/提取器注册表

registry.rs 是整个自定义数据体系的"调度中心"。它通过OnceLock初始化静态注册状态,内部用DashMap存储三类处理器:

  • JSON 反序列化器:以type_name为键,签名是Fn(serde_json::Value) -> Result<Arc<dyn CustomDataTrait>>
  • Arrow schema/编码器/解码器三元组:以type_name为键,编码器负责&[Arc<dyn CustomDataTrait>] -> RecordBatch,解码器负责(metadata, RecordBatch) -> Vec<Data>
  • Python 提取器(PyExtractor):把 Python 对象转换为Arc<dyn CustomDataTrait>
  • Rust 提取器工厂(RustExtractorFactory):为同二进制类型生产 Python 提取器。

注册操作使用原子化的DashMap::entry()语义,文档特别强调:register_*ensure_*并发调用时不会在抢占条目上产生竞态。二者的行为差异在源码中有明确体现并配有测试:

  • register_json_deserializer/register_arrow/register_py_extractor采用Entry::Occupied检查,重复注册会直接bail!("... already registered ...")(对应测试register_json_deserializer_fails_on_duplicate);
  • ensure_json_deserializer_registered/ensure_arrow_registered/ensure_py_extractor_registered采用or_insert_with幂等且不覆盖已注册处理器,适合在模块初始化等可能重复执行的路径中调用(对应测试ensure_json_deserializer_registered_is_idempotentensure_arrow_registered_is_idempotent)。

关键设计:不把任何类型硬编码进主二进制,而是在运行时根据DataType中保存的type_name(以及 Parquet 元数据中的type_name)动态解析处理器。这意味着注册表可在不重新编译引擎的情况下扩展新类型。

4.2CustomData包装器:跨 FFI 边界的统一容器

custom.rs 中定义了外层 PyO3 包装器CustomData

pub struct CustomData { pub data: Arc<dyn CustomDataTrait>, // 内部自定义 payload pub data_type: DataType, // 数据类型身份 }

构造签名遵循"先DataType、后 payload"的顺序:CustomData(data_type, data)。它包含:

  • 一个DataType
  • 一个实现CustomDataTrait的内部 payload(包在Arc<dyn CustomDataTrait>中,Arc 克隆是 O(1) 的,向 Python 传值时代价极低)。

时间戳委托ts_eventts_init委托给内部CustomDataTrait实现,并在包装器上以属性暴露。

相等性语义:Rust 侧PartialEq先比较DataType,再委托eq_arc比较内部 payload。Python 侧实现了__eq____repr__。实例故意不可哈希——避免哈希与 payload 比较语义不一致(可哈希对象要求a == b蕴含hash(a) == hash(b),而 payload 相等性由 trait 动态决定)。

两个构造入口CustomData::from_arc(arc)从内部类型名推导DataTypeCustomData::new(arc, data_type)显式传入DataType——后者用于从外部元数据(如 Parquet)恢复数据类型的场景。

4.2.1 统一的 JSON 信封(envelope)

CustomData被序列化为 JSON 时——无论是to_json_bytes/from_json_bytes、SQL 缓存还是 Redis——都使用同一个规范化信封(源码中为CustomDataEnvelope,见 custom.rs):

  • type:自定义类型名(来自CustomDataTrait::type_name);
  • data_type:一个对象,包含type_namemetadata,以及可选的identifier
  • payload仅内部 payload(即CustomDataTrait::to_json解析后的值)。

这个信封的意义在于:反序列化不依赖用户 payload 的字段名。注册的 JSON 反序列化器只接收payload值(见 registry.rs 的parse_envelope_payload),因此用户结构体可以随意使用字段名——包括valuetype这类可能与包装元数据冲突的名字——而不会干扰。

配套测试test_custom_data_json_roundtrip验证了带metadataidentifierDataType经过序列化往返后,type_namemetadataidentifier与内部 payload 均保持一致。

4.3DataType:路由与持久化的身份标识

DataType(见 mod.rs)是自定义数据路由与持久化的身份模型,构造签名:DataType(type_name, metadata=None, identifier=None)

字段角色:

  • type_name:处理器的查找键,也是 Parquet 路径的一部分;
  • metadata(可选):参与相等性、哈希与 topic 派生;
  • identifier(可选)不影响路由、相等性、哈希,仅用于持久化路径与缓存数据库查找。

源码揭示了topichash的预计算机制:DataType::new在构造时用type_name与 metadata 拼接出 topic(如type_name.key1=value1.key2=value2,见params_to_topic_suffix),并对 topic 预计算 hash 缓存。由此可以推断:

  • 两个type_namemetadata相同、但identifier不同的DataType比较相等、哈希相同、发布到同一个消息总线 topic
  • identifier决定目录路径,即data/custom/<type_name>/<identifier...>,并参与 PostgreSQL 与 Redis 的过滤。

持久化时完整保存DataTypeto_persistence_json只序列化type_namemetadataidentifier,不含 topic/hash),查询时恢复;而处理器查找只使用type_name。因此同一逻辑类型可以携带不同 metadata 或 identifier,仍通过同一个注册处理器解码——这正是"动态类型注册"得以成立的身份基础。

五、注册架构:从 Python 对象到 Rust trait 对象的桥

注册是衔接用户类型与引擎数据管道的桥梁,两条路径的汇合关系如下:

5.1 纯 Python 注册:register_custom_data_class(MyType)

当 Python 代码调用register_custom_data_class(MyType)时,依次发生:

  1. Rust 保留该类引用,用于后续 JSON 重建(源码中register_python_data_class将类存入进程级DashMap<String, Py<PyAny>>);
  2. Rust 注册调用类回调的 JSON 与 Arrow 处理器;
  3. 构造CustomData时,若注册的原生提取器接受该对象则用之;否则将该对象包装进PythonCustomDataWrapper

需要特别说明:此路径上的 JSON 与 Arrow 回调都在 Python GIL 下执行——这是纯 Python 模式的主要开销来源,也是原生 Rust 模式存在的原因。

5.2 同二进制 Rust 注册:#[custom_data]ensure_*系列

对编译进进程的 Rust 类型:

  1. #[custom_data]#[custom_data(pyo3)]过程宏生成CustomDataTrait与 JSON 实现,默认还生成 Arrow 实现;
  2. ensure_custom_data_registered::<T>()把原生 schema/编码器/解码器插入进程级注册表;
  3. ensure_rust_extractor_registered::<T>()注册一个提取器工厂(源码中RustExtractorFactory是一个Fn() -> PyExtractor,惰性构建)。一旦通过 Python 类注册激活,该提取器可以恢复出具体 Rust 类型,而不再退回 Python 包装器。

该路径的编码/解码全程停留在 Rust 原生侧,不涉及 GIL 与 Python 回调。

5.3 注册优先级

register_custom_data_class(...)解析处理器时按以下顺序:

  1. 优先:若存在已注册的原生提取器及原生 JSON/Arrow 处理器,则使用原生路径;
  2. 兜底:否则使用 Python 包装器与回调处理器。

由于ensure_*注册是幂等的、且不覆盖既有原生处理器,同一种类型无论被注册多少次,只要原生处理器先于 Python 回调注册,就始终走原生路径。

六、包装器后端:两种 payload 实现

外层CustomData内部可以持有不同的 payload 实现。

6.1PythonCustomDataWrapper

用于纯 Python 自定义数据(custom.rs),职责:

  • 持有 Python 对象引用;
  • 缓存ts_eventts_inittype_name——源码注释明确说明这是为了避免在热路径(数据排序、消息路由)频繁获取 GIL;type_name还通过intern_type_name_static字符串驻留,只保留每类一个静态副本;
  • 实现CustomDataTrait
  • 支持在 GIL 下调用 Python 的 JSON 与 Arrow 回调路径(to_json优先调用对象的to_json()方法,否则回退到json.dumps(obj.__dict__))。

它是构造时的默认回退:当没有注册的提取器接受对象时使用;Python 侧 JSON/Arrow 解码器也直接产出该包装器。eq_arc按 Python 对象身份 +==比较,避免两个不同 Python 对象因同名同时间戳被误判相等。

6.2 原生同二进制 Rust payload

对编译进进程的 Rust 类型,内部 payload 就是具体 Rust 值,可直接从Arc<dyn CustomDataTrait>向下转型(downcast)。序列化与解码不需要任何 Python 回调路径,性能与内置类型一致。

七、持久化架构:动态 Arrow 注册与 ParquetDataCatalog

7.1 为什么需要动态 Arrow 注册

内置类型的 schema 与编码器对 Rust 二进制是静态已知的,自定义数据则不然。因此持久化层通过已注册的type_name动态解析自定义数据——这决定了注册必须发生在任何写入/查询之前,也解释了第 4.1 节注册表存在的必要性。

7.2 Catalog 写入流

ParquetDataCatalog要求自定义数据以CustomData值的形式写入。写入路径(custom.rs 中的编排逻辑):

  1. 从内部 payload 取type_name,从首个值DataTypemetadataidentifier
  2. 在进程级注册表中查找 Arrow 编码器;
  3. 把值编码为RecordBatch
  4. 追加data_type——每行写入该DataType的持久化 JSON(schema_with_data_type_column/augment_batch_with_data_type_column完成列的追加与 schema 元数据合并);
  5. type_name与 metadata 附加到 Arrow schema 元数据;
  6. 将 batch 写入data/custom/<type_name>/<identifier...>路径下的 Parquet 文件。

identifier在成为路径段之前会被规范化。由于 metadata/identifier 取自首个值,同一批写入要求类型一致,这也是文档明确写出"first value's DataType"的原因。

7.3 Catalog 读取流

查询时:

  1. 读取匹配的 Parquet 文件;
  2. 从 schema 元数据中提取type_name
  3. 向进程级注册表请求解码器;
  4. RecordBatch解码为Vec<Data>
  5. 用原始DataType重建CustomDatadata_type列每行保存了持久化 JSON,可直接恢复)。

读取与写入时的注册查找完全对称。此外,文档还说明:当把 Feather 流转换为 Parquet(如回测之后)时,自定义数据分支被设计为直接变换 Arrow batch 并写入对应的自定义数据路径。

7.4 已知限制:Streaming Feather 暂不支持

文档明确给出警告:

Streaming Feather 持久化目前不可用。Python 的StreamingFeatherWriter会以OSError拒绝CustomDataconvert_stream_to_data也不会把自定义数据 Feather 流转换为 Parquet。这将在未来版本中支持。在此期间,请直接用ParquetDataCatalog.write_custom_data将自定义数据写入目录。

这是当前版本的真实边界,接入时应使用write_custom_data直写目录,而不是走流式 Feather 管道。

八、Arrow C FFI 桥:零拷贝穿越 Python/Rust 边界

纯 Python 自定义数据没有原生 Rust Arrow 编码逻辑。为此,NautilusTrader 借助Arrow C FFI 接口在 Python 与 Rust 之间传递RecordBatch,全程不经过 JSON 或二进制序列化,避免数据复制与格式转换开销。

8.1 纯 Python 编码路径

  1. Rust 获取 GIL;
  2. 对第一个 Python payload 调用encode_record_batch_py(...)
  3. Python 把对象转换为pyarrow.RecordBatch
  4. Python 通过_export_to_c把 batch 导出为 Arrow C FFI 结构体(FFI_ArrowArray+FFI_ArrowSchema);
  5. Rust 从 FFI 结构体重构原生RecordBatch并写出。

8.2 纯 Python 解码路径

反向流程:

  1. Rust 把自身RecordBatch转为 Arrow C FFI 结构体;
  2. Python 通过RecordBatch._import_from_c导入;
  3. Python 在类上调用decode_record_batch_py(metadata, batch)
  4. Rust 把返回的 Python 对象包进PythonCustomDataWrapper

8.3 原生路径不走桥

同二进制 Rust 自定义数据不使用Arrow C FFI 桥,它们直接使用注册在进程中的原生 Rust 编码/解码处理器。

九、查询时的数据重建

从目录加载自定义数据时,重建方式取决于后端:

  • 同二进制 Rust 类型:直接解码为原生 Rust 值;
  • 纯 Python 类型:通过注册类的decode_record_batch_py(...)回调重建。

无论哪种后端,调用方在 PyO3 API 边界拿到的都是同一个外层CustomData包装器——这正是"统一边界包装器"设计目标在查询阶段的落地。

十、运行时集成:消息总线、actor 与策略

自定义数据并非只在持久化层生效,它完整参与 NautilusTrader 的运行时路由。文档给出的集成点与仓库源码对应如下:

  • crates/data/src/engine/mod.rs:数据引擎通过消息总线发布CustomData
  • crates/common/src/msgbus/switchboard.rs:从DataType派生自定义 topic——结合第 4.3 节可知,topic 由type_name+metadata预计算,因此同一类型的不同 identifier 会路由到同一 topic;
  • crates/common/src/actor/:把自定义数据路由进 actor 订阅;
  • crates/trading/src/python/strategy.rs:向 Python 策略的on_data暴露自定义数据;
  • crates/backtest/src/engine.rs:把Data::Custom视为数据引擎投递的输入,而非交易所路由数据——这意味着自定义数据在回测中可以与内置数据一样按时序进入引擎。

结论:一个注册过的自定义类型可以被持久化、查询、订阅、消费,走的全是与内置数据族相同的运行时接口。

十一、缓存数据库集成:PostgreSQL 与 Redis

自定义数据同样接入缓存数据库:

  • PostgreSQL:存储在custom表中,记录包含data_typemetadataidentifier与完整 JSON payload;读取时通过CustomData::from_json_bytes(...)重建;Python SQL 绑定暴露add_custom_dataload_custom_data(对应 schema/sql/tables.sql 中的表结构设计);
  • Redis:存储在custom:<ts_init_020>:<uuid>键下,值为完整CustomDataJSON;add_custom_dataload_custom_dataDataType(type_name、metadata、identifier)过滤,并按ts_init排序返回,通过 PyO3RedisCacheDatabaseAPI 暴露。

值得注意的是 JSON 信封在此复用:无论是 SQL 缓存还是 Redis,序列化走的都是第 4.2.1 节描述的同一套{type, data_type, payload}信封,反序列化统一由注册表按type_name分发。

十二、实践启示

文档在结尾给出的核心结论值得反复强调:"纯 Python 编写"与"原生 Rust 编解码"不是两套割裂的功能集,而是同一套自定义数据概念系统的两种后端。据此可以总结出以下接入建议:

  1. 原型阶段:用纯 Python 类 +register_custom_data_class快速定义数据,配合 Arrow C FFI 桥即可直接写入ParquetDataCatalog
  2. 性能敏感场景:把类型编译进进程,使用#[custom_data]宏与ensure_*注册,全程免 GIL;
  3. 混合部署:利用注册优先级——原生处理器先注册、Python 类再注册,即可让同一种类型在 Python 侧保持简单、在引擎侧走原生路径;
  4. 注意边界:Streaming Feather 流式持久化当前不可用,回测后如需落盘请直写write_custom_data
  5. 身份设计:把参与路由的维度放进metadata,把仅用于区分存储路径的维度放进identifier,二者职责不可混淆。

延伸阅读

  • 概念总览:overview.md、data/index.md
  • 事件与数据模型:value_types.md、custom_data.md
  • 引擎与策略路由:message_bus.md、strategies.md
  • 持久化相关:persistence.md
  • 核心源码:custom.rs、registry.rs、mod.rs、backend/custom.rs

【免费下载链接】nautilus_traderProduction-grade Rust-native trading engine with deterministic event-driven architecture项目地址: https://gitcode.com/GitHub_Trending/na/nautilus_trader

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询