【问题标题】:How to handle offset commit failures with enable.auto.commit disabled in Spark Streaming with Kafka?如何在使用 Kafka 的 Spark Streaming 中禁用 enable.auto.commit 来处理偏移提交失败?
【发布时间】: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


【解决方案1】:

我对这种情况的理解是你使用了Kafka Direct Stream集成(使用Spark Streaming + Kafka Integration Guide (Kafka broker version 0.10.0 or higher)中描述的spark-streaming-kafka-0-10_2.11模块)。

如错误信息中所说:

提交无法完成,因为该组已经重新平衡并将分区分配给另一个成员。

Kafka 管理消费者使用的主题分区,因此 Direct Stream 将创建一个消费者池(在单个消费者组内)。

与任何消费者组一样,您应该期待重新平衡(引用 Kafka: The Definitive Guide 的第 4 章“Kafka 消费者 - 从 Kafka 读取数据”):

消费者组中的消费者共享他们订阅的主题中的分区的所有权。当我们向组中添加一个新的消费者时,它开始消费来自之前被另一个消费者消费的分区的消息。当消费者关闭或崩溃时,也会发生同样的事情,它会离开组,它曾经消费的分区将被剩余的消费者之一消费。当消费者组正在消费的主题被修改时,例如如果管理员添加了新分区,也会将分区重新分配给消费者。

在相当多的情况下,可能会发生重新平衡并且应该是可以预料的。而你做到了。

你问:

我该如何恢复?这不是关于我如何逃避这个错误,而是一旦发生如何处理它?

我的回答是使用CanCommitOffsets的另一种方法:

def commitAsync(offsetRanges: Array[OffsetRange], callback: OffsetCommitCallback): Unit

这使您可以访问 Kafka 的 OffsetCommitCallback:

OffsetCommitCallback 是一个回调接口,用户可以实现它以在提交请求完成时触发自定义操作。回调可以在任何调用 poll() 的线程中执行。

我认为onComplete 可以让您了解异步提交如何完成并采取相应措施。

我无法帮助您的是如何在无法提交某些偏移量时恢复 Spark Streaming 应用程序中的更改。我认为这需要跟踪偏移量并接受某些偏移量无法提交和重新处理的情况。

【讨论】:

  • 这也是我尝试过的,问题是我可以识别出错误,我可以捕捉到错误,但是一旦我掌握了错误,我该怎么做才能恢复,因为我希望仍然能够提交
  • 鉴于主题的分区不再属于您,即提交的代码,我认为您不允许提交。由于尚未提交偏移量,因此其他一些消费者会处理它。我的理解是,您必须以某种方式回滚由于(未提交的)偏移而做出的更改。
  • 那么我可以停止并重新开始我的直播吗?无需停止流上下文?
  • 你会如何“停止和开始”你的直播?您可以停止流式传输上下文,但这需要一个检查点来跟踪您的偏移量。
  • 我玩过它,似乎 commitAsync 不是最好的前进方式,计划在 kafka 之外手动保存偏移量......
猜你喜欢
  • 2017-02-06
  • 2021-05-22
  • 2018-09-22
  • 2020-09-03
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-06-08
  • 2017-06-22
相关资源
最近更新 更多