news 2026/10/9 1:46:52

Apache Storm Distributed RPC 完整实战:DRPC 原理、拓扑构建与集群部署

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Storm Distributed RPC 完整实战:DRPC 原理、拓扑构建与集群部署
  • 后端
  • 大数据

【免费下载链接】storm

Apache Storm

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

Distributed RPC(DRPC)是 Apache Storm 提供的一种以流式拓扑形式并行执行高强度函数计算的模式:客户端像调用普通 RPC 一样提交函数名与参数,Storm 集群则在拓扑中并行完成计算并把结果返回给等待的客户端。本文以仓库 docs/Distributed-RPC.md 为主线,结合 LinearDRPCTopologyBuilder、DRPCSpout、DRPC 等源码与 storm-starter 示例,讲解 DRPC 的整体架构、客户端调用、拓扑构建、本地/远程模式部署,以及 reach 这类需要真正并行计算能力的复杂用例,帮助你从零搭建并运行一个可用的 DRPC 服务。

DRPC 工作流程

DRPC 的核心思想与整体架构

DRPC 的出发点是用 Storm 实时并行化那些计算量巨大的函数:拓扑以「函数参数流」作为输入,并输出每条函数调用的结果流。严格来说,DRPC 与其说是 Storm 的一个独立特性,不如说它是基于 Storm 的流(streams)、spout、bolt 和拓扑等原语组合出的一种模式(pattern)——它本可以被打包成独立库,但因为它太常用了,Storm 直接将其内置。

整个协调工作由一个DRPC server完成,Storm 自带该实现。DRPC server 负责四件事:

  1. 接收客户端发来的 RPC 请求;
  2. 把请求(函数调用)投递给 Storm 拓扑;
  3. 从拓扑接收计算结果;
  4. 把结果回传给正在等待的客户端。

从客户端视角看,一次分布式 RPC 调用与普通 RPC 几乎没有区别。下面这段代码演示如何计算函数reach在参数"http://twitter.com"上的结果:

Config conf = new Config(); conf.put("storm.thrift.transport", "org.apache.storm.security.auth.plain.PlainSaslTransportPlugin"); conf.put(Config.STORM_NIMBUS_RETRY_TIMES, 3); conf.put(Config.STORM_NIMBUS_RETRY_INTERVAL, 10); conf.put(Config.STORM_NIMBUS_RETRY_INTERVAL_CEILING, 20); DRPCClient client = new DRPCClient(conf, "drpc-host", 3772); String result = client.execute("reach", "http://twitter.com");

其中3772是 DRPC server 接收客户端请求的默认端口(drpc.port,见 conf/defaults.yaml)。如果你不想手工指定主机,也可以使用预配置客户端:它会从配置的 DRPC server 列表中随机挑选一台主机,若该主机不可用,则依次遍历所有已配置主机寻找可用者(对应 DRPCClient.getConfiguredClient 中Collections.shuffle(servers)后逐个尝试建立连接的逻辑):

DRPCClient client = DRPCClient.getConfiguredClient(conf); String result = client.execute("reach", "http://twitter.com");

注意:getConfiguredClient会把Config.DRPC_PORT(默认 3772)作为连接端口,并要求配置中必须存在drpc.servers,否则会抛出IllegalStateException。

请求如何穿过拓扑:DRPCSpout 与 ReturnResults

一次 DRPC 调用的完整数据流如下:

  1. 客户端把「函数名 + 参数」发给 DRPC server;
  2. 实现该函数的拓扑通过DRPCSpout从 DRPC server 读取函数调用流——DRPC server 会为每次函数调用打上唯一 id;
  3. 拓扑完成计算,位于拓扑末端的ReturnResultsbolt 连接 DRPC server,把「函数调用 id + 结果」交还给它;
  4. DRPC server 依据 id 匹配到正在等待的客户端,解除其阻塞并把结果返回。

这个「先入队、被拓扑取走、再回填结果」的过程在服务端由 DRPC.java 实现:请求先进入按函数名分组的queues(等待被拓扑 fetch),被fetchRequest取走后转入requests(等待结果回填),returnResult通过 id 找到对应的OutstandingRequest并唤醒阻塞中的客户端。

