【问题标题】:kafka-console-consumer running in docker container stops consuming messages在 docker 容器中运行的 kafka-console-consumer 停止消费消息
【发布时间】:2018-06-11 05:36:37
【问题描述】:

我有融合的 kafka (v4.0.0) 代理和 3 个在 docker-compose 中运行的动物园管理员。创建了一个包含 10 个分区且复制因子为 3 的测试主题。 当一个控制台消费者创建时没有传递 --group 选项(其中 group.id 将被自动分配),即使在代理被杀死并且代理重新上线后,它也可以持续消费消息。

但是,如果我使用 --group 选项('console-group')创建一个控制台消费者,则在终止 kafka 代理后会停止消息消费。

$ docker run --net=host confluentinc/cp-kafka:4.0.0 kafka-console-consumer --bootstrap-server localhost:19092,localhost:29092,localhost:39092 --topic starcom.status --from-beginning --group console-group
<< some messages consumed >>
<< broker got killed >> 
[2017-12-31 18:34:05,344] WARN [Consumer clientId=consumer-1, groupId=console-group] Connection to node -1 could not be established. Broker may not be available. (org.apache.kafka.clients.NetworkClient)
<< no message after this >> 

即使在代理重新上线后,消费者组也不会再消费任何消息。

奇怪的是,当我使用 kafka-consumer-groups 工具检查时,该消费者组没有滞后。换句话说,该消费者群体的消费者补偿正在增加。没有其他消费者使用 group.id 运行,所以出了点问题。

根据日志,该组似乎已经稳定。

kafka-2_1           | [2017-12-31 17:35:40,743] INFO [GroupCoordinator 2]: Preparing to rebalance group console-group with old generation 0 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 17:35:43,746] INFO [GroupCoordinator 2]: Stabilized group console-group generation 1 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 17:35:43,765] INFO [GroupCoordinator 2]: Assignment received from leader for group console-group for generation 1 (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 17:54:30,228] INFO [GroupCoordinator 2]: Preparing to rebalance group console-group with old generation 1 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 17:54:31,162] INFO [GroupCoordinator 2]: Stabilized group console-group generation 2 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 17:54:31,173] INFO [GroupCoordinator 2]: Assignment received from leader for group console-group for generation 2 (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 17:57:25,273] INFO [GroupCoordinator 2]: Preparing to rebalance group console-group with old generation 2 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 17:57:28,256] INFO [GroupCoordinator 2]: Stabilized group console-group generation 3 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 17:57:28,267] INFO [GroupCoordinator 2]: Assignment received from leader for group console-group for generation 3 (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 17:57:53,594] INFO [GroupCoordinator 2]: Preparing to rebalance group console-group with old generation 3 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 17:57:55,322] INFO [GroupCoordinator 2]: Stabilized group console-group generation 4 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 17:57:55,336] INFO [GroupCoordinator 2]: Assignment received from leader for group console-group for generation 4 (kafka.coordinator.group.GroupCoordinator)
kafka-3_1           | [2017-12-31 18:15:07,953] INFO [GroupCoordinator 3]: Preparing to rebalance group console-group-2 with old generation 0 (__consumer_offsets-22) (kafka.coordinator.group.GroupCoordinator)
kafka-3_1           | [2017-12-31 18:15:10,987] INFO [GroupCoordinator 3]: Stabilized group console-group-2 generation 1 (__consumer_offsets-22) (kafka.coordinator.group.GroupCoordinator)
kafka-3_1           | [2017-12-31 18:15:11,044] INFO [GroupCoordinator 3]: Assignment received from leader for group console-group-2 for generation 1 (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 18:08:59,087] INFO [GroupCoordinator 2]: Loading group metadata for console-group with generation 4 (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 18:09:02,453] INFO [GroupCoordinator 2]: Preparing to rebalance group console-group with old generation 4 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 18:09:03,309] INFO [GroupCoordinator 2]: Stabilized group console-group generation 5 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 18:09:03,471] INFO [GroupCoordinator 2]: Assignment received from leader for group console-group for generation 5 (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 18:10:32,010] INFO [GroupCoordinator 2]: Preparing to rebalance group console-group with old generation 5 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 18:10:34,006] INFO [GroupCoordinator 2]: Stabilized group console-group generation 6 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 18:10:34,040] INFO [GroupCoordinator 2]: Assignment received from leader for group console-group for generation 6 (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 18:12:02,014] INFO [GroupCoordinator 2]: Preparing to rebalance group console-group with old generation 6 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 18:12:09,449] INFO [GroupCoordinator 2]: Stabilized group console-group generation 7 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 18:12:09,466] INFO [GroupCoordinator 2]: Assignment received from leader for group console-group for generation 7 (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 18:16:29,277] INFO [GroupCoordinator 2]: Preparing to rebalance group console-group with old generation 7 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 18:16:31,924] INFO [GroupCoordinator 2]: Stabilized group console-group generation 8 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 18:16:31,945] INFO [GroupCoordinator 2]: Assignment received from leader for group console-group for generation 8 (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 18:17:54,813] INFO [GroupCoordinator 2]: Preparing to rebalance group console-group with old generation 8 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 18:18:01,256] INFO [GroupCoordinator 2]: Stabilized group console-group generation 9 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 18:18:01,278] INFO [GroupCoordinator 2]: Assignment received from leader for group console-group for generation 9 (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 18:33:47,316] INFO [GroupCoordinator 2]: Preparing to rebalance group console-group with old generation 9 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 18:33:49,709] INFO [GroupCoordinator 2]: Stabilized group console-group generation 10 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 18:33:49,745] INFO [GroupCoordinator 2]: Assignment received from leader for group console-group for generation 10 (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 18:34:05,484] INFO [GroupCoordinator 2]: Preparing to rebalance group console-group with old generation 10 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 18:34:07,845] INFO [GroupCoordinator 2]: Stabilized group console-group generation 11 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 18:34:07,865] INFO [GroupCoordinator 2]: Assignment received from leader for group console-group for generation 11 (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 19:34:16,436] INFO [GroupCoordinator 2]: Preparing to rebalance group console-group with old generation 11 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 19:34:18,221] INFO [GroupCoordinator 2]: Stabilized group console-group generation 12 (__consumer_offsets-33) (kafka.coordinator.group.GroupCoordinator)
kafka-2_1           | [2017-12-31 19:34:18,248] INFO [GroupCoordinator 2]: Assignment received from leader for group console-group for generation 12 (kafka.coordinator.group.GroupCoordinator) 

