NautilusTrader 如何配置 Redis 作为缓存数据库与消息总线后端?
【免费下载链接】nautilus_traderProduction-grade Rust-native trading engine with deterministic event-driven architecture项目地址: https://gitcode.com/GitHub_Trending/na/nautilus_trader
NautilusTrader 中 Redis 是可选组件,只有当你要把缓存数据库(Cache backing)或消息总线(MessageBus backing)配置为 Redis 后端时才需要它。本文的任务就是完成这两项配置:让LiveNode用 Redis 持久化订单、持仓、账户等缓存记录以便重启后恢复,同时让消息总线通过 Redis Streams 把序列化消息转发给其他节点。适用前提是已安装nautilus_trader,并且 Redis 版本不低于 6.2(消息总线依赖 Streams 功能,最低支持版本 6.2,见 安装文档)。
准备 Redis 实例
文档推荐用 Docker 容器快速起一个本地 Redis,安装文档给出的命令是:
docker run -d --name redis -p 6379:6379 redis:latest该命令会从 Docker Hub 拉取最新 Redis 镜像(若本地已有则跳过),以分离模式(-d)运行一个名为redis的容器,并把 6379 端口映射到本机。容器日常操作用docker start redis启动、docker stop redis停止。
仓库的 .docker/docker-compose.yml 中也包含一个redis服务(容器名nautilus-redis,镜像public.ecr.aws/docker/library/redis,端口绑定在127.0.0.1:6379),适合同时需要 Postgres 的场景。
文档建议用 Redis Insight 这类 GUI 工具来可视化和调试 Redis 中的数据,后文验证时会用到。
配置缓存数据库后端
缓存数据库的分工是:CacheConfig控制缓存行为(编码、缓冲、容量等),连接参数放在RedisCacheConfig里。Python 侧通过LiveNodeBuilder.with_cache_database_factory注入,节点启动时自行构造并持有适配器,因此连接只在节点run()时才真正建立:
from nautilus_trader.common import Environment from nautilus_trader.infrastructure import RedisCacheConfig from nautilus_trader.live import LiveNode from nautilus_trader.model import TraderId node = ( LiveNode.builder("LiveNode", TraderId("TRADER-001"), Environment.LIVE) .with_cache_database_factory(RedisCacheConfig(host="localhost", port=6379)) .build() ) try: node.run() finally: node.dispose()RedisCacheConfig来自nautilus_trader.infrastructure,host和port替换为你实际部署的 Redis 地址;文档示例使用本机localhost:6379。Rust 原生调用方则构造RedisCacheConfig(含host、port、username、password、connection_timeout、response_timeout等字段)并通过CacheDatabaseFactorytrait 的create(trader_id, instance_id, config)生成适配器,再在启动前用node.set_cache_database(cache_database)?;挂到节点上,完整示例见 缓存概念文档。
几个影响运行行为的配置点:
- 后端是恢复机制,不是完整事件归档,也不是同步的分布式缓存。重启后可恢复的记录包括 general data、currencies、instruments、instrument closes、accounts、orders 和 positions;有界的市场数据历史和运行中的进程不会被恢复。
- 在默认的
LiveExecutionEngineConfig.load_cache = true下,节点会在连接客户端、做执行状态核对之前先恢复持久化缓存状态并重建派生索引;把CacheConfig.flush_on_start设为true则是相反操作——清空后端而不是恢复。 - 必须用
run()运行带缓存数据库后端的节点,run_async()会拒绝这种会阻塞宿主事件循环的后端。传入其他类型对象会抛NotImplementedError。 - 文档明确要求节点结束时调用
dispose():它会关闭后端并冲刷CacheConfig.buffer_interval_ms缓冲中尚未写入的数据,直接从run()返回可能丢掉最后这批写入。
配置消息总线后端
消息总线的分工类似:MessageBusConfig控制总线行为,RedisMessageBusConfig持有 Redis 连接参数并实现MessageBusBackingFactory。注意MessageBusConfig单独使用不会安装任何后端,必须再配合 factory 调用,否则会走默认路径而没有 Redis 转发:
from nautilus_trader.common import Environment from nautilus_trader.common import MessageBusConfig from nautilus_trader.infrastructure import RedisMessageBusConfig from nautilus_trader.live import LiveNode from nautilus_trader.model import TraderId trader_id = TraderId("TRADER-001") message_bus = MessageBusConfig( external_streams=["external-stream"], stream_per_topic=False, ) redis_config = RedisMessageBusConfig( host="localhost", port=6379, ) node = ( LiveNode.builder("LiveNode", trader_id, Environment.LIVE) .with_msgbus_config(message_bus) .with_external_msgbus_factory(redis_config) .build() ) node.run()Rust 侧等价做法是LiveNodeBuilder::with_external_msgbus_factory(Box::new(redis_config)),已有代码也可以继续传RedisMessageBusFactory(redis_config)。factory 总是会安装外部 egress;external_streams配置了哪些流名,run()时就会消费哪些流。注意内置 ingress 从当前时间戳开始消费每个流,节点启动前已存在于流中的条目不会被重放;需要启动前重放时应使用缓存恢复或 event store。
常用行为选项(均在MessageBusConfig中):
stream_per_topic:Redis 不支持通配符流主题,文档建议设为False,此时所有消息写同一个基础流 key,topic 保留为消息字段;设为True则按 topic 追加到流 key 上。- 流 key 结构由
use_trader_prefix、use_trader_id、use_instance_id、streams_prefix控制;默认选项下基础流 key 为trader-{trader_id}:{streams_prefix},完整形式为trader-{trader_id}:{instance_id}:{streams_prefix}。 types_filter:传入要排除的 payload 类型名列表,防止高频 quote 等数据灌满流,例如MessageBusConfig(types_filter=["QuoteTick", "TradeTick"])。autotrim_mins与autotrim_maxlen:分别按时间窗口和条目数修剪流;两者同时设置时超出任一阈值即修剪。Redis 侧每条流每分钟最多修剪一次,条目可能比窗口多留约一分钟。- 编码:默认
json;关注 payload 体积和序列化性能时用msgpack。Redis 缓存 payload 路径只支持 MessagePack 和 JSON,SBE 与 Cap'n Proto 是外部 egress 的 schema 编码,不适用于 Redis 缓存。
生产者/消费者节点转发行情(可选分支)
如果目标是让一个节点把行情转发给下游节点,文档给出一个完整示例:生产者节点把streams_prefix设为"binance"、use_trader_id/use_trader_prefix/use_instance_id全部设为false,得到简单可预测的流 key;消费者节点用MessageBusConfig(external_streams=Some(vec!["binance".to_string()]))接收同一流,并把RedisMessageBusConfig传给with_external_msgbus_factory。完整 Rust 代码见 消息总线概念文档。
消费者侧还有一步配套配置:在LiveDataEngineConfig里把外部客户端 ID 加入external_clients(示例中为"BINANCE_EXT"),这样DataEngine不会向这个客户端发送订阅命令,RustDataEngine跳过外部客户端订阅时会为对应 payload 类型注册入站重发。外部生产者若直接写 Redis 流,每条消息必须包含topic、type、payload字段,缺少type的条目会被接收节点跳过。
启动、验证与限制
按文档给出的行为,验证配置是否生效的路径如下:
- 连接成功:带缓存数据库后端的节点在数据库连接失败时会直接使
run()失败。因此节点能正常跑起来,说明 Redis 连接已建立;起不来的情况先检查容器是否在运行(docker start redis)、host/port是否指向同一个实例。 - 数据落库:用 Redis Insight 之类的 GUI 查看 Redis 中的数据,确认缓存后端和消息流按预期写入。
- 重启恢复:重启节点后,检查 orders、positions、accounts、instruments 等恢复记录是否按 缓存文档 描述的范围恢复;确认
load_cache未被置为false、flush_on_start未被设为true。 - 跨节点消费:消费者节点收到外部流消息后会把反序列化的 payload 发布到其内部消息总线,可以通过内部订阅确认消息到达。
文档明确给出的限制,配置前需要知道:
- 缓存后端指向同一个数据库命名空间不会让多个节点的缓存保持一致,每个节点各自持有内存缓存。
- actor 与策略的状态持久化(
with_load_state/with_save_state)只支持 Redis 后端;Postgres 适配器只支持缓存状态,两者不要混用。 - 状态持久化不是连续检查点:内核每次运行最多保存一次状态,
SIGKILL或崩溃会丢失上次保存后的所有变更,所以节点结束时要用dispose()而不是直接返回。 - 一个进程内只支持一个
LiveNode,多个节点需分进程运行。
更完整的参数参考见 配置实盘交易节点 和 消息总线概念文档。
【免费下载链接】nautilus_traderProduction-grade Rust-native trading engine with deterministic event-driven architecture项目地址: https://gitcode.com/GitHub_Trending/na/nautilus_trader
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考