Kubernetes 上的 Kafka:Strimzi Operator 生产实践

从有状态服务的视角讲透如何在 Kubernetes 上跑 Kafka:StatefulSet 的稳定网络标识与 Operator 模式的价值、Strimzi 的 Cluster Operator 与 Entity Operator 架构、核心 CRD 模型、Kafka CR 的 replicas 与 listeners 与 config 与 jvmOptions 与 storage 字段、StorageClass 与 PVC 持久化与 JBOD 与在线扩容、podAntiAffinity 与 topologySpreadConstraints 与机架感知的跨 AZ 高可用、滚动升级与扩容与分区再均衡与优雅停机、JMX Prometheus Exporter 与 PodMonitor 与 Cruise Control 监控运维,以及一整套真实踩坑清单

把 Kafka 搬上 Kubernetes 从来不是「写个 Deployment 挂个卷」这么简单。Kafka 是有状态服务里最难伺候的一类:每个 broker 要绑定固定的 broker.id 与稳定的网络标识,数据目录不能随 Pod 漂移,副本分布要跨故障域,重启要能优雅下线让 ISR 收敛。这些约束恰好是 StatefulSet 与 Operator 模式擅长解决的,也是本文要展开的部分。

1. 为什么在 Kubernetes 上跑 Kafka

1.1 有状态服务的三个硬约束

无状态服务上 K8s 是「换 Pod 即换机器」的简单模型,Kafka 不行。它有三条硬约束:

  • 身份稳定:broker.id 与 Pod 一一绑定,换名即换身份,会破坏副本与 ISR 的既有关系。
  • 数据落地:日志段写在本地盘,Pod 重建后必须挂回同一份卷,否则等于数据丢失。
  • 顺序敏感:broker 启动、下线、扩容都讲究次序,不能并发乱来,否则分区 leader 频繁切换。

K8s 原生对象里,只有 StatefulSet 同时提供「稳定网络标识 + 稳定存储 + 有序操作」这三件套,所以 Kafka 的编排底座必然是它。

1.2 StatefulSet 提供的稳定标识

StatefulSet 为每个副本生成确定的 Pod 名:kafka-0、kafka-1、kafka-2。这个确定名字带来两件事:一是 DNS 记录 kafka-0.kafka-headless.ns.svc 长期稳定,二是可以按序号推导 broker.id(Strimzi 默认就是「序号即 broker.id」)。

# StatefulSet 的稳定标识
# Pod 名: kafka-0 / kafka-1 / kafka-2(固定,重建后不变)
# 卷名:  data-kafka-0(与 Pod 绑定,随 Pod 重建复用)
# 启动次序: 0 -> 1 -> 2(默认 OrderedReady)
# 结论: 序号是 Kafka 的「身份锚点」,扩容只能从尾部加

正因为如此,Kafka 集群只能扩容不能缩容(缩容会改变副本与分区的归属,风险极高),这条铁律在 K8s 上同样成立。

1.3 Operator 模式的价值

纯手写 StatefulSet + ConfigMap 能跑起来,但「改一个配置触发滚动重启」「换证书」「加一个 listener」「扩一个主题分区」这些日常操作都要人肉编排。Operator 把这些领域知识写成代码:你只声明「我要 5 个 broker、每个 500Gi 盘、开一个 TLS listener」,Cluster Operator 负责把它翻译成 StatefulSet、Service、ConfigMap、Secret、PVC,并在后续持续调谐(reconcile)到目标状态。

2. Strimzi 架构与 CRD 模型

2.1 控制平面:三个 Operator 的职责

Strimzi 的控制平面由三类 Operator 组成,职责清晰分层:

Operator部署形态职责
Cluster Operator单副本 Deployment调谐 Kafka/KafkaConnect/KafkaMirrorMaker 等顶层 CR,管理 StatefulSet 与配置
Topic Operator与 Kafka 集群同命名空间把 KafkaTopic CR 同步为真实主题,双向对账
User Operator与 Kafka 集群同命名空间把 KafkaUser CR 同步为 SCRAM 凭证与 ACL

