尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Storm Java API实战:从拓扑构建到Kafka实时管线调优

Storm Java API实战:从拓扑构建到Kafka实时管线调优 说句实话刚接到实时订单监控这个需求时我第一反应是准备直接上 Flink。但翻了翻团队现有的大数据组件、运维脚本和线上部署环境最后老老实实用回了 Storm。原因并不复杂Storm 的 Java API 足够直接一套 Topology 从本地模式到集群提交几乎无缝衔接团队也有现成的 Storm ZooKeeper 环境可以依赖。这篇文章就围绕 Storm Java API 实战来写从零开始搭一个能跑起来的实时数据处理应用覆盖 Spout/Bolt 编写、分组策略、Ack 机制、Kafka 接入、上集群前的调优和监控。适合刚接触 Storm、有一点 Java 基础、想把实时数据处理真正落地的人参考。1. 为什么还在用 Storm老框架的适用边界与选型逻辑1.1 Storm到底解决了什么问题Storm 是 Twitter 开源、后来成为 Apache 顶级项目的分布式实时计算框架。它把数据处理抽象成一张 Topology——你可以把它想象成一条流水线数据从源头 Topic 进来经过若干处理节点最后落到目标系统。和 Spark Streaming 那类微批模式不同Storm 是真正的逐条流式处理数据到达一条就处理一条从 Spout 发射到 Bolt 完成处理的链路延迟通常在毫秒到秒级。我遇到的那个需求最初是轮询订单数据库每隔几秒扫一次新增记录再把结果刷新到大屏上。逻辑不复杂但问题很明显轮询有固定的延迟窗口数据量一涨数据库压力也跟着涨。把逻辑迁移到 Storm 之后订单事件通过 Kafka 实时推送进来Storm 负责过滤、统计和告警整个链路从“隔几分钟看一眼”变成了“秒出结果”。这就是 Storm 的核心价值——让数据在流动过程中被处理而不是等数据落库之后再去扫。1.2 和Flink/Spark Streamin比该怎么选我知道很多人会问2024 年了为什么不学 Flink我不否认 Flink 在现代流处理领域的领先地位。原生状态管理、窗口机制、Event Time 语义、精确一次性保证这些确实是 Flink 的强项。但回到真实项目里选型往往不是“哪个最强选哪个”而是“哪个最适合现有环境”。Storm 的优势在于三件事部署运维链路短Nimbus、Supervisor、UI 三个角色配合 ZooKeeper 做协调旧集群资料一大堆排障经验也成熟。API 足够简单核心就 Spout 和 Bolt 两个抽象业务逻辑写起来快团队上手成本低。逐条处理延迟低没有微批的攒批过程对延迟极其敏感的告警场景反而更合适。当然如果是从零搭建新系统并且有复杂状态计算、CEP、实时 SQL 这类需求我建议认真评估 Flink。但如果你维护的是存量 Storm 项目、或者需求主要是过滤、分流、简单聚合和告警那 Storm 完全够用而且它不是学完就废——Topology 里的分组策略、消息保障语义、并行度的思路在整个流处理领域都是通用的。1.3 学Storm真的不亏吗不亏。就算你最后用 FlinkStorm 里面几个核心概念——Spout/Bolt 的责任划分、Stream Grouping 的数据路由规则、Ack 机制下的消息重发语义——这些在 Flink 里都能找到对应物。理解消息在分布式流水线里是怎么流转的、怎么被确认的、失败了会怎样比单纯学会 API 更有价值。而且 Storm 的代码量小非常适合拿来理解“分布式流处理”背后最基本的机制。2. 环境搭建与Topology骨架先让数据在本地流动起来2.1 依赖引入和开发模式选择新建 Maven 项目后第一步是引入 Storm 依赖。以 1.2.x 版本为例dependency groupIdorg.apache.storm/groupId artifactIdstorm-core/artifactId version1.2.3/version scopeprovided/scope /dependency这里需要注意provided这个 scope。集群模式下Storm 集群本身已经带了 storm-core 的 jar 包你的应用 jar 里如果再打一份很容易触发类冲突和奇怪的 NoSuchMethodError。所以打包时排除掉 Storm 依赖只提交业务代码。本地开发时情况不太一样。如果你直接在 IDE 里跑 main 方法provided作用域在运行期可能不包含依赖导致本地模式启动失败。我这里的做法是开发阶段用一个专门的 profile 引入 runtime 依赖打成包含依赖的 fat jar 用于本地测试提交集群时再用provided。Storm 官方给的 maven-shade 插件配置里面也体现了这个思路——把 Storm 依赖从最终的 bundle jar 里排除。2.2 一个最快跑通的Topology三件套一个最简单的 Storm Topology 通常由三部分组成Spout 负责发射数据Bolt 负责处理和输出TopologyBuilder 负责把它们串联起来。先看 Spoutpublic class OrderSpout implements IRichSpout { private SpoutOutputCollector collector; private AtomicInteger seq new AtomicInteger(0); Override public void open(Map conf, TopologyContext context, SpoutOutputCollector collector) { this.collector collector; } Override public void nextTuple() { String orderId order- seq.incrementAndGet(); long userId ThreadLocalRandom.current().nextLong(10000); double amount ThreadLocalRandom.current().nextDouble(10, 5000); String status PAID; collector.emit(new Values(orderId, userId, amount, status), orderId); } Override public void declareOutputFields(OutputFieldsDeclarer declarer) { declarer.declare(new Fields(orderId, userId, amount, status)); } Override public void ack(Object msgId) { } Override public void fail(Object msgId) { } }nextTuple是 Spout 的核心方法框架会不断循环调用它来拉取新数据。这里有个我在新手阶段踩过的坑nextTuple必须快速返回绝不能在里面做数据库轮询、sleep 等阻塞操作。一旦阻塞整个 Spout task 就卡死了后面所有流程全部停滞。正确的做法是从内存队列或 Kafka 这样的外部系统尽量快地把数据emit出去。Bolt 端需要实现IRichBolt核心是execute方法public class FilterBolt implements IRichBolt { private OutputCollector collector; Override public void prepare(Map conf, TopologyContext context, OutputCollector collector) { this.collector collector; } Override public void execute(Tuple tuple) { String status tuple.getStringByField(status); if (PAID.equals(status)) { collector.emit(tuple, new Values( tuple.getStringByField(orderId), tuple.getLongByField(userId), tuple.getDoubleByField(amount) )); } collector.ack(tuple); } Override public void declareOutputFields(OutputFieldsDeclarer declarer) { declarer.declare(new Fields(orderId, userId, amount)); } }注意到execute里ack是显式调用的。很多第一次写的人会忘记这一步结果就是数据明明处理完了但 Spout 端的ack回调一直不来最后触发超时重发。这个在第四章会展开。最后是 Topology 的组装和启动public class OrderTopology { public static void main(String[] args) throws Exception { TopologyBuilder builder new TopologyBuilder(); builder.setSpout(order-spout, new OrderSpout(), 2); builder.setBolt(filter-bolt, new FilterBolt(), 2) .shuffleGrouping(order-spout); builder.setBolt(count-bolt, new CountBolt(), 4) .fieldsGrouping(filter-bolt, new Fields(userId)); Config config new Config(); config.setNumWorkers(3); LocalCluster cluster new LocalCluster(); cluster.submitTopology(order-topology, config, builder.createTopology()); Thread.sleep(60000); cluster.shutdown(); } }setSpout和setBolt的第二个参数是并行度表示启动几个实例shuffleGrouping和fieldsGrouping是数据路由规则第三章会重点说。本地模式的好处是不需要真实集群直接可以把整条拓扑跑起来看日志。2.3 确认数据真的流过了我习惯在本地开发时每层 Bolt 入口打一条简短的 debug 日志确认数据路径通没通。比如 FilterBolt 的execute第一行就打印订单号CountBolt 打印收到的数量。直接看 UI 和生产数据可能被各种指标干扰但本地日志能最快告诉你“Spout 有没有发出来、Bolt 有没有收到、最后的输出有没有落下去”。还有一个实用小技巧本地跑通后先cluster.shutdown()一次看看进程能不能干净退出。很多 Topology 在提交后无法正常关闭多半是close方法里的资源没释放干净这个习惯能帮你提前发现连接泄漏这类问题。3. 分组策略与并行度实时管线里最容易被忽略的“方向盘”3.1 六类Stream Grouping的行为差异与选型Stream Grouping 决定了数据从上游发射之后到底路由到下游的哪一个 task。我见过不少项目拓扑能跑、指标也正常但结果就是不对——十有八九是分组方式选得不对。以下是 Storm 里几类核心分组策略的行为差异分组方式路由规则典型场景需要警惕的点Shuffle Grouping随机轮流分发负载均衡、无状态过滤不能保序不能保证同一 key 到同一 taskFields Grouping按指定字段哈希按用户/订单聚合、保序字段选错等于白选Global Grouping全部送到下游第 0 个 task全局合并、排序单 task 会成为瓶颈All Grouping广播给所有 task配置同步、缓存刷新数据量会被放大 N 倍Direct Grouping由发射端指定接收 task手动精细控制必须配合 emitDirect 使用Local or Shuffle Grouping优先 worker 内随机减少跨进程网络传输只在有多个 worker 时有意义我个人最常用的就是 Fields Grouping 和 Shuffle Grouping。做按用户维度的统计和告警必须用 Fields Grouping做纯过滤、格式转换这种无状态逻辑Shuffle Grouping 最合适它能把负载均匀摊到下游所有并行实例上避免某个 task 过热而其他 task 闲置。Global Grouping 一定要慎用。全量数据汇聚到单个 task系统吞吐上限就被那一个 task 卡死了。如果确实需要全量聚合建议先考虑用 Fields Grouping 做分桶聚合最后再单独用一个 Global Grouping 只做归并这样至少前段压力是分摊的。3.2 Worker、Executor、Task三层并行度Storm 的并行度不是单个概念而是分层的Worker运行在 Supervisor 节点上的 JVM 进程Config.setNumWorkers 控制整拓扑的 Worker 数。ExecutorWorker 内部的线程setSpout/setBolt 的第二个参数控制的是 Executor 数量。TaskExecutor 内部的最小调度单元setNumTasks 可以设置默认为 1 Executor 对应 1 Task。很多人有个误解设置了setNumTasks(100)就以为并行度提高了。实际上真正决定执行并发的是 Executor 数也就是线程数。Task 增大通常意味着同一个 Executor 要轮询处理更多 Task 的数据并不会带来处理能力的提升。调整并行度时先想清楚 CPU 核数和数据量一个 Worker 内的 Executor 数量不要超过 CPU 核数太多否则线程切换的开销反而吃掉性能。另一个经验并行度并不是越大越好。我的一个告警拓扑最初给下游统计 Bolt 开了 16 个 Executor结果因为是按用户 ID 做 Fields Grouping用户分布不均出现严重的数据倾斜部分 Executor 忙到饱和部分一直闲着。后来把 Executor 数降回 8数据分布反而更均匀。这里的教训是分组策略决定了数据怎么分桶并行度决定了有多少个桶两者要一起调单独动哪一个都可能出问题。3.3 一个因为分组选错导致的乱序问题直接说一个我实际处理过的故障。当时有一个实时订单流水处理拓扑下游要按用户维度输出“最近 5 笔订单”的聚合结果。第一个版本为了负载均衡把处理 Bolt 的上游连接配成了 Shuffle Grouping。结果线上很快出现反馈同一个用户的多笔订单到达下游的顺序经常是乱掉的比如订单 1、订单 2 发的最后聚合结果里反而订单 2 排在前面。排查链路是这样的先看 UI所有 Executor 的负载都很正常没有倾斜再看日志发现同一个 userId 的订单会随机落到不同的 task 上每个 task 各自维护自己的本地队列顺序自然没法保证。最后把 Shuffle Grouping 改成fieldsGrouping(userId)同一个用户的所有订单被哈希到同一个 Executor顺序才稳定下来。这个场景也说明了另一个问题一旦同一 key 的数据被路由到同一个 Executor并行度就受限了——同一个 key 在同一时刻只能被一个 task 处理。所以 Fields Grouping 的并行度上限是 key 的分布数量设计时要对流量做充分评估。4. 消息可靠性Ack机制与“不丢不重”的权衡4.1 锚定与Ack/Fail回调消息树是怎么收拢的Storm 的可靠性保障核心是 Ack 机制。当 Spout 发射一条 tuple 并带上 messageId 的时候它就开启了一棵“消息树”。下游 Bolt 每次emit(tuple, newValues)时新 tuple 会和旧 tuple 建立锚定关系当所有锚定关系都走到尾部并且被成功 ack整棵树才算完成Spout 的ack(msgId)才会被回调。如果任何一个环节调用fail(tuple)或者超时Spout 的fail(msgId)就会被触发。这个机制听起来简单实际写 Bolt 时很容易漏掉。尤其是过滤场景不少新手会把collector.ack(tuple)写在 filter 判断的 if 分支里不满足条件的数据就没有 ack。这个行为本身不会报错但 Spout 就会一直等这条消息的树最后超时之后触发 fail 重发。如果重发又没处理整条链路就会出现持续的数据膨胀和重复。所以我在项目里定了两条写 Bolt 的纪律execute方法里无论走哪个分支最终都必须调用 ack 或 fail。如果 Bolt 不再向外发射数据说明这条分支的使命已经完成必须 ack 掉表示“我这里处理完了”。4.2 超时设置与重发窗口Storm 通过Config.TOPOLOGY_MESSAGE_TIMEOUT_SECS控制一条消息从 Spout 发出到整棵树完成 ack 的超时时间默认值是 30 秒。如果链路处理链路偏长比如要写文件、调外部接口30 秒很容易超时。超时之后 Spout 的fail会被触发如果没有妥善处理最直接的表现就是消息被重新发射下游看到重复。我当时的告警拓扑要把结果写入数据库多次实测后发现默认超时时间不够。我的做法是先做了一次全链路耗时统计算出 P99 延迟然后把超时时间设为 P99 的三到四倍留出足够的抖动空间。同时把 Spout 的fail回调实现成“重发一次”并配合幂等写入兜底而不是无限重试。这里有个值得注意的点超时时间设得越长消息出问题后重新被处理的间隔就越长设得太短正常处理稍微慢一点就误判为失败成批重发反而放大重复。这个参数最好基于真实链路的延迟数据来调整不要拍脑袋。4.3 开Ack的代价和幂等兜底开启 Ack 是有代价的Storm 需要记录和回溯 tuple 之间的锚定关系。实测下来在简单过滤和字段转换场景里开启 Ack 比关闭 Ack 的吞吐大约低 20% 到 30%。如果下游能接受少量重复并且你的处理逻辑是纯幂等的确实可以通过关闭 Ack 来提升吞吐但前提是你能接受“可能有消息丢”。我在生产环境里的处理方式是“至少一次 幂等兜底”。也就是说开启 Ack保证每条消息至少被处理一次同时下游写入端做好幂等设计。以写 MySQL 为例订单统计表以 orderId 作为唯一键用insert ... on duplicate key update保证同一订单重复插入时只更新一次。Redis 写告警状态用 SETNX只有当 key 不存在时才真正写入。这样就算 Storm 因为超时或 fail 重发也不会产生重复的统计和重复的告警。可靠性这层不光是框架的事还得靠应用自己的幂等设计兜底。5. 与Kafka联动实时订单告警管线的完整实战5.1 接入KafkaSpout的配置细节先引入 storm-kafka-client 依赖dependency groupIdorg.apache.storm/groupId artifactIdstorm-kafka-client/artifactId version1.2.3/version scopeprovided/scope /dependency然后创建一个 KafkaSpoutKafkaSpoutConfig.BuilderString, String kafkaBuilder KafkaSpoutConfig.builder(kafka1:9092,kafka2:9092, order-topic); kafkaBuilder.setGroupId(storm-order-group); kafkaBuilder.setFirstPollOffsetStrategy( KafkaSpoutConfig.FirstPollOffsetStrategy.UNCOMMITTED_EARLIEST ); kafkaBuilder.setProcessingGuarantee( KafkaSpoutConfig.ProcessingGuarantee.AT_LEAST_ONCE ); KafkaSpoutConfigString, String kafkaConfig kafkaBuilder.build(); KafkaSpoutString, String kafkaSpout new KafkaSpout(kafkaConfig); builder.setSpout(kafka-spout, kafkaSpout, 2);FirstPollOffsetStrategy这里要解释一下很多第一次接入的人在这里吃过亏。它控制的是“第一次从无提交状态的地方开始消费时从哪里读”EARLIEST从头开始读历史数据适合需要回放或算历史指标的场景。LATEST只读最新数据适合只要实时新增的场景。UNCOMMITTED_EARLIEST优先读未提交的偏移如果没有就把最早的偏移给它。UNCOMMITTED_LATEST同理优先未提交偏移其次是最新偏移。我的经验是告警系统这种场景选UNCOMMITTED_EARLIEST最稳既能消费到还没提交但可能已经处理过的数据保证至少一次又不至于把整个历史全拉一遍。如果你选了LATEST而中间有数据还没被提交重启之后这部分数据会直接跳过告警就漏了。5.2 订单统计与窗口的实现很多人在学 Storm 时都有一个疑问它不是有窗口功能吗这里需要说清楚Storm 原生 API 里并没有像 Flink 那样的内置窗口抽象。如果你用的是基础 API非 Trident窗口逻辑基本都要自己手写。好在按分钟、按小时这类基础窗口并不复杂。我在订单告警拓扑里的统计 Bolt 是这样的逻辑输入Kafka 里解析出来的订单事件字段包括 orderId、userId、amount、status、eventTime。处理按 userId 做 Fields Grouping同一个用户的事件落到同一个 Bolt 实例。窗口自己维护一个ConcurrentHashMapLong, CountAndAmountkey 对应当前分钟时间戳value 存订单量和金额合计。每次execute先判断事件所属分钟和当前维护的分钟是否一致不一致就把上一分钟的统计输出并清理。这个大致的伪代码写法如下Override public void execute(Tuple tuple) { long eventTs tuple.getLongByField(eventTs); long minute eventTs / 60000L; if (minute ! currentMinute) { emitAndCleanup(currentMinute); currentMinute minute; } MapString, Object stat minuteStats.computeIfAbsent(minute, k - new HashMap()); stat.put(amount, (double) stat.getOrDefault(amount, 0.0) amount); stat.put(count, (int) stat.getOrDefault(count, 0) 1); collector.ack(tuple); }这里还要提醒一句窗口的清理逻辑不能只依赖事件时间。如果某个用户在某分钟内一直没新订单那这一分钟永远不会被推进上一分钟的统计就永远不输出。所以我额外加了一个基于当前系统时间的定时触发每 20 秒扫一次凡是落后于当前时间超过 90 秒的窗口直接输出并清理。这样才能保证那些低谷时段的数据也能被及时产出。5.3 告警怎么输出才能不打扰人告警拓扑的输出一开始我天真地直接 log 出来结果告警每 10 秒来一条收件人直接把我拉黑了。后来我做了一个收敛策略同一 userId 在 5 分钟内重复触发同一规则的告警自动合并成一条只更新触发次数。告警内容里带上实时统计数据和当前系统时间方便核对。写入 Redis 的时候用 SETNX 做去重key 设计成alarm:{ruleId}:{userId}:{minute}。这样做之后告警量下降了一个数量级而且每一条告警都还有业务参考价值。很多流处理项目会忽略“输出端的体验设计”但真实系统里“报警对不对、烦不烦人”往往比“处理得快不快”更影响口碑。6. 上集群前的调优与监控本地能跑不代表集群能扛6.1 Kryo序列化和并发处理的隐形坑Storm 的消息传递默认走 Kryo 序列化。如果你在 tuple 里传的是 Java 对象必须保证这个类实现了 Serializable 接口并且在提交拓扑之前把它注册到 Kryo 里否则很可能会在运行期遇到序列化异常。更稳妥的做法是在 tuple 里只传基本的 String、Long、Double 类型下游需要对象的时候自己把 JSON 字符串反序列化出来。这样既减少了 Kryo 序列化出问题的概率也让 tuple 的结构在调试时一眼能看清。另一个容易踩的坑在并发处理。一个 Bolt 的多个 Executor 是多个线程如果你的 Bolt 代码里持有了一个非线程安全的成员变量比如 SimpleDateFormat在并发调用format时会得到错误的结果甚至直接抛异常。我在一个日切日志解析拓扑里就遇见过时间字段偶尔差几个小时排查到最后发现是 SimpleDateFormat 的并发问题。后来改用 LocalDateTime 和线程安全的设计再没出过问题。6.2 UI指标教你判断瓶颈在哪Storm UI 是你的第一排查工具。我一般按这个顺序看指标Complete latency一条 tuple 从 Spout 发出到全部处理完成的耗时。如果这个值持续上涨链路里有明显慢节点。CapacityExecutor 的繁忙程度接近 1 说明这个处理线程几乎没闲着再往上就是瓶颈所在。Failed列如果出现持续性的 failed优先检查 Bolt 的 ack/fail 分支和超时设置。execute latency和process latency前者是 execute 方法的内部耗时后者包含排队等待时间。两者差距过大说明任务排队严重需要调大并行度或者引入背压。集群模式下动态调整拓扑并行度不需要重新提交代码直接在命令行执行storm rebalance order-topology -n 5 -e filter-bolt3这条命令把 worker 数调整为 5并单独把 filter-bolt 的并行度调整为 3。这个能力在压测调优阶段非常有用不用反复打包提交。另外还有一个参数值得单独说Config.setMaxSpoutPending(int)。它可以限制 Spout 未 ack 的最大 tuple 数相当于给整条流水线一个背压机制。当下游处理不过来时Spout 会主动放慢发射速度防止数据在内存里无限堆积。压测时我会从几百起步逐步调大这个值观察完整链路延迟和吞吐的变化找到当前集群条件下的甜点值。6.3 从本地压测到集群提交的操作习惯本地模式跑通和集群模式稳定运行是两码事。本地模式忽略了很多集群特有的因素网络延迟、多 Worker 间数据传输、Executor 调度开销。所以我自己定的操作流程是第一步本地用 LocalCluster 跑通全链路确认业务逻辑正确。第二步在测试集群用真实数据量做一次短时间压测重点观察 UI 里的 Complete latency 和 Failed。第三步根据压测结果调整并行度、超时时间和 maxSpoutPending再跑一轮。第四步确定参数后用storm jar提交到生产集群并通过 UI 观察启动后的前 15 分钟表现。提交命令本身很简单storm jar order-topology.jar com.example.OrderTopology prod提交流程里还有一个容易踩的坑拓扑名冲突。如果同名拓扑已经运行在集群上提交会失败或者覆盖老拓扑的行为让你措手不及。我的习惯是每次提交之前先storm kill topologyName清理旧拓扑再提交新版本同时拓扑名里带上版本号比如order-topology-v1.4回滚也方便。最后分享一个我自己的习惯每次调参都只改一个变量。并行度、超时时间、maxSpoutPending这三者相互影响如果一次动两个以上你根本没法判断指标变化到底是谁引起的。老老实实一次调一个记录下改动前后的 Complete latency、Capacity、吞吐量调优这件事就会变得有据可查而不是玄学。
返回列表