You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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: .*

可能原因分析

  1. Kraft控制器选举时序延迟:控制器切换过程中,旧控制器节点未及时感知新控制器当选,导致转发到旧节点的请求失败。
  2. 请求转发竞态条件:Strimzi的请求转发逻辑在集群状态同步不及时时,错误地将请求发往非活跃控制器。
  3. 版本兼容bug:Strimzi v0.45.0与Kafka 3.9.0的组合存在已知的控制器状态同步问题。
  4. 网络稳定性问题: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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.14 12:34:49