Topic 与 User Operator 合称 Entity Operator,通常作为 Kafka CR 里的一个 entityOperator 字段声明,由 Cluster Operator 拉起一个 Pod 承载。这样设计的好处是:集群级资源由 Cluster Operator 全权负责,命名空间级资源(主题、用户)交给 Entity Operator,避免权限过大。

2.2 核心 CRD 一览

Strimzi 的核心 CRD 覆盖了从集群到主题的完整声明面:

CRD作用典型字段
Kafka定义一套 Kafka 集群replicas / listeners / storage / config
KafkaNodePool定义节点池(KRaft 模式下划分角色)replicas / roles / storage
KafkaTopic声明一个主题partitions / replicas / config
KafkaUser声明一个客户端身份authentication / authorization
KafkaConnect声明 Connect 集群replicas / bootstrapServers
KafkaMirrorMaker2声明跨集群复制sourceCluster / targetCluster

其中 KafkaNodePool 是 Strimzi 0.36+ 引入的重要抽象:把「broker」拆成「节点池」,可以一个池只做 controller、一个池只做 broker,或按机型分层(大内存池放热数据、大磁盘池放冷数据)。它也让 KRaft 模式下的角色划分变得自然,配合 https://plumephp.com/kafka-kraft/ 里讲的 controller quorum 机制,可以平滑摆脱 ZooKeeper。

2.3 部署 Cluster Operator

生产环境推荐用 Helm 或 Operator Lifecycle Manager 安装,快速验证则可以直接 apply 官方 install 清单:

# 1) 创建命名空间
kubectl create namespace kafka

# 2) 安装 Cluster Operator(以 Helm 为例)
helm install strimzi-cluster-operator strimzi/strimzi-kafka-operator \
  --namespace kafka \
  --set watchAnyNamespace=false \
  --set replicas=1

# 3) 确认 Operator 就绪
kubectl -n kafka get deploy strimzi-cluster-operator
kubectl -n kafka logs deploy/strimzi-cluster-operator --tail=50

watchAnyNamespace=false 时 Operator 只监听自身命名空间,多租户场景下建议每个业务命名空间各装一个,或改用 watchNamespaces 显式列表,避免一个 Operator 管全集群带来的爆炸半径。

3. Kafka 集群 CR 与节点配置

3.1 Kafka CR 的整体骨架

一个生产可用的 Kafka CR 大致长这样,字段不多但每个都关键:

apiVersion: kafka.strimzi.io/v1beta2
kind: Kafka
metadata:
  name: prod-kafka
  namespace: kafka
spec:
  kafka:
    version: 3.7.0
    replicas: 3
    listeners:
      - name: plain
        port: 9092
        type: internal
        tls: false
      - name: tls
        port: 9093
        type: internal
        tls: true
        authentication:
          type: tls
    config:
      # 1) 默认副本数,按需覆盖
      default.replication.factor: 3
      min.insync.replicas: 2
      # 2) 关闭自动建主题,避免脏主题
      auto.create.topics.enable: false
      # 3) 日志保留策略
      log.retention.hours: 168
      num.partitions: 6
    storage:
      type: jbod
      volumes:
        - id: 0
          type: persistent-claim
          size: 500Gi
          class: fast-ssd
          deleteClaim: false
    resources:
      requests:
        cpu: "4"
        memory: 16Gi
      limits:
        cpu: "8"
        memory: 16Gi
  entityOperator:
    topicOperator: {}
    userOperator: {}

3.2 replicas 与 listeners

replicas 决定了 broker 数量,直接决定最大副本因子与跨故障域的分布能力。生产底线是 3 副本起步,跨 3 个可用区。注意 replicas 在 K8s 上只增不减。

listeners 是外部访问的入口设计,几种 type 的取舍很关键:

type暴露方式适用场景
internal集群内 headless Service同集群客户端
cluster-ip每 broker 一个 ClusterIP Service跨命名空间但不出集群
nodeport每 broker 一个 NodePort临时外部访问
loadbalancer每 broker 一个 LB云上外部访问
ingress七层网关按 SNI 路由共享入口、省 LB 成本