LinearDRPCTopologyBuilder:一站式构建 DRPC 拓扑

Storm 提供了 LinearDRPCTopologyBuilder 来自动化 DRPC 几乎所有的编排步骤,包括:

  1. 搭建 spout;
  2. 把结果返回给 DRPC server;
  3. 为 bolt 提供对一组 tuple 做「有限聚合(finite aggregation)」的能力。

最简单的例子:ExclaimBolt

下面是一个完整的 DRPC 拓扑实现,其函数行为是给输入参数追加一个"!":

public static class ExclaimBolt extends BaseBasicBolt { public void execute(Tuple tuple, BasicOutputCollector collector) { String input = tuple.getString(1); collector.emit(new Values(tuple.getValue(0), input + "!")); } public void declareOutputFields(OutputFieldsDeclarer declarer) { declarer.declare(new Fields("id", "result")); } } public static void main(String[] args) throws Exception { LinearDRPCTopologyBuilder builder = new LinearDRPCTopologyBuilder("exclamation"); builder.addBolt(new ExclaimBolt(), 3); // ... }

这段代码的核心约定(也是 LinearDRPCTopologyBuilder.createTopology 源码中的硬性校验):

  • 构造LinearDRPCTopologyBuilder时传入的是拓扑对应的 DRPC 函数名。单个 DRPC server 可以协调多个函数,函数名用于彼此区分(服务端DRPC也按函数名分队列维护请求);
  • 第一个声明的 bolt 接收二元组,第一个字段是请求 id,第二个字段是该请求的参数;
  • 最后一个 bolt 必须输出[id, result]形式的二元组——builder 在源码中会通过OutputFieldsGetter检查最后一个 bolt 恰好只声明一条输出流且恰好包含两个字段,否则抛出RuntimeException;
  • 所有中间 tuple 都必须把请求 id 放在第一个字段。

本例中ExclaimBolt只是给 tuple 的第二个字段追加"!",其余与 DRPC server 的连接、结果回传全部由LinearDRPCTopologyBuilder代劳。仓库 BasicDRPCTopology.java 给出了该示例的可运行版本(函数名、拓扑名均可作为命令行参数传入)。

本地模式 DRPC

过去,在本地模式使用 DRPC 需要手工创建一个特殊的LocalDRPC实例。这一用法在编写测试时仍然保留,但在当前版本的 Storm 中,本地模式下会自动创建LocalDRPC实例,任何新建的DRPCClient都会自动链接到它而不是外部世界。因此,与旧版LocalDRPC一样,你想测试的任何交互都必须包含在启动拓扑的那段脚本里。

从源码看,本地模式的桥接通过ServiceRegistry完成:DRPCSpout 在本地分支里以localDrpcId为键从服务注册表中取出DistributedRPCInvocations.Iface并调用fetchRequest;DRPCClient 内部维护一个静态的localOverrideClient,一旦存在本地覆盖,新客户端就会直接使用该本地实例而不再走 Thrift 网络。测试代码可通过DRPCClient.LocalOverride(实现了AutoCloseable)临时注入本地 DRPC。

远程模式 DRPC:在真实集群上运行

在真实集群上使用 DRPC 同样直接,只需三步:

