求助:KafkaConsumer陷入加入组循环的问题排查
Kafka消费者陷入“加入组”循环问题
问题详情
我有一个KafkaConsumer陷入了“加入组”循环。这是一个简单的Python Kafka消费者脚本:
from kafka import KafkaConsumer import logging import sys root = logging.getLogger('kafka') root.setLevel(logging.INFO) handler = logging.StreamHandler(sys.stdout) handler.setLevel(logging.INFO) formatter = logging.Formatter( '%(asctime)s - %(name)s - %(levelname)s - %(message)s') handler.setFormatter(formatter) root.addHandler(handler) bootstrap_servers_sasl = ['node1.dev.company.local:9092', 'node2.dev.company.local:9092', 'node3.dev.company.local:9092'] topicName = 'test_sasl' consumer = KafkaConsumer( topicName, bootstrap_servers = bootstrap_servers_sasl, security_protocol = 'SASL_PLAINTEXT', sasl_mechanism = 'SCRAM-SHA-512', sasl_plain_username = 'test_user', sasl_plain_password = 't3st_us3r', group_id = 'test_group' ) try: for message in consumer: if message: print(f"Received message: {message.value.decode('utf-8')}") except Exception as e: print(f"An exception occurred: {e}") finally: consumer.close()
当创建KafkaConsumer时指定group_id后,日志会反复输出以下内容,消费者始终无法消费目标主题的消息:
2023-09-13 08:35:44,102 - kafka.cluster - INFO - Group coordinator for test_group is BrokerMetadata(nodeId='coordinator-0', host='node1.dev.company.local', port=9092, rack=None) 2023-09-13 08:35:44,102 - kafka.coordinator - INFO - Discovered coordinator coordinator-0 for group test_group 2023-09-13 08:35:44,102 - kafka.coordinator - INFO - (Re-)joining group test_group 2023-09-13 08:35:44,104 - kafka.coordinator - WARNING - Marking the coordinator dead (node coordinator-0) for group test_group: [Error 16] NotCoordinatorForGroupError.
如果不指定group_id,一切工作正常。
Kafka Broker是Confluent Kafka,协议版本为3.4-IV0。
服务器端controller.log包含以下内容,不确定是否存在问题:
[2023-09-13 08:56:03,345] INFO [Controller id=0] Processing automatic preferred replica leader election (kafka.controller.KafkaController) [2023-09-13 08:56:03,345] TRACE [Controller id=0] Checking need to trigger auto leader balancing (kafka.controller.KafkaController) [2023-09-13 08:56:03,345] DEBUG [Controller id=0] Topics not in preferred replica for broker 0 HashMap() (kafka.controller.KafkaController) [2023-09-13 08:56:03,345] TRACE [Controller id=0] Leader imbalance ratio for broker 0 is 0.0 (kafka.controller.KafkaController) [2023-09-13 08:56:03,345] DEBUG [Controller id=0] Topics not in preferred replica for broker 1 HashMap() (kafka.controller.KafkaController) [2023-09-13 08:56:03,345] TRACE [Controller id=0] Leader imbalance ratio for broker 1 is 0.0 (kafka.controller.KafkaController) [2023-09-13 08:56:03,345] DEBUG [Controller id=0] Topics not in preferred replica for broker 2 HashMap() (kafka.controller.KafkaController) [2023-09-13 08:56:03,345] TRACE [Controller id=0] Leader imbalance ratio for broker 2 is 0.0 (kafka.controller.KafkaController)
补充配置与测试信息
Broker配置(node1的server.properties,集群共3节点)
broker.id=0 listeners=LISTENER_ONE://:9092,LISTENER_TWO://:9096 inter.broker.listener.name=LISTENER_ONE sasl.mechanism.inter.broker.protocol=PLAIN sasl.enabled.mechanisms=PLAIN,SCRAM-SHA-512 authorizer.class.name=kafka.security.authorizer.AclAuthorizer super.users=User:admin allow.everyone.if.no.acl.found=true security.protocol=SASL_PLAINTEXT advertised.listeners=LISTENER_ONE://node1.dev.company.local:9092,LISTENER_TWO://node1.dev.company.local:9096 listener.security.protocol.map=LISTENER_ONE:SASL_PLAINTEXT,LISTENER_TWO:PLAINTEXT num.network.threads=3 num.io.threads=8 socket.send.buffer.bytes=102400 socket.receive.buffer.bytes=102400 socket.request.max.bytes=104857600 log.dirs=/var/lib/kafka num.partitions=1 num.recovery.threads.per.data.dir=1 offsets.topic.replication.factor=3 transaction.state.log.replication.factor=3 transaction.state.log.min.isr=2 log.retention.hours=24 log.segment.bytes=1073741824 log.retention.check.interval.ms=300000 zookeeper.connect=node1.dev.company.local:2181,node2.dev.company.local:2181,node3.dev.company.local:2181 zookeeper.connection.timeout.ms=6000 group.initial.rebalance.delay.ms=3 auto.create.topics.enable=false inter.broker.protocol.version=3.4-IV0
控制台消费者测试
在Broker机器上运行以下命令:
kafka-console-consumer --bootstrap-server node1.dev.company.local:9092,node2.dev.company.local:9092,node3.dev.company.local:9092 --group test_group --topic test_sasl --consumer.config consumer.config
对应的consumer.config内容:
security.protocol=SASL_PLAINTEXT sasl.mechanism=SCRAM-SHA-512 sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required \ username="test_user" \ password="t3st_us3r";
测试无错误输出,但也不显示处理生产的测试消息;不使用--group参数运行时,同样无消息显示,这与Python客户端表现不同。
Topic信息查询
运行以下命令查询topic详情:
kafka-topics --bootstrap-server node1.dev.company.local:9092,node2.dev.company.local:9092,node3.dev.company.local:9092 --describe --topic test_sasl --command-config admin-plaintext.config
对应的admin-plaintext.config内容:
security.protocol=SASL_PLAINTEXT sasl.mechanism=PLAIN sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required \ username="admin-username" \ password="admin-password";
查询结果:
Topic: test_sasl TopicId: Y3hiju-ZQsOvdRgLk4vpZw PartitionCount: 1 ReplicationFactor: 2 Configs: cleanup.policy=delete,segment.bytes=1073741824,retention.ms=86400000,unclean.leader.election.enable=true Topic: test_sasl Partition: 0 Leader: 2 Replicas: 2,0 Isr: 0,2
内容的提问来源于stack exchange,提问作者jceddy
相关产品推荐
相关产品推荐