最容易被忽视的是 advertised host:nodeport/loadbalancer 类型必须配 overrides.bootstrap 与 overrides.brokers,否则客户端拿到的是 Pod IP 或内部 DNS,连不上。跨网络访问时,Kafka 返回给客户端的地址必须是「客户端视角可达的地址」,这是 Kafka 网络模型的老坑。

3.3 config 与 jvmOptions

spec.kafka.config 里放的是原生 server.properties 键值,Strimzi 会合并进 broker 配置。几条生产必设:

  • default.replication.factor: 3 与 min.insync.replicas: 2 组合,保证 acks=all 时至少两份落盘,配合 https://plumephp.com/kafka-cluster-ha/ 里的 ISR 机制,才能既容灾又不丢数据。
  • auto.create.topics.enable: false:生产必须关闭,主题走 KafkaTopic CR 声明,避免拼错名字就凭空建出一个单副本主题。
  • log.segment.bytes 与 log.retention.*:按吞吐与磁盘容量算,别用默认。

jvmOptions 单独控制 JVM 堆与 GC,Kafka 对堆外内存(page cache)依赖极大:

    jvmOptions:
      # 1) 堆不要给太大,把内存留给 page cache
      -Xms: 6g
      -Xmx: 6g
      # 2) G1 是 3.x 默认,低延迟场景可评估 ZGC
      -XX:
        UseG1GC: true
        MaxGCPauseMillis: 50

经验法则:堆内存不超过容器内存的一半。给 16Gi 容器配 6~8Gi 堆,其余留给操作系统 page cache,才能让日志段的读写命中缓存。

3.4 resources 与 storage

resources 的 requests 决定调度可行性,limits 决定 cgroup 上限。Kafka 是「吃满 CPU 也不嫌多」的服务,建议 CPU 不设 limit 或设得宽松(如 requests 4 / limits 8),内存 limit 等于 requests 防止被 OOM 前频繁换页。

storage 见下一章详述。这里只强调一点:storage 字段一旦创建就不能改类型,从 ephemeral 切到 persistent-claim 需要重建集群,规划阶段就要定死。

4. 存储类与 PVC 持久化

4.1 StorageClass 的选择

StorageClass 决定了卷的性能、可用区约束与扩容能力。生产选型的三个维度:

维度关注点建议
介质云盘类型(gp3/io2/ESSD)日志密集用高 IOPS 盘
绑定模式WaitForFirstConsumer 或 Immediate必须 WaitForFirstConsumer
扩容allowVolumeExpansion必须为 true

WaitForFirstConsumer 至关重要:它让 PVC 在 Pod 调度之后才绑定,从而把卷创建在 Pod 实际落地的可用区。若用 Immediate,卷可能先建在 AZ-a,而 Pod 被调度到 AZ-b,导致跨 AZ 挂载失败或性能暴跌。

4.2 ephemeral、persistent-claim 与 JBOD

Strimzi 的 storage 有三种类型:

  • ephemeral:用 emptyDir,Pod 重建数据全丢。仅用于 CI/演示,生产绝对禁用。
  • persistent-claim:单卷,最简单的生产形态。
  • jbod:多个卷组成一个逻辑存储,Kafka 会在卷之间轮转写入日志段,适合「多块小盘凑成大容量」或「SSD+HDD 分层」。

JBOD 的关键配置是 deleteClaim:设 false 时删除 Kafka 集群不会删 PVC,数据保留;设 true 会连卷一起删。生产一律 false,把删除权交给人。

    storage:
      type: jbod
      volumes:
        - id: 0
          type: persistent-claim
          size: 1Ti
          class: fast-ssd
          deleteClaim: false
        - id: 1
          type: persistent-claim
          size: 4Ti
          class: standard-hdd
          deleteClaim: false

注意 JBOD 下 Kafka 按「卷内剩余空间」策略选择落盘位置,容量不均衡会导致热卷。加盘时要成对加,别只给一个 broker 加。

