首次部署集成Apache Kafka的Storm-Core拓扑及报错排查求助
咱们一步步来搞定这个部署,流程其实很清晰:
先搞定依赖配置
不管你用Maven还是Gradle,都得把Storm-Core和Storm-Kafka的依赖加进项目里。重点注意版本匹配!比如Storm 2.4.0就对应Storm-Kafka 2.4.0,Kafka选2.8.x左右的版本就稳。给你个Maven的示例:<dependencies> <dependency> <groupId>org.apache.storm</groupId> <artifactId>storm-core</artifactId> <version>2.4.0</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.storm</groupId> <artifactId>storm-kafka-client</artifactId> <version>2.4.0</version> </dependency> </dependencies>配置Kafka Spout
接下来创建KafkaSpoutConfig,把Kafka集群地址、消费组ID、目标Topic这些关键信息填进去,还要设置初始偏移量策略(比如从头消费用EARLIEST,取最新消息用LATEST)。简单的代码示例:import org.apache.storm.kafka.spout.KafkaSpout; import org.apache.storm.kafka.spout.KafkaSpoutConfig; public class KafkaStormTopology { public static void main(String[] args) throws Exception { KafkaSpoutConfig<String, String> spoutConfig = KafkaSpoutConfig.builder("kafka-broker1:9092,kafka-broker2:9092", "myfirsttopic") .setGroupId("storm-consumer-group") .setOffsetCommitPeriodMs(10000) .setFirstPollOffsetStrategy(KafkaSpoutConfig.FirstPollOffsetStrategy.EARLIEST) .build(); KafkaSpout<String, String> kafkaSpout = new KafkaSpout<>(spoutConfig); // 后续构建拓扑... } }组装并提交拓扑
用TopologyBuilder把Spout和处理数据的Bolt串起来,然后用StormSubmitter提交到集群。比如:TopologyBuilder builder = new TopologyBuilder(); builder.setSpout("kafka-spout", kafkaSpout, 2); // 设置并行度为2 builder.setBolt("data-processing-bolt", new DataProcessingBolt(), 4) .shuffleGrouping("kafka-spout"); // 和Spout做随机分组 Config config = new Config(); config.setDebug(false); if (args != null && args.length > 0) { // 提交到集群 StormSubmitter.submitTopology(args[0], config, builder.createTopology()); } else { // 本地调试用 LocalCluster cluster = new LocalCluster(); cluster.submitTopology("local-kafka-topology", config, builder.createTopology()); Thread.sleep(60000); cluster.shutdown(); }提交到Storm集群
打包成Jar包后,用Storm的命令行工具提交:storm jar your-topology-jar-with-dependencies.jar com.your.package.KafkaStormTopology my-kafka-storm-topology验证部署
打开Storm UI(默认地址是http://storm-nimbus:8080),看看拓扑状态:Spout的发射条数、Bolt的处理条数是不是在增长,再去Kafka那边用kafka-consumer-groups.sh检查消费偏移量有没有推进,确保数据正常流动。
先看你这个报错:java.lang.IllegalArgumentException: zk-node '/kafka-cluster-1/brokers/topics/myfirsttopic/aadb3eb4-2224-4c18-b8ad-6959a1c9f607' dose not exists,问题出在Storm的旧版偏移量监控工具在ZK里找不到对应节点,咱们一步步排查:
先搞清楚偏移量存在哪
如果你用的是新版KafkaSpout(基于Kafka Consumer API的那种),消费者偏移量是存在Kafka内部的__consumer_offsets主题里,根本不会存在ZK里!但Storm UI的旧监控逻辑还在去ZK找,自然会报错。切换Storm的偏移量监控方式
去Storm集群配置文件storm.yaml里加这两行,让Storm用Kafka读取偏移量,而非ZK:storm.kafka.monitor.offset.storage: kafka storm.kafka.monitor.kafka.bootstrap.servers: "kafka-broker1:9092,kafka-broker2:9092"改完后重启Storm UI服务,应该就能正常获取滞后值了。
如果是旧版ZK-based消费者的情况
要是你真的在用旧版Storm-Kafka(比如基于ZkHosts的那种),那得检查ZK路径:- 用
zkCli.sh连接ZK,执行ls /kafka-cluster-1/brokers/topics/myfirsttopic,看看实际的分区节点是什么(正常分区是数字,比如0、1,你报错里的UUID大概率是路径配置错了)。 - 检查Storm拓扑里的Kafka配置,是不是把ZK根路径或Topic名称写错了?比如是不是把消费者组路径和Topic路径搞混了?
- 确认Kafka的
myfirsttopic是否存在、分区状态正常。用Kafka命令行工具查一下:kafka-topics.sh --describe --topic myfirsttopic --bootstrap-server kafka-broker:9092
- 用
版本兼容性要注意
最后再确认下Storm和Kafka的版本是否匹配,比如Storm 1.x对应Kafka 0.10.x,Storm 2.x对应Kafka 2.x,版本不兼容也会出现各种奇怪问题。
内容的提问来源于stack exchange,提问作者Yu Watanabe

