- 流处理
- 后端
- 大数据
【免费下载链接】storm
Apache 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。
在当前仓库中,这一演进体现在两个层面:
- 队列层:disruptor.clj 封装了
DisruptorQueue,并提供block/yield/sleep/spin四种等待策略以及multi-threaded/single-threaded两种 claim 策略的映射(见该文件 L27-L51); - 网络层:消息传输通过可插拔的
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 之间的连接,遵循原文档所述的两个触发条件:
- 定时触发:由
task.refresh.poll.secs(默认 10 秒,defaults.yaml L142)驱动的定时器周期调度; - 事件触发: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]后内存路由给实际任务——在当前实现中以更明确的形态保留:
- Worker 启动时通过
msg-loader/launch-receive-thread!(loader.clj L62-L84)拉起若干接收线程,数量由topology.worker.receiver.thread.count(默认 1,defaults.yaml L139)决定; - 每个接收线程循环调用
(.recv socket 0 thread-id)从绑定端口批量拉取消息(loader.clj L27-L55); - 收到
task == -1的特殊消息时视为关闭通知,关闭 socket 并退出线程; - 普通消息聚成
[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.transport | backtype.storm.messaging.netty.Context | 传输插件类名,经TransportFactory反射加载 |
topology.transfer.buffer.size | 1024 | 每 Worker transfer queue 容量(按批计) |
topology.executor.receive.buffer.size | 1024 | 每个 executor receive queue 容量(按批计,需为 2 的幂,见 Config.java L1275-L1276) |
topology.disruptor.wait.strategy | com.lmax.disruptor.BlockingWaitStrategy | Disruptor 等待策略,可选 block/yield/sleep/spin |
topology.disruptor.wait.timeout.millis | 1000 | Disruptor 消费等待超时(毫秒) |
topology.worker.receiver.thread.count | 1 | 每 Worker 接收线程数,可提升入站吞吐 |
task.refresh.poll.secs | 10 | 连接刷新定时周期(秒),assignment 变化时也会触发 |
topology.testing.always.try.serialize | false | 开启后所有 tuple 先做序列化断言,仅建议测试环境使用 |
topology.stats.sample.rate | 0.05 | emitted/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 的完整旅程如下:
- Spout/Bolt 发射:Task 调用
tasks-fn(regular 走分组函数、direct 走目标任务 id 校验)得到目标 task 集合(task.clj); - Task 提交:通过 Worker 注入的
transfer-fn提交[task, tuple]对(task.clj L107-L123); - Worker 分流:本地目标直接
disruptor/publish进目标 executor 的 receive queue;远端目标经 Kryo 序列化后写入共享 transfer queue(worker.clj L117-L149); - 单线程发送:transfer thread 批量排空 transfer queue,
TransferDrainer按节点/端口分组,经 Netty 连接批量发出(worker.clj L334-L350); - 远端接收:目标 Worker 的接收线程从绑定端口批量读取,按 task id 分发到对应 executor 的 receive queue(loader.clj L27-L55);
- 消费执行: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
相关推荐
Apache Storm 消息传递实现深度解析:从 Tuple 发射到跨 Worker 传输的完整链路
Apache Storm 消息传递实现深度解析:从 Tuple 发射到跨 Worker 传输的完整链路 本文以 Apache Storm 官方文档 Messag
大数据流处理后端gte-small安全与隐私考虑:企业级文本嵌入部署的最佳实践
gte small安全与隐私考虑:企业级文本嵌入部署的最佳实践 在当今数据驱动的商业环境中,文本嵌入技术作为连接自然语言与机器学习系统的关键桥梁,其安全与隐私保
突破消息传递瓶颈:nats-server智能路由决策算法深度解析
突破消息传递瓶颈:nats server智能路由决策算法深度解析 你是否在分布式系统中遇到过消息延迟飙升、网络带宽浪费或节点负载不均的问题?作为NATS(高性能
后端消息队列消息路由通信
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考