主题复制一切正常。

$ docker run --net=host confluentinc/cp-kafka:4.0.0 kafka-topics --zookeeper localhost:22181 --topic starcom.status --describe
Topic:starcom.status    PartitionCount:10       ReplicationFactor:3     Configs:
        Topic: starcom.status   Partition: 0    Leader: 3       Replicas: 3,1,2 Isr: 2,3,1
        Topic: starcom.status   Partition: 1    Leader: 1       Replicas: 1,2,3 Isr: 3,2,1
        Topic: starcom.status   Partition: 2    Leader: 2       Replicas: 2,3,1 Isr: 3,2,1
        Topic: starcom.status   Partition: 3    Leader: 3       Replicas: 3,2,1 Isr: 3,2,1
        Topic: starcom.status   Partition: 4    Leader: 1       Replicas: 1,3,2 Isr: 3,2,1
        Topic: starcom.status   Partition: 5    Leader: 2       Replicas: 2,1,3 Isr: 3,2,1
        Topic: starcom.status   Partition: 6    Leader: 3       Replicas: 3,1,2 Isr: 2,3,1
        Topic: starcom.status   Partition: 7    Leader: 1       Replicas: 1,2,3 Isr: 3,2,1
        Topic: starcom.status   Partition: 8    Leader: 2       Replicas: 2,3,1 Isr: 3,2,1
        Topic: starcom.status   Partition: 9    Leader: 3       Replicas: 3,2,1 Isr: 3,2,1

$ docker run --net=host confluentinc/cp-kafka:4.0.0 kafka-topics --zookeeper localhost:22181 --topic starcom.status --describe
Topic:starcom.status    PartitionCount:10       ReplicationFactor:3     Configs:
        Topic: starcom.status   Partition: 0    Leader: 3       Replicas: 3,1,2 Isr: 2,3
        Topic: starcom.status   Partition: 1    Leader: 2       Replicas: 1,2,3 Isr: 3,2
        Topic: starcom.status   Partition: 2    Leader: 2       Replicas: 2,3,1 Isr: 3,2
        Topic: starcom.status   Partition: 3    Leader: 3       Replicas: 3,2,1 Isr: 3,2
        Topic: starcom.status   Partition: 4    Leader: 3       Replicas: 1,3,2 Isr: 3,2
        Topic: starcom.status   Partition: 5    Leader: 2       Replicas: 2,1,3 Isr: 3,2
        Topic: starcom.status   Partition: 6    Leader: 3       Replicas: 3,1,2 Isr: 2,3
        Topic: starcom.status   Partition: 7    Leader: 2       Replicas: 1,2,3 Isr: 3,2
        Topic: starcom.status   Partition: 8    Leader: 2       Replicas: 2,3,1 Isr: 3,2
        Topic: starcom.status   Partition: 9    Leader: 3       Replicas: 3,2,1 Isr: 3,2

