Kraft模式下Kafka集群频繁出现NotControllerException问题求助
Kafka Kraft模式下NotControllerException频繁出现的解决方案
问题背景
我们在Kubernetes环境中通过Strimzi v0.45.0部署了Kafka 3.9.0集群,采用Kraft模式,包含3个控制器节点和3个Broker节点。集群每日会出现5-10次NotControllerException: The active controller appears to be node xx错误,所有控制器节点均会触发该错误。已尝试将Broker与控制器节点合并或分离部署,问题仍未解决。
相关日志
[ControllerApis nodeId=11] Unexpected error handling request Forwarded request: RequestContext(header=RequestHeader(apiKey=ENVELOPE, apiVersion=0, clientId=1, correlationId=127, headerVersion=2), connectionId='10.42.8.188:9090-10.42.4.228:36804-7', clientAddress=/10.42.4.228, principal=User:CN=kafka-multi-kafka,O=io.strimzi, listenerName=ListenerName(CONTROLPLANE-9090), securityProtocol=SSL, clientInformation=ClientInformation(softwareName=apache-kafka-java, softwareVersion=3.9.0), fromPrivilegedListener=false, principalSerde=Optional[org.apache.kafka.common.security.authenticator.DefaultKafkaPrincipalBuilder@60dd23f0]) RequestHeader(apiKey=ALTER_USER_SCRAM_CREDENTIALS, apiVersion=0, clientId=adminclient-1, correlationId=5925, headerVersion=2) -- AlterUserScramCredentialsRequestData(deletions=[], upsertions=[ScramCredentialUpsertion("This is removed because it contains password information") with context RequestContext(header=RequestHeader(apiKey=ALTER_USER_SCRAM_CREDENTIALS, apiVersion=0, clientId=adminclient-1, correlationId=5925, headerVersion=2), connectionId='10.42.8.188:9090-10.42.4.228:36804-7', clientAddress=/10.42.5.10, principal=User:CN=kafka-multi-entity-user-operator,O=io.strimzi, listenerName=ListenerName(CONTROLPLANE-9090), securityProtocol=SSL, clientInformation=ClientInformation(softwareName=unknown, softwareVersion=unknown), fromPrivilegedListener=false, principalSerde=Optional.empty) (kafka.server.ControllerApis) [quorum-controller-11-event-handler] java.util.concurrent.CompletionException: org.apache.kafka.common.errors.NotControllerException: The active controller appears to be node 12. at java.base/java.util.concurrent.CompletableFuture.encodeThrowable(CompletableFuture.java:332) at java.base/java.util.concurrent.CompletableFuture.completeThrowable(CompletableFuture.java:347) at java.base/java.util.concurrent.CompletableFuture$UniApply.tryFire(CompletableFuture.java:636) at java.base/java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:510) at java.base/java.util.concurrent.CompletableFuture.completeExceptionally(CompletableFuture.java:2162) at org.apache.kafka.controller.QuorumController$ControllerWriteEvent.complete(QuorumController.java:884) at org.apache.kafka.controller.QuorumController$ControllerWriteEvent.handleException(QuorumController.java:875) at org.apache.kafka.queue.KafkaEventQueue$EventContext.completeWithException(KafkaEventQueue.java:153) at org.apache.kafka.queue.KafkaEventQueue$EventContext.run(KafkaEventQueue.java:142) at org.apache.kafka.queue.KafkaEventQueue$EventHandler.handleEvents(KafkaEventQueue.java:215) at org.apache.kafka.queue.KafkaEventQueue$EventHandler.run(KafkaEventQueue.java:186) at java.base/java.lang.Thread.run(Thread.java:840) Caused by: org.apache.kafka.common.errors.NotControllerException: The active controller appears to be node 12.
集群CRD配置
apiVersion: kafka.strimzi.io/v1beta2 kind: Kafka metadata: annotations: strimzi.io/kraft: enabled strimzi.io/node-pools: enabled labels: argocd.argoproj.io/instance: development2-kafka name: kafka-multi namespace: development2 spec: cruiseControl: config: default.goals: > com.linkedin.kafka.cruisecontrol.analyzer.goals.MinTopicLeadersPerBrokerGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.ReplicaCapacityGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.DiskCapacityGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.NetworkInboundCapacityGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.NetworkOutboundCapacityGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.CpuCapacityGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.ReplicaDistributionGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.PotentialNwOutGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.DiskUsageDistributionGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.NetworkInboundUsageDistributionGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.NetworkOutboundUsageDistributionGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.CpuUsageDistributionGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.TopicReplicaDistributionGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.LeaderReplicaDistributionGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.LeaderBytesInDistributionGoal goals: > com.linkedin.kafka.cruisecontrol.analyzer.goals.MinTopicLeadersPerBrokerGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.ReplicaCapacityGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.DiskCapacityGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.NetworkInboundCapacityGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.NetworkOutboundCapacityGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.CpuCapacityGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.ReplicaDistributionGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.PotentialNwOutGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.DiskUsageDistributionGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.NetworkInboundUsageDistributionGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.NetworkOutboundUsageDistributionGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.CpuUsageDistributionGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.TopicReplicaDistributionGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.LeaderReplicaDistributionGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.LeaderBytesInDistributionGoal hard.goals: > com.linkedin.kafka.cruisecontrol.analyzer.goals.MinTopicLeadersPerBrokerGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.ReplicaCapacityGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.DiskCapacityGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.NetworkInboundCapacityGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.NetworkOutboundCapacityGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.CpuCapacityGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.ReplicaDistributionGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.PotentialNwOutGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.DiskUsageDistributionGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.NetworkInboundUsageDistributionGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.NetworkOutboundUsageDistributionGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.CpuUsageDistributionGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.TopicReplicaDistributionGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.LeaderReplicaDistributionGoal, com.linkedin.kafka.cruisecontrol.analyzer.goals.LeaderBytesInDistributionGoal image: 'sosconreg.azurecr.io/kafka:0.45.0-kafka-3.9.0' resources: limits: cpu: 1 memory: 1Gi requests: cpu: 200m memory: 512Mi entityOperator: topicOperator: resources: limits: cpu: 1 memory: 512Mi requests: cpu: 200m memory: 512Mi userOperator: resources: limits: cpu: 1 memory: 512Mi requests: cpu: 200m memory: 512Mi secretPrefix: kafka- kafka: authorization: superUsers: - super-user1-user type: simple config: auto.create.topics.enable: false default.replication.factor: 2 log.retention.minutes: 30240 min.insync.replicas: 1 num.partitions: 1 offsets.retention.minutes: 30240 offsets.topic.replication.factor: 2 transaction.state.log.min.isr: 1 transaction.state.log.replication.factor: 2 image: 'sosconreg.azurecr.io/kafka:0.45.0-kafka-3.9.0' listeners: - authentication: type: scram-sha-512 name: plain port: 9092 tls: false type: internal - authentication: type: scram-sha-512 name: tls port: 9093 tls: false type: internal livenessProbe: initialDelaySeconds: 15 timeoutSeconds: 5 logging: loggers: kafka.root.logger.level: DEBUG log4j.logger.kafka.authorizer.logger: INFO log4j.logger.kafka.controller: DEBUG log4j.logger.kafka.coordinator.transaction: DEBUG log4j.logger.kafka.request.logger: INFO log4j.rootLogger: 'DEBUG, CONSOLE' type: inline metricsConfig: type: jmxPrometheusExporter valueFrom: configMapKeyRef: key: kafka-metrics-config.yml name: kafka-metrics rack: topologyKey: topology.kubernetes.io/zone readinessProbe: initialDelaySeconds: 15 timeoutSeconds: 5 template: clusterRoleBinding: metadata: annotations: argocd.argoproj.io/compare-options: IgnoreExtraneous version: 3.9.0 kafkaExporter: groupRegex: .* image: 'sosconreg.azurecr.io/kafka:0.45.0-kafka-3.9.0' resources: limits: cpu: 1 memory: 512Mi requests: cpu: 200m memory: 512Mi topicRegex: .*
可能原因分析
- Kraft控制器选举时序延迟:控制器切换过程中,旧控制器节点未及时感知新控制器当选,导致转发到旧节点的请求失败。
- 请求转发竞态条件:Strimzi的请求转发逻辑在集群状态同步不及时时,错误地将请求发往非活跃控制器。
- 版本兼容bug:Strimzi v0.45.0与Kafka 3.9.0的组合存在已知的控制器状态同步问题。
- 网络稳定性问题:Kubernetes环境中控制器节点间的网络延迟或分区,导致状态同步不及时。
解决方案
1. 调整Kraft控制器同步配置
在Kafka CR的kafka.config中添加以下参数,缩短控制器集群的超时时间,加快状态同步:
kafka: config: controller.quorum.timeout.ms: 3000 controller.quorum.election.timeout.ms: 5000 controller.quorum.fetch.timeout.ms: 1000
2. 升级Strimzi与Kafka版本
升级至Strimzi v0.46.0+和Kafka 3.9.1+,这两个版本修复了多个Kraft控制器相关的竞态条件问题。
3. 优化Entity Operator请求重试逻辑
修改User Operator配置,增加请求重试次数与延迟,避免在控制器切换时立即发送请求:
entityOperator: userOperator: config: retry.backoff.ms: 200 retry.max.attempts: 5
4. 检查Kubernetes网络稳定性
验证控制器节点间的网络延迟,确保无持续网络分区或高延迟:
# 在控制器节点间执行ping测试 kubectl exec -it kafka-multi-controller-0 -- ping -c 10 kafka-multi-controller-1 kubectl exec -it kafka-multi-controller-0 -- ping -c 10 kafka-multi-controller-2
5. 提升控制器资源配置
确保控制器节点有足够资源处理状态同步,避免因资源不足导致延迟:
kafka: controllerNodePool: resources: requests: cpu: 500m memory: 1Gi limits: cpu: 2 memory: 2Gi
验证步骤
- 监控集群日志,统计
NotControllerException的出现频率是否下降。 - 查看控制器选举次数,确认是否存在频繁切换。
- 检查User/Topic Operator的请求成功率,确保功能正常。
内容的提问来源于stack exchange,提问作者Casper T
相关产品推荐
相关产品推荐

