news 2026/9/13 12:21:09

Neon 存储消息传递架构:从 Safekeeper Gossip 到集中式 Storage Broker 的设计与落地

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Neon 存储消息传递架构:从 Safekeeper Gossip 到集中式 Storage Broker 的设计与落地

Neon 存储消息传递架构:从 Safekeeper Gossip 到集中式 Storage Broker 的设计与落地

【免费下载链接】neonNeon: Serverless Postgres. We separated storage and compute to offer autoscaling, code-like database branching, and scale to zero.项目地址: https://gitcode.com/GitHub_Trending/ne/neon

本篇技术指南以 Neon 仓库中的设计文档 docs/rfcs/015-storage-messaging.md 为主体,系统讲解 Neon 存储层(safekeeper / pageserver)如何通过一个集中式的消息/状态中枢来解决 WAL 裁剪、S3 上传决策、WAL 转发与连接生命周期管理等核心问题;同时结合仓库中storage_brokercrate、safekeeper 与 pageserver 的实际实现,帮助读者理解该 RFC 从设计到落地的完整脉络,掌握其中的核心概念与工程权衡。

1. 背景与动机:存储层要解决的四个问题

在 Neon 的"存储与计算分离"架构中,一条 timeline 的 WAL 会同时被写入多个 safekeeper(副本,用于容错),而 pageserver 需要消费这些 WAL 来更新页面数据。为了保证数据一致性、避免存储无限增长,系统必须解决以下四个问题:

  1. 在 safekeeper 上裁剪(Trim)WAL:已确认持久化的 WAL 可以被删除,防止磁盘被历史日志占满;
  2. 决定哪个 safekeeper 负责把 WAL 上传到 S3:避免多个副本重复上传,同时保证至少有一个在做;
  3. 决定哪个 safekeeper 把 WAL 转发给 pageserver:pageserver 需要一个"最新、最健康"的 WAL 源,并且在当前源不可用时快速切换;
  4. 决定何时关闭 safekeeper ↔ pageserver 连接:当 compute 断开、或 safekeeper 落后/宕机时,及时释放资源并切换数据源。

该 RFC(创建于 2022-01-19,最初由 @kelvich 提出)给出了比早期 014-safekeepers-gossip RFC 更通用、更易维护的解法:用集中式的"存储消息中枢"取代节点间的 P2P gossip。正如 RFC 中明确指出的,与 gossip 方案不同,该方案并不试图解决 safekeeper 之间的直接恢复,而是把 gossip 方案中纠缠在一起的两类问题拆开、分别处理。

该 RFC 同时声明了非目标(Non-goals):它并不打算把"compute→pageserver"和"compute→safekeeper"的映射关系移出 console——console 仍然是集群中该信息的唯一持久化来源,每个 pageserver / safekeeper 仍须确切知道它服务哪些 timeline。至于新 pageserver 如何发现映射信息,则属于本 RFC 范围之外。

2. 核心思路:集中式状态中枢取代 P2P gossip

RFC 的总体设计是:

不再使用点对点 gossip,而是引入一个集中的 broker,所有存储节点(safekeeper、pageserver)向它汇报各自的 per-timeline 状态,同时订阅自己关心的状态集合,从而在本地维护一份一致性的视图。

为此,每个存储节点都应提供一个--broker-url=1.2.3.4形式的 CLI 参数,用于指定 broker 的地址。

相比 gossip,这种集中式方案解决了 gossip 固有的一个棘手问题——"对端地址失忆症"(peer address amnesia):在 gossip 方案中,如果多个 safekeeper 同时重启,在下一个 compute 连接到来之前,它们彼此不知道对方的地址;而集中式 broker 中所有节点都能随时查询/订阅到其他节点的地址信息,彻底消除了这一盲区。

RFC 中给出了两个候选实现路径,并最终决定采用 etcd 思路:

  • 方案 A(最初提案,已否决):在 console 中实现一个自定义 gRPC 消息代理;
  • 方案 B(最终选择):部署 etcd,在其中维护 per-timeline 的状态树。

2.1 已否决的方案:console 内的自定义消息代理

RFC 用一个折叠块保留了最初的提议:在 console 内新增一个 gRPC 服务充当消息代理(broker 可以忽略 payload、只做消息转发),消息格式为:

{sender, destination, payload}

其中 destination 有两种:

  • sk_#{tenant}_#{timeline}—— 广播给负责该 timeline 的所有 safekeeper;
  • pserver_#{tenant}_#{timeline}—— 广播给负责该 timeline 的所有 pageserver。

sender 也有两种:

  • sk_#{sk_id}
  • pserver_#{pserver_id}

