☰
CouchDB 集群内部 RPC 机制深度解析:rexi 分布式调用框架的原理、源码与配置
2026/10/9 2:17:04 网站建设 项目流程
  • 数据库
  • 文档数据库
  • 后端

【免费下载链接】couchdb

Seamless multi-primary syncing database with an intuitive HTTP/JSON API, designed for reliability

项目地址:https://gitcode.com/gh_mirrors/co/couchdb
点击查看免费下载

本文基于 Apache CouchDB 仓库中的src/rexi应用及其 README 展开,完整解读这套为分布式集群定制的 Erlang RPC 框架:它为什么被设计出来、如何工作、怎样在 Fabric 等上层模块中使用,以及相关的配置项与可观测性指标。读完本文,你将掌握 rexi 的 API 用法、消息流时序、故障处理模型和调优要点。

一、为什么 CouchDB 需要一套自己的 RPC:从 rex 到 rexi

CouchDB 以 HTTP/JSON API 闻名,但集群内部节点之间的通信依赖 Erlang 分布式能力。src/rexi/README.md开篇即阐明 rexi 的定位:它是专为把 CouchDB 操作发送到集群中其他节点而量身定制的 RPC 服务器应用,在 BigCouch 架构中充当远程过程调用的载体,让fabric模块的函数得以在远程集群节点上执行。

OTP 自带一个 RPC 服务器rex(rpc模块的后端),但 rexi 针对 BigCouch 分布式数据存储的需求做了取舍——去掉rex中不必要的开销,专门优化"批量派生远程进程"这一核心场景。关键设计差异体现在三点:

  1. 消息只送到远端一个常驻服务器:cast消息从发起方(origin)发送到远端节点的 rexi_server,由它在本地spawn工作进程执行任务;
  2. 本地派生而非远程派生:相比从发起方直接远程spawn进程,本地派生开销小得多,能支撑大量并发远程任务;
  3. 监控语义保持,但调用方不被阻塞:远程进程仍然被监控,但处理请求的进程不会卡在"尝试连接一个过载或宕机节点"上,rexi_DOWN消息最终会到达客户端。这种"延迟与故障检测"的组合,是 rexi 相比传统 RPC 的显著优势。

rexi 通常与 Fabric(同为集群应用,位于 src/fabric)配合使用,但也可独立运行。其核心源码位于 src/rexi/src,模块职责如下:

模块职责
rexi.erl对外 API:cast 系列、kill_all、reply、流式接口
rexi_server.erl远端常驻服务器:接收任务、派生并监控 worker、转发错误
rexi_buffer.erl发送缓冲:目标节点不可达时暂存消息,超限丢弃
rexi_monitor.erl批量化监控进程,产出rexi_DOWN通知
rexi_server_mon.erl按集群拓扑管理各节点上的 server/buffer 生命周期
rexi_sup.erl顶层 supervisor,组织上述组件

二、一次远程调用的完整生命周期

以rexi:cast/2为例,一条远程调用的典型路径如下:

%% 在发起方节点上: Ref = rexi:cast(Node, {fabric_rpc, compact, [Name]}),

cast/2等价于cast_ref(make_ref(), Node, self(), MFA)(见 rexi.erl),即默认把当前进程作为 Caller,并自动创建一个 reference 作为本次调用的唯一标识。消息体构造为:

Msg = cast_msg({doit, {Caller, Ref}, get(nonce), MFA}) %% cast_msg/1 实际包装为 {'$gen_cast', Msg}

随后通过rexi_utils:send(rexi_utils:server_pid(Node), Msg)发送。server_pid/1(rexi_utils.erl)返回{rexi_server_<Node>, Node}注册名——每个集群节点上都有一个以自身节点名命名的 rexi_server。

远端rexi_server收到{doit, From, Nonce, MFA}后(rexi_server.erl):

  1. 执行spawn_monitor(?MODULE, init_p, [From, MFA, Nonce]),在本地派生一个 worker 进程并监控它;
  2. 把#job{}记录同时写入workers与clients两张 ETS 表(分别以 worker 的 monitor ref 和 client ref 为键),用于双向索引;
  3. worker 进程init_p/3把rexi_from = {CallerPid, Ref}和 nonce 写入进程字典,然后apply(M, F, A)执行目标函数(rexi_server.erl)。

目标函数内部若要回传结果,调用:

rexi:reply(Reply) %% 等价于:{Caller, Ref} = get(rexi_from), erlang:send(Caller, {Ref, Reply})

于是发起方收到形如{Ref, Reply}的消息,与自己的 ref 匹配即可收敛结果。

2.1 调用方如何聚合多节点结果:rexi_utils:recv

集群调用通常同时向多个分片节点发起cast,然后统一收消息。rexi_utils:recv/6(rexi_utils.erl)提供带总超时与单消息超时的接收循环:传入Refs(ref 列表)、Keypos(匹配键位置)、处理函数Fun、累加器Acc0,循环消费消息直到Fun返回{stop, Acc}或超时。它能识别{Ref, Msg}、{Ref, From, Msg}、{rexi, Ref, Msg}、{rexi, Ref, From, Msg}、{rexi, '$rexi_ping'}以及{rexi_DOWN, _, _, _}等消息形态,并对不匹配 ref 的消息静默忽略。这是 Fabric 各模块收敛分片响应的基础设施。

