【发布时间】: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