4.3 PVC 回收策略与在线扩容

云盘 StorageClass 的 reclaimPolicy 决定 PVC 删除后底层盘的去向。Delete 会自动删盘(省钱但危险),Retain 保留盘(安全但要人工清理)。生产建议 Retain,配合 deleteClaim: false 形成双保险。

在线扩容 PVC 依赖两件事:StorageClass 的 allowVolumeExpansion: true,以及底层 CSI 驱动支持在线扩容。操作时直接改 Kafka CR 里的 size:

# 把每个 broker 的卷从 500Gi 扩到 1Ti
kubectl -n kafka patch kafka prod-kafka --type merge \
  -p '{"spec":{"kafka":{"storage":{"volumes":[{"id":0,"type":"persistent-claim","size":"1Ti","class":"fast-ssd","deleteClaim":false}]}}}}'

# 观察 PVC 扩容进度
kubectl -n kafka get pvc -l strimzi.io/cluster=prod-kafka -w

扩容本身不触发 broker 重启(部分 CSI 需要文件系统扩展,可能触发一次 Pod 重启),但扩容后必须手动触发分区再均衡,否则新空间不会自动被利用——这一点在第六章讲。

5. 反亲和、拓扑与高可用

5.1 podAntiAffinity 强制分散

默认情况下 K8s 可能把三个 broker 塞到同一台机器上,一台机器挂了整个集群完蛋。必须用反亲和性把 broker 打散到不同节点:

    template:
      pod:
        affinity:
          podAntiAffinity:
            requiredDuringSchedulingIgnoredDuringExecution:
              - labelSelector:
                  matchLabels:
                    strimzi.io/cluster: prod-kafka
                    strimzi.io/name: prod-kafka-kafka
                topologyKey: kubernetes.io/hostname

requiredDuringScheduling 是硬约束:节点不够时 Pod 会 Pending,宁可调度失败也不挤在一起。这在 3 副本集群 + 3 节点时是刚需。若节点数紧张,可退化为 preferredDuringScheduling 的软约束,但要清楚这是拿可用性换调度弹性。

5.2 topologySpreadConstraints 更精细的分布

反亲和只能保证「不同节点」,管不了「不同可用区」。要跨 AZ 均匀分布,用拓扑分布约束:

        topologySpreadConstraints:
          - maxSkew: 1
            topologyKey: topology.kubernetes.io/zone
            whenUnsatisfiable: DoNotSchedule
            labelSelector:
              matchLabels:
                strimzi.io/name: prod-kafka-kafka

maxSkew: 1 表示任意两个可用区的 broker 数量差不超过 1,DoNotSchedule 表示不满足就不调度。这一条加上 hostname 反亲和,就实现了「跨 AZ + 跨节点」的双层打散,是生产高可用的标配组合。

5.3 机架感知 rack awareness 与跨 AZ

K8s 层面的打散只解决「Pod 落在哪」,Kafka 层面的副本分布还要靠 rack awareness:给每个 broker 声明它属于哪个机架(可用区),Kafka 建副本时会尽量把同一分区的副本放到不同机架。Strimzi 从 Pod 的拓扑标签自动推导机架:

    rack:
      topologyKey: topology.kubernetes.io/zone

开启后,num.partitions 里的每个分区副本会跨 AZ 分布,单个 AZ 整体故障时分区仍能在其他 AZ 选出 leader。代价是跨 AZ 写入的同步延迟增加(同区可能 1ms,跨区可能 5~10ms),所以 acks=all 的延迟要按跨区算。吞吐敏感的场景可以按主题区分:核心交易主题跨 AZ 保安全,日志类主题单 AZ 求性能。

分区与副本的规划直接影响这套机制的效果,具体策略可参考 https://plumephp.com/kafka-topic-design/。

6. 滚动升级与扩容

6.1 谁触发了滚动重启

