【问题标题】:Flink not commiting offsets to kafkaFlink 没有向 kafka 提交偏移量
【发布时间】:2017-06-24 01:38:58
【问题描述】:

我有一个 flink 流作业,它正在从 kafka 读取数据并记录它。我已启用检查点。

我在 kafka 中看不到已提交的偏移量,而是出现了错误。

非常感谢任何帮助。

{$KAFKA_HOME/bin/kafka-consumer-groups.sh --new-consumer --bootstrap-server localhost:9092 --describe --group flink-consumer-group
Error while executing consumer group command Group flink-consumer-group with protocol type '' is not a valid consumer group
java.lang.IllegalArgumentException: Group flink-consumer-group with protocol type '' is not a valid consumer group
at kafka.admin.AdminClient.describeConsumerGroup(AdminClient.scala:152)
at kafka.admin.ConsumerGroupCommand$KafkaConsumerGroupService.describeGroup(ConsumerGroupCommand.scala:308)
at kafka.admin.ConsumerGroupCommand$ConsumerGroupService$class.describe(ConsumerGroupCommand.scala:89)
at kafka.admin.ConsumerGroupCommand$KafkaConsumerGroupService.describe(ConsumerGroupCommand.scala:296)
at kafka.admin.ConsumerGroupCommand$.main(ConsumerGroupCommand.scala:68)
at kafka.admin.ConsumerGroupCommand.main(ConsumerGroupCommand.scala)}

版本

kafka_2.11-0.10.1.0 (server with) flink-connector-kafka-0.10_2.11

【问题讨论】:

  • 我观察到同样的问题,想知道您是否能够解决它?

标签: java scala streaming apache-kafka apache-flink


【解决方案1】:

Flink 自己处理偏移量。提交给 kafka(或旧版本或设置中的 zookeeper)的偏移量或多或少只是为了您的信息或监控目的。

您的错误看起来像是您混淆了不同的 kafka 版本(代理版本与客户端版本)。也许你可以仔细检查一下。

【讨论】:

  • 感谢 TobiSh 的回复。我知道 flink 不使用提交的偏移量,但是在没有保存点的情况下重新启动应用程序时,它应该使用提交的偏移量,而在我的情况下不会发生这种情况。我正在使用 kafka_2.11-0.10.1.0(带有服务器)flink-connector-kafka-0.10_2.11
  • 嗨@mukh007 你有一些我们可以重现这种行为的源代码吗? flink 集群的配置和 kafka 属性也很有用。
  • 还有一件事。您是否在日志文件中找到类似的内容:分区 {} 没有初始偏移;消费者有位置 {},因此初始偏移量将设置为 {}
【解决方案2】:

所以我认为 flink 在检查点默认情况下确实将偏移量提交给 kafka,因为值 FlinkKafkaConsumer#setCommitOffsetsOnCheckpoints 默认为 true。

很遗憾,这些偏移量在 kafka-offset 检查器 cli 中是不可见的。

我们实现了一个 scala kafka 消费者,它使用相同的消费者组连接到 kafka,但不订阅主题以从 kafka 获取偏移量。

注意:从 Kafka 版本 0.9 开始,消费者 Flink Kafka 导出所有标准指标,请参阅 documentation

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-04-06
    • 2022-10-04
    • 2020-01-27
    • 2021-11-14
    • 2017-03-17
    • 2018-07-02
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多