【问题标题】:storm.kafka.UpdateOffsetException - Issue with Opaque Trident Kafka Spoutstorm.kafka.UpdateOffsetException - 不透明 Trident Kafka Spout 的问题
【发布时间】:2016-03-13 21:24:10
【问题描述】:

我正在使用带有 OpaqueTridentKafkaSpout 的三叉戟拓扑。

我正在使用的 TridentKafkaConfig 的代码 sn-p :-

OpaqueTridentKafkaSpout kafkaSpout = null;
TridentKafkaConfig spoutConfig = new TridentKafkaConfig(new ZkHosts("xxx.x.x.9:2181,xxx.x.x.1:2181,xxx.x.x.2:2181"), "topic_name");
spoutConfig.scheme = new SchemeAsMultiScheme(new StringScheme());
spoutConfig.fetchSizeBytes = 147483600;
kafkaSpout = new OpaqueTridentKafkaSpout(spoutConfig);

我从一名工人那里得到了这个运行时异常:-

java.lang.RuntimeException:storm.kafka.UpdateOffsetException 在 backtype.storm.utils.DisruptorQueue.consumeBatchToCursor(DisruptorQueue.java:135) 在 backtype.storm.utils.DisruptorQueue.consumeBatchWhenAvailable(DisruptorQueue.java:106) 在 backtype.storm.disruptor$consume_batch_when_available.invoke(disruptor.clj:80) 在 backtype.storm.daemon.executor$fn_5694$fn5707$fn5758.invoke(executor.clj:819) 在 backtype.storm.util$async_loop$fn545.invoke(util.clj:479) 在 clojure.lang.AFn.run(AFn.java:22) 在 java.lang.Thread.run(Thread.java:745) 原因: storm.kafka.UpdateOffsetException 在 storm.kafka.KafkaUtils.fetchMessages(KafkaUtils.java:186) 在 storm.kafka.trident.TridentKafkaEmitter.fetchMessages(TridentKafkaEmitter.java:132) 在 storm.kafka.trident.TridentKafkaEmitter.doEmitNewPartitionBatch(TridentKafkaEmitter.java:113) 在 storm.kafka.trident.TridentKafkaEmitter.failFastEmitNewPartitionBatch(TridentKafkaEmitter.java:72) 在 storm.kafka.trident.TridentKafkaEmitter.emitNewPartitionBatch(TridentKafkaEmitter.java:79) 在 storm.kafka.trident.TridentKafkaEmitter.access$000(TridentKafkaEmitter.java:46) 在 storm.kafka.trident.TridentKafkaEmitter$1.emitPartitionBatch(TridentKafkaEmitter.java:204) 在 storm.kafka.trident.TridentKafkaEmitter$1.emitPartitionBatch(TridentKafkaEmitter.java:194) 在 storm.trident.spout.OpaquePartitionedTridentSpoutExecutor$Emitter.emitBatch(OpaquePartitionedTridentSpoutExecutor.java:127) 在 storm.trident.spout.TridentSpoutExecutor.execute(TridentSpoutExecutor.java:82) 在 storm.trident.topology.TridentBoltExecutor.execute(TridentBoltExecutor.java:370) 在 backtype.storm.daemon.executor$fn5694$tuple_action_fn5696.invoke(executor.clj:690) 在 backtype.storm.daemon.executor$mk_task_receiver$fn5615.invoke(executor.clj:436) 在 backtype.storm.disruptor$clojure_handler$reify_5189.onEvent(disruptor.clj:58) 在 backtype.storm.utils.DisruptorQueue.consumeBatchToCursor(DisruptorQueue.java:127) ... 6 更多

根据一些帖子,我尝试设置 spoutConfig :- spoutConfig.maxOffsetBehind = Long.MAX_VALUE; spoutConfig.startOffsetTime = kafka.api.OffsetRequest.EarliestTime(); 我的 Kafka 保留时间是默认值 - 128 小时,即 7 天,并且 kafka 生产者每秒向 Storm/Trident 拓扑发送 6800 条消息。我浏览了大部分帖子,但似乎没有一个能解决这个问题。处理此问题的最佳方法是什么?

【问题讨论】:

    标签: clojure apache-kafka apache-storm


    【解决方案1】:

    我仍然不知道是什么导致了这个问题。但基本上我们没有正确关闭storm、zookeeper和kafka。这导致风暴拓扑失败,我们不得不拆除整个集群并重新构建它。更新到 Storm 0.10.0 有助于解决其他一些问题。

    【讨论】:

      猜你喜欢
      • 2018-09-19
      • 2014-12-12
      • 2018-11-03
      • 2014-08-22
      • 2019-03-14
      • 2015-09-03
      • 2010-11-23
      • 1970-01-01
      • 2011-06-04
      相关资源
      最近更新 更多