☰
Apache Storm 消息传递实现深度解析:Worker 传输与 Task 路由的完整链路
2026/10/9 5:03:45 网站建设 项目流程
  • 流处理
  • 后端
  • 大数据

【免费下载链接】storm

Apache Storm

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

导读

本文基于 Storm 仓库中的 Message-passing-implementation.md 展开,梳理 Apache Storm 中 tuple 从"被 emit"到"被目标任务接收"的完整链路。原文档面向 0.7.x 版本撰写,并明确标注"0.8.0 已用 Disruptor 重构消息传递基础设施",因此本文将结合当前仓库(基于 Disruptor + Netty 的实现)重新走查这一机制。读完本文,你将掌握:Worker 与 Task 在消息传递中的职责边界、transfer queue 与 receive queue 的角色、direct stream 与 regular stream 的路由差异,以及消息传输相关的全部关键配置参数。

一、整体职责划分:Worker 负责传输,Task 负责路由

Storm 的消息传递实现遵循一个清晰的分工原则:

  • Worker(进程)负责消息传输:管理到其他 Worker 的网络连接、执行序列化、维护每 Worker 唯一的"transfer queue"、单线程批量发送;
  • Task(线程内逻辑单元)负责消息路由:决定一个 tuple 应该发给哪些目标任务,然后调用 Worker 提供的 transfer 函数完成实际投递。

这条主线在原文档中就已确立,在 Disruptor 重构后依然成立,只是底层队列与网络层实现被整体替换。理解这一分工是读懂下文所有细节的前提。

二、从 ZeroMQ 到 Disruptor + Netty:一次基础设施重构

原文档开篇即给出重要提示:"this walkthrough is out of date as of 0.8.0. 0.8.0 revamped the message passing infrastructure to be based on the Disruptor"。这意味着:

  • 0.7.x 时代,分布式模式的消息发送走 ZeroMQ(zmq.clj),本地模式走内存 Java 队列(local.clj);
  • 0.8.0 之后,队列层全面切换到 LMAX Disruptor(环形缓冲队列),网络层则由 Netty 取代 ZeroMQ。

在当前仓库中,这一演进体现在两个层面:

  1. 队列层:disruptor.clj 封装了DisruptorQueue,并提供block/yield/sleep/spin四种等待策略以及multi-threaded/single-threaded两种 claim 策略的映射(见该文件 L27-L51);
  2. 网络层:消息传输通过可插拔的IContext插件实现,默认插件为 Netty 实现(见下文第五节)。

因此,本文后续内容均以当前仓库的实现为准,原文档中关于 ZeroMQ 与virtual_port.clj的细节仅作为历史背景提及。

三、Worker 侧:发送路径的四个关键构件

3.1 每 Worker 唯一的 Transfer Queue

worker-data在 Worker 启动时为每个 Worker 创建一个全局唯一的传输队列transfer-queue,其缓冲大小、等待超时与等待策略分别由topology.transfer.buffer.size、topology.disruptor.wait.timeout.millis、topology.disruptor.wait.strategy三个配置控制(worker.clj L197-L199)。

与此同时,每个 executor(线程)拥有一个独立的 receive queue(mk-receive-queue-map,worker.clj L151-L159),容量由topology.executor.receive.buffer.size控制。这样便形成了"一条共享出站队列 + 多条按 executor 隔离的入站队列"的拓扑结构。

3.2 Transfer Function:序列化与分流

Worker 通过mk-transfer-fn向各 executor 提供统一的传输函数(worker.clj L117-L149)。它对一批[task, tuple]对执行如下逻辑:

  • 若目标 task 属于本 Worker(local-tasks),直接收集到local列表,走local-transfer路径,不经序列化;
  • 若目标 task 在远端,则使用线程安全的KryoTupleSerializer将 tuple 序列化为字节,封装成TaskMessage,放入remoteMap;
  • 本地批调用local-transfer分发,远端批通过disruptor/publish写入transfer-queue。

