Golang Kafka生产者断连后分区数自动变更的原因及解决方法
我有一个Golang Kafka生产者,曾临时失去与Broker的连接。每个topic都配置了10个分区,断连后,生产者进程日志显示如下:
%3|1690891270.524|FAIL|rdkafka#producer-1| [thrd:my_ip:9092/bootstrap]: my_ip:9092/bootstrap: Connect to ipv4#my_ip:9092 failed: Connection refused (after 0ms in state CONNECT, 10 identical error(s) suppressed) %3|1690891273.524|FAIL|rdkafka#producer-1| [thrd:my_ip:9093/bootstrap]: my_ip:9093/bootstrap: Connect to ipv4#my_ip:9093 failed: Connection refused (after 0ms in state CONNECT, 7 identical error(s) suppressed) %3|1690891273.584|FAIL|rdkafka#producer-1| [thrd:hostname:9094/1001]: hostname:9094/1001: Connect to ipv4#my_ip:9094 failed: Connection refused (after 0ms in state CONNECT, 4 identical error(s) suppressed) %4|1690891277.639|CLUSTERID|rdkafka#producer-1| [thrd:main]: Broker my_ip:9092/bootstrap reports different ClusterId "TkRgxModQH-mlCkT9mr3lQ" than previously known "o4j9GSQrQ5Svuwvq4uiWQQ": a client must not be simultaneously connected to multiple clusters %5|1690891278.529|PARTCNT|rdkafka#producer-1| [thrd:main]: Topic my_topic partition count changed from 10 to 1
我的生产者配置非常基础:
config := kafka.ConfigMap{ "security.protocol": "plaintext", "bootstrap.servers": my_ip, } producer, err := kafka.NewProducer(&config)
我的Kafka Docker Compose配置如下:
zoo1: image: "${IMAGE_ZOOKEEPER}" hostname: zoo1 container_name: zoo1 ports: - "127.0.0.1:2181:2181" environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_SERVER_ID: 1 ZOOKEEPER_SERVERS: zoo1:2888:3888;zoo2:2888:3888;zoo3:2888:3888 networks: - internal_network zoo2: image: "${IMAGE_ZOOKEEPER}" hostname: zoo2 container_name: zoo2 ports: - "127.0.0.1:2182:2182" environment: ZOOKEEPER_CLIENT_PORT: 2182 ZOOKEEPER_SERVER_ID: 2 ZOOKEEPER_SERVERS: zoo1:2888:3888;zoo2:2888:3888;zoo3:2888:3888 networks: - internal_network zoo3: image: "${IMAGE_ZOOKEEPER}" hostname: zoo3 container_name: zoo3 ports: - "127.0.0.1:2183:2183" environment: ZOOKEEPER_CLIENT_PORT: 2183 ZOOKEEPER_SERVER_ID: 3 ZOOKEEPER_SERVERS: zoo1:2888:3888;zoo2:2888:3888;zoo3:2888:3888 networks: - internal_network kafka1: image: "${IMAGE_KAFKA}" hostname: kafka1 container_name: kafka1 ports: - "9092:9092" - "29092:29092" environment: KAFKA_LISTENERS: INTERNAL://0.0.0.0:29092,EXTERNAL://0.0.0.0:9092 KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka1:29092,EXTERNAL://niro:9092 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL KAFKA_ZOOKEEPER_CONNECT: "zoo1:2181,zoo2:2182,zoo3:2183" depends_on: - zoo1 - zoo2 - zoo3 networks: - internal_network kafka2: image: "${IMAGE_KAFKA}" hostname: kafka2 container_name: kafka2 ports: - "9093:9093" - "29093:29093" environment: KAFKA_LISTENERS: INTERNAL://0.0.0.0:29093,EXTERNAL://0.0.0.0:9093 KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka2:29093,EXTERNAL://niro:9093 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL KAFKA_ZOOKEEPER_CONNECT: "zoo1:2181,zoo2:2182,zoo3:2183" depends_on: - zoo1 - zoo2 - zoo3 networks: - internal_network kafka3: image: "${IMAGE_KAFKA}" hostname: kafka3 container_name: kafka3 ports: - "9094:9094" - "29094:29094" environment: KAFKA_LISTENERS: INTERNAL://0.0.0.0:29094,EXTERNAL://0.0.0.0:9094 KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka3:29094,EXTERNAL://niro:9094 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL KAFKA_ZOOKEEPER_CONNECT: "zoo1:2181,zoo2:2182,zoo3:2183" depends_on: - zoo1 - zoo2 - zoo3 networks: - internal_network kafka-topic-create: image: "${IMAGE_KAFKA}" hostname: kafka-topic-create container_name: kafka-topic-create depends_on: - kafka1 - kafka2 - kafka3 entrypoint: [ '/bin/sh', '-c' ] command: | " # blocks until kafka is reachable kafka-topics --bootstrap-server kafka1:29092 --list echo -e 'Creating kafka topics' kafka-topics --bootstrap-server kafka1:29092 --create --if-not-exists --topic my_topic --replication-factor 1 --partitions 10 echo -e 'Successfully created the following topics for Kafka1:' kafka-topics --bootstrap-server kafka1:29092 --list echo -e 'Kafka2:' kafka-topics --bootstrap-server kafka2:29093 --list echo -e 'Kafka3:' kafka-topics --bootstrap-server kafka3:29094 --list " environment: KAFKA_BROKER_ID: ignored KAFKA_ZOOKEEPER_CONNECT: ignored networks: - internal_network
使用的镜像为confluentinc/cp-kafka:7.3.2和confluentinc/cp-zookeeper:7.3.2
请问Kafka为何会自动修改分区数?如何阻止这种情况发生?
分区数被修改的原因
从日志里的ClusterId不一致提示就能定位核心问题:你的生产者实际上连接到了两个不同的Kafka集群,具体场景如下:
- 原有集群不可用时,生产者意外连接到了一个新的Kafka实例(可能是Docker重启后ZooKeeper数据丢失导致集群重建,或者本地存在其他测试集群)。
- 这个新集群里没有预先创建
my_topic,而librdkafka默认开启了auto.create.topics.enable=true,生产者尝试发消息时触发了自动创建topic逻辑,而Kafka默认创建的topic分区数为1。 - 当生产者重新连接回原有集群时,就会出现ClusterId冲突的日志,同时感知到topic分区数从原有集群的10变成了新集群的1。
另外你的Docker Compose里,Kafka的EXTERNAL监听器使用了niro主机名,如果本地DNS解析异常,也可能导致生产者连接到错误的节点,进而触发上述问题。
阻止分区数异常的解决方案
1. 关闭自动创建topic功能
在生产者配置里明确禁用自动创建:
config := kafka.ConfigMap{ "security.protocol": "plaintext", "bootstrap.servers": my_ip, "auto.create.topics.enable": false, // 禁用自动创建topic }
同时在每个Kafka Broker的Docker环境变量中也关闭该功能:
KAFKA_AUTO_CREATE_TOPICS_ENABLE: false
这样当生产者尝试往不存在的topic发消息时会直接报错,不会自动创建不符合预期的topic。
2. 强制生产者只连接目标集群
- 检查
bootstrap.servers配置,确保只包含目标集群的Broker地址,不要混入其他节点。 - 验证
niro主机名在生产者所在机器能正确解析到Docker宿主机IP,避免连接错误地址。 - 可以在生产者配置里添加
cluster.id参数,强制匹配目标集群的ClusterId:
"cluster.id": "o4j9GSQrQ5Svuwvq4uiWQQ", // 替换为你的原有集群ClusterId
这样如果连接到ClusterId不匹配的集群,生产者会直接拒绝连接,避免出现分区数异常。
3. 保证集群数据持久化
你的Docker Compose没有配置数据卷,容器重启后ZooKeeper和Kafka的数据会丢失,导致集群重建生成新的ClusterId。给每个ZooKeeper和Kafka服务添加数据卷:
# 在每个zoo服务下添加 volumes: - zoo1-data:/var/lib/zookeeper/data # 在每个kafka服务下添加 volumes: - kafka1-data:/var/lib/kafka/data # 在Compose文件末尾定义全局卷 volumes: zoo1-data: zoo2-data: zoo3-data: kafka1-data: kafka2-data: kafka3-data:
这样容器重启后数据不会丢失,集群ClusterId保持一致,不会出现重建情况。
4. 恢复与验证分区数
集群恢复后,先检查topic当前状态:
kafka-topics --bootstrap-server <你的BrokerIP>:9092 --describe --topic my_topic
如果分区数确实是1,可以手动扩容回10个:
kafka-topics --bootstrap-server <你的BrokerIP>:9092 --alter --topic my_topic --partitions 10
内容的提问来源于stack exchange,提问作者Omri. B