Strimzi 的 Cluster Operator 在检测到「需要重建 Pod 的变更」时会自动执行滚动重启。常见触发源:

  • 修改 spec.kafka.config(部分参数需重启生效)
  • 修改 resources、jvmOptions、storage
  • 升级 Kafka 版本或 Strimzi 版本
  • 修改 template 里的 Pod 相关字段
  • 证书轮换(TLS 证书到期自动续期)

滚动重启的次序由 Strimzi 编排:默认从序号最大的 broker 开始,逐个重启,每个重启后等它重新加入 ISR 再动下一个。这个过程会自动尊重 podDisruptionBudget 与 min.insync.replicas,不会一次性把副本打光。

6.2 滚动升级的参数控制

大集群的滚动重启可能持续数小时,Strimzi 提供了两个关键旋钮控制节奏:

参数位置作用
maxUnavailablespec.kafka.template.podDisruptionBudget允许同时不可用的 Pod 数
brokerRackAwareRebalancespec.kafka 相关注解按机架分批重启
STRIMZI_OPERATOR 的调谐间隔Operator 环境变量控制对账频率
    template:
      podDisruptionBudget:
        maxUnavailable: 1

maxUnavailable: 1 保证任意时刻最多一个 broker 下线,配合 3 副本 + min.insync.replicas: 2,滚动期间仍能满足 acks=all 写入。若把 maxUnavailable 设成 2,写入会因 ISR 不足而失败。

6.3 扩容 broker 与分区再均衡

扩容只需改 replicas,但扩容不等于负载均衡:

# 1) 3 个 broker 扩到 5 个
kubectl -n kafka patch kafka prod-kafka --type merge \
  -p '{"spec":{"kafka":{"replicas":5}}}'

# 2) 等待新 Pod 就绪
kubectl -n kafka get pods -l strimzi.io/cluster=prod-kafka -w

新 broker 加入后,既有主题的分区不会自动迁过去,必须触发分区再均衡。Strimzi 的推荐做法是配好 Cruise Control 后提交一个 rebalance 提案:

# 用 Cruise Control 生成再均衡提案
kubectl -n kafka exec -it deploy/cruise-control -- \
  curl -X POST http://localhost:9090/kafka-cruise-control/add_broker \
  -H 'Content-Type: application/json' \
  -d '{"dryrun":true,"goals":["RackAwareGoal","ReplicaCapacityGoal"]}'

dryrun: true 先看提案是否安全(不违反机架感知、不超过副本容量上限),确认后再提交执行。再均衡本身是大量分区副本迁移,会占用网络与磁盘 IO,务必在业务低峰做,并盯住 broker 的网络出口带宽。

6.4 优雅停机与 PDB

K8s 删除 Pod 时默认先发 SIGTERM,Strimzi 会拦截这个信号并触发 Kafka 的受控关闭:broker 停止接受新请求、把 leader 角色交出去、等待副本同步、最后退出。这个过程由 terminationGracePeriodSeconds 控制窗口:

    template:
      kafkaContainer:
        ...
      pod:
        terminationGracePeriodSeconds: 300

给足 300s 让 leader 平滑迁移,避免直接 kill 导致分区短暂不可用。PodDisruptionBudget 则保护节点维护(如 drain)时的主动驱逐,两者一个管「被动删除」一个管「主动驱逐」,缺一不可。

7. 监控、运维与常见坑

7.1 JMX Prometheus Exporter 与 PodMonitor

Strimzi 内置了 JMX Prometheus Exporter,只需在 Kafka CR 里声明 metricsConfig 指向一份 ConfigMap,Exporter 就会以 sidecar 形式注入并暴露 /metrics:

  kafka:
    metricsConfig:
      type: jmxPrometheusExporter
      valueFrom:
        configMapKeyRef:
          name: kafka-metrics
          key: kafka-metrics-config.yml

抓取侧用 PodMonitor(Prometheus Operator 环境)自动发现:

apiVersion: monitoring.coreos.com/v1
kind: PodMonitor
metadata:
  name: kafka-resources-metrics
  namespace: kafka
spec:
  selector:
    matchLabels:
      strimzi.io/cluster: prod-kafka
  podMetricsEndpoints:
    - path: /metrics
      port: tcp-prometheus
  namespaceSelector:
    matchNames:
      - kafka