原文档强调"serializer 是线程安全的"(指向KryoTupleSerializer),这在当前实现中依旧成立:多个 executor 线程可以并发调用同一个 transfer 函数而无需额外加锁。

另外,topology.testing.always.try.serialize(默认false,见 defaults.yaml L199)打开后,会先对所有 tuple 做一次序列化断言(assert-can-serialize),用于在本地模式或调试阶段提前暴露序列化问题,线上环境应保持关闭。

3.3 Refresh Connections:连接的生命周期管理

mk-refresh-connections(worker.clj L268-L320)负责维护与其他 Worker 之间的连接,遵循原文档所述的两个触发条件:

  1. 定时触发:由task.refresh.poll.secs(默认 10 秒,defaults.yaml L142)驱动的定时器周期调度;
  2. 事件触发:ZooKeeper 中的 assignment 版本变化(通过assignment-version回调感知)。

其核心工作是计算"本 Worker 出站任务所需的目标节点/端口"(worker-outbound-tasks结合 assignment 映射得出),对比当前已建连接集合,增量地建立新连接、关闭已废弃连接,并更新cached-task->node+port与cached-node+port->socket两个缓存。这正是原文档所说"维护 task -> worker 映射"的现行实现。

一个值得注意的细节是连接就绪门控:activate-worker-when-all-connections-ready(worker.clj L367-L381)会持续检查所有出站连接是否就绪(ConnectionWithStatus.Status.Ready),全部就绪后才把worker-active-flag置为 true,Spout/Bolt 才会被激活,从而避免拓扑在连接未建立时就开始发射导致消息丢失。

3.4 Transfer Thread:单线程排空与批量发送

原文档指出"Worker 用单线程排空 transfer queue 并发送消息"。当前实现由mk-transfer-tuples-handler(worker.clj L334-L350)完成:它作为 Disruptor 的 EventHandler 挂在 transfer queue 上,使用TransferDrainer聚合一批消息,在批量结束时一次性按task->node+port分组、通过对应的 socket 发送,随后清空 drainer。批量化显著降低了小消息场景下的网络开销。

四、Task 侧:路由决策的完整逻辑

4.1 tasks-fn:从分组函数到目标 Task ID 列表

原文档描述的 routing map({stream id} -> {component id} -> {stream grouping function})在当前实现中对应 executor 数据中的stream->component->grouper。mk-tasks-fn(task.clj L125-L175)针对两种发射方式分别实现:

  • Regular emit:遍历当前 stream 上每个下游组件的 grouping 函数,调用(grouper task-id values)得到目标 task 集合,汇总为out-tasks。若某 stream 声明为 direct 却做 regular emit,会抛出IllegalArgumentException;
  • Direct emit:指定目标任务out-task-id,先查出该任务所属组件及其分组,若该组件对该 stream 采用 regular grouping 而非 direct grouping,同样抛异常——这就是原文档所说的"direct stream 只发给订阅它的 bolt"在代码层面的强制约束。

4.2 send-unanchored:路由结果驱动实际传输

send-unanchored(task.clj L107-L123)是 Task 侧调用 transfer 的入口:它构造TupleImpl(携带 values、源 task id、stream id),然后对tasks-fn返回的每个目标任务调用由 Worker 注入的transfer-fn。也就是说,Task 只负责"算出来发给谁","怎么发出去"完全交给 Worker,与第三节的职责划分一一对应。

4.3 度量与调试钩子

在路由过程中,mk-tasks-fn还内建了统计采样:通过emit-sampler(由topology.stats.sample.rate控制采样率)触发stats/emitted-tuple!与stats/transferred-tuples!计数,并在topology.debug开启时打印每次发射的目标与数值。这些数据最终汇入 UI 展示的 emitted/transferred 指标。

五、消息协议层:可插拔传输插件

5.1 插件机制:TransportFactory

