Spring Boot Kafka消费者无法维持与3节点Kafka集群的连接
问题描述
通过docker-compose部署3节点Kafka集群,配置如下:
services: kafka1: image: confluentinc/cp-kafka:latest hostname: kafka1 container_name: kafka1 ports: - "9092:9092" environment: CLUSTER_ID: "7bes6xCzSMeUR-RRas3gMQ" KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT' KAFKA_ADVERTISED_LISTENERS: 'PLAINTEXT://kafka1:29092,PLAINTEXT_HOST://localhost:9092' KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3 KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 5000 KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2 KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3 KAFKA_JMX_PORT: 9997 KAFKA_JMX_OPTS: '-Dcom.sun.management.jmxremote -Dcom.sun.management.jmxremote.authenticate=false -Dcom.sun.management.jmxremote.ssl=false -Djava.rmi.server.hostname=kafka1 -Dcom.sun.management.jmxremote.rmi.port=9997' KAFKA_PROCESS_ROLES: 'broker,controller' KAFKA_NODE_ID: 1 KAFKA_CONTROLLER_QUORUM_VOTERS: '1@kafka1:29093,2@kafka2:29093,3@kafka3:29093' KAFKA_LISTENERS: 'PLAINTEXT://kafka1:29092,CONTROLLER://kafka1:29093,,PLAINTEXT_HOST://0.0.0.0:9092' KAFKA_INTER_BROKER_LISTENER_NAME: 'PLAINTEXT' KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER' KAFKA_LOG_DIRS: '/tmp/kraft-combined-logs' volumes: - kafka1-data:/var/lib/kafka/data - ./scripts/update_run.sh:/tmp/update_run.sh command: "bash -c 'if [ ! -f /tmp/update_run.sh ]; then echo \"ERROR: Did you forget the update_run.sh file that came with this docker-compose.yml file?\" && exit 1 ; else /tmp/update_run.sh && /etc/confluent/docker/run ; fi'" kafka2: image: confluentinc/cp-kafka:latest hostname: kafka2 container_name: kafka2 ports: - "9093:9092" environment: CLUSTER_ID: "7bes6xCzSMeUR-RRas3gMQ" KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT' KAFKA_ADVERTISED_LISTENERS: 'PLAINTEXT://kafka2:29092,PLAINTEXT_HOST://localhost:9093' KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3 KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 5000 KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2 KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3 KAFKA_JMX_PORT: 9997 KAFKA_JMX_OPTS: '-Dcom.sun.management.jmxremote -Dcom.sun.management.jmxremote.authenticate=false -Dcom.sun.management.jmxremote.ssl=false -Djava.rmi.server.hostname=kafka2 -Dcom.sun.management.jmxremote.rmi.port=9998' KAFKA_PROCESS_ROLES: 'broker,controller' KAFKA_NODE_ID: 2 KAFKA_CONTROLLER_QUORUM_VOTERS: '1@kafka1:29093,2@kafka2:29093,3@kafka3:29093' KAFKA_LISTENERS: 'PLAINTEXT://kafka2:29092,CONTROLLER://kafka2:29093,PLAINTEXT_HOST://0.0.0.0:9093' KAFKA_INTER_BROKER_LISTENER_NAME: 'PLAINTEXT' KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER' KAFKA_LOG_DIRS: '/tmp/kraft-combined-logs' volumes: - kafka2-data:/var/lib/kafka/data - ./scripts/update_run.sh:/tmp/update_run.sh command: "bash -c 'if [ ! -f /tmp/update_run.sh ]; then echo \"ERROR: Did you forget the update_run.sh file that came with this docker-compose.yml file?\" && exit 1 ; else /tmp/update_run.sh && /etc/confluent/docker/run ; fi'" kafka3: image: confluentinc/cp-kafka:latest hostname: kafka3 container_name: kafka3 ports: - "9094:9092" environment: CLUSTER_ID: "7bes6xCzSMeUR-RRas3gMQ" KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT' KAFKA_ADVERTISED_LISTENERS: 'PLAINTEXT://kafka3:29092,PLAINTEXT_HOST://localhost:9094' KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3 KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 5000 KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2 KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3 KAFKA_JMX_PORT: 9997 KAFKA_JMX_OPTS: '-Dcom.sun.management.jmxremote -Dcom.sun.management.jmxremote.authenticate=false -Dcom.sun.management.jmxremote.ssl=false -Djava.rmi.server.hostname=kafka3 -Dcom.sun.management.jmxremote.rmi.port=9999' KAFKA_PROCESS_ROLES: 'broker,controller' KAFKA_NODE_ID: 3 KAFKA_CONTROLLER_QUORUM_VOTERS: '1@kafka1:29093,2@kafka2:29093,3@kafka3:29093' KAFKA_LISTENERS: 'PLAINTEXT://kafka3:29092,CONTROLLER://kafka3:29093,PLAINTEXT_HOST://0.0.0.0:9094' KAFKA_INTER_BROKER_LISTENER_NAME: 'PLAINTEXT' KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER' KAFKA_LOG_DIRS: '/tmp/kraft-combined-logs' volumes: - kafka3-data:/var/lib/kafka/data - ./scripts/update_run.sh:/tmp/update_run.sh command: "bash -c 'if [ ! -f /tmp/update_run.sh ]; then echo \"ERROR: Did you forget the update_run.sh file that came with this docker-compose.yml file?\" && exit 1 ; else /tmp/update_run.sh && /etc/confluent/docker/run ; fi'"
Spring Boot应用无法维持集群连接,日志如下:
2024-03-29T20:00:23.440-06:00 [Consumer clientId=consumer-transaction-group-2, groupId=transaction-group] Discovered group coordinator localhost:9093 (id: 2147483645 rack: null) 2024-03-29T20:00:23.440-06:00 [Consumer clientId=consumer-transaction-group-1, groupId=transaction-group] Discovered group coordinator localhost:9093 (id: 2147483645 rack: null) 2024-03-29T20:00:23.440-06:00 [Consumer clientId=consumer-transaction-group-2, groupId=transaction-group] Group coordinator localhost:9093 (id: 2147483645 rack: null) is unavailable or invalid due to cause: coordinator unavailable. isDisconnected: false. Rediscovery will be attempted. 2024-03-29T20:00:23.440-06:00 [Consumer clientId=consumer-transaction-group-1, groupId=transaction-group] Group coordinator localhost:9093 (id: 2147483645 rack: null) is unavailable or invalid due to cause: coordinator unavailable. isDisconnected: false. Rediscovery will be attempted. 2024-03-29T20:00:23.440-06:00 [Consumer clientId=consumer-transaction-group-2, groupId=transaction-group] Requesting disconnect from last known coordinator localhost:9093 (id: 2147483645 rack: null) 2024-03-29T20:00:23.440-06:00 [Consumer clientId=consumer-transaction-group-1, groupId=transaction-group] Requesting disconnect from last known coordinator localhost:9093 (id: 2147483645 rack: null) 2024-03-29T20:00:23.547-06:00 [Consumer clientId=consumer-transaction-group-2, groupId=transaction-group] Discovered group coordinator localhost:9093 (id: 2147483645 rack: null)
应用Kafka配置:
kafka: bootstrap-servers: localhost:9092,localhost:9093,localhost:9094 listener: ack-mode: manual properties: schema.registry.url: http://localhost:8185/ consumer: group-id: transaction-group enable-auto-commit: false auto-offset-reset: earliest key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: io.confluent.kafka.serializers.protobuf.KafkaProtobufDeserializer producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: io.confluent.kafka.serializers.protobuf.KafkaProtobufSerializer transaction-id-prefix: tx- properties: enable-idempotence: true
单节点集群时问题消失,请问是什么导致连接中断?
原因分析与解决方案
1. Kafka1节点监听器配置存在语法错误
查看docker-compose中kafka1的KAFKA_LISTENERS配置:
KAFKA_LISTENERS: 'PLAINTEXT://kafka1:29092,CONTROLLER://kafka1:29093,,PLAINTEXT_HOST://0.0.0.0:9092'
这里出现连续两个逗号(,,),属于无效配置格式,会导致Kafka节点启动时无法正确解析监听器列表,进而破坏集群内节点通信和对外服务能力。
修复方式:删除多余的逗号,修改为:
KAFKA_LISTENERS: 'PLAINTEXT://kafka1:29092,CONTROLLER://kafka1:29093,PLAINTEXT_HOST://0.0.0.0:9092'
2. 事务日志ISR配置与集群启动状态不匹配
当前配置中KAFKA_TRANSACTION_STATE_LOG_MIN_ISR=2,事务日志副本数KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR=3。如果集群启动过程中某个节点未完全就绪,事务日志的同步副本数(ISR)可能达不到2,导致协调节点无法正常对外提供服务,引发消费者连接失败。
优化建议:
- 确保集群所有节点完全启动、同步完成后,再启动Spring Boot应用。
- 测试环境可临时将
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR调整为1,降低ISR要求(生产环境不建议此操作)。
3. 验证集群节点状态
修复配置后,使用以下命令检查集群事务日志的副本和ISR状态:
docker exec -it kafka1 kafka-topics --bootstrap-server localhost:9092 --describe --topic __transaction_state
确认所有副本处于同步状态,集群节点均正常加入。
内容的提问来源于stack exchange,提问作者phoenix
相关产品推荐
相关产品推荐

