KRaft模式Kafka集群添加副本时停滞问题求助
Kafka KRaft集群副本同步故障求助
集群环境
- 3节点KRaft模式集群,每个节点同时作为Broker和Controller
- Broker间通信采用SSL加密,生产者/消费者暂用明文
- 测试过Bitnami Kafka 3.3.2、3.4.1版本镜像
问题现象
集群启动后节点间可互相识别,但创建副本因子为3的主题时,副本同步陷入停滞:
创建主题命令
/opt/bitnami/kafka/bin/kafka-topics.sh --bootstrap-server=ofd-kafka-0:9092 --create --topic containers --partitions 3 --replication-factor 3
主题状态信息
Topic: containers TopicId: ZrOENsloQXuqTn2jt4NbNA PartitionCount: 3 ReplicationFactor: 1 Configs: segment.bytes=1073741824,retention.ms=2160000,max.message.bytes=15728640,retention.bytes=-1 Topic: containers Partition: 0 Leader: 1 Replicas: 0,1,2 Isr: 1 Adding Replicas: 0,2 Removing Replicas: Topic: containers Partition: 1 Leader: 2 Replicas: 0,1,2 Isr: 2 Adding Replicas: 0,1 Removing Replicas: Topic: containers Partition: 2 Leader: 0 Replicas: 2,1,0 Isr: 0 Adding Replicas: 1,2 Removing Replicas:
节点日志错误
所有节点日志均出现连接超时警告:
[2023-06-21 14:02:31,827] INFO [ReplicaFetcher replicaId=2, leaderId=1, fetcherId=0] Disconnecting from node 1 due to socket connection setup timeout. The timeout value is 8458 ms. (org.apache.kafka.clients.NetworkClient) [2023-06-21 14:02:31,827] INFO [ReplicaFetcher replicaId=2, leaderId=1, fetcherId=0] Client requested connection close from node 1 (org.apache.kafka.clients.NetworkClient) [2023-06-21 14:02:31,827] INFO [ReplicaFetcher replicaId=2, leaderId=1, fetcherId=0] Error sending fetch request (sessionId=INVALID, epoch=INITIAL) to node 1: (org.apache.kafka.clients.FetchSessionHandler) java.net.SocketTimeoutException: Failed to connect within 30000 ms at kafka.server.BrokerBlockingSender.sendRequest(BrokerBlockingSender.scala:109) at kafka.server.RemoteLeaderEndPoint.fetch(RemoteLeaderEndPoint.scala:78) at kafka.server.AbstractFetcherThread.processFetchRequest(AbstractFetcherThread.scala:309) at kafka.server.AbstractFetcherThread.$anonfun$maybeFetch$3(AbstractFetcherThread.scala:124) at kafka.server.AbstractFetcherThread.$anonfun$maybeFetch$3$adapted(AbstractFetcherThread.scala:123) at scala.Option.foreach(Option.scala:407) at kafka.server.AbstractFetcherThread.maybeFetch(AbstractFetcherThread.scala:123) at kafka.server.AbstractFetcherThread.doWork(AbstractFetcherThread.scala:106) at kafka.server.ReplicaFetcherThread.doWork(ReplicaFetcherThread.scala:97) at kafka.utils.ShutdownableThread.run(ShutdownableThread.scala:96) [2023-06-21 14:02:31,827] WARN [ReplicaFetcher replicaId=2, leaderId=1, fetcherId=0] Error in response for fetch request (type=FetchRequest, replicaId=2, maxWait=500, minBytes=1, maxBytes=10485760, fetchData={containers-1=PartitionData(topicId=wA7p8AUSRCmvg-Nc4qL9hA, fetchOffset=0, logStartOffset=0, maxBytes=1048576, currentLeaderEpoch=Optional[1], lastFetchedEpoch=Optional.empty)}, isolationLevel=READ_UNCOMMITTED, removed=, replaced=, metadata=(sessionId=INVALID, epoch=INITIAL), rackId=) (kafka.server.ReplicaFetcherThread) java.net.SocketTimeoutException: Failed to connect within 30000 ms at kafka.server.BrokerBlockingSender.sendRequest(BrokerBlockingSender.scala:109) at kafka.server.RemoteLeaderEndPoint.fetch(RemoteLeaderEndPoint.scala:78) at kafka.server.AbstractFetcherThread.processFetchRequest(AbstractFetcherThread.scala:309) at kafka.server.AbstractFetcherThread.$anonfun$maybeFetch$3(AbstractFetcherThread.scala:124) at kafka.server.AbstractFetcherThread.$anonfun$maybeFetch$3$adapted(AbstractFetcherThread.scala:123) at scala.Option.foreach(Option.scala:407) at kafka.server.AbstractFetcherThread.maybeFetch(AbstractFetcherThread.scala:123) at kafka.server.AbstractFetcherThread.doWork(AbstractFetcherThread.scala:106) at kafka.server.ReplicaFetcherThread.doWork(ReplicaFetcherThread.scala:97) at kafka.utils.ShutdownableThread.run(ShutdownableThread.scala:96)
StatefulSet核心配置
apiVersion: apps/v1 kind: StatefulSet ... spec: template: spec: containers: - args: - -ec - | export KAFKA_CFG_NODE_ID="$(echo "$MY_POD_NAME" | grep -o -E '[0-9[]*$')" /opt/bitnami/scripts/kafka/entrypoint.sh /opt/bitnami/scripts/kafka/run.sh command: - /bin/bash env: - name: BITNAMI_DEBUG value: "true" - name: KAFKA_KRAFT_CLUSTER_ID value: "Y2MyZmRlNDY3MjM0NGU4Yj" - name: KAFKA_CFG_PROCESS_ROLES value: controller,broker - name: KAFKA_CFG_CONTROLLER_QUORUM_VOTERS value: 0@kafka-inst-0:9093,1@kafka-inst-1:9093,2@kafka-inst-2:9093 - name: KAFKA_CFG_LISTENERS value: PLAINTEXT://:9092,CONTROLLER://:9093,INTERNAL://:9094 - name: KAFKA_CFG_CONTROLLER_LISTENER_NAMES value: CONTROLLER - name: KAFKA_HEAP_OPTS value: -Xmx1G -Xms1G - name: ALLOW_PLAINTEXT_LISTENER value: "yes" - name: KAFKA_CFG_MIN_INSYNC_REPLICAS value: "2" - name: KAFKA_CFG_DEFAULT_REPLICATION_FACTOR value: "3" - name: MY_POD_NAME valueFrom: fieldRef: apiVersion: v1 fieldPath: metadata.name - name: KAFKA_CFG_ADVERTISED_LISTENERS value: INTERNAL://$(MY_POD_NAME).test.svc.cluster.local:9094,PLAINTEXT://$(MY_POD_NAME).test.svc.cluster.local:9092 - name: KAFKA_INTER_BROKER_LISTENER_NAME value: INTERNAL - name: KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP value: INTERNAL:SSL,CONTROLLER:SSL,PLAINTEXT:PLAINTEXT - name: KAFKA_TLS_TYPE value: PEM - name: KAFKA_CFG_INITIAL_BROKER_REGISTRATION_TIMEOUT_MS value: "240000" image: bitnami/kafka:3.3.2-debian-11-r49 imagePullPolicy: IfNotPresent ...
补充信息
- SSL证书:通过Helm genCA生成自签名CA,Broker使用PEM格式证书,挂载配置如下:
volumeMounts: - name: cacert mountPath: /opt/bitnami/kafka/config/certs/kafka.truststore.pem subPath: tls.crt - name: storage mountPath: /opt/bitnami/kafka/config/certs/kafka.keystore.pem subPath: certs/kafka.keystore.pem - name: storage mountPath: /opt/bitnami/kafka/config/certs/kafka.keystore.key subPath: certs/kafka.keystore.key - 添加
-Djavax.net.debug=all参数后,日志显示SSL通信正常;创建单副本主题时,所有Broker均可在本地创建分区 - 集群用途:作为消息负载均衡器缓存突发消息,暂不考虑使用Kafka Operator
寻求排查线索或解决思路,感谢!
内容的提问来源于stack exchange,提问作者Adavan
相关产品推荐
相关产品推荐

