【发布时间】:2017-10-16 15:56:16
【问题描述】:
我使用 Spark 2.0.0 和 Kafka 0.10.2。
我有一个应用程序正在处理来自 Kafka 的消息,并且是一项长期运行的工作。
我有时会在日志中看到以下消息。我知道如何增加超时时间和一切,但我想知道的是我确实有这个错误,我该如何从中恢复?
错误 ConsumerCoordinator:偏移量提交失败。 org.apache.kafka.clients.consumer.CommitFailedException:
由于该组已经重新平衡并将分区分配给另一个成员,因此无法完成提交。
这意味着后续调用 poll() 之间的时间比配置的 session.timeout.ms 长,这通常意味着轮询循环花费了太多时间处理消息。
您可以通过增加会话超时或使用 max.poll.records 减少 poll() 中返回的批处理的最大大小来解决此问题。
这不是我如何逃避这个错误,而是一旦它发生后如何处理它
背景:在正常情况下,我不会看到提交错误,但如果我确实遇到了错误,我应该能够从中恢复。我正在使用AT_LEAST_ONCE 设置,所以我对重新处理一些消息非常满意。
我正在运行 Java 并使用带有手动提交的 DirectKakfaStreams。
创建流:
JavaInputDStream<ConsumerRecord<String, String>> directKafkaStream =
KafkaUtils.createDirectStream(
jssc,
LocationStrategies.PreferConsistent(),
ConsumerStrategies.<String, String>Subscribe(topics, kafkaParams));
提交偏移量
((CanCommitOffsets) directKafkaStream.inputDStream()).commitAsync(offsetRanges);
【问题讨论】:
-
Subscribe对应的kafkaParams是什么? -
enable.auto.commit 设置为 false,rest 没什么特别的,服务器详情,groupid 等
标签: java apache-kafka spark-streaming apache-spark-2.0