TransportFactory.makeContext(TransportFactory.java L29-L56)根据配置项storm.messaging.transport(Config.java L58)反射加载传输插件:若插件类实现IContext则直接实例化并调用prepare,否则要求其提供makeContext(Map)静态工厂方法。这为替换传输层(如历史 ZeroMQ、默认 Netty)保留了标准扩展点。

5.2 核心接口:IContext / IConnection / TaskMessage

传输抽象由三个接口/类构成:

  • IConnection(IConnection.java):定义recv(flags, clientId)(flags 0 表示阻塞、1 表示非阻塞)与send(单条或批量)以及close;
  • IContext:定义prepare / bind / connect / term,是每 Worker 一个的传输上下文;
  • TaskMessage(TaskMessage.java L22-L53):封装task(短整型)与message(字节数组),其serialize()输出格式为short task + payload——这正是原文档所述"虚拟端口接收 [task id, message] 二元组"的线上格式。

5.3 默认 Netty 实现与批处理编码

默认配置下(defaults.yaml L42)传输插件为backtype.storm.messaging.netty.Context:

  • Context.bind(storm-id, port)创建Server监听单一 TCP 端口,Context.connect(...)建立到远端 Worker 的Client连接(Context.java L73-L92);
  • 出站侧由MessageBuffer(MessageBuffer.java L25-L57)累积TaskMessage到MessageBatch,批满即返回待发批次;
  • 每个MessageBatch(MessageBatch.java L83-L118)按task(short 2B) + len(int 4B) + payload逐条编码,末尾追加EOB_MESSAGE结束标记;
  • 入站侧由MessageDecoder(MessageDecoder.java L29-L144)在单次调用中尽可能解码多条消息,并区分控制消息(负编码)、SASL 令牌(编码 -500)与普通 TaskMessage(task >= 0)。

此外,Netty 客户端还受storm.messaging.netty.buffer.size(消息批大小)、storm.messaging.netty.max.retries、storm.messaging.netty.min/max.sleep.ms(重连退避窗口)等参数约束(见 Client.java L135-L150 的读取逻辑)。

六、接收路径:从虚拟端口到 Executor 接收队列

原文档描述的"虚拟端口"模型——每个 Worker 监听单一 TCP 端口,收到[task id, message]后内存路由给实际任务——在当前实现中以更明确的形态保留:

  1. Worker 启动时通过msg-loader/launch-receive-thread!(loader.clj L62-L84)拉起若干接收线程,数量由topology.worker.receiver.thread.count(默认 1,defaults.yaml L139)决定;
  2. 每个接收线程循环调用(.recv socket 0 thread-id)从绑定端口批量拉取消息(loader.clj L27-L55);
  3. 收到task == -1的特殊消息时视为关闭通知,关闭 socket 并退出线程;
  4. 普通消息聚成[task, message]批次后交给transfer-local-fn,由其按task -> short-executor映射分组,最终disruptor/publish到对应 executor 的 receive queue(worker.clj L99-L110)。

也就是说,当前实现用"接收线程 + task 到 executor 的队列映射"取代了旧版virtual_port.clj与内存 ZeroMQ 端口,但"单端口接入、按 task id 分流"的核心语义一脉相承。

七、本地模式:纯内存队列实现

本地模式无需任何网络设施即可运行,对应 local.clj:

  • LocalContext.prepare创建全局queues-map(storm-id-port为键、LinkedBlockingQueue为值)与锁;
  • bind返回持有该队列的LocalConnection,connect返回不带队列的"发送端";
  • send直接put一个TaskMessage到目标队列,recv支持阻塞(flags 0)与非阻塞(flags 1,poll)两种取法(local.clj L32-L54)。

原文档将其动机概括为"本地使用 Storm 无需安装 ZeroMQ";在现版本中,本地模式还让开发者可以脱离真实网络环境调试拓扑逻辑。Worker 在本地模式下同样会创建 Disruptor 队列与接收线程,只是底层 socket 换成了内存队列(见worker-data中mq-context的构造与mk-local-context的注入路径 loader.clj L24-L25)。