  1. 启动 DRPC server(s);
  2. 配置 DRPC server 的位置;
  3. 把 DRPC 拓扑提交到 Storm 集群。

启动 DRPC server

用storm脚本启动 DRPC server,与启动 Nimbus 或 UI 一样简单:

bin/storm drpc

配置 DRPC server 位置

接下来需要让 Storm 集群知道 DRPC server 的位置——这是DRPCSpout知道从哪里读取函数调用的依据。配置可以写在storm.yaml或拓扑配置中;同时应把storm.thrift.transport属性配成与DRPCClient一致的值。在storm.yaml中大致如下:

drpc.servers: - "drpc1.foo.com" - "drpc2.foo.com" drpc.http.port: 8081 storm.thrift.transport: "org.apache.storm.security.auth.plain.PlainSaslTransportPlugin"

注意drpc.http.port: 8081只是文档示例;当前仓库 conf/defaults.yaml 中 DRPC 相关默认值如下,供你按实际集群调整:

配置项默认值作用
drpc.port3772DRPC server 接收客户端 Thrift 请求的端口
drpc.invocations.port3773DRPC 拓扑(spout/bolt)收发函数调用与结果的端口
drpc.http.port3774DRPC HTTP 服务端口
drpc.https.port-1HTTPS 端口(-1 表示关闭)
drpc.request.timeout.secs600服务端请求超时时间(秒),超时后按SERVER_TIMEOUT失败
drpc.worker.threads64DRPC Thrift server 工作线程数
drpc.queue.size128DRPC Thrift 请求队列大小
drpc.max_buffer_size1048576DRPC 最大缓冲区字节数
drpc.invocations.threads64invocations 通道线程数
drpc.childopts-Xmx768mDRPC 进程 JVM 参数
drpc.authorizer.acl.filenamedrpc-auth-acl.yaml基于 ACL 的授权器使用的文件名
drpc.disable.http.bindingtrue是否禁用 DRPC HTTP 绑定

关于storm.thrift.transport:文档中的PlainSaslTransportPlugin是不做认证的明文插件。若要启用认证,仓库还提供了 drpc-auth-acl.yaml.example、jaas_digest.conf 等示例,并内置DRPCSimpleACLAuthorizer(见 DRPCSimpleACLAuthorizer.java)与 DRPCAuthorizerBase.java;服务端通过drpc.authorizer配置授权器(定义于 DaemonConfig.java 的DRPC_AUTHORIZER)。

提交 DRPC 拓扑

最后,用StormSubmitter像提交任何普通拓扑一样提交 DRPC 拓扑。上面的例子在远程模式下这样运行:

StormSubmitter.submitTopology("exclamation-drpc", conf, builder.createRemoteTopology());

createRemoteTopology()用于生成适合 Storm 集群的拓扑(对应源码中createTopology(new DRPCSpout(function)));与之相对,本地模式应使用createLocalTopology(ILocalDRPC)。

三种调用 DRPC 函数的方式

假设拓扑正在监听exclaim函数,可以有多种调用方式:

编程方式调用:

Config conf = new Config(); try (DRPCClient drpc = DRPCClient.getConfiguredClient(conf)) { //User the drpc client String result = drpc.execute("exclaim", "argument"); }

通过 curl(HTTP 调用):

curl http://hostname:8081/drpc/exclaim/argument

该 HTTP 端点由 DRPCResource.java 提供:@Path("/drpc/")下暴露/{func}/{args}的 GET 接口与/{func}的 POST 接口,内部均调用drpc.executeBlocking(func, args)同步等待结果。URL 中的端口对应drpc.http.port。

通过命令行:

bin/storm drpc-client exclaim argument

复杂示例:并行计算 URL 的 Twitter reach

前面的 exclamation 只是演示概念的玩具例子。下面看一个真正需要 Storm 集群并行能力的函数:计算某个 URL 在 Twitter 上的reach——即有多少独立用户接触过该 URL。

计算 reach 需要四步:

  1. 获取所有发过该 URL 的推文作者(tweeters);
  2. 获取这些作者的全部粉丝(followers);
  3. 对粉丝集合去重;
  4. 统计去重后的粉丝数量。

单次 reach 计算可能涉及上千次数据库调用和数千万条粉丝记录,是名副其实的重计算。在单机上可能需要数分钟,而在 Storm 集群上即使是最困难的 URL 也能在几秒内算完。仓库中的 ReachTopology.java 以内存 HashMap 模拟了真实的 tweeters/followers 数据库,其拓扑定义如下:

LinearDRPCTopologyBuilder builder = new LinearDRPCTopologyBuilder("reach"); builder.addBolt(new GetTweeters(), 3); builder.addBolt(new GetFollowers(), 12) .shuffleGrouping(); builder.addBolt(new PartialUniquer(), 6) .fieldsGrouping(new Fields("id", "follower")); builder.addBolt(new CountAggregator(), 2) .fieldsGrouping(new Fields("id"));

(仓库中实际并行度配置为GetTweeters×4、GetFollowers×12、PartialUniquer×6、CountAggregator×3,并设置了conf.setNumWorkers(6)。)

拓扑按四个步骤执行:

  1. GetTweeters:获取发布该 URL 的用户。把[id, url]输入流转换为[id, tweeter]输出流,每个urltuple 会映射出多个tweetertuple。
  2. GetFollowers:获取这些用户的粉丝。把[id, tweeter]转换为[id, follower]。由于同一个人可能关注了多个发过该 URL 的作者,跨所有 task 看,follower tuple 会出现重复。
  3. PartialUniquer:按 follower id 对流分组,使同一个 follower 始终进入同一个 task,这样每个 task 收到的是互不重叠的 follower 子集。当它收到针对某个请求 id 的全部 follower tuple 后,发出自己子集的去重数量。
  4. CountAggregator:汇总各PartialUniquertask 的部分计数,完成 reach 计算。

PartialUniquer是理解「批次聚合」的关键,其实现基于BaseBatchBolt:

public class PartialUniquer extends BaseBatchBolt { BatchOutputCollector _collector; Object _id; Set<String> _followers = new HashSet<String>(); @Override public void prepare(Map conf, TopologyContext context, BatchOutputCollector collector, Object id) { _collector = collector; _id = id; } @Override public void execute(Tuple tuple) { _followers.add(tuple.getString(1)); } @Override public void finishBatch() { _collector.emit(new Values(_id, _followers.size())); } @Override public void declareOutputFields(OutputFieldsDeclarer declarer) { declarer.declare(new Fields("id", "partial-count")); } }

PartialUniquer通过继承BaseBatchBolt实现了IBatchBolt接口。批处理 bolt(batch bolt)提供了一等公民的 API,把「一批 tuple」作为具体单元来处理:每个请求 id 都会创建一个新的 batch bolt 实例,Storm 会在合适时机负责清理这些实例。execute方法把收到的 follower 加入该请求 id 对应的内部HashSet;当该 task 处理完这个批次的所有 tuple 后,finishBatch回调被触发,PartialUniquer发出包含其 follower 子集去重数量的单个 tuple。

在底层,CoordinatedBolt负责检测某个 bolt 是否已收到某个请求 id 的全部 tuple,它利用**直连流(direct streams)**来完成这种协调。reach 拓扑的每一步都是并行执行的——这正是定义 DRPC 拓扑极其简单、却能获得横向扩展能力的根源。

非线性 DRPC 拓扑

LinearDRPCTopologyBuilder只处理线性DRPC 拓扑——即计算被表达为一系列步骤(像 reach 那样)。不难想象有些函数需要带分支、合并的更复杂拓扑。当前阶段要实现这类拓扑,需要直接使用CoordinatedBolt手工搭建(CoordinatedBolt位于 storm-client,LinearDRPCTopologyBuilder内部也正是用它包装每个 bolt)。官方建议在邮件列表中讨论你的非线性用例,以推动 DRPC 更通用抽象的建设。

LinearDRPCTopologyBuilder 的工作原理

结合 LinearDRPCTopologyBuilder.createTopology 的源码,builder 实际构造的拓扑由以下部件构成:

  • DRPCSpout:发射[args, return-info]。其中return-info是 DRPC server 的主机、端口以及 server 生成的请求 id——查看 DRPCSpout.nextTuple 可见它通过DRPCInvocationsClient.fetchRequest(function)轮询取请求,并把{"id":..., "host":..., "port":...}以 JSON 字符串形式放入 return-info 字段;
  • PrepareRequest:为请求生成一个请求 id,并分出三条流——参数流(ARGS_STREAM)、返回信息流(RETURN_STREAM)、id 流(ID_STREAM),见 PrepareRequest.java(请求 id 用rand.nextLong()生成);
  • CoordinatedBolt包装与直连分组(direct groupings):除最后一个 bolt 外,每个 bolt 都被包装成CoordinatedBolt以感知「某个请求 id 的批次是否收齐」,并通过Constants.COORDINATED_STREAM_ID直连流传递协调信号;
  • JoinResult:把计算结果与 return-info 按请求 id 拼接(对结果流按结果首字段fieldsGrouping、对返回信息流按"request"字段fieldsGrouping);
  • ReturnResults:连接 DRPC server 并返回结果。查看 ReturnResults.java,它会解析 return-info JSON,在本地模式从ServiceRegistry取服务、远程模式则缓存/新建DRPCInvocationsClient调用result(id, result),并对TException做最多 3 次重连重试。

LinearDRPCTopologyBuilder因此是「在 Storm 原语之上构建更高层抽象」的一个优秀范例——它只依赖公开的 spout、bolt、grouping 与CoordinatedBolt,却把 DRPC 端到端编排完全封装掉了。

高级主题

  • KeyedFairBolt:用于在同一时刻交错处理多个请求,见 KeyedFairBolt.java。它以 tuple 的第一个字段(即请求 id)为键,把进入的 tuple 按 key 放入KeyedRoundRobinQueue,由后台线程轮询取走交棒给委托 bolt,从而公平地推进多个请求的批次处理,避免单个大请求饿死其他请求。
  • 如何直接使用CoordinatedBolt:当线性 builder 无法表达你的拓扑时,可手工把业务 bolt 包进CoordinatedBolt,指定SourceArgs(single()/all())与IdStreamSpec(id 检测流),再自行实现FinishedCallback.finishedId完成批次聚合,具体可对照LinearDRPCTopologyBuilder.createTopology的装配方式。
  • 服务端请求生命周期与超时:DRPC server 以定时任务(timer.scheduleAtFixedRate,间隔为超时的一半)扫描所有OutstandingRequest,超过drpc.request.timeout.secs未完成的请求会以SERVER_TIMEOUT异常失败并清理(见 DRPC.java)。拓扑侧的DRPCSpout.fail也会调用failRequest通知服务端该请求失败。

小结

DRPC 将 Storm 的流式计算能力包装成同步的 RPC 语义:客户端无感、拓扑侧只需遵循「id 打头、末 bolt 输出 [id, result]」的约定。无论是 toy 级的 exclamation 拓扑,还是 reach 这种涉及海量数据去重统计的重计算,LinearDRPCTopologyBuilder都让开发者能用极少的代码获得 Storm 集群的横向并行能力;需要更复杂拓扑时,则可以直接站在CoordinatedBolt、KeyedFairBolt等底层抽象之上自行编排。部署侧只需启动bin/storm drpc、配置drpc.servers与storm.thrift.transport、再通过StormSubmitter提交拓扑,即可通过 Java 客户端、curl 或storm drpc-client三种方式对外提供服务。

  • 后端
  • 大数据

【免费下载链接】storm

Apache Storm

项目地址:https://gitcode.com/gh_mirrors/storm22/storm
点击查看免费下载
上一篇:Starscream路线图展望:未来WebSocket功能发展预测
下一篇:Exchange Calendar未来路线图:即将推出的5大令人期待的新功能

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

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

Pi编程智能体实战:概念厘清、skill导入与subagent编排全解析

不整虚的&#xff0c;直接聊点真东西。最近“pi”这个词热度很诡异&#xff0c;你搜出来一堆结果&#xff0c;有讲树莓派的&#xff0c;有讲圆周率的&#xff0c;还有讲控制器的PI参数的。但真正在开发者圈子里炸开锅的&#xff0c;是那个叫 Pi 的 Agent——也就是大家口中的 p…

作者头像 李华
网站建设 2026/10/9 1:46:14

AI Native团队落地手册:从CLAUDE.md到多Agent编排的完整SDLC实践

1. 从“用AI写代码”到“AI Native团队”&#xff1a;差的不是工具&#xff0c;是整套协作骨架很多团队嘴上说着“我们已经 AI Native 了”&#xff0c;实际干的事无非是给每个人开了个 AI 编程助手的账号&#xff0c;然后继续用三年前那套需求评审、排期、联调、提测的流程。结…

作者头像 李华
网站建设 2026/10/9 1:45:44

Agent Memory架构设计与落地:为大模型打造外挂大脑

做过 Agent 的同学&#xff0c;应该都踩过同一个坑&#xff1a;Agent 明明能理解复杂指令&#xff0c;可你让它处理完跟用户的整个对话流程&#xff0c;或者隔几天再回来看它&#xff0c;发现它什么都不记得了。我也在这上面翻过车&#xff0c;最后发现根子不在模型能力上&…

作者头像 李华