针对最初的四个问题,该方案的交互设计是:

  • WAL 裁剪:每个 safekeeper 周期性广播(write_lsn, commit_lsn)给同一 timeline 的所有对端 safekeeper;
  • 决定谁推 S3:每个 safekeeper 周期性广播i_am_alive_#{current_timestamp}给对端,从而维护一个"存活对端"向量(宽松的、允许假阴性的集合),id 最小的存活 safekeeper负责推送数据到 S3;
  • 决定谁转发 WAL 给 pageserver:每个 safekeeper 周期性向相关 pageserver 发送(write_lsn, commit_lsn, compute_connected),pageserver 据此维护 safekeeper 状态视图、连接随机一个、并在检测到落后或宕机时重连;pageserver 通过 console 将#sk_#{sk_id}解析为真实 IP(如4.5.6.7:6400);
  • 何时关闭 SK↔PS 连接:pageserver 拥有足够信息自行判断。

该方案的可用性论证也很直接:broker 随 console 一起由 k8s 保活;即使 console 宕机导致 WAL 无法裁剪、pageserver 无法切换 safekeeper,此时系统本来也无法接受新的 compute 连接、无法启动已停止的 compute,因此"只是变得更糟了一点,而非灾难性的"。

不过 RFC 最终决定放弃该方案——理由见第 5 节的权衡分析。

2.2 最终方案:etcd 集中式状态存储

最终选定的方案是部署 etcd,并在其中维护如下数据结构(以某条tenant_timeline为例):

"compute_#{tenant}_#{timeline}" => { safekeepers => { "sk_#{sk_id}" => { write_lsn: "0/AEDF130", commit_lsn: "0/AEDF100", compute_connected: true, last_updated: 1642621138, }, } }

由于 etcd 不支持嵌套对象内的字段级更新,该结构实际对应为一组扁平 key:

"compute_#{tenant}_#{timeline}/safekeepers/sk_#{sk_id}/write_lsn", "compute_#{tenant}_#{timeline}/safekeepers/sk_#{sk_id}/commit_lsn", ...

每个存储节点订阅与其相关的 key 集合,并在本地维护这份结构的视图。数据流本质上与方案 A 一致,但避免了自行实现消息代理,也消除了存储节点对 console 的运行时依赖(本地开发时不必构建并运行依赖 Postgres 的 console,只需后台跑一个 etcd)。

Safekeeper 地址发现

在启动时,safekeeper 将其监听的地址以{"sk_#{sk_id}" => ip_address}的形式发布到 etcd;pageserver 即可将sk_#{sk_id}解析为真实地址。这种方式在本地与云上部署均适用。为此 safekeeper 需要增加--advertised-addressCLI 选项,以便"监听 0.0.0.0 但对外公布更实用的地址"。

Safekeeper 行为

对每一条 timeline,safekeeper 周期性广播compute_#{tenant}_#{timeline}/safekeepers/sk_#{sk_id}/*字段,并订阅compute_#{tenant}_#{timeline}的变化——这样它就能掌握同一 timeline 对端 safekeeper 的信息,足以正确裁剪 WAL。至于谁推 S3,safekeeper 可以使用 etcd lease,或广播时间戳来追踪对端存活状态。

Pageserver 行为

pageserver 为其拥有的每个 tenant 订阅compute_#{tenant}_#{timeline}。基于这些信息,它可以:

  • 维护 safekeeper 状态视图,连接随机一个 safekeeper;
  • 检测到某个 safekeeper 停止推进或宕机时,重连到另一个;
  • 通过 state changecompute_connected: false -> true触发连接(从而不再需要 "call me maybe"机制)。

RFC 还给出了compute_connected的一个替代方案:追踪从 compute 到达 safekeeper 的最新消息时间戳。compute 通常每秒向所有 safekeeper 广播一次 KeepAlive,因此连接正常时该时间戳每秒更新;当它数秒未更新时即视为连接断开。这个方案能更快地发现两类问题:

  1. compute 已失败,但 TCP 连接仍存活直到超时(通常约一分钟);
  2. safekeeper 已失败,却没有把compute_connected置为 false。

另一种兼容本 RFC 的做法是:pageserver 把(write_lsn, commit_lsn, compute_connected)当作 KeepAlive 处理,当某个sk_id一段时间内没有任何消息时即判定异常。

3. 关键工作流时序

RFC 用三张时序图(参与者包括 Compute、SK1-SK3、PS1-PS2、Orchestrator、Metadata Service——即 etcd/broker)完整描述了系统在典型场景下的协作过程。

3.1 集群启动(Cluster startup)

