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 二进制中静态已知;而用户自定义数据在编译期不可预知。文档明确列出该架构要满足的五项需求,这也决定了后续每一处设计取舍:
- 纯 Python 可定义:用户无需编写 Rust 代码即可定义自定义数据;
- Rust 原生路径:Rust 侧定义的自定义数据使用原生 Rust JSON 与 Arrow 处理器;
- 统一边界包装器:在 PyO3 边界只保留一个面向用户的
CustomData包装器; - 动态注册持久化:
ParquetDataCatalog通过动态类型注册(而非硬编码 schema)支持持久化; - 完整路由能力:自定义数据可走与内置数据完全相同的数据引擎、actor、策略订阅流程。
从源码结构看,这套需求被拆解为三部分实现:进程级注册表(registry.rs)、统一包装器CustomData与 trait(custom.rs)、持久化编排层(custom.rs)。
二、高层模型:两种编写模式的统一
自定义数据支持两种编写形式,它们最终都汇入同一个外层CustomData包装器与同一个DataType身份模型:
| 模式 | 编写形式 | 注册路径 | 编码/解码路径 | 包装器后端 |
|---|---|---|---|---|
| 纯 Python | 带 JSON 与 Arrow 方法的类 | register_custom_data_class(...) | Python 回调 + Arrow C FFI | PythonCustomDataWrapper |
| 同二进制 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_idempotent、ensure_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_event、ts_init委托给内部CustomDataTrait实现,并在包装器上以属性暴露。
相等性语义:Rust 侧PartialEq先比较DataType,再委托eq_arc比较内部 payload。Python 侧实现了__eq__与__repr__。实例故意不可哈希——避免哈希与 payload 比较语义不一致(可哈希对象要求a == b蕴含hash(a) == hash(b),而 payload 相等性由 trait 动态决定)。
两个构造入口:CustomData::from_arc(arc)从内部类型名推导DataType;CustomData::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_name、metadata,以及可选的identifier;payload:仅内部 payload(即CustomDataTrait::to_json解析后的值)。
这个信封的意义在于:反序列化不依赖用户 payload 的字段名。注册的 JSON 反序列化器只接收payload值(见 registry.rs 的parse_envelope_payload),因此用户结构体可以随意使用字段名——包括value、type这类可能与包装元数据冲突的名字——而不会干扰。
配套测试test_custom_data_json_roundtrip验证了带metadata与identifier的DataType经过序列化往返后,type_name、metadata、identifier与内部 payload 均保持一致。
4.3DataType:路由与持久化的身份标识
DataType(见 mod.rs)是自定义数据路由与持久化的身份模型,构造签名:DataType(type_name, metadata=None, identifier=None)。
字段角色:
type_name:处理器的查找键,也是 Parquet 路径的一部分;metadata(可选):参与相等性、哈希与 topic 派生;identifier(可选):不影响路由、相等性、哈希,仅用于持久化路径与缓存数据库查找。
源码揭示了topic与hash的预计算机制:DataType::new在构造时用type_name与 metadata 拼接出 topic(如type_name.key1=value1.key2=value2,见params_to_topic_suffix),并对 topic 预计算 hash 缓存。由此可以推断:
- 两个
type_name与metadata相同、但identifier不同的DataType,比较相等、哈希相同、发布到同一个消息总线 topic; identifier决定目录路径,即data/custom/<type_name>/<identifier...>,并参与 PostgreSQL 与 Redis 的过滤。
持久化时完整保存DataType(to_persistence_json只序列化type_name、metadata、identifier,不含 topic/hash),查询时恢复;而处理器查找只使用type_name。因此同一逻辑类型可以携带不同 metadata 或 identifier,仍通过同一个注册处理器解码——这正是"动态类型注册"得以成立的身份基础。
五、注册架构:从 Python 对象到 Rust trait 对象的桥
注册是衔接用户类型与引擎数据管道的桥梁,两条路径的汇合关系如下:
5.1 纯 Python 注册:register_custom_data_class(MyType)
当 Python 代码调用register_custom_data_class(MyType)时,依次发生:
- Rust 保留该类引用,用于后续 JSON 重建(源码中
register_python_data_class将类存入进程级DashMap<String, Py<PyAny>>); - Rust 注册调用类回调的 JSON 与 Arrow 处理器;
- 构造
CustomData时,若注册的原生提取器接受该对象则用之;否则将该对象包装进PythonCustomDataWrapper。
需要特别说明:此路径上的 JSON 与 Arrow 回调都在 Python GIL 下执行——这是纯 Python 模式的主要开销来源,也是原生 Rust 模式存在的原因。
5.2 同二进制 Rust 注册:#[custom_data]与ensure_*系列
对编译进进程的 Rust 类型:
#[custom_data]或#[custom_data(pyo3)]过程宏生成CustomDataTrait与 JSON 实现,默认还生成 Arrow 实现;ensure_custom_data_registered::<T>()把原生 schema/编码器/解码器插入进程级注册表;ensure_rust_extractor_registered::<T>()注册一个提取器工厂(源码中RustExtractorFactory是一个Fn() -> PyExtractor,惰性构建)。一旦通过 Python 类注册激活,该提取器可以恢复出具体 Rust 类型,而不再退回 Python 包装器。
该路径的编码/解码全程停留在 Rust 原生侧,不涉及 GIL 与 Python 回调。
5.3 注册优先级
register_custom_data_class(...)解析处理器时按以下顺序:
- 优先:若存在已注册的原生提取器及原生 JSON/Arrow 处理器,则使用原生路径;
- 兜底:否则使用 Python 包装器与回调处理器。
由于ensure_*注册是幂等的、且不覆盖既有原生处理器,同一种类型无论被注册多少次,只要原生处理器先于 Python 回调注册,就始终走原生路径。
六、包装器后端:两种 payload 实现
外层CustomData内部可以持有不同的 payload 实现。
6.1PythonCustomDataWrapper
用于纯 Python 自定义数据(custom.rs),职责:
- 持有 Python 对象引用;
- 缓存
ts_event、ts_init、type_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 中的编排逻辑):
- 从内部 payload 取
type_name,从首个值的DataType取metadata与identifier; - 在进程级注册表中查找 Arrow 编码器;
- 把值编码为
RecordBatch; - 追加
data_type列——每行写入该DataType的持久化 JSON(schema_with_data_type_column/augment_batch_with_data_type_column完成列的追加与 schema 元数据合并); - 把
type_name与 metadata 附加到 Arrow schema 元数据; - 将 batch 写入
data/custom/<type_name>/<identifier...>路径下的 Parquet 文件。
identifier在成为路径段之前会被规范化。由于 metadata/identifier 取自首个值,同一批写入要求类型一致,这也是文档明确写出"first value's DataType"的原因。
7.3 Catalog 读取流
查询时:
- 读取匹配的 Parquet 文件;
- 从 schema 元数据中提取
type_name; - 向进程级注册表请求解码器;
- 把
RecordBatch解码为Vec<Data>; - 用原始
DataType重建CustomData(data_type列每行保存了持久化 JSON,可直接恢复)。
读取与写入时的注册查找完全对称。此外,文档还说明:当把 Feather 流转换为 Parquet(如回测之后)时,自定义数据分支被设计为直接变换 Arrow batch 并写入对应的自定义数据路径。
7.4 已知限制:Streaming Feather 暂不支持
文档明确给出警告:
Streaming Feather 持久化目前不可用。Python 的
StreamingFeatherWriter会以OSError拒绝CustomData,convert_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 编码路径
- Rust 获取 GIL;
- 对第一个 Python payload 调用
encode_record_batch_py(...); - Python 把对象转换为
pyarrow.RecordBatch; - Python 通过
_export_to_c把 batch 导出为 Arrow C FFI 结构体(FFI_ArrowArray+FFI_ArrowSchema); - Rust 从 FFI 结构体重构原生
RecordBatch并写出。
8.2 纯 Python 解码路径
反向流程:
- Rust 把自身
RecordBatch转为 Arrow C FFI 结构体; - Python 通过
RecordBatch._import_from_c导入; - Python 在类上调用
decode_record_batch_py(metadata, batch); - 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_type、metadata、identifier与完整 JSON payload;读取时通过CustomData::from_json_bytes(...)重建;Python SQL 绑定暴露add_custom_data与load_custom_data(对应 schema/sql/tables.sql 中的表结构设计); - Redis:存储在
custom:<ts_init_020>:<uuid>键下,值为完整CustomDataJSON;add_custom_data与load_custom_data按DataType(type_name、metadata、identifier)过滤,并按ts_init排序返回,通过 PyO3RedisCacheDatabaseAPI 暴露。
值得注意的是 JSON 信封在此复用:无论是 SQL 缓存还是 Redis,序列化走的都是第 4.2.1 节描述的同一套{type, data_type, payload}信封,反序列化统一由注册表按type_name分发。
十二、实践启示
文档在结尾给出的核心结论值得反复强调:"纯 Python 编写"与"原生 Rust 编解码"不是两套割裂的功能集,而是同一套自定义数据概念系统的两种后端。据此可以总结出以下接入建议:
- 原型阶段:用纯 Python 类 +
register_custom_data_class快速定义数据,配合 Arrow C FFI 桥即可直接写入ParquetDataCatalog; - 性能敏感场景:把类型编译进进程,使用
#[custom_data]宏与ensure_*注册,全程免 GIL; - 混合部署:利用注册优先级——原生处理器先注册、Python 类再注册,即可让同一种类型在 Python 侧保持简单、在引擎侧走原生路径;
- 注意边界:Streaming Feather 流式持久化当前不可用,回测后如需落盘请直写
write_custom_data; - 身份设计:把参与路由的维度放进
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),仅供参考