资讯详情

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

📅 2026/10/9 1:46:55 | 华诺云谱 👁 阅读
Apache Storm Distributed RPC 完整实战:DRPC 原理、拓扑构建与集群部署
后端大数据【免费下载链接】stormApache Storm项目地址https://gitcode.com/gh_mirrors/storm22/storm点击查看免费下载Distributed RPCDRPC是 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 负责四件事接收客户端发来的 RPC 请求把请求函数调用投递给 Storm 拓扑从拓扑接收计算结果把结果回传给正在等待的客户端。从客户端视角看一次分布式 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 调用的完整数据流如下客户端把「函数名 参数」发给 DRPC server实现该函数的拓扑通过DRPCSpout从 DRPC server 读取函数调用流——DRPC server 会为每次函数调用打上唯一 id拓扑完成计算位于拓扑末端的ReturnResultsbolt 连接 DRPC server把「函数调用 id 结果」交还给它DRPC server 依据 id 匹配到正在等待的客户端解除其阻塞并把结果返回。这个「先入队、被拓扑取走、再回填结果」的过程在服务端由 DRPC.java 实现请求先进入按函数名分组的queues等待被拓扑 fetch被fetchRequest取走后转入requests等待结果回填returnResult通过 id 找到对应的OutstandingRequest并唤醒阻塞中的客户端。LinearDRPCTopologyBuilder一站式构建 DRPC 拓扑Storm 提供了 LinearDRPCTopologyBuilder 来自动化 DRPC 几乎所有的编排步骤包括搭建 spout把结果返回给 DRPC server为 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并调用fetchRequestDRPCClient 内部维护一个静态的localOverrideClient一旦存在本地覆盖新客户端就会直接使用该本地实例而不再走 Thrift 网络。测试代码可通过DRPCClient.LocalOverride实现了AutoCloseable临时注入本地 DRPC。远程模式 DRPC在真实集群上运行在真实集群上使用 DRPC 同样直接只需三步启动 DRPC server(s)配置 DRPC server 的位置把 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); }通过 curlHTTP 调用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 需要四步获取所有发过该 URL 的推文作者tweeters获取这些作者的全部粉丝followers对粉丝集合去重统计去重后的粉丝数量。单次 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)。拓扑按四个步骤执行GetTweeters获取发布该 URL 的用户。把[id, url]输入流转换为[id, tweeter]输出流每个urltuple 会映射出多个tweetertuple。GetFollowers获取这些用户的粉丝。把[id, tweeter]转换为[id, follower]。由于同一个人可能关注了多个发过该 URL 的作者跨所有 task 看follower tuple 会出现重复。PartialUniquer按 follower id 对流分组使同一个 follower 始终进入同一个 task这样每个 task 收到的是互不重叠的 follower 子集。当它收到针对某个请求 id 的全部 follower tuple 后发出自己子集的去重数量。CountAggregator汇总各PartialUniquertask 的部分计数完成 reach 计算。PartialUniquer是理解「批次聚合」的关键其实现基于BaseBatchBoltpublic class PartialUniquer extends BaseBatchBolt { BatchOutputCollector _collector; Object _id; SetString _followers new HashSetString(); 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接口。批处理 boltbatch 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-clientLinearDRPCTopologyBuilder内部也正是用它包装每个 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字段fieldsGroupingReturnResults连接 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指定SourceArgssingle()/all()与IdStreamSpecid 检测流再自行实现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三种方式对外提供服务。赞分享后端大数据【免费下载链接】stormApache Storm项目地址https://gitcode.com/gh_mirrors/storm22/storm点击查看免费下载相关推荐Apache Storm 分布式 RPCDRPC实战指南原理、LinearDRPCTopologyBuilder 与拓扑开发Apache Storm 分布式 RPCDRPC实战指南原理、LinearDRPCTopologyBuilder 与拓扑开发 Apache Storm 的大数据流处理后端Apache Storm UI REST API 完全指南集群监控、拓扑管理与 DRPC 调用实战Apache Storm UI REST API 完全指南集群监控、拓扑管理与 DRPC 调用实战 本文以 Apache Storm 官方文档 docs/ST大数据流处理后端Apache Storm 集群搭建完全指南从 ZooKeeper 到 DRPC 的实战部署手册Apache Storm 集群搭建完全指南从 ZooKeeper 到 DRPC 的实战部署手册 Apache Storm 是一个分布式实时计算系统本指南以大数据流处理后端上一篇Starscream路线图展望未来WebSocket功能发展预测下一篇Exchange Calendar未来路线图即将推出的5大令人期待的新功能创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
📝

华诺云谱内容团队

资深建站顾问 · 行业研究员

10年+企业数字化服务经验,专注智能建站、SEO优化与品牌营销,持续输出建站技巧、行业洞察与营销干货,已帮助5000+企业实现数字化增长。

你可能需要的服务

订阅华诺云谱资讯周报

每周一封,精选建站技巧、SEO与营销干货,直达邮箱。已有 8,000+ 企业主订阅,助你少走弯路。

↑