要点:pageserver 先订阅 timeline N 的状态;各 safekeeper 持续把当前 LSN 汇报给元数据服务;当新 compute 出现时,元数据服务把"谁在哪个 LSN"推给 pageserver,pageserver 选择 LSN 最大(最新)的 safekeeper 作为复制源

3.2 典型运行中的行为(Typical operations)

三个典型场景:

  1. pageserver checkpoint:pageserver 把数据上传到 S3 后,在元数据服务更新remote consistent lsn,元数据服务将其广播给所有 safekeeper,各 safekeeper 据此把 WAL 裁剪到该 LSN;
  2. safekeeper 发现自身落后:当MAX(对端LSN) - 自己LSN > THRESHOLD时,落后的 safekeeper 直接从最新的对端拉取 WAL 增量补数据(对应仓库中 safekeeper 的pull_timeline/ 从对端补 WAL 的机制);
  3. pageserver 检测到源 safekeeper 失联:连接断开或 30 秒无消息时,pageserver 依据视图选择次新的 safekeeper 重新开始复制。

3.3 Timeline 迁移(Timeline relocation)

迁移流程的关键步骤:orchestrator 让新 pageserver(PS2)attach timeline(S3 中存在则返回 202 Accepted)→ 下载数据 → 注册 timeline 并从元数据服务获取 safekeeper 列表、订阅变化→ 从最新 safekeeper 开始复制追赶 → 追上后重启 compute 并切换 pageserver 地址 → 旧 pageserver(PS1)detach。图中还覆盖了三个失败场景的处理策略:attach 失败可安全重试(超阈值换另一台 pageserver)、attach 成功但下载/复制失败可等待超时后换机(并限制尝试次数)、detach 失败可重试(持续失败可能导致 S3 数据重复)。

4. 权衡分析(Pros/Cons)

RFC 在结尾给出了两轮决策依据,这正是理解该架构取舍的关键。

4.1 集中式 broker/etcd vs gossip

  • Gossip 的唯一优点:允许存储层在没有 console 或 etcd 的情况下独立运行;
  • 集中式 broker/etcd 的优点
    1. 更简单;
    2. 一并解决了 "call me maybe" 问题;
    3. 避免 gossip 在"不把 safekeeper 预分组为三元组"情况下的 N-to-N 连接问题。

4.2 Console broker vs etcd

RFC 作者坦言最初想避免引入 etcd 依赖,理由是 ClickHouse 对 ZooKeeper 的依赖带来了沉重的配置与运维负担(作者还特别提到 ClickHouse 甚至为此重实现了内嵌的 ZooKeeper)。但对 Neon 的场景,etcd 有两个关键差异使风险可控:

  1. 存入 etcd 的数据不需要持久性与强一致性保证(只是一份可随时重建的状态快照);
  2. etcd 使用 gRPC 协议,消息非常简单

因此,"未来如果想实现一个带 etcd 接口的内存存储,是顺理成章的事";而现在完全不必实现它——本地跑 Neon 只需要在后台运行一个 etcd,而不必构建并运行依赖 Postgres 的 console。

5. 仓库中的实际落地:storage_broker crate

该 RFC 最终在仓库中落地为独立的 storage_broker crate(对应Cargo.toml中的storage_broker包),以及 docs 下的 docs/storage_broker.md。它不是 etcd,而是一个自研的、以 gRPC(tonic)为协议的集中式 broker——本质上正是 RFC 所设想的"与 etcd 接口等价"的服务端实现。

5.1 Broker 服务与消息协议

服务定义位于 storage_broker/proto/broker.proto,BrokerService暴露四个 RPC:

