【问题标题】:Cannot Restart Kafka Consumer Application, Failing due to OffsetOutOfRangeException无法重新启动 Kafka 消费者应用程序,由于 OffsetOutOfRangeException 而失败
【发布时间】:2019-07-19 08:01:17
【问题描述】:

目前,我的 Kafka Consumer 流应用程序正在手动将偏移量提交到 Kafka,并将 enable.auto.commit 设置为 false。 当我尝试重新启动它时应用程序失败并抛出以下异常:

org.apache.kafka.clients.consumer.OffsetOutOfRangeException: Offsets out of range with no configured reset policy for partitions:{partition-12=155555555}

假设上述错误是由于保留期导致消息不存在/分区被删除,我尝试了以下方法:

我禁用了手动提交并启用了自动提交(enable.auto.commit=trueauto.offset.reset=earliest) 仍然失败并出现同样的错误

org.apache.kafka.clients.consumer.OffsetOutOfRangeException: Offsets out of range with no configured reset policy for partitions:{partition-12=155555555}

请建议重新启动作业的方法,以便它可以成功读取存在消息/分区的正确偏移量

【问题讨论】:

标签: apache-kafka offset kafka-consumer-api hadoop-streaming


【解决方案1】:

您正在尝试从主题partition 的分区12 中读取偏移量155555555,但 - 很可能 - 由于您的保留策略,它可能已经被删除。

您可以使用 Kafka Streams Application Reset Tool 来重置 Kafka Streams 应用程序的内部状态,以便它可以从头开始重新处理其输入数据

$ bin/kafka-streams-application-reset.sh

Option (* = required)         Description
---------------------         -----------
* --application-id <id>       The Kafka Streams application ID (application.id)
--bootstrap-servers <urls>    Comma-separated list of broker urls with format: HOST1:PORT1,HOST2:PORT2
                                (default: localhost:9092)
--intermediate-topics <list>  Comma-separated list of intermediate user topics
--input-topics <list>         Comma-separated list of user input topics
--zookeeper <url>             Format: HOST:POST
                                (default: localhost:2181)

或使用新的消费者组 ID 启动您的消费者。

【讨论】:

  • 使用这些选项会起作用。但我不明白为什么不应用auto.offset.reset 策略。
  • @ValBonn 仅当组没有提交的偏移量时才应用此配置。这就是我提到使用新组 ID 的原因。
  • 谢谢@Giorgos,但这不是我理解文档的方式kafka.apache.org/documentationauto.offset.reset What to do when there is no initial offset in Kafka or if the current offset does not exist any more on the server (e.g. because that data has been deleted): 而@user1326784 似乎是第二种情况。
  • @ValBonn 我猜@user1326784 也应该尝试使用auto.offset.reset=largest 看看会发生什么 - 以确保重置策略在这种情况下无效。
【解决方案2】:

我遇到了同样的问题,我在我的应用程序中使用包 org.apache.spark.streaming.kafka010。一开始,我怀疑 auto.offset.reset 策略无效,但是当我阅读了描述对象KafkaUtils中的fixKafkaParams方法,我发现配置被覆盖了。我猜它为executor调整配置ConsumerConfig.AUTO_OFFSET_RESET_CONFIG的原因是为了保持驱动程序和执行程序获得的偏移量一致。

【讨论】:

    猜你喜欢
    • 2017-08-16
    • 2021-07-18
    • 1970-01-01
    • 1970-01-01
    • 2021-04-07
    • 2023-01-18
    • 2021-05-15
    • 2018-01-27
    • 1970-01-01
    相关资源
    最近更新 更多