资讯动态

Apache Storm 消息传递实现深度解析:Worker 传输与 Task 路由的完整链路

发布时间:2026/10/9 5:03:47 来源:尧图企业网站定制
流处理后端大数据【免费下载链接】stormApache 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 时代分布式模式的消息发送走 ZeroMQzmq.clj本地模式走内存 Java 队列local.clj0.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 Queueworker-data在 Worker 启动时为每个 Worker 创建一个全局唯一的传输队列transfer-queue其缓冲大小、等待超时与等待策略分别由topology.transfer.buffer.size、topology.disruptor.wait.timeout.millis、topology.disruptor.wait.strategy三个配置控制worker.clj L197-L199。与此同时每个 executor线程拥有一个独立的 receive queuemk-receive-queue-mapworker.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 属于本 Workerlocal-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-connectionsworker.clj L268-L320负责维护与其他 Worker 之间的连接遵循原文档所述的两个触发条件定时触发由task.refresh.poll.secs默认 10 秒defaults.yaml L142驱动的定时器周期调度事件触发ZooKeeper 中的 assignment 版本变化通过assignment-version回调感知。其核心工作是计算本 Worker 出站任务所需的目标节点/端口worker-outbound-tasks结合 assignment 映射得出对比当前已建连接集合增量地建立新连接、关闭已废弃连接并更新cached-task-nodeport与cached-nodeport-socket两个缓存。这正是原文档所说维护 task - worker 映射的现行实现。一个值得注意的细节是连接就绪门控activate-worker-when-all-connections-readyworker.clj L367-L381会持续检查所有出站连接是否就绪ConnectionWithStatus.Status.Ready全部就绪后才把worker-active-flag置为 trueSpout/Bolt 才会被激活从而避免拓扑在连接未建立时就开始发射导致消息丢失。3.4 Transfer Thread单线程排空与批量发送原文档指出Worker 用单线程排空 transfer queue 并发送消息。当前实现由mk-transfer-tuples-handlerworker.clj L334-L350完成它作为 Disruptor 的 EventHandler 挂在 transfer queue 上使用TransferDrainer聚合一批消息在批量结束时一次性按task-nodeport分组、通过对应的 socket 发送随后清空 drainer。批量化显著降低了小消息场景下的网络开销。四、Task 侧路由决策的完整逻辑4.1 tasks-fn从分组函数到目标 Task ID 列表原文档描述的 routing map{stream id} - {component id} - {stream grouping function}在当前实现中对应 executor 数据中的stream-component-grouper。mk-tasks-fntask.clj L125-L175针对两种发射方式分别实现Regular emit遍历当前 stream 上每个下游组件的 grouping 函数调用(grouper task-id values)得到目标 task 集合汇总为out-tasks。若某 stream 声明为 direct 却做 regular emit会抛出IllegalArgumentExceptionDirect emit指定目标任务out-task-id先查出该任务所属组件及其分组若该组件对该 stream 采用 regular grouping 而非 direct grouping同样抛异常——这就是原文档所说的direct stream 只发给订阅它的 bolt在代码层面的强制约束。4.2 send-unanchored路由结果驱动实际传输send-unanchoredtask.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 插件机制TransportFactoryTransportFactory.makeContextTransportFactory.java L29-L56根据配置项storm.messaging.transportConfig.java L58反射加载传输插件若插件类实现IContext则直接实例化并调用prepare否则要求其提供makeContext(Map)静态工厂方法。这为替换传输层如历史 ZeroMQ、默认 Netty保留了标准扩展点。5.2 核心接口IContext / IConnection / TaskMessage传输抽象由三个接口/类构成IConnectionIConnection.java定义recv(flags, clientId)flags 0 表示阻塞、1 表示非阻塞与send单条或批量以及closeIContext定义prepare / bind / connect / term是每 Worker 一个的传输上下文TaskMessageTaskMessage.java L22-L53封装task短整型与message字节数组其serialize()输出格式为short task payload——这正是原文档所述虚拟端口接收 [task id, message] 二元组的线上格式。5.3 默认 Netty 实现与批处理编码默认配置下defaults.yaml L42传输插件为backtype.storm.messaging.netty.ContextContext.bind(storm-id, port)创建Server监听单一 TCP 端口Context.connect(...)建立到远端 Worker 的Client连接Context.java L73-L92出站侧由MessageBufferMessageBuffer.java L25-L57累积TaskMessage到MessageBatch批满即返回待发批次每个MessageBatchMessageBatch.java L83-L118按task(short 2B) len(int 4B) payload逐条编码末尾追加EOB_MESSAGE结束标记入站侧由MessageDecoderMessageDecoder.java L29-L144在单次调用中尽可能解码多条消息并区分控制消息负编码、SASL 令牌编码 -500与普通 TaskMessagetask 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默认 1defaults.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 queueworker.clj L99-L110。也就是说当前实现用接收线程 task 到 executor 的队列映射取代了旧版virtual_port.clj与内存 ZeroMQ 端口但单端口接入、按 task id 分流的核心语义一脉相承。七、本地模式纯内存队列实现本地模式无需任何网络设施即可运行对应 local.cljLocalContext.prepare创建全局queues-mapstorm-id-port为键、LinkedBlockingQueue为值与锁bind返回持有该队列的LocalConnectionconnect返回不带队列的发送端send直接put一个TaskMessage到目标队列recv支持阻塞flags 0与非阻塞flags 1poll两种取法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-L1276topology.disruptor.wait.strategycom.lmax.disruptor.BlockingWaitStrategyDisruptor 等待策略可选 block/yield/sleep/spintopology.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 承载 tasktuple 的传输以 Worker 为单位、路由以 task 为单位。九、消息生命周期全景回顾将前述各节串联一条 tuple 的完整旅程如下Spout/Bolt 发射Task 调用tasks-fnregular 走分组函数、direct 走目标任务 id 校验得到目标 task 集合task.cljTask 提交通过 Worker 注入的transfer-fn提交[task, tuple]对task.clj L107-L123Worker 分流本地目标直接disruptor/publish进目标 executor 的 receive queue远端目标经 Kryo 序列化后写入共享 transfer queueworker.clj L117-L149单线程发送transfer thread 批量排空 transfer queueTransferDrainer按节点/端口分组经 Netty 连接批量发出worker.clj L334-L350远端接收目标 Worker 的接收线程从绑定端口批量读取按 task id 分发到对应 executor 的 receive queueloader.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 等测试用例验证上述行为。赞分享流处理后端大数据【免费下载链接】stormApache Storm项目地址https://gitcode.com/gh_mirrors/storm26/storm点击查看免费下载相关推荐Apache Storm 消息传递实现深度解析从 Tuple 发射到跨 Worker 传输的完整链路Apache Storm 消息传递实现深度解析从 Tuple 发射到跨 Worker 传输的完整链路 本文以 Apache Storm 官方文档 Messag大数据流处理后端gte-small安全与隐私考虑企业级文本嵌入部署的最佳实践gte small安全与隐私考虑企业级文本嵌入部署的最佳实践 在当今数据驱动的商业环境中文本嵌入技术作为连接自然语言与机器学习系统的关键桥梁其安全与隐私保突破消息传递瓶颈nats-server智能路由决策算法深度解析突破消息传递瓶颈nats server智能路由决策算法深度解析 你是否在分布式系统中遇到过消息延迟飙升、网络带宽浪费或节点负载不均的问题作为NATS高性能后端消息队列消息路由通信上一篇如何3步完成《艾尔登法环》角色存档迁移终极免费工具完整指南下一篇让经典游戏手柄重获新生XOutput协议转换工具的终极指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