ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

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

NautilusTrader 自定义数据(Custom Data)全解析:Python/Rust 双模式注册、Parquet 持久化与运行时路由实战 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与 traitcustom.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) - ResultArcdyn CustomDataTraitArrow schema/编码器/解码器三元组以type_name为键编码器负责[Arcdyn CustomDataTrait] - RecordBatch解码器负责(metadata, RecordBatch) - VecDataPython 提取器PyExtractor把 Python 对象转换为Arcdyn CustomDataTraitRust 提取器工厂RustExtractorFactory为同二进制类型生产 Python 提取器。注册操作使用原子化的DashMap::entry()语义文档特别强调register_*与ensure_*并发调用时不会在抢占条目上产生竞态。二者的行为差异在源码中有明确体现并配有测试register_json_deserializer/register_arrow/register_py_extractor采用Entry::Occupied检查重复注册会直接bail!(... already registered ...)对应测试register_json_deserializer_fails_on_duplicateensure_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 包装器CustomDatapub struct CustomData { pub data: Arcdyn CustomDataTrait, // 内部自定义 payload pub data_type: DataType, // 数据类型身份 }构造签名遵循先DataType、后 payload的顺序CustomData(data_type, data)。它包含一个DataType一个实现CustomDataTrait的内部 payload包在Arcdyn 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)从内部类型名推导DataTypeCustomData::new(arc, data_type)显式传入DataType——后者用于从外部元数据如 Parquet恢复数据类型的场景。4.2.1 统一的 JSON 信封envelope当CustomData被序列化为 JSON 时——无论是to_json_bytes/from_json_bytes、SQL 缓存还是 Redis——都使用同一个规范化信封源码中为CustomDataEnvelope见 custom.rstype自定义类型名来自CustomDataTrait::type_namedata_type一个对象包含type_name、metadata以及可选的identifierpayload仅内部 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, metadataNone, identifierNone)。字段角色type_name处理器的查找键也是 Parquet 路径的一部分metadata可选参与相等性、哈希与 topic 派生identifier可选不影响路由、相等性、哈希仅用于持久化路径与缓存数据库查找。源码揭示了topic与hash的预计算机制DataType::new在构造时用type_name与 metadata 拼接出 topic如type_name.key1value1.key2value2见params_to_topic_suffix并对 topic 预计算 hash 缓存。由此可以推断两个type_name与metadata相同、但identifier不同的DataType比较相等、哈希相同、发布到同一个消息总线 topicidentifier决定目录路径即data/custom/type_name/identifier...并参与 PostgreSQL 与 Redis 的过滤。持久化时完整保存DataTypeto_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将类存入进程级DashMapString, PyPyAnyRust 注册调用类回调的 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——源码注释明确说明这是为了避免在热路径数据排序、消息路由频繁获取 GILtype_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 值可直接从Arcdyn CustomDataTrait向下转型downcast。序列化与解码不需要任何 Python 回调路径性能与内置类型一致。七、持久化架构动态 Arrow 注册与 ParquetDataCatalog7.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的持久化 JSONschema_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 values DataType的原因。7.3 Catalog 读取流查询时读取匹配的 Parquet 文件从 schema 元数据中提取type_name向进程级注册表请求解码器把RecordBatch解码为VecData用原始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 编码路径Rust 获取 GIL对第一个 Python payload 调用encode_record_batch_py(...)Python 把对象转换为pyarrow.RecordBatchPython 通过_export_to_c把 batch 导出为 Arrow C FFI 结构体FFI_ArrowArrayFFI_ArrowSchemaRust 从 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数据引擎通过消息总线发布CustomDatacrates/common/src/msgbus/switchboard.rs从DataType派生自定义 topic——结合第 4.3 节可知topic 由type_namemetadata预计算因此同一类型的不同 identifier 会路由到同一 topiccrates/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键下值为完整CustomDataJSONadd_custom_data与load_custom_data按DataTypetype_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),仅供参考
返回列表