Apache Storm Trident与Kafka集成时OpaqueTridentKafkaSpout启动报错求助
我在使用OpaqueTridentKafkaSpout消费Kafka消息时,使用了以下代码:
TridentKafkaConfig tridentKafkaConfig = new TridentKafkaConfig(hosts,properties.getProperty("topic", "mytopic")); tridentKafkaConfig.scheme = new SchemeAsMultiScheme(new StringScheme()); OpaqueTridentKafkaSpout kafkaSpout = new OpaqueTridentKafkaSpout(tridentKafkaConfig);
由于max spout pending配置会导致同一条Kafka消息出现在多个批次中,因此我未配置该参数。但Kafka Spout启动时会抛出一次如下错误,后续运行却正常:
2018-05-29 09:47:21.703 o.a.s.util Thread-9-spout-myspout-Spout-executor[33 33] [ERROR] Async loop died! java.lang.RuntimeException: java.lang.NullPointerException at org.apache.storm.utils.DisruptorQueue.consumeBatchToCursor(DisruptorQueue.java:522) ~[storm-core-1.2.1.jar:1.2.1] at org.apache.storm.utils.DisruptorQueue.consumeBatchWhenAvailable(DisruptorQueue.java:487) ~[storm-core-1.2.1.jar:1.2.1] at org.apache.storm.disruptor$consume_batch_when_available.invoke(disruptor.clj:74) ~[storm-core-1.2.1.jar:1.2.1] at org.apache.storm.daemon.executor$fn__5043$fn__5056$fn__5109.invoke(executor.clj:861) ~[storm-core-1.2.1.jar:1.2.1] at org.apache.storm.util$async_loop$fn__557.invoke(util.clj:484) [storm-core-1.2.1.jar:1.2.1] at clojure.lang.AFn.run(AFn.java:22) [clojure-1.7.0.jar:?] at java.lang.Thread.run(Thread.java:748) [?:1.8.0_171] Caused by: java.lang.NullPointerException at org.apache.storm.kafka.spout.trident.KafkaTridentSpoutEmitter.seek(KafkaTridentSpoutEmitter.java:193) ~[stormjar.jar:?] at org.apache.storm.kafka.spout.trident.KafkaTridentSpoutEmitter.emitPartitionBatch(KafkaTridentSpoutEmitter.java:127) ~[stormjar.jar:?] at org.apache.storm.kafka.spout.trident.KafkaTridentSpoutEmitter.emitPartitionBatch(KafkaTridentSpoutEmitter.java:51) ~[stormjar.jar:?] at org.apache.storm.trident.spout.OpaquePartitionedTridentSpoutExecutor$Emitter.emitBatch(OpaquePartitionedTridentSpoutExecutor.java:141) ~[storm-core-1.2.1.jar:1.2.1] at org.apache.storm.trident.spout.TridentSpoutExecutor.execute(TridentSpoutExecutor.java:82) ~[storm-core-1.2.1.jar:1.2.1] at org.apache.storm.trident.topology.TridentBoltExecutor.execute(TridentBoltExecutor.java:383) ~[storm-core-1.2.1.jar:1.2.1] at org.apache.storm.daemon.executor$fn__5043$tuple_action_fn__5045.invoke(executor.clj:739) ~[storm-core-1.2.1.jar:1.2.1] at org.apache.storm.daemon.executor$mk_task_receiver$fn__4964.invoke(executor.clj:468) ~[storm-core-1.2.1.jar:1.2.1] at org.apache.storm.disruptor$clojure_handler$reify__4475.onEvent(disruptor.clj:41) ~[storm-core-1.2.1.jar:1.2.1] at org.apache.storm.utils.DisruptorQueue.consumeBatchToCursor(DisruptorQueue.java:509) ~[storm-core-1.2.1.jar:1.2.1] ... 6 more
针对这个问题,我整理了几个可行的解决建议:
显式配置
maxSpoutPending参数:虽然你担心重复消息,但未配置该参数时,内部初始化逻辑可能触发空指针。尝试设置一个合理的值(比如100),触发正确的初始化流程:tridentKafkaConfig.maxSpoutPending = 100;实际运行中,只要你的后续处理逻辑是幂等的,少量重复消息不会影响业务。
显式设置偏移量策略:空指针可能是因为初始化时偏移量未正确赋值,你可以在配置中明确指定消费的起始偏移量:
// 从最新偏移量开始消费 tridentKafkaConfig.startOffsetTime = kafka.api.OffsetRequest.LatestTime(); // 或者从最早偏移量开始 // tridentKafkaConfig.startOffsetTime = kafka.api.OffsetRequest.EarliestTime();检查Kafka元数据加载时机:这个错误也可能是Spout启动时Kafka集群的元数据还未完全加载导致的。你可以在提交拓扑前,先通过Kafka客户端手动获取一次目标topic的元数据,确保连接正常后再启动拓扑。
升级Storm Kafka组件版本:你使用的Storm 1.2.1对应的Kafka Spout可能存在初始化阶段的已知bug,尝试升级到该Storm版本兼容的最新storm-kafka组件版本,大概率能修复这个问题。
内容的提问来源于stack exchange,提问作者phaigeim