service BrokerService { // Subscribe to safekeeper updates. rpc SubscribeSafekeeperInfo(SubscribeSafekeeperInfoRequest) returns (stream SafekeeperTimelineInfo) {}; // Publish safekeeper updates. rpc PublishSafekeeperInfo(stream SafekeeperTimelineInfo) returns (google.protobuf.Empty) {}; // Subscribe to all messages, limited by a filter. rpc SubscribeByFilter(SubscribeByFilterRequest) returns (stream TypedMessage) {}; // Publish one message. rpc PublishOne(TypedMessage) returns (google.protobuf.Empty) {}; }

SafekeeperTimelineInfo消息完整覆盖了 RFC 中提到的状态字段,并扩展了生产所需的更多元数据:

  • safekeeper_idtenant_timeline_id:标识谁在汇报哪条 timeline;
  • termlast_log_term:safekeeper 的任期信息(与 WAL 一致性协议相关);
  • flush_lsn:已刷盘的最新 LSN;
  • commit_lsn:safekeeper 认为已提交的 LSN;
  • backup_lsn:已备份到 S3 的 LSN;
  • remote_consistent_lsn:pageserver 最近一次 checkpoint 上传到 S3 的 LSN(对应时序图中 WAL 裁剪的依据);
  • peer_horizon_lsnlocal_start_lsnstandby_horizon:用于 WAL 保留与对端补数的边界;
  • safekeeper_connstrhttp_connstrhttps_connstr:连接字符串(对应 RFC 中"地址发现"的诉求——把sk_id解析为真实地址);
  • availability_zone:可用区信息。

此外 proto 中还定义了订阅过滤器SubscribeByFilterRequest(可按MessageType与 tenant/timeline 过滤)和MessageType枚举(SAFEKEEPER_TIMELINE_INFOSAFEKEEPER_DISCOVERY_REQUESTSAFEKEEPER_DISCOVERY_RESPONSE),以及用于"发现"流程的SafekeeperDiscoveryRequest/SafekeeperDiscoveryResponse(后者是SafekeeperTimelineInfo的短版本,只含 WAL 下载所需的字段:commit_lsnsafekeeper_connstravailability_zonestandby_horizon)。这与 RFC 中"pageserver 把sk_#{sk_id}解析为 IP"的设计一一对应,只是把地址发现也搬进了 broker 协议。

客户端封装在 storage_broker/src/lib.rs:storage_broker::connect()负责建立到 broker 的 gRPC 通道,支持基于 URL scheme 的明文/ TLS 切换(https://前缀启用 rustls 加密)、HTTP/2 keepalive(默认DEFAULT_KEEPALIVE_INTERVAL"5000 ms")与 5 秒连接超时,并提供parse_proto_ttid()完成 protobuf 与内部TenantTimelineId类型的转换。

5.2 Safekeeper 侧的三个后台循环

safekeeper/src/broker.rs 实现了 safekeeper 与 broker 的全部交互,对应 RFC 中"Safekeeper 周期性广播状态并订阅对端"的行为,其task_main以 1 秒(RETRY_INTERVAL_MSEC)为周期托管三个任务并在失败后自动重启:

  1. push_loop(推送):每 1 秒(PUSH_INTERVAL_MSEC)把本节点所有活跃 timelineSafekeeperTimelineInfo通过流式publish_safekeeper_info推给 broker;若一轮推送耗时超过 push 间隔的一半会记录告警日志。配置项disable_periodic_broker_push为 true 时该循环被禁用。
  2. pull_loop(拉取/订阅):通过subscribe_safekeeper_info订阅(当前实现是订阅全部 timeline,代码中有 TODO 注明后续可只订阅本地 timeline),收到每条更新后调用tli.record_safekeeper_info(msg)写入本地 timeline 状态。值得注意的是代码注释明确说明:pull 也会收到本节点自己推送的信息,这被用作"与 broker 连接仍然存活"的指示——这正是 RFC 中"把状态更新当作 KeepAlive"思想的落地。
  3. discover_loop(发现响应):独立订阅SAFEKEEPER_DISCOVERY_REQUEST类型的消息,当收到针对本节点所服务 timeline 的发现请求时,构造SafekeeperDiscoveryResponse并通过publish_one回复,从而避免影响正常的 push/pull 主循环。

safekeeper/src/bin/safekeeper.rssafekeeper/src/lib.rs中可以看到对应配置项:broker_endpoint(broker 地址)、broker_keepalive_interval(默认 5 秒)等,支持--broker-endpoint--broker-keepalive-interval这类 CLI 参数——正是 RFC 中--broker-url设想的实现。此外该目录下还有task_stats辅助任务:若超过 10 秒未从 broker 收到任何更新,会周期性输出告警日志,用于暴露 broker 链路故障。

5.3 Pageserver 侧的连接管理:WAL Receiver Connection Manager

RFC 中"pageserver 订阅 timeline 状态、选择最合适的 safekeeper、并在连接失效时切换"的行为,在 pageserver 侧落地为 pageserver/src/tenant/timeline/walreceiver/connection_manager.rs。其模块注释直接点明了设计意图:

为了实现这一点使用了 storage broker:safekeepers 在 broker 中传播其 timeline 状态,manager 订阅这些变化并积累状态,以查询 LSN 最大的那个进行连接;当前连接状态也会被跟踪,以确保它不会过期。

该 manager 的核心职责包括:

  • 通过 broker 客户端订阅指定 timeline 的更新(subscribe_for_timeline_updates),且多个 connection manager 共享同一条底层 TCP 连接上的多个流;
  • 每次收到 broker 更新或连接状态变化后,重新评估并决定是否需要(重)连接,选择LSN 最大(最新)的 safekeeper作为复制源;
  • 当没有可用候选时,周期性地发送SafekeeperDiscoveryRequest(间隔复用lagging_wal_timeout配置)来探测对端——对应 proto 中的发现流程;
  • 维护连接重试历史:一旦从某 safekeeper 收到并处理了 WAL(推进了last_record_lsn),即清除该 safekeeper 的重试记录,允许立刻重连。

pageserver 的 broker 配置集中在 pageserver/src/config.rs:broker_endpoint(可配置多个端点)与broker_keepalive_interval,与 safekeeper 侧保持一致。这条链路把 RFC 中"按最新 LSN 选择复制源""检测落后/失联并切换""通过状态变化触发连接"全部落实为可运行的代码。

5.4 与 RFC 设计的一一对应

RFC 设计点仓库落地位置
集中式 broker,节点上报 per-timeline 状态storage_broker/proto/broker.proto 的BrokerServiceSafekeeperTimelineInfo
safekeeper 周期性广播(write_lsn, commit_lsn, ...)safekeeper/src/broker.rs 的push_loop(1s 周期)
safekeeper 订阅对端状态用于 WAL 裁剪同上pull_loop+record_safekeeper_info;裁剪依据remote_consistent_lsn
地址发现(sk_id → 真实地址)proto 中SafekeeperDiscoveryRequest/Response+safekeeper_connstr字段
pageserver 订阅状态、选最新 safekeeper、失效切换pageserver/src/tenant/timeline/walreceiver/connection_manager.rs
状态更新作为 KeepAlive,替代 "call me maybe"safekeeper/src/broker.rs pull_loop 中"收到自己的信息即代表 broker 链路存活"的注释与实现

6. 总结

015-storage-messaging这篇 RFC 的价值在于:它把 Neon 存储层"WAL 裁剪、S3 上传决策、WAL 转发源选择、SK↔PS 连接生命周期"四个看似独立的问题,统一收敛为**"所有存储节点向集中式状态中枢汇报 per-timeline 状态、并订阅所需状态集合"**这一个模式,并用清晰的时序图规定了各组件在启动、运行与迁移场景下的协作契约。它对 gossip 与集中式两种路线的诚实权衡(尤其是"为何 etcd 依赖可以接受"的分析),也为后来者在同类分布式系统中的依赖选型提供了可复用的思考框架。

从仓库现状看,该设计最终被实现为一个自研的 gRPCstorage_broker(而非直接部署 etcd),safekeeper 与 pageserver 围绕它分别建立了 push/pull/discover 与连接管理逻辑,RFC 中的消息字段、地址发现、KeepAlive 语义与"选择最新 safekeeper"的决策规则均能在源码中找到一一对应。对想深入 Neon 存储层协作机制的读者,建议按此顺序阅读:015-storage-messaging.md → storage_broker/proto/broker.proto → safekeeper/src/broker.rs → pageserver/src/tenant/timeline/walreceiver/connection_manager.rs,即可从设计到实现完整贯通。

【免费下载链接】neonNeon: Serverless Postgres. We separated storage and compute to offer autoscaling, code-like database branching, and scale to zero.项目地址: https://gitcode.com/GitHub_Trending/ne/neon

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/13 12:20:45

LunaTranslator快速上手:3分钟让日文视觉小说游戏不再卡壳

LunaTranslator快速上手:3分钟让日文视觉小说游戏不再卡壳 【免费下载链接】LunaTranslator 视觉小说翻译器 / Visual Novel Translator 项目地址: https://gitcode.com/GitHub_Trending/lu/LunaTranslator LunaTranslator是一款开源的视觉小说游戏翻译工具&…

作者头像 李华
网站建设 2026/9/13 12:19:13

PMSM无感矢量控制:滑模观测器SMO原理与工程实现

简介:本资源是一套面向电机控制方向研究生、工程师及高阶学习者的三相永磁同步电机(PMSM)矢量控制MATLAB/Simulink仿真实践包,聚焦无模型控制与无感矢量控制两大前沿策略,解决传统FOC依赖精确模型和位置传感器带来的鲁…

作者头像 李华
网站建设 2026/9/13 12:18:36

Megatron Core 多模态模型实战指南:从 LLaVA、NVLM 到 MIMO 框架

Megatron Core 多模态模型实战指南:从 LLaVA、NVLM 到 MIMO 框架 【免费下载链接】Megatron-LM Ongoing research training transformer models at scale 项目地址: https://gitcode.com/GitHub_Trending/me/Megatron-LM Megatron Core(本仓库 me…

作者头像 李华