八、关键配置参数速查

以下参数共同决定消息传递的吞吐与延迟特性,均以仓库 defaults.yaml 与 Config.java 为准:

配置项默认值作用
storm.messaging.transportbacktype.storm.messaging.netty.Context传输插件类名,经TransportFactory反射加载
topology.transfer.buffer.size1024每 Worker transfer queue 容量(按批计)
topology.executor.receive.buffer.size1024每个 executor receive queue 容量(按批计,需为 2 的幂,见 Config.java L1275-L1276)
topology.disruptor.wait.strategycom.lmax.disruptor.BlockingWaitStrategyDisruptor 等待策略,可选 block/yield/sleep/spin
topology.disruptor.wait.timeout.millis1000Disruptor 消费等待超时(毫秒)
topology.worker.receiver.thread.count1每 Worker 接收线程数,可提升入站吞吐
task.refresh.poll.secs10连接刷新定时周期(秒),assignment 变化时也会触发
topology.testing.always.try.serializefalse开启后所有 tuple 先做序列化断言,仅建议测试环境使用
topology.stats.sample.rate0.05emitted/transferred 统计采样率
storm.messaging.netty.buffer.size—Netty 客户端消息批大小(见 Client.java L135)

其中等待策略的选择值得注意:block(默认)在低吞吐场景下 CPU 占用低,但disruptor.clj注释提醒 block 策略在 Trident 单批处理场景下需要配合超时机制避免消费者假阻塞(disruptor.clj L43-L46);追求极致延迟可换用yield或spin,代价是更高的 CPU 占用。

上图展示消息传递涉及的三个层级:Worker 进程承载 executor 线程,executor 承载 task,tuple 的传输以 Worker 为单位、路由以 task 为单位。

九、消息生命周期全景回顾

将前述各节串联,一条 tuple 的完整旅程如下:

  1. Spout/Bolt 发射:Task 调用tasks-fn(regular 走分组函数、direct 走目标任务 id 校验)得到目标 task 集合(task.clj);
  2. Task 提交:通过 Worker 注入的transfer-fn提交[task, tuple]对(task.clj L107-L123);
  3. Worker 分流:本地目标直接disruptor/publish进目标 executor 的 receive queue;远端目标经 Kryo 序列化后写入共享 transfer queue(worker.clj L117-L149);
  4. 单线程发送:transfer thread 批量排空 transfer queue,TransferDrainer按节点/端口分组,经 Netty 连接批量发出(worker.clj L334-L350);
  5. 远端接收:目标 Worker 的接收线程从绑定端口批量读取,按 task id 分发到对应 executor 的 receive queue(loader.clj L27-L55);
  6. 消费执行:executor 从 receive queue 取批并执行 Bolt/Spout 逻辑,完成一次消息传递闭环。

十、结语

从 0.7.x 的 ZeroMQ 时代到当前的 Disruptor + Netty 架构,Storm 消息传递的"Worker 管传输、Task 管路由"设计骨架始终未变,变的只是队列与网络的具体载体。理解 transfer queue 与 receive queue 的对称设计、direct/regular 两种路由语义以及task.refresh.poll.secs等参数的作用,是诊断拓扑延迟、吞吐瓶颈与消息丢失问题的前提。感兴趣的读者可以继续深入 worker.clj、task.clj 以及 netty 目录 下的实现,并结合 messaging_test.clj 等测试用例验证上述行为。

  • 流处理
  • 后端
  • 大数据

【免费下载链接】storm

Apache Storm

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

相关推荐

上一篇:如何3步完成《艾尔登法环》角色存档迁移?终极免费工具完整指南
下一篇:让经典游戏手柄重获新生:XOutput协议转换工具的终极指南

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

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

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

立即咨询