必看的四组指标:kafka_server_replicamanager_underreplicatedpartitions(ISR 落后分区数,非 0 即告警)、kafka_controller_kafkacontroller_activecontrollercount(控制器唯一性)、kafka_server_brokertopicmetrics_bytesinpersec(入流量)、kafka_log_log_size(磁盘水位)。完整的监控体系与告警阈值设计,可参考 https://plumephp.com/kafka-monitoring-operations/。

7.2 Cruise Control 自动再均衡

Cruise Control 是一个独立的负载均衡服务,通过分析 broker 的资源利用率生成分区迁移提案。Strimzi 把它作为 Kafka CR 的一个子字段部署:

  cruiseControl:
    brokerCapacity:
      inboundNetwork: 10000MB/s
      outboundNetwork: 10000MB/s
    config:
      # 1) 目标权重:磁盘优先于 CPU
      hard.goals: >
        com.linkedin.kafka.cruisecontrol.analyzer.goals.RackAwareGoal,
        com.linkedin.kafka.cruisecontrol.analyzer.goals.ReplicaCapacityGoal

hard.goals 是绝不能违反的约束(如机架感知、副本容量上限),default.goals 是优化目标(如磁盘均衡)。Cruise Control 还支持 KafkaRebalance CR 声明式触发,配合 autoRebalanceEnabled 可以在扩容后自动收敛。

7.3 常见坑清单

  • ephemeral 存储上生产:Pod 一漂移数据全没,重建集群才发现,血的教训。
  • StorageClass 用 Immediate 绑定:卷与 Pod 跨 AZ,挂载失败或延迟暴涨,必须 WaitForFirstConsumer。
  • 没开反亲和:三个 broker 全在同一节点,节点故障即全集群不可用。
  • resources 没设 requests:调度器把 broker 塞进资源紧张的节点,OOM 与 CPU 抢占随之而来。
  • 堆内存给太大:堆吃掉容器内存,page cache 被挤没,读性能断崖。
  • 缩容 replicas:Kafka 在 K8s 上只能扩不能缩,缩容会导致分区数据丢失。
  • auto.create.topics.enable 没关:客户端拼错主题名就建出单副本主题,埋下数据丢失隐患。
  • listeners 的 advertised host 没配:客户端拿到 Pod 内网地址,连接超时。
  • 扩容后忘了再均衡:新 broker 空转,老 broker 磁盘打满,容量假象。
  • PDB 与 maxUnavailable 冲突:驱逐时被 PDB 拦住,节点维护卡死,需协调两者取值。
  • 证书自动轮换期间客户端不重载:TLS 证书换了客户端还在用旧的,连接被拒,需客户端支持热加载。
  • terminationGracePeriodSeconds 太短:broker 来不及交接 leader 就被 kill,分区短暂不可用。

8. 总结

在 Kubernetes 上跑 Kafka,本质是「用声明式编排驯服有状态服务」:底座靠 StatefulSet 的稳定标识与稳定存储,控制面靠 Strimzi 的三类 Operator 把领域知识代码化,配置面靠 Kafka CR 把 replicas、listeners、config、storage 一次声明到位,可用性靠 podAntiAffinity 与 topologySpreadConstraints 与 rack awareness 三层打散,变更靠 Cluster Operator 编排的滚动重启与 Cruise Control 的再均衡。真正决定成败的不是「能不能跑起来」,而是那些细节:WaitForFirstConsumer 的绑定模式、deleteClaim 的取值、堆与 page cache 的比例、maxUnavailable 与 min.insync.replicas 的配合。把这些细节钉死,Kafka 在 K8s 上就能既享受云原生的编排红利,又守住有状态服务的可靠性底线。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「kafka」更多文章

  1. ksqlDB 流式 SQL:流表模型、窗口聚合与生产运维
  2. Kafka 消费延迟诊断:Lag 定位、分区倾斜与治理
  3. Kafka 分层存储:KIP-405 冷热数据卸载与对象存储实践