ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

用 Rust 编写 Nautilus Trader 数据 Actor:从 SpreadMonitor 到生产级组件

用 Rust 编写 Nautilus Trader 数据 Actor:从 SpreadMonitor 到生产级组件 用 Rust 编写 Nautilus Trader 数据 Actor从 SpreadMonitor 到生产级组件【免费下载链接】nautilus_traderProduction-grade Rust-native trading engine with deterministic event-driven architecture项目地址: https://gitcode.com/GitHub_Trending/na/nautilus_trader这篇指南面向希望用 Rust 为 Nautilus Trader 编写自定义数据组件的开发者。文中以SpreadMonitor价差监控器为完整示例逐步讲解如何定义结构体、构造DataActorCore、接入nautilus_actor!宏、实现DataActortrait 的数据与生命周期处理器并最终把 Actor 注册到BacktestEngine或LiveNode。读完后你将掌握 Rust Actor 的完整编写范式、门面facade与原生native访问 API 的边界以及回调中安全使用ActorRef的运行时约束。Nautilus Trader 中的 Actor数据 Actor是一个接收行情数据、自定义数据/信号和系统事件但不管理订单的组件。与之相对Strategy策略在 Actor 能力之上增加了完整的订单管理能力。Actor 常用于监控、指标计算、数据聚合、信号生成等旁路任务例如订阅报价并计算买卖价差、跟踪订单簿失衡、监听队列压力与 Socket 连接状态等。Actor 与 Strategy 的完整能力对照可参考 Actors 概念指南 与 Strategies 指南。1. 定义结构体持有 DataActorCore 与业务状态一个 Rust Actor 本质上是一个普通结构体它必须持有DataActorCore存放运行时身份与状态并根据需要携带自己的业务字段。DataActorCore在 crates/common/src/actor/data_actor.rs 中定义从源码可见它集中管理了actor_idActor 标识符configDataActorConfig配置trader_id/clock/cache注册前为None在注册时由运行时接线wire upstateComponentState生命周期状态一组按数据类型组织的订阅处理器映射quote_handlers、bar_handlers、book_handlers等indicatorsActor 内注册的技术指标集合。用户代码通常不直接操作DataActorCore的内部字段而是通过DataActortrait 提供的门面方法facade methods访问运行时状态例如clock()cache()config()actor_id()trader_id()各类订阅方法以指南中的SpreadMonitor为例它订阅某合约的报价并记录买卖价差业务状态只需一个instrument_iduse nautilus_common::{nautilus_actor, actor::{DataActor, DataActorConfig, DataActorCore}}; use nautilus_model::{data::QuoteTick, identifiers::{ActorId, InstrumentId}}; pub struct SpreadMonitor { core: DataActorCore, instrument_id: InstrumentId, }从源码结构看DataActorCore采用RcRefCell_保存clock与cache并且所有类型化订阅处理器都按数据种类拆分存储data_actor.rs这决定了 Actor 在回调中按类型分发数据的底层机制。2. 实现构造函数DataActorConfig 与默认值Actor 的构造函数负责创建DataActorConfig并交给DataActorCore::new。DataActorConfig在 data_actor.rs 中定义共有三个字段字段类型默认值说明actor_idOptionActorIdNone自定义 Actor 标识符log_eventsbooltrue是否记录事件日志log_commandsbooltrue是否记录命令日志由于除actor_id外的字段都有Default实现构造函数中通常只显式指定actor_id其余用..Default::default()补齐impl SpreadMonitor { pub fn new(instrument_id: InstrumentId) - Self { let config DataActorConfig { actor_id: Some(ActorId::from(SPREAD_MON-001)), ..Default::default() }; Self { core: DataActorCore::new(config), instrument_id, } } }关于actor_id有两个重要事实均可在源码中确认未配置actor_id时的默认值不是类型名。DataActorCore::new通过config.actor_id.unwrap_or_else(Self::default_actor_id)回退而default_actor_id返回ActorId::from(stringify!(DataActor))即字符串DataActordata_actor.rs 与 data_actor.rs。也就是说无论 Actor 是什么类型只要不显式指定 ID注册时都会是DataActor。重复 ID 会在注册时被拒绝。当同一运行环境存在多个同类型 Actor 实例时必须为每个实例显式指定唯一的actor_id。另外DataActorCore::new的源码data_actor.rs显示trader_id、clock、cache初始均为None/空全部依赖注册阶段注入——这也解释了为什么在注册前调用clock()、cache()会 panic。3. 接入运行时nautilus_actor! 宏与 Debug 实现nautilus_actor!宏负责把 Actor 结构体中的DataActorCore字段与运行时契约连接起来。它的定义位于 crates/common/src/macros.rs逻辑非常直接为该类型实现DataActorNativetrait 的core()与core_mut()两个方法默认委托给名为core的字段若字段名不同可传入第二个参数指定// 字段名为 core 时的默认形式 nautilus_actor!(SpreadMonitor); // 字段名为 actor_core 时的显式形式 nautilus_actor!(MyActor, actor_core);宏展开后的DataActorNative实现为impl nautilus_common::actor::DataActorNative for SpreadMonitor { fn core(self) - DataActorCore { self.core } fn core_mut(mut self) - mut DataActorCore { mut self.core } }随后为 Actor 实现Debug。可以手动实现只打印关键字段避免把庞大的订阅映射带进日志也可以直接派生nautilus_actor!(SpreadMonitor); impl std::fmt::Debug for SpreadMonitor { fn fmt(self, f: mut std::fmt::Formatter_) - std::fmt::Result { f.debug_struct(SpreadMonitor).finish() } }需要强调的是普通回调代码不要调用宏生成的core()/core_mut()原生访问器而应使用DataActor门面方法详见第 5 节。运行时的注册与组件装配由 blanket 的Actor、Componenttrait 实现统一提供宏只负责最底层的字段接线。4. 实现 DataActor trait生命周期与数据处理DataActortrait 定义在 data_actor.rs。它包含生命周期回调、数据处理器和订阅/请求方法三大部分所有处理器都有默认实现多数为空操作on_start/on_stop/on_resume/on_reset的默认实现会输出log::warn!提示未重写因此只需要覆写自己关心的那部分每个处理器返回anyhow::Result()。SpreadMonitor覆写了on_start发起订阅与on_quote计算价差impl DataActor for SpreadMonitor { fn on_start(mut self) - anyhow::Result() { self.subscribe_quotes(self.instrument_id, None, None); Ok(()) } fn on_quote(mut self, quote: QuoteTick) - anyhow::Result() { let spread quote.ask_price.as_f64() - quote.bid_price.as_f64(); log::info!(Spread: {spread:.5}); Ok(()) } }subscribe_quotes直接通过DataActortrait 暴露在self上。从源码看它的底层实现在DataActorCore中data_actor.rs先构造一个DataCommand::Subscribe(SubscribeCommand::Quotes(...))命令自动携带instrument_id的 venue 作为venue字段再通过add_quote_subscription在本地注册类型化处理器成功后把命令发给消息总线。订阅失败如重复订阅同一 topic不会发出命令而是记录log::warn!——这一点对调试很有价值。生命周期处理器Actor 在注册、启动、停止、恢复、降级、故障、重置、销毁之间流转完整状态机见 Actors 概念指南 的 Lifecycle 小节。常用生命周期钩子如下方法调用时机on_start()Actor 启动时适合在此订阅数据on_stop()Actor 停止时清理 Actor 拥有的资源on_resume()Actor 从 stopped/degraded 恢复时on_reset()Actor 重置时包括回测轮次之间钩子成功后释放保留的数据订阅on_degrade()Actor 进入降级状态、仅提供部分功能时on_fault()Actor 遭遇故障进入 faulted 状态时on_dispose()Actor 被销毁、必须释放剩余资源时以on_reset为例默认实现会警告“重置指标和其他状态”提示典型用法是复位 Actor 内维护的指标与累积状态以便在下一轮回测中干净启动。数据处理器总览DataActortrait 中可覆写的数据处理器覆盖了 Nautilus Trader 支持的全部市场数据类型方法签名均可从 data_actor.rs 核对行情流on_quote报价、on_trade成交、on_barK 线、on_book/on_book_deltas/on_book_depth订单簿快照/增量/深度衍生价格on_mark_price标记价、on_index_price指数价、on_funding_rate资金费率合约与期权on_instrument合约定义、on_instrument_status/on_instrument_close合约状态/收盘、on_option_greeks/on_option_chain期权希腊值与链切片自定义与信号on_data自定义数据、on_signal信号系统事件on_time_event定时器/提醒、on_queue_state运行队列压力、on_socket_stateSocket 传输状态历史数据on_historical_bars/on_historical_quotes/on_historical_trades/on_historical_book_deltas/on_historical_book_depth/on_historical_data等。一个关键的运行时语义是请求响应与订阅更新走不同的回调通过request_bars()、request_quotes()等请求获取的历史数据进入on_historical_*系列处理器而通过subscribe_bars()、subscribe_quotes()等订阅的持续数据进入on_bar()、on_quote()等处理器。完整映射表见 Actors 概念指南 的 Callback handlers 小节。调试数据流问题时先确认数据来自请求还是订阅避免在错误的处理器中等待数据。DataActor还提供handle_*系列内部方法如handle_quote、handle_data、handle_signal它们先记录接收日志、检查not_running()状态再分发到对应的on_*处理器并捕获错误data_actor.rs 附近。这也解释了“组件未运行时数据会被静默丢弃并打日志”的行为。5. 原生运行时访问门面 vs DataActorNativeActor 的运行时访问分两层明确区分它们能避免写出无法跨 Python 边界移植的代码DataActor门面默认选择所有用户代码应优先使用DataActortrait 上的门面方法包括只读属性config()、actor_id()、trader_id()、is_registered()以及clock()、cache()、shutdown_system()、publish_data()、publish_signal()、add_synthetic()等能力方法。门面方法在 data_actor.rs只读属性与同文件后续章节能力方法中定义。DataActorNative仅在需要显式原生访问时使用定义于 data_actor.rs提供core()、core_mut()、clock_mut()、clock_rc()、cache_ref()、cache_rc()等返回原生借用/引用计数的接口。这些类型不能跨越 Python 边界因此只应在与引擎同构的同一原生二进制中的性能敏感路径或宿主集成代码中使用面向 Python 或插件编写面plug-in authoring surface的策略/ Actor 代码不要导入它。注意clock()、cache()等门面方法在 Actor 注册前调用会 panic源码中通过unwrap_or_else(|| panic!(...))实现见 data_actor.rs。所以“注册前只构造、不访问运行时状态”是编写构造函数时的硬约束。Native traits 的适用性矩阵与DataActorNative方法完整表格可参考 Rust 概念指南 的 Native traits 与 DataActorNative methods 小节。6. 注册 ActorBacktestEngine 与 LiveNodeActor 编写完成后注册方式取决于运行载体回测环境BacktestEnginelet actor SpreadMonitor::new(instrument_id); engine.add_actor(actor)?;实盘环境LiveNodelet actor SpreadMonitor::new(instrument_id); node.add_actor(actor)?;从DataActorCore的字段设计可以推断注册阶段发生的事运行时把trader_id、clock、cache注入 core并把 Actor 纳入组件注册表。因此add_actor必须在引擎/节点启动前调用且同一 trader 下不允许注册重复的actor_idActors 概念指南 明确说明 Python 侧重复 ID 会抛RuntimeErrorRust 侧同样在注册时拒绝。完整可运行的回测搭建示例可参考 run_rust_backtest 指南 与 run_rust_live_trading 指南。7. Guard 安全回调中访问其他 Actor 的规则当系统向你的 Actor 分发消息时它会从注册表获取一个短生命周期的ActorRefguard。这些 guard 由运行时管理你不需要也不应该直接创建它们。但如果你在回调代码中需要访问其他 Actor必须遵守以下三条规则来自原指南并可在运行时注册表设计中印证每次按 ID 查找 Actor绝不缓存ActorRef在作用域结束前释放 guard绝不把它存进字段绝不让 guard 跨越.await点存活。DataActorCore上的订阅方法天然满足这些约束它们在闭包内捕获 actor ID 并执行查找从而在正确的作用域内使用 guard。完整的线程模型与注册表语义参见 Rust 开发者指南 的 Runtime invariants 小节。8. 完整示例从 SpreadMonitor 到 BookImbalanceActorSpreadMonitor展示了一个最小 Actor 的全部骨架。仓库中提供了更完整的生产级示例——BookImbalanceActor订单簿失衡监控 Actor位于 crates/trading/src/examples/actors/imbalance/actor.rs。它与SpreadMonitor的差异恰好展示了进阶模式多合约状态管理用states: AHashMapInstrumentId, ImbalanceState按合约维护累积状态ImbalanceState记录update_count、bid_volume_total、ask_volume_total并提供imbalance()方法计算累计买卖量失衡pub fn imbalance(self) - f64 { let total self.bid_volume_total self.ask_volume_total; if total 0.0 { (self.bid_volume_total - self.ask_volume_total) / total } else { 0.0 } }构造函数参数化new(instrument_ids, log_interval, actor_id)支持多个合约、可配置的日志间隔与可选actor_idNone时回退到BOOK_IMBALANCE-001并提供了from_config(config)工厂方法对接配置对象。on_start批量订阅在on_start中对每个instrument_id调用subscribe_book_deltas(instrument_id, BookType::L2_MBP, None, None, false, None)订阅 L2 订单簿增量。on_book_deltas增量聚合遍历deltas.deltas按delta.order.side汇总买卖挂单量累加到对应合约的状态并按log_interval周期性打印进度。on_stop输出汇总在停止时调用print_summary()按合约 ID 排序打印更新次数、累计买卖量与失衡值。自定义Debug只打印instrument_ids与log_interval避免把大状态表带入调试输出。impl DataActor for BookImbalanceActor { fn on_start(mut self) - anyhow::Result() { let ids self.instrument_ids.clone(); for instrument_id in ids { self.subscribe_book_deltas( instrument_id, BookType::L2_MBP, None, // depth None, // client_id false, // managed None, // params ); } Ok(()) } fn on_stop(mut self) - anyhow::Result() { self.print_summary(); Ok(()) } fn on_book_deltas(mut self, deltas: OrderBookDeltas) - anyhow::Result() { let mut bid_volume 0.0; let mut ask_volume 0.0; for delta in deltas.deltas { let size delta.order.size.as_f64(); match delta.order.side { Some(OrderSide::Buy) bid_volume size, Some(OrderSide::Sell) ask_volume size, None {} } } let state self.states.entry(deltas.instrument_id).or_default(); state.update_count 1; state.bid_volume_total bid_volume; state.ask_volume_total ask_volume; if self.log_interval 0 state.update_count.is_multiple_of(self.log_interval) { println!( [{}] update #{}: batch bid{:.2} ask{:.2} cumulative imbalance{:.4}, deltas.instrument_id, state.update_count, bid_volume, ask_volume, state.imbalance(), ); } Ok(()) } }该示例的测试位于 crates/trading/src/examples/actors/imbalance/tests.rs是验证 Actor 行为订阅、数据累积、汇总输出的参考实现。9. 进阶能力与相关指南Rust Actor 还可以利用以下能力相关实现均可从源码与文档继续深入定时器与提醒通过clock()门面调度周期定时器与一次性提醒set_timer/set_time_alert并用on_time_event或显式回调处理TimeEvent定时器名与提醒名共享时钟命名空间同名注册会替换旧定时器见 Actors 概念指南。队列压力与 Socket 状态监控subscribe_queue_state(Some(50))/subscribe_socket_state(Some(50))可订阅运行队列压力与 Socket 传输状态变化处理函数on_queue_state/on_socket_state优先级参数控制同主题订阅者之间的投递顺序数值越大越先执行重复订阅不会改变已有优先级需要先退订再以新优先级订阅。SocketState.CONNECTED只表示传输层可用不代表认证、订阅回放或适配器就绪。Socket 重连请求实盘环境下reconnect_socket(client_id, endpoint)可请求恢复单个端点但它是 fire-and-observe 语义——返回成功只表示命令通过本地校验并入队需通过SocketStateChanged事件确认恢复结果。消息发布publish_data()按数据类型派生的 topic 发布自定义数据publish_signal()发布轻量信号。进一步阅读Actors 概念指南Actor 能力、生命周期状态机、数据处理器映射表、队列/Socket 事件语义Rust 概念指南handler 方法总表、native traits 适用性矩阵与DataActorNative方法表Rust 开发者指南线程模型与运行时不变量Runtime invariants用 Rust 编写策略在 Actor 能力之上增加订单管理运行 Rust 回测 与 运行 Rust 实盘交易BacktestEngine与LiveNode的完整搭建流程。【免费下载链接】nautilus_traderProduction-grade Rust-native trading engine with deterministic event-driven architecture项目地址: https://gitcode.com/GitHub_Trending/na/nautilus_trader创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表