$ docker run --net=host confluentinc/cp-kafka:4.0.0 kafka-topics --zookeeper localhost:22181 --topic starcom.status --describe
Topic:starcom.status    PartitionCount:10       ReplicationFactor:3     Configs:
        Topic: starcom.status   Partition: 0    Leader: 3       Replicas: 3,1,2 Isr: 2,3,1
        Topic: starcom.status   Partition: 1    Leader: 1       Replicas: 1,2,3 Isr: 3,2,1
        Topic: starcom.status   Partition: 2    Leader: 2       Replicas: 2,3,1 Isr: 3,2,1
        Topic: starcom.status   Partition: 3    Leader: 3       Replicas: 3,2,1 Isr: 3,2,1
        Topic: starcom.status   Partition: 4    Leader: 1       Replicas: 1,3,2 Isr: 3,2,1
        Topic: starcom.status   Partition: 5    Leader: 2       Replicas: 2,1,3 Isr: 3,2,1
        Topic: starcom.status   Partition: 6    Leader: 3       Replicas: 3,1,2 Isr: 2,3,1
        Topic: starcom.status   Partition: 7    Leader: 1       Replicas: 1,2,3 Isr: 3,2,1
        Topic: starcom.status   Partition: 8    Leader: 2       Replicas: 2,3,1 Isr: 3,2,1
        Topic: starcom.status   Partition: 9    Leader: 3       Replicas: 3,2,1 Isr: 3,2,1

这是(融合)kafka 控制台消费者的限制吗?基本上,我试图通过运行这个较小的测试来确保我真正的 Java Kafka 消费者能够在代理停机期间幸存下来。

任何帮助将不胜感激。

编辑(2018 年!)

我完全重新创建了我的 docker(-compose) 环境并且能够重现它。这次我创建了“新组”消费者组,并且控制台消费者在代理重新启动后抛出错误。从那时起,消息不再被消费。同样,根据消费者组工具,消费者抵消正在向前发展。

[2018-01-01 19:18:32,935] ERROR [Consumer clientId=consumer-1, groupId=new-group] Offset commit failed on partition starcom.status-4 at offset 0: This is not the correct coordinator. (org.apache.kafka.clients.consumer.internals.ConsumerCoordinator)
[2018-01-01 19:18:32,936] WARN [Consumer clientId=consumer-1, groupId=new-group] Asynchronous auto-commit of offsets {starcom.status-4=OffsetAndMetadata{offset=0, metadata=''}, starcom.status-5=OffsetAndMetadata{offset=0, metadata=''}, starcom.status-6=OffsetAndMetadata{offset=2, metadata=''}} failed: Offset commit failed with a retriable exception. You should retry committing offsets. The underlying error was: This is not the correct coordinator. (org.apache.kafka.clients.consumer.internals.ConsumerCoordinator)

【问题讨论】:

  • 控制台消费者只是Java API的一个包装器,但它无法重新建立与新分区领导者的连接?如果是这样,似乎是网络/配置问题。如果您只是杀死一个 docker 容器,它会丢失该代理的所有相关数据
  • 感谢您的评论。但是,正如我上面所说,如果我不通过 --group 选项,控制台消费者可以在不正常的代理中断中幸存 --- 这是一件非常好的事情。这让我相信这不是网络/配置问题。
  • 在上面提供了更多错误信息
  • 你是如何产生消息的?您确定它们已分发给其他经纪人吗?重新平衡一个主题真的需要 2 个小时吗?
  • 感谢您的后续问题。然而,我想我已经深究了。请看我的回答。

标签: docker apache-kafka docker-compose confluent-platform


【解决方案1】:

原来是 docker newbie 错误。

当我在 kafka-console-consumer shell 上 ctrl+c 时,容器(group.id:“console-group”)被置于分离模式。直到我运行了docker ps [-n | -a] 命令,我才知道这一点。当我使用相同的命令 (docker run --net=host confluentinc/cp-kafka:4.0.0 kafka-console-consumer --bootstrap-server localhost:19092,localhost:29092,localhost:39092 --topic starcom.status --from-beginning --group console-group) 启动另一个控制台消费者时,消费者加入了相同的“控制台组”。这就是为什么后续消息(显然我正在生成具有相同分区键的消息)被在后台运行的第一个消费者消费并给我错误的印象,即消息正在丢失。这就是为什么 consumer-groups 命令显示正确的偏移量提升。在不同窗口中将原始消费者重新附加到前台 (docker attach &lt;&lt;container-id&gt;&gt;) 后,现在我看到所有生成的消息都根据分区分配在两个不同的控制台中使用。一切都按预期工作。很抱歉误报,但希望遇到同样问题的人能从中得到一些提示。

【讨论】:

  • 你可能应该用过docker run --rm
  • 是的。 --rm 选项肯定会帮助清理停止容器。但是,我认为罪魁祸首更多的是没有提供 -i 和 -t 选项。
【解决方案2】:

总而言之,如果我想消费一些消息,那么在 docker 环境中设置 kafka-console-consumer 的正确方法应该是

docker run --net=host --rm -i -t \ 
    confluentinc/cp-kafka:4.0.0 \
      kafka-console-consumer --bootstrap-server localhost:19092,localhost:29092,localhost:39092 --topic foo.bar

注意有 --rm、-i、-t 选项。如果您传递了--max-messages,则不需要'-i 和-t',在这种情况下,控制台将正常退出,停止并拆除容器。

【讨论】:

    猜你喜欢
    • 2019-12-06
    • 2017-09-23
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-11-29
    • 2018-08-04
    相关资源
    最近更新 更多