- 后端
- 大数据
【免费下载链接】storm
Apache 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 负责四件事:
- 接收客户端发来的 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并调用fetchRequest;DRPCClient 内部维护一个静态的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.port | 3772 | DRPC server 接收客户端 Thrift 请求的端口 |
drpc.invocations.port | 3773 | DRPC 拓扑(spout/bolt)收发函数调用与结果的端口 |
drpc.http.port | 3774 | DRPC HTTP 服务端口 |
drpc.https.port | -1 | HTTPS 端口(-1 表示关闭) |
drpc.request.timeout.secs | 600 | 服务端请求超时时间(秒),超时后按SERVER_TIMEOUT失败 |
drpc.worker.threads | 64 | DRPC Thrift server 工作线程数 |
drpc.queue.size | 128 | DRPC Thrift 请求队列大小 |
drpc.max_buffer_size | 1048576 | DRPC 最大缓冲区字节数 |
drpc.invocations.threads | 64 | invocations 通道线程数 |
drpc.childopts | -Xmx768m | DRPC 进程 JVM 参数 |
drpc.authorizer.acl.filename | drpc-auth-acl.yaml | 基于 ACL 的授权器使用的文件名 |
drpc.disable.http.binding | true | 是否禁用 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 需要四步:
- 获取所有发过该 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是理解「批次聚合」的关键,其实现基于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
相关推荐
Apache Storm 分布式 RPC(DRPC)实战指南:原理、LinearDRPCTopologyBuilder 与拓扑开发
Apache Storm 分布式 RPC(DRPC)实战指南:原理、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 是一个分布式实时计算系统,本指南以
大数据流处理后端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考