ARTICLE DETAIL

资讯详情

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

NautilusTrader 基础设施解密:nautilus-infrastructure 的 Redis 缓存、消息总线与 PostgreSQL 存储架构

NautilusTrader 基础设施解密:nautilus-infrastructure 的 Redis 缓存、消息总线与 PostgreSQL 存储架构 NautilusTrader 基础设施解密nautilus-infrastructure 的 Redis 缓存、消息总线与 PostgreSQL 存储架构【免费下载链接】nautilus_traderProduction-grade Rust-native trading engine with deterministic event-driven architecture项目地址: https://gitcode.com/GitHub_Trending/na/nautilus_trader本指南深入剖析 NautilusTrader 开源交易引擎中的基础设施层组件——nautilus-infrastructurecrate。它以 Redis 与 PostgreSQL 为后端为引擎提供缓存数据库Cache Database与消息总线Message Bus两类核心基础设施支撑系统从开发环境平滑扩展到生产部署。读完本文你将掌握该 crate 的 feature flag 编译开关、RedisCacheConfig/RedisMessageBusConfig/PostgresCacheConfig的完整配置字段与默认值、Redis 键设计与流修剪机制以及它与nautilus-common中CacheConfig、MessageBusConfig的协作关系。一、nautilus-infrastructure 在引擎架构中的定位NautilusTrader 是一个开源、生产级、Rust 原生的多资产多交易所交易引擎在单一事件驱动架构中横跨研究、确定性模拟与实盘执行三个阶段提供研究到实盘research-to-live的语义一致性。nautilus-infrastructure正是为这套架构提供地基的 crate它实现了引擎在运行期持久化状态与对外通信所需的底层后端而引擎本身并不关心后端是 Redis 还是 PostgreSQL——抽象通过nautilus-common中的 trait 完成。从 crates/infrastructure/src/lib.rs 的 crate 级文档可以看到其定位与能力清单Redis 集成基于 Redis 实现的缓存数据库cache database与消息总线message bus后端PostgreSQL 集成基于 SQL 的缓存数据库带完整数据模型full data models连接管理带重试逻辑与健康监控的连接处理序列化选项支持 JSON 与 MessagePack 两种编码格式Python 绑定通过 PyO3 提供 Python 互操作。crate 通过 feature flag 支持多种数据库后端使用者可针对自身部署需求与规模选择合适的基础设施组件。模块的组织方式也印证了这一点见 lib.rs#[cfg(feature python)] pub mod python; #[cfg(feature redis)] pub mod redis; #[cfg(feature postgres)] pub mod sql;redis模块下包含cache缓存数据库、msgbus消息总线、queries查询封装与stream_fields流字段定义sql模块下则包含cache、pg连接选项与modelsaccounts、data、enums、general、instruments、orders、positions、types 等完整数据模型。二、Feature flags按需裁剪编译内容nautilus-infrastructure提供 4 个 feature flag 来控制编译期源码包含范围见 crates/infrastructure/Cargo.tomlFeature作用redis默认启用 Redis 缓存数据库与消息总线后端实现依赖dep:redispostgres启用 PostgreSQL SQLx 模型与缓存数据库后端依赖dep:sqlxpython启用基于 PyO3 的 Python 绑定extension-module以 Python 扩展模块cdylib方式构建三个开关之间存在依赖关系python会连带启用nautilus-common/python、nautilus-core/python、nautilus-model/python、pyo3、pyo3-async-runtimes、pyo3-stub-genextension-module则进一步要求nautilus-common/extension-module、nautilus-core/extension-module、nautilus-model/extension-module以及pyo3/extension-module。需要注意两个细节默认特性是redisCargo.toml 中的注释明确说明目前nautilus_trader默认需要 redisdefault [redis] # redis needed by nautilus_trader by default for now即引擎主 crate 默认依赖 Redis 后端docs.rs 文档构建同时启用postgres与redis见 Cargo.toml 的 metadata 段并注明 Python 文档依赖本地 pyo3-stub-gen patch发布到 crates.io 的 crate 并不携带该补丁。从 python/mod.rs 的模块初始化代码可以看到infrastructurePython 模块会按编译特性注册对应类启用redis时注册RedisCacheConfig、RedisCacheDatabase、PyRedisMessageBusBacking、PyRedisMessageBusFactory、RedisMessageBusConfig并注册缓存/消息总线工厂启用postgres时注册PostgresCacheConfig、PostgresCacheDatabase、PostgresConnectOptions并注册 Postgres 缓存数据库工厂。三、Redis 缓存数据库RedisCacheDatabase3.1 配置字段与默认值RedisCacheConfig定义在 crates/infrastructure/src/redis/cache.rs要求Redis 6.2 或更高版本才能正确运行。所有字段均带#[serde(default, deny_unknown_fields)]即可以从 JSON 反序列化且拒绝未知字段字段类型默认值说明hostOptionStringNone即127.0.0.1Redis 主机地址portOptionu16None即6379Redis 端口usernameOptionStringNoneRedis 账户用户名passwordOptionStringNoneRedis 账户密码sslboolfalse是否使用 SSL 加密连接rediss://connection_timeoutu1620建立新连接的超时秒response_timeoutu1620等待响应的超时秒number_of_retriesusize100连接重试次数指数退避exponent_baseu642指数退避的底数max_delayu641000重试之间的最大延迟秒factoru642重试延迟计算的乘法因子该结构体实现了RedisConnectionConfigtrait见 crates/infrastructure/src/redis/mod.rs该 trait 是所有 Redis 相关配置缓存与消息总线共享的连接参数抽象。3.2 连接管理URL 构造、版本检查与指数退避get_redis_url函数redis/mod.rs根据配置构造连接 URL并返回一份密码脱敏版用于日志与监控保留密码首尾各两个字符中间用...代替。其认证矩阵如下用户名密码生成的 user-info 部分非空非空user:pass空非空:pass空空省略一个重要的约束提供了用户名却没有密码会直接 panic报错信息提示请提供密码或省略用户名Redis config error: username supplied without password。ssl: true时 scheme 切换为rediss://。create_redis_connectionredis/mod.rs完成实际连接建立其重连策略为重试number_of_retries次延迟按factor * (exponent_base ^ 当前尝试次数)指数增长上限为max_delay新连接的操作超时为response_timeout每次连接尝试超时为connection_timeout。连接成功后还会执行INFO命令读取redis_version与最低版本6.2.0常量REDIS_MIN_VERSION比对并记录日志——低于最低版本时输出 error 级别日志。单元测试覆盖了 URL 构造的各种组合见 redis/mod.rs 测试段例如仅密码时生成redis://:secretpwexample.com:6380与脱敏版redis://:se...pwexample.com:6380带 SSL 时生成rediss://user:passexample.com:6380。3.3 双连接架构READ 与 WRITE 分离RedisCacheDatabase采用双连接架构见 cache.rs 模块文档READ 连接self.con用于同步查询keys、read、load_all归主结构体所有WRITE 连接由后台任务cache-process持有在get_runtime()上创建通过无界tokio::sync::mpsc通道接收命令。所有写操作insert、update、delete、flush都经由命令通道路由到 WRITE 连接执行。这样设计的原因是避免跨 runtime 的 I/O 问题——WRITE 连接始终在 Nautilus runtime 上创建。同步调用方close、flushdb_sync使用std::sync::mpsc应答通道阻塞等待后台任务完成当从 Nautilus runtime 自身调用时blocking_recv会自动使用block_in_place避免阻塞 worker 线程见 cache.rs。写操作还支持缓冲批量写入CacheConfig.buffer_interval_ms控制相邻流水线事务之间的缓冲间隔毫秒为 0 时每条消息立即刷新非零时后台任务按定时器周期性调用flush_buffer并将缓冲中的命令通过原子 Redis Pipelinepipe.atomic()一次性执行见drain_buffercache.rs。3.4 键设计集合、索引与 trader 命名空间缓存数据在 Redis 中以集合collection为单位组织每个集合键使用:常量REDIS_DELIMITER分隔。集合键常量cache.rs包括index、general、currencies、instruments、instrument_closes、synthetics、accounts、orders、positions、actors、strategies、snapshots、health、custom。订单与持仓还维护了一组索引键cache.rsindex:order_ids、index:order_position、index:order_client、index:orders、index:orders_open、index:orders_closed、index:orders_emulated、index:orders_inflight、index:positions、index:positions_open、index:positions_closed。get_index_keyredis/mod.rs可从完整键中提取索引部分例如trader-id:uuid:index:order_position→index:order_position。所有键都带有trader 命名空间前缀。get_trader_keycache.rs根据CacheConfig构造默认前缀为trader-use_trader_prefix随后是trader_id若use_instance_id为 true 则追加:instance_id。不同集合映射到不同的Redis 数据类型见insert函数的分发逻辑cache.rs集合Redis 命令/类型index集合索引用SADDSet哈希索引用HSETHashgeneral、currencies、instruments、instrument_closes、synthetics、actors、strategies、health、customSETStringaccounts、orders、positions、snapshotsRPUSHList事件追加型账户、订单、持仓以事件列表list of events形式追加存储replace_list则先用DEL删除再RPUSH用于add_order/add_position的初始写入。删除操作会同步清理相关索引delete_order会依次清理订单本身、6 个订单集合索引与 2 个哈希索引见 cache.rs。订单状态更新UpdateOrder操作会先读取完整事件列表、反序列化重建订单、追加新事件后再更新索引update_order_event_logcache.rs。3.5 序列化JSON 与 MessagePackCacheConfig.encoding控制缓存操作的序列化编码默认 JSON。DatabaseQueries::serialize_payload/deserialize_payloadqueries.rs支持MsgPackrmp_serde与Json两种编码并统一做了时间戳转换处理SBE 与 Capn Proto 编码在 Redis 缓存场景明确不支持会直接报错。若CacheConfig.timestamps_as_iso8601为 true时间戳以 ISO 8601 字符串持久化否则以 UNIX 纳秒存储。批量读取方面read_bulk使用MGET单次网络操作取回多键queries.rsread_bulk_batched则按bulk_read_batch_size分块执行MGET以规避部分 Redis 提供商的请求大小限制queries.rs。键扫描使用SCAN游标循环每批COUNT 5000queries.rs。3.6 加载与健康检查RedisCacheDatabaseAdapter通过CacheDatabaseAdaptertrait定义于nautilus-common向引擎暴露统一接口RedisCacheConfig实现CacheDatabaseFactorycache.rs。load_all使用tokio::try_join!并行加载 currencies、instruments、instrument_closes、synthetics、accounts、orders、positions、greeks、yield_curves 九类数据cache.rs。加载类功能支持load_currency、load_instrument、load_order、load_position、load_actor读取actors:{actor_id}:state、load_strategy读取strategies:{strategy_id}:state等细粒度查询。需要注意该适配器的能力边界市场数据quote/trade/bar/funding rate/order book、signal 与订单/持仓快照的加载明确不支持返回anyhow::bail!市场数据持久化属于其他组件如nautilus-persistence的职责。heartbeat方法将格式化后的时间戳写入health:heartbeat键构成健康监控的一部分。四、Redis 消息总线RedisMessageBusBacking4.1 配置连接参数 总线语义消息总线同样要求Redis 6.2。RedisMessageBusConfigcrates/infrastructure/src/redis/msgbus.rs的 11 个连接字段与RedisCacheConfig完全一致含相同的默认值并实现MessageBusBackingFactory用于创建RedisMessageBusBacking。总线自身的语义由nautilus-common中的MessageBusConfig决定crates/common/src/msgbus/config.rs关键字段如下字段默认值说明encodingJson外部发布负载的默认编码encoding_market_data/encoding_builtinNone市场数据 / 内置账户、组合、订单、持仓负载的编码覆盖timestamps_as_iso8601false时间戳是否以 ISO 8601 字符串持久化buffer_interval_msNone流水线批量事务的缓冲间隔ms建议范围[10, 1000]100ms 是良好折中autotrim_minsNone流自动修剪的窗口分钟实际窗口可能比设定值多延长最多 1 分钟至少每分钟修剪一次需 Redis 6.2autotrim_maxlenNone每条流近似保留的最大条目数近似修剪use_trader_prefixtrue流名是否使用trader-前缀use_trader_idtrue流名是否包含 trader IDuse_instance_idfalse流名是否包含实例 IDstreams_prefixstream外部发布流名的前缀stream_per_topictruetrue 时按主题写独立流false 时全部写入同一流external_streamsNone总线监听的、用于回灌反序列化消息的外部流键types_filterNone不对外发布的负载类型名列表heartbeat_interval_secsNone心跳间隔秒get_stream_keyredis/mod.rs按上述开关拼装流键例如use_trader_prefix use_trader_id use_instance_id streams_prefix生成trader-tester-123:{instance_id}:stream全部关闭时仅剩stream。4.2 三个后台任务发布、流读取、心跳RedisMessageBusBacking::newmsgbus.rs在 Nautilus runtime 上创建最多三个后台任务每个任务持有自己的 Redis 连接发布任务msgbus-publish接收无界通道中的BusMessage按stream_per_topic决定流键stream_key:topic用XADD写入 Redis 流publish_messages。支持两种自动修剪策略autotrim_maxlen时使用XADD ... MAXLEN ~近似长度修剪autotrim_mins时以 60 秒为缓冲阈值周期性执行XTRIM ... MINID常量REDIS_XTRIM、REDIS_MINID、TRIM_BUFFER_SECS见 msgbus.rs。流读取任务msgbus-stream仅当配置了external_streams时创建使用XREAD ... BLOCK 100从当前时间戳开始阻塞读取外部流通过容量 100_000 的有界通道回灌到引擎内部stream_messages。关闭时通过ArcAtomicBool信号终止。心跳任务msgbus-heartbeat按heartbeat_interval_secs周期向health:heartbeat主题发布心跳消息run_heartbeat。注意heartbeat_interval_secs Some(0)是非法的直接报错。close是同步方法使用block_in_place桥接进异步关闭路径close_async且必须在任何current_threadTokio runtime 之外调用。处理句柄以OptionJoinHandle存储保证关闭幂等。另外crate 还附带一个基于 criterion 的基准redis_msgbus见 Cargo.toml bench 段 与 benches/redis_msgbus.rs用于度量消息总线的吞吐与延迟。五、PostgreSQL 缓存数据库PostgresCacheDatabase5.1 连接选项参数、环境变量与默认值PostgresConnectOptionscrates/infrastructure/src/sql/pg.rs包含 5 个字段host、port、username、password、database。其默认值default_administrator为localhost:5432、用户nautilus、密码pass、数据库nautilus——这是本地开发/测试的约定配置。get_postgres_connect_optionspg.rs按显式参数 → 环境变量 → 默认值的优先级合并配置支持的环境变量包括POSTGRES_HOSTPOSTGRES_PORTPOSTGRES_USERNAMEPOSTGRES_PASSWORDPOSTGRES_DATABASEconnect_pg通过sqlx的PgPool::connect_with建立连接池连接字符串构造提供connection_string()与脱敏版connection_string_masked()密码替换为***。转换为PgConnectOptions时调用了disable_statement_logging()以避免在日志中泄露语句。此外validate_sql_identifier强制表名等标识符仅含 ASCII 字母数字与下划线escape_sql_string对字符串做单引号转义体现对 SQL 注入的基础防护。5.2 缓存配置与数据模型PostgresCacheConfigcrates/infrastructure/src/sql/cache.rs字段与连接选项一致host、port、username、password、database均为Option缺省时从环境变量与内置默认值解析实现了CacheDatabaseFactory以接入引擎缓存。其 SQL 查询封装在sql/queries.rs数据库命令处理同样采用后台任务 缓冲 Pipeline 的模式。完整数据模型体现在 crates/infrastructure/src/sql/models/ 目录accounts.rs、data.rs、enums.rs、general.rs、instruments.rs、orders.rs、positions.rs、types.rs。对应的建表、分区与函数脚本位于仓库根目录的 schema/sql/tables.sql表结构、types.sql自定义类型、partitions.sql分区、functions.sql函数。值得注意Postgres 后端要求先执行nautilus database init迁移常量SCHEMA_MIGRATION_COMMAND的提示文案为run nautilus database init to migrate见 sql/cache.rs迁移命令由 crates/cli 提供。5.3 能力边界与 Redis 缓存适配器类似Postgres 适配器同样聚焦账户事件、订单事件、持仓事件、通用键值状态与快照的持久化市场数据quotes、trades、bars、funding rates、order book与 signal 的保存/加载不在其职责范围内。选择 Postgres 的意义在于 SQL 查询能力、事务性与长期归档场景而 Redis 则更侧重低延迟读写。六、Python 绑定与使用方式nautilus-infrastructure通过 PyO3 暴露 Python 接口统一挂载在nautilus_trader.infrastructure模块下见 python/mod.rs。Python 侧可用的类包括RedisCacheConfig、RedisCacheDatabase、RedisMessageBusConfig、RedisMessageBusFactory、PyRedisMessageBusBacking以及PostgresCacheConfig、PostgresCacheDatabase、PostgresConnectOptions。Python 用户构建扩展模块时需启用extension-module并连带python特性Rust-only 构建则无需启用任何 Python 特性。类型桩由pyo3-stub-gen自动生成与 python/nautilus_trader/infrastructure/ 目录下的.py/.pyi文件对应Python 层对配置对象的校验逻辑同样复用 Rust 侧的deny_unknown_fields与默认值语义。七、测试、文档与工程规范该 crate 的测试资产集中在 crates/infrastructure/tests/5 个集成测试文件覆盖 Redis/Postgres 后端的端到端行为TESTS.md记录了相关测试说明。dev-dependencies 引入了nautilus-live启用node特性、nautilus-persistence、nautilus-serialization与nautilus-model启用test-support用于构造真实的引擎运行环境验证基础设施组件。源码工程规范方面crate 启用了严格的 lint#![deny(unsafe_code)]、#![warn(clippy::pedantic)]、#![deny(missing_debug_implementations)]、#![deny(clippy::missing_errors_doc)]、#![deny(clippy::missing_panics_doc)]与#![deny(rustdoc::broken_intra_doc_links)]见 lib.rs并且为similar_names等 clippy 规则给出了 domain 层面的豁免理由如trader_id/trade_id、price_precision/size_precision是有意为之的平行命名。crate 依赖的核心 Nautilus 组件包括nautilus-common启用live特性、nautilus-core、nautilus-cryptography与nautilus-model。八、许可证与归属nautilus-infrastructure的源代码以 GNU Lesser General Public License v3.0 授权版权归 Nautech Systems Pty Ltd2015-2026所有。NautilusTrader 商标及其扩展组件由 Nautech Systems 持续开发维护使用本软件前请阅读项目根目录的 LICENSE 文件了解完整条款。延伸阅读若想进一步理解引擎如何消费这些基础设施可继续阅读nautilus-common中的CacheConfigcrates/common/src/cache/config.rs含tick_capacity/bar_capacity范围校验[1, 1_000_000]与persist_account_events等字段与MessageBusConfigcrates/common/src/msgbus/config.rs了解数据持久化管线可参考 crates/persistence运行时如何组织连接与任务可参考 docs/concepts/event_sourcing.md。【免费下载链接】nautilus_traderProduction-grade Rust-native trading engine with deterministic event-driven architecture项目地址: https://gitcode.com/GitHub_Trending/na/nautilus_trader创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表