2.2 主动终止:kill_all

当调用方决定放弃尚未完成的远程任务(例如请求超时、客户端断开),可调用rexi:kill_all/1批量发送异步 kill 信号(rexi.erl):

%% NodeRefs = [{Node, Ref}, ...] rexi:kill_all([{Node1, Ref1}, {Node2, Ref2}, ...])

该函数先用maps:groups_from_list按节点分组,再对每个节点只发送一条{kill_all, Refs}消息,避免逐条发送。远端 server 收到后对每个 ref 执行kill_worker/2(rexi_server.erl):先demonitor再exit(Pid, kill),并从 ETS 表中移除记录。被 kill 的 worker不会再发rexi_EXIT消息。

三、故障处理模型:rexi_DOWN 与 rexi_EXIT

rexi 的故障语义分两层:

  • 进程崩溃:worker 执行apply(M, F, A)抛出异常时,init_p/3捕获Class:Reason:Stack,构造#error{}记录(见 rexi.hrl,包含 timestamp、reason、mfa、nonce、stack 字段),然后以该记录作为退出原因退出。远端 server 收到{'DOWN', ...}后调用notify_caller,向调用方发送:
{Ref, {rexi_EXIT, Reason}}

同时把错误存入内存队列(受error_limit约束,可通过set_error_limit与get_errors等 call 查询,见 rexi_server.erl)。进程字典中的 nonce 会被带回错误记录,便于调用方关联上下文。

  • 节点故障/不可达:rexi 的核心卖点是"调用方进程不会被卡住"。当目标节点宕机或连接异常时,rexi_utils:send/2(rexi_utils.erl)先用erlang:send(Dest, Msg, [noconnect, nosuspend])尝试立即发送;若返回非ok(如noconnect、nosuspend),则转交rexi_buffer缓冲,调用方进程自身立刻返回,绝不阻塞等待。

同时,rexi_monitor:start/1(rexi_monitor.erl)会派生一个进程批量监控一组目标(pid 或{Name, Node})。对不可达/不在集群的节点立即以noconnect理由通知父进程;对可达目标则erlang:monitor(process, P)并循环等待 DOWN。任何被监控进程退出,父进程都会收到:

{rexi_DOWN, MonitoringPid, DeadPid, Reason}

rexi_monitor:stop/1则负责优雅关闭监控进程并清空邮箱中积压的rexi_DOWN消息。Fabric 各模块正是靠这套机制在节点宕机时迅速完成降级处理,例如 fabric_db_create.erl 对{rexi_DOWN, _, {_, Node}, _}的分支处理。

四、发送缓冲与背压:rexi_buffer

rexi_buffer是 rexi 应对"目标节点暂时不可达"的关键组件。每个节点对应一个rexi_buffer_<Node>进程,行为如下(rexi_buffer.erl):

  • 收到{deliver, Dest, Msg}时把消息入队,并通过counters记录缓冲量;
  • 当缓冲数达到max_count(配置rexi.buffer_count,默认 2000,见 rexi_buffer.erl)时执行should_drop/1,丢弃队尾消息(queue:drop(Q2))并递增[rexi, dropped]统计;
  • 通过timeout消息驱动排空:能直接发送则立即发送,遇到持续延迟则spawn_monitor(erlang, send, [Dest, Msg])派一个发送进程,避免 buffer 进程自身被阻塞;消息全部发出后进入hibernate以回收内存(见 rexi_buffer.erl)。

由此形成明确的背压语义:宁可丢弃积压消息,也不能让调用链卡死。可调用rexi_buffer:get_buffered_count/1查询缓冲量,erase_buffer/1清空缓冲。

五、流式传输协议:stream2 的握手与确认

视图查询、变更流(changes feed)等场景需要 worker 向协调者(coordinator)持续推送大量结果,rexi:stream2/1,2,3为此定义了带确认的流式协议。协议时序在 rexi.erl 的注释中有完整描述:

Coordinator Worker ---------- -------- cast/2,3,4 -> {doit, ...} -> rexi_server: spawn_monitor worker worker 首次调用 stream2 时执行 init_stream/1 <- rexi_STREAM_INIT <- 同步发送并等待回复 部分 worker 收到 rexi_STREAM_START 继续, 其余收到 rexi_STREAM_CANCEL 停止 -> rexi_STREAM_START -> worker 持续调用 rexi:stream2 推送数据 <- Caller ! {Ref, self(), Msg} <- 协调者必须确认(ack) -> {rexi_ack, 1} -> 发送最后一条消息(无需 ack) <- Caller ! Msg <-

对应 API 语义:

  • rexi:stream2(Msg)默认Limit=5、Timeout=300000(配置rexi.stream_limit可调,见 rexi.erl);
  • worker 在途未确认消息数达到Limit后停止发送,等待{rexi_ack, N}(wait_for_ack/2带超时,超时则递增[rexi, streams, timeout, wait_for_ack]统计并exit(timeout));
  • 协调者用rexi:stream_ack(Client)发送{rexi_ack, 1}(rexi.erl);
  • rexi:stream_start({Pid, Tag})/rexi:stream_cancel({Pid, Tag})基于gen_server:reply机制回复rexi_STREAM_START/rexi_STREAM_CANCEL,分别让 worker 继续或退出(exit(normal));
  • 最后一条消息用rexi:stream_last(Msg)发送,不携带 worker pid 也不需要 ack(rexi.erl);
  • 长时间运行的任务可用rexi:ping()向协调者发送{rexi, '$rexi_ping'},rexi_utils:recv会静默忽略这类心跳,避免协调者误判超时。

rexi_tests.erl中的t_stream2、t_stream2_acks、t_stream2_cancel等用例(src/rexi/test/rexi_tests.erl)验证了这条协议的完整路径。

六、集群拓扑感知:每节点一个 server

rexi_server_mon实现mem3_cluster行为(rexi_server_mon.erl),通过mem3_cluster:start_link监听集群稳定性事件:

  • 集群不稳定时(有节点加入/离开),立即启动缺失节点的 server("可以启动,但不立即停掉旧节点");
  • 集群稳定后,补齐缺失 server 并停止多余 server(start_servers/1+stop_servers/1,见 rexi_server_mon.erl)。

也就是说,rexi 的 server 集合是动态跟随集群成员变化的。rexi_sup.erl(src/rexi/src/rexi_sup.erl)以rest_for_one策略组织rexi_server、rexi_server_sup、rexi_server_mon、rexi_buffer_sup、rexi_buffer_mon。aggregate_server_queue_len/0与aggregate_buffer_queue_len/0(rexi.erl)可跨节点聚合队列长度,用于观察系统负载。

七、配置项与运维观测

7.1 配置参数

rexi 的配置位于[rexi]配置节,默认值见 rel/overlay/etc/default.ini:

[rexi] ;buffer_count = 2000 ;stream_limit = 5 ;shard_split_timeout_msec = 600000 ;shard_split_topoff_batch_size = 500

官方配置文档 src/docs/src/config/cluster.rst 对核心参数的解释:

参数默认值作用
rexi.buffer_count2000本地 RPC server 在远端节点不可达时最多缓冲多少条消息,超出后开始丢弃。对应 rexi_buffer.erl 中的config:get_integer("rexi", "buffer_count", 2000)
rexi.stream_limit5流式操作(视图、变更流)中,worker 不等待协调者确认即可发送的最大消息数。过大时协调者可能被 worker 消息淹没,反而降低整体吞吐(CouchDB 2.x 中该值曾硬编码)。对应 rexi.erl

7.2 可观测性指标

rexi 通过 priv/stats_descriptions.cfg 暴露以下计数器:

指标说明
rexi.buffered被缓冲的 rexi 消息数量
rexi.dropped从缓冲中丢弃的 rexi 消息数量
rexi.down已处理的rexi_DOWN消息数量
rexi.streams.timeout.init_stream流初始化超时次数
rexi.streams.timeout.stream流消息发送超时次数
rexi.streams.timeout.wait_for_ack等待确认时超时次数

这些指标可帮助运维判断集群是否出现节点不可达(dropped/buffered上升)或协调者处理能力不足(stream系列超时上升)。

八、测试与验证

rexi 的测试集中在 src/rexi/test/rexi_tests.erl:

  • t_cast、t_cast_ref、t_cast_call:验证rexi:cast各种变体把 MFA 送到远端执行并返回结果;
  • t_stream2、t_stream2_acks、t_stream2_cancel:验证流式协议的初始化、确认与取消;
  • rpc_test_fun系列:模拟远端 worker 依次调用stream2与stream_last的行为。

配合 src/rexi/test/rexi_buffer_tests.erl 对缓冲/丢弃逻辑的覆盖,构成了 rexi 的完整测试面。

九、构建与许可

README 记载 rexi 需要R13B03 或更高版本的 Erlang,可使用仓库内捆绑的 rebar 构建。rexi 以Apache 2.0许可发布(与仓库根目录 LICENSE 一致)。该应用当前作为 CouchDB 集群的核心依赖运行——例如 fabric.app.src 将rexi列为运行依赖,fabric.erl 直接以rexi:cast(Node, {fabric_rpc, compact, [Name]})的形式下发压缩任务,读者可顺着这些调用点继续深入理解 CouchDB 的分布式执行层。

  • 数据库
  • 文档数据库
  • 后端

【免费下载链接】couchdb

Seamless multi-primary syncing database with an intuitive HTTP/JSON API, designed for reliability

项目地址:https://gitcode.com/gh_mirrors/co/couchdb
点击查看免费下载

相关推荐

上一篇:如何快速掌握 Substrate 开发?Awesome Substrate 提供的完整学习路线图
下一篇:PhotoGIMP完整教程:3分钟让GIMP界面变成Photoshop的终极指南

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

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询