文章

Kafka 与 Zookeeper 部署实践与排障

Kafka 与 Zookeeper 部署实践与排障

概述

本文覆盖 Kafka + Zookeeper 在 Kubernetes 环境下的部署实践、原理机制、常见问题排障、安全加固、监控集成、以及 KRaft 模式迁移路径。适用于 Kafka 3.x(含 KRaft GA)+ ZK 3.6+ 的生产集群。

结论先行:新部署建议直接上 KRaft 模式(Kafka 3.3+ GA),省掉 ZK 运维成本;存量 ZK 模式集群按计划迁移。两种模式的部署实践本文都会覆盖。


目录

  • [[#架构总览]]
  • [[#一、Zookeeper 原理与 K8s 部署]]
  • [[#二、Kafka 原理与 K8s 部署]]
  • [[#三、常见部署问题排障]]
  • [[#四、生产环境最佳实践]]
  • [[#五、安全加固:SASL + TLS + ACL]]
  • [[#六、监控集成:Prometheus + Grafana]]
  • [[#七、快速排查 Checklist]]
  • [[#八、ZK → KRaft 迁移路径]]
  • [[#九、Strimzi Operator 方案简介]]

架构总览

传统 ZK 模式

┌─────────────────────────────────────────────────────┐
│              Zookeeper 集群 (3/5 节点)                 │
│         元数据:Topic / Partition / ISR / Controller   │
│         端口:2181(client) 2888(follower) 3888(election)│
└──────────────┬──────────────────────────────────────┘
               │ 注册 & 心跳 (session.timeout.ms)
    ┌──────────┼──────────┐
    ↓          ↓          ↓
┌────────┐ ┌────────┐ ┌────────┐
│Broker 0│ │Broker 1│ │Broker 2│   Kafka 集群
│:9092   │ │:9092   │ │:9092   │   (log.dirs → PVC)
└────────┘ └────────┘ └────────┘
    ↑          ↑          ↑
    └──────────┼──────────┘

     Producer / Consumer (acks=all)

KRaft 模式(3.3+ GA)

┌─────────────────────────────────────────────────────┐
│         Kafka 集群 (Controller + Broker 合一)          │
│  ┌────────────┐  ┌────────────┐  ┌────────────┐      │
│  │Broker 0    │  │Broker 1    │  │Broker 2    │      │
│  │+ Controller│  │+ Controller│  │+ Controller│      │
│  │:9092 data  │  │:9092 data  │  │:9092 data  │      │
│  │:9093 ctrl  │  │:9093 ctrl  │  │:9093 ctrl  │      │
│  └────────────┘  └────────────┘  └────────────┘      │
│  元数据通过 Raft 协议同步,无需 ZK                      │
│  controller.quorum.voters=1@b0:9093,2@b1:9093,3@b2:9093│
└─────────────────────────────────────────────────────┘

     Producer / Consumer (acks=all)

模式对比

维度ZK 模式KRaft 模式
元数据存储ZK 集群(外部 znode)Kafka 内部 Raft 元数据日志
组件数量ZK(3+) + Kafka(3+)Kafka(3+) only
运维复杂度高(多维护一套 ZK 集群)
元数据操作延迟较高(经 ZK 网络读写)低(本地日志 + Raft)
分区上限~20 万(ZK watch 瓶颈)数百万
Controller 故障恢复秒级(需 ZK 重新选举)毫秒级(Raft 预选主)
推荐版本Kafka ≤ 3.x(兼容期)Kafka ≥ 3.3(生产可用)
部署 YAML 复杂度高(两套 StatefulSet)中(一套 StatefulSet)
迁移路径→ 逐步迁移到 KRaft新部署首选

一、Zookeeper 原理与 K8s 部署

1.1 ZAB 协议与 Leader 选举

Zookeeper 使用 ZAB (Zookeeper Atomic Broadcast) 协议保证数据一致性,类似于 Paxos 但做了简化优化。

阶段 1: 发现 (Discovery)        阶段 2: 同步 (Synchronization)
┌─────────┐                    ┌─────────┐
│ 各节点交换最新 epoch          │ Leader 将最新数据同步给所有 Follower │
│ 选出包含最多已提交事务的节点    │ Follower 追齐 Leader 的数据          │
│ 作为 Leader 候选              │                                      │
└─────────┘                    └─────────┘
       ↓                              ↓
阶段 3: 广播 (Broadcast)
┌─────────────────────────────────────────┐
│ Leader 收到写请求 → 生成 Proposal (ZXID) │
│ → 广播给所有 Follower                    │
│ → 过半 Follower ACK → Commit → 响应客户端 │
└─────────────────────────────────────────┘

ZXID(事务 ID)结构

ZXID (64 bit)
┌────────────────┬────────────────┐
│  epoch (32bit)  │ counter (32bit) │
└────────────────┴────────────────┘
epoch: Leader 任期号,每次选举递增
counter: 该 epoch 内的事务序号,单调递增

类比:epoch 就像”朝代号”,counter 就像”圣旨编号”。换了皇帝(Leader),朝代号变,但圣旨编号从 0 重新开始。通过 epoch + counter 可以全局排序所有事务。

Leader 选举触发条件

  • Leader 节点崩溃 / 网络不可达
  • Follower 在 session.timeout 内未收到 Leader 心跳
  • 集群启动时初始化选举

选举过程

  1. 每个节点提出自己为 Leader 候选(投票格式:(myid, ZXID)
  2. 收到其他节点投票后比较:先比 ZXID(大的优先),再比 myid(大的优先)
  3. 过半节点同意 → 选出 Leader → 其他节点变为 Follower

1.2 事务日志与快照机制

ZK 的持久化由两部分组成,必须分盘存放

类型目录内容写入时机文件名格式
事务日志dataLogDir每条写操作的 WAL 日志每次写操作同步刷盘log.<ZXID>
快照dataDir内存数据的定期全量序列化snapCount 次事务触发snapshot.<ZXID>
事务日志 (dataLogDir)          快照 (dataDir)
┌──────────────────┐          ┌────────────────────┐
│ log.100          │          │ snapshot.500       │  ← 内存数据全量
│ log.200          │          │ snapshot.800       │
│ log.300          │          │ snapshot.1200      │
│ ...              │          │                    │
│ log.1500 (active)│          │                    │
└──────────────────┘          └────────────────────┘
       ↑                              ↑
  每次写操作同步写入             每 snapCount 次事务触发
  (性能瓶颈!必须用 SSD)         可后台异步生成

恢复流程:加载最新快照 → 重放快照点之后的事务日志 → 恢复到最新状态

自动清理配置

# zoo.cfg
autopurge.snapRetainCount=3      # 保留最近 3 个快照及其对应日志
autopurge.purgeInterval=1        # 每小时自动清理一次

踩坑:如果不配 autopurge,事务日志会无限增长直到磁盘满。ZK 进程会因为无法写入事务日志而 OOM 或 hang,直接导致 Kafka 集群不可用。

1.3 核心配置项详解

参数推荐值说明调优场景
tickTime2000基本时间单元(ms),心跳 = 1 tick网络延迟高的环境可调大到 3000
initLimit10follower 初始连接 leader 的最大 tick 数(= 20s)首次启动慢或数据量大时调大
syncLimit5follower 与 leader 同步的最大 tick 数(= 10s)网络抖动频繁时适当调大
maxClientCnxns60单个客户端 IP 最大并发连接数防止单客户端耗尽连接
autopurge.snapRetainCount3保留的快照数量最少保留 3 个用于恢复
autopurge.purgeInterval1自动清理间隔(小时)设为 0 则禁用自动清理(不推荐)
dataDir/data/zookeeper/data快照目录(可放 HDD)-
dataLogDir/data/zookeeper/log事务日志目录(必须 SSDIO 敏感,决定 ZK 写入性能
snapCount100000每多少次事务触发一次快照事务频繁时调大减少快照开销
maxSessionTimeout40000最大 session 超时(ms)限制客户端 session 上限
4lw.commands.whitelistsrvr,stat,mntr允许的 4 字母命令白名单监控需要 mntr,排查需要 stat
JVM 堆 (-Xmx)2-4G不要超过 4GGC 停顿会阻塞 Kafka ISR 变更

1.4 StatefulSet 完整部署 YAML

---
# Headless Service - Pod 间通过域名互访
apiVersion: v1
kind: Service
metadata:
  name: zk-headless
  namespace: middleware
  labels:
    app: zookeeper
spec:
  clusterIP: None
  selector:
    app: zookeeper
  ports:
    - name: client
      port: 2181
      targetPort: 2181
    - name: follower
      port: 2888
      targetPort: 2888
    - name: leader-election
      port: 3888
      targetPort: 3888
  publishNotReadyAddresses: true   # ZK 集群组建期间需要 DNS 解析

---
# ConfigMap - zoo.cfg + 4lw 白名单 + JVM 参数
apiVersion: v1
kind: ConfigMap
metadata:
  name: zk-config
  namespace: middleware
data:
  zoo.cfg: |
    tickTime=2000
    initLimit=10
    syncLimit=5
    dataDir=/data/zookeeper/data
    dataLogDir=/data/zookeeper/log
    maxClientCnxns=60
    autopurge.snapRetainCount=3
    autopurge.purgeInterval=1
    snapCount=100000
    # 4 字母命令白名单(监控 + 排障用)
    4lw.commands.whitelist=srvr,stat,mntr,ruok,conf,envi
    # StatefulSet Pod 域名:zk-{ordinal}.zk-headless.middleware.svc.cluster.local
    server.1=zk-0.zk-headless.middleware.svc.cluster.local:2888:3888
    server.2=zk-1.zk-headless.middleware.svc.cluster.local:2888:3888
    server.3=zk-2.zk-headless.middleware.svc.cluster.local:2888:3888

  jvm.env: |
    JVMFLAGS="-Xms2g -Xmx2g -XX:MetaspaceSize=128m
    -XX:+UseG1GC
    -XX:MaxGCPauseMillis=50
    -XX:InitiatingHeapOccupancyPercent=40
    -XX:+ExplicitGCInvokesConcurrent
    -Djute.maxbuffer=1048576
    -XX:+HeapDumpOnOutOfMemoryError
    -XX:HeapDumpPath=/data/zookeeper/heapdump"

---
# PodDisruptionBudget - 保证至少 2 个 ZK 可用
apiVersion: policy/v1
kind: PodDisruptionBudget
metadata:
  name: zk-pdb
  namespace: middleware
spec:
  minAvailable: 2
  selector:
    matchLabels:
      app: zookeeper

---
# StatefulSet
apiVersion: apps/v1
kind: StatefulSet
metadata:
  name: zk
  namespace: middleware
  labels:
    app: zookeeper
spec:
  serviceName: zk-headless
  replicas: 3
  podManagementPolicy: OrderedReady        # 有序启动,确保 myid 递增
  updateStrategy:
    type: RollingUpdate
    rollingUpdate:
      partition: 0                          # 滚更时从 ordinal 最大开始替换
  selector:
    matchLabels:
      app: zookeeper
  template:
    metadata:
      labels:
        app: zookeeper
    spec:
      affinity:
        podAntiAffinity:                   # 反亲和,Pod 打散到不同物理节点
          requiredDuringSchedulingIgnoredDuringExecution:
            - labelSelector:
                matchLabels:
                  app: zookeeper
              topologyKey: kubernetes.io/hostname
      terminationGracePeriodSeconds: 120   # 给 ZK 足够时间做 leader 重新选举
      initContainers:
        - name: init-myid                  # 动态生成 myid 文件
          image: zookeeper:3.8
          command:
            - sh
            - -c
            - |
              # 从 Pod 名提取序号,+1 作为 myid(myid 从 1 开始)
              ORDINAL=${HOSTNAME##*-}
              MY_ID=$((ORDINAL + 1))
              echo "Writing myid=${MY_ID} to /data/zookeeper/data/myid"
              mkdir -p /data/zookeeper/data /data/zookeeper/log
              echo "${MY_ID}" > /data/zookeeper/data/myid
              # 如果数据目录已有 myid,检查是否一致
              EXISTING=$(cat /data/zookeeper/data/myid 2>/dev/null || echo "")
              if [ -n "${EXISTING}" ] && [ "${EXISTING}" != "${MY_ID}" ]; then
                echo "WARNING: myid changed from ${EXISTING} to ${MY_ID}!"
              fi
          volumeMounts:
            - name: zk-data
              mountPath: /data/zookeeper/data
            - name: zk-log
              mountPath: /data/zookeeper/log
      containers:
        - name: zookeeper
          image: zookeeper:3.8
          ports:
            - containerPort: 2181
              name: client
            - containerPort: 2888
              name: follower
            - containerPort: 3888
              name: leader-election
          env:
            - name: ZOO_CFG_DIR
              value: /conf
          command:
            - sh
            - -c
            - |
              # 加载 JVM 参数
              source /conf/jvm.env
              export JVMFLAGS
              exec zkServer.sh start-foreground
          resources:
            requests:
              cpu: 500m
              memory: 2Gi
            limits:
              cpu: 1
              memory: 4Gi
          readinessProbe:                  # 就绪检查:ZK 能响应客户端请求
            exec:
              command:
                - sh
                - -c
                - 'echo ruok | nc localhost 2181 | grep imok'
            initialDelaySeconds: 10
            periodSeconds: 10
            timeoutSeconds: 5
            failureThreshold: 3
          livenessProbe:                   # 存活检查:进程是否存活
            tcpSocket:
              port: 2181
            initialDelaySeconds: 30
            periodSeconds: 30
            timeoutSeconds: 5
            failureThreshold: 3
          startupProbe:                     # 启动检查:给 ZK 足够启动时间
            tcpSocket:
              port: 2181
            failureThreshold: 30
            periodSeconds: 10
          lifecycle:
            preStop:                        # 优雅关闭:主动断开客户端连接
              exec:
                command:
                  - sh
                  - -c
                  - "zkServer.sh stop"
          volumeMounts:
            - name: zk-data
              mountPath: /data/zookeeper/data
            - name: zk-log
              mountPath: /data/zookeeper/log
            - name: config
              mountPath: /conf/zoo.cfg
              subPath: zoo.cfg
            - name: config
              mountPath: /conf/jvm.env
              subPath: jvm.env
      volumes:
        - name: config
          configMap:
            name: zk-config
  volumeClaimTemplates:
    - metadata:
        name: zk-data                       # 快照目录
      spec:
        accessModes: ["ReadWriteOnce"]
        storageClassName: fast-ssd
        resources:
          requests:
            storage: 20Gi
    - metadata:
        name: zk-log                        # 事务日志目录(独立 PVC)
      spec:
        accessModes: ["ReadWriteOnce"]
        storageClassName: fast-ssd          # 必须用 SSD!
        resources:
          requests:
            storage: 20Gi

1.5 ZK 集群扩缩容注意事项

ZK 集群扩容(如 3 → 5 节点)不像 Kafka 那样简单加 Pod 即可,需要:

  1. 更新 zoo.cfg:添加新的 server.N 配置
  2. 逐个重启:所有现有节点需要加载新的 server.N 列表
  3. Reconfig(ZK 3.5+):使用动态重新配置命令,不需要重启
# ZK 3.5+ 动态添加节点(不需要重启现有节点)
kubectl exec zk-0 -n middleware -- zkCli.sh -server localhost:2181 \
  reconfig -add server.4=zk-3.zk-headless.middleware.svc.cluster.local:2888:3888:participant

# 动态移除节点
kubectl exec zk-0 -n middleware -- zkCli.sh -server localhost:2181 \
  reconfig -remove server.4

关键:ZK 集群节点数必须是奇数(3/5/7),过半机制要求。偶数节点没有额外容错收益反而增加选举复杂度。

1.6 ZK 4lw 命令速查

命令作用示例
ruok检查 ZK 是否存活(返回 imokecho ruok | nc localhost 2181
stat查看 ZK 状态(版本、mode、连接数)echo stat | nc localhost 2181
srvr查看 Server 信息(zxid、模式、延迟)echo srvr | nc localhost 2181
mntr监控指标输出(Prometheus 格式友好)echo mntr | nc localhost 2181
conf查看当前配置echo conf | nc localhost 2181
envi查看环境信息echo envi | nc localhost 2181
dump列出未处理的 session 和临时节点(仅 Leader)echo dump | nc localhost 2181
wchs列出 watch 摘要echo wchs | nc localhost 2181
wchc按 session 列出 watch 详情echo wchc | nc localhost 2181
wchp按 path 列出 watch 详情echo wchp | nc localhost 2181
isro是否只读模式(返回 RORWecho isro | nc localhost 2181
crst重置连接/客户端统计(仅 Leader)echo crst | nc localhost 2181

安全提示4lw.commands.whitelist 默认只允许 srvr。生产环境按需开放,crst 等危险命令不要加白名单。


二、Kafka 原理与 K8s 部署

2.1 分区 / 副本 / ISR 机制

Topic: order-events (partitions=3, replication.factor=3)

Partition 0:                    Partition 1:                    Partition 2:
┌────────────┐                  ┌────────────┐                  ┌────────────┐
│ Broker 0   │ ← Leader         │ Broker 1   │ ← Leader         │ Broker 2   │ ← Leader
│ Broker 1   │ ← Follower (ISR) │ Broker 2   │ ← Follower (ISR) │ Broker 0   │ ← Follower (ISR)
│ Broker 2   │ ← Follower (ISR) │ Broker 0   │ ← Follower (ISR) │ Broker 1   │ ← Follower (ISR)
└────────────┘                  └────────────┘                  └────────────┘
     ↕                               ↕                               ↕
  ISR = {0,1,2}                  ISR = {0,1,2}                  ISR = {0,1,2}

核心概念

概念说明
Partition(分区)Topic 的并行单元,每个分区是一个有序的、不可变的消息序列
Replica(副本)每个分区有 N 个副本分布在不同 Broker 上,其中一个为 Leader
Leader分区的主副本,负责所有读写请求
Follower分区的从副本,从 Leader 拉取数据保持同步
ISR (In-Sync Replicas)与 Leader 保持同步的副本集合(含 Leader 自身)
OSR (Out-of-Sync Replicas)落后于 Leader 的副本,不在 ISR 中
AR (All Replicas)ISR + OSR = 分区所有副本
HW (High Watermark)已被所有 ISR 副本确认的最大 offset,消费者只能看到 HW 之前的消息
LEO (Log End Offset)每个副本日志的下一条写入 offset

ISR 扩缩机制

Follower 持续从 Leader 拉取数据

Follower 的 LEO 落后 Leader 超过 replica.lag.time.max.ms (默认 30s)

Leader 将该 Follower 从 ISR 中移除 (Shrinking ISR)

Follower 追上 Leader 的 LEO

Leader 将该 Follower 重新加入 ISR (Expanding ISR)

关键参数

参数默认值说明
replica.lag.time.max.ms30000Follower 落后超过此时间则移出 ISR
min.insync.replicas1配合 acks=all,ISR 少于此值时 Broker 拒绝写入
unclean.leader.election.enablefalse是否允许非 ISR 副本成为 Leader(数据丢失风险)
default.replication.factor1默认副本数(生产建议 3)

生产建议replication.factor=3 + min.insync.replicas=2 + acks=all + unclean.leader.election.enable=false。这样在 1 个 Broker 故障时仍可正常读写,且不丢数据。

2.2 日志存储与清理策略

Kafka 的日志存储结构:

log.dirs=/data/kafka/logs/
├── order-events-0/              # Topic-partition 目录
│   ├── 00000000000000000000.log  # segment 数据文件
│   ├── 00000000000000000000.index # offset 索引
│   ├── 00000000000000000000.timeindex # timestamp 索引
│   ├── 00000000000005367851.log  # 下一个 segment
│   ├── 00000000000005367851.index
│   ├── 00000000000005367851.timeindex
│   └── leader-epoch-checkpoint   # leader epoch 信息
├── order-events-1/
├── order-events-2/
└── __consumer_offsets-0/         # 内置 Topic:消费偏移量

两种清理策略

策略配置值适用场景原理
delete(默认)cleanup.policy=delete普通消息 Topic按时间/大小删除旧 segment
compactcleanup.policy=compactKV 状态 Topic(如 __consumer_offsets每个 key 只保留最新值
delete+compactcleanup.policy=delete,compact同时需要两种策略先 compact 再按时间删除

delete 策略参数

# 按时间删除(任一条件满足即删除)
log.retention.hours=168              # 保留 7 天
log.retention.bytes=10737418240      # 或保留 10GB(先到先删)

# segment 滚动
log.segment.bytes=1073741824         # 单个 segment 1GB
log.segment.ms=604800000            # 或 7 天强制滚动新 segment

# 清理线程
log.retention.check.interval.ms=300000  # 每 5 分钟检查一次
num.io.threads=8                     # IO 线程数

compact 策略参数

# compact 相关
min.cleanable.dirty.ratio=0.5        # 脏数据占比超过 50% 触发压缩
delete.retention.ms=86400000         # tombstone(null 值)保留 24 小时
segment.ms=300000                    # compact segment 滚动间隔

2.3 核心配置项详解

Broker 级配置

参数推荐值说明调优场景
broker.id0, 1, 2…集群内唯一,StatefulSet 从 Pod 序号生成-
log.dirs/data/kafka/logs日志目录(可逗号分隔多目录,跨磁盘)多磁盘环境分散 IO
num.partitions3+新建 Topic 的默认分区数按吞吐量需求调整
default.replication.factor3默认副本数生产环境最少 3
min.insync.replicas2ISR 最小副本数配合 acks=all
log.retention.hours168日志保留时间(7 天)按业务需求调整
log.retention.bytes-1日志保留大小(-1 不限制)磁盘有限时设置上限
log.segment.bytes1073741824单个 segment 大小(1GB)小 segment 更快清理
zookeeper.connectzk-0.zk-headless:2181,…ZK 地址(ZK 模式)KRaft 模式不需要
advertised.listenersPLAINTEXT://<Pod域名>:9092对外暴露地址不能用 0.0.0.0
auto.create.topics.enablefalse生产环境关闭自动创建防止误创建 Topic
unclean.leader.election.enablefalse禁止非 ISR 副本成为 Leader防止数据丢失
replica.lag.time.max.ms30000Follower 落后超时移出 ISR网络差时适当调大
num.network.threads3网络请求处理线程数高并发时调大
num.io.threads8磁盘 IO 线程数多磁盘时调大
socket.send.buffer.bytes102400TCP 发送缓冲区高吞吐场景调大
socket.receive.buffer.bytes102400TCP 接收缓冲区高吞吐场景调大
log.flush.interval.messagesLong.MAX每多少条消息强制刷盘默认由 OS 管理,不建议改
log.flush.interval.msnull每多少 ms 强制刷盘强制刷盘会降低性能
group.initial.rebalance.delay.ms3000Consumer Group 首次 rebalance 延迟减少 rebalance 期间消息延迟

Producer 关键配置

参数推荐值说明
acksall等待所有 ISR 副本确认(最安全)
retries2147483647无限重试(配合 delivery.timeout.ms
delivery.timeout.ms120000消息发送的总超时(含重试)
max.in.flight.requests.per.connection5单连接未确认请求上限(保持顺序需 ≤ 5)
enable.idempotencetrue幂等 Producer(防重复消息)
compression.typelz4压缩算法(lz4 性能最好)
batch.size16384批量发送大小(字节)
linger.ms10等待批量发送的时间(ms)
buffer.memory33554432Producer 缓冲区大小(32MB)
max.request.size1048576单条消息最大大小(1MB)

Consumer 关键配置

参数推荐值说明
enable.auto.commitfalse关闭自动提交偏移量(手动提交更安全)
auto.offset.resetearliest无偏移量时从最早开始消费
session.timeout.ms45000Consumer 心跳超时
heartbeat.interval.ms15000心跳间隔(建议 = session.timeout / 3)
max.poll.records500单次 poll 最大消息数
max.poll.interval.ms300000两次 poll 之间的最大间隔(处理慢时调大)
fetch.min.bytes1最小拉取字节数
fetch.max.wait.ms500最大等待时间
isolation.levelread_committed只读取已提交的事务消息

2.4 Kafka StatefulSet 部署 — ZK 模式(完整 YAML)

---
apiVersion: v1
kind: Service
metadata:
  name: kafka-headless
  namespace: middleware
  labels:
    app: kafka
spec:
  clusterIP: None
  selector:
    app: kafka
  ports:
    - name: client
      port: 9092
      targetPort: 9092
    - name: controller
      port: 9093
      targetPort: 9093
    - name: jmx
      port: 9999
      targetPort: 9999
  publishNotReadyAddresses: true

---
apiVersion: v1
kind: ConfigMap
metadata:
  name: kafka-config
  namespace: middleware
data:
  server.properties: |
    # Broker 基础
    broker.id=0
    log.dirs=/bitnami/kafka/data
    # 网络配置
    listeners=PLAINTEXT://0.0.0.0:9092
    # advertised.listeners 在启动脚本中动态生成
    # Zookeeper
    zookeeper.connect=zk-0.zk-headless:2181,zk-1.zk-headless:2181,zk-2.zk-headless:2181
    zookeeper.connection.timeout.ms=18000
    # Topic 默认配置
    num.partitions=3
    default.replication.factor=3
    min.insync.replicas=2
    # 日志保留
    log.retention.hours=168
    log.segment.bytes=1073741824
    log.retention.check.interval.ms=300000
    # 性能调优
    num.network.threads=3
    num.io.threads=8
    socket.send.buffer.bytes=102400
    socket.receive.buffer.bytes=102400
    socket.request.max.bytes=104857600
    # 安全
    auto.create.topics.enable=false
    unclean.leader.election.enable=false
    replica.lag.time.max.ms=30000
    # 内置 Topic
    offsets.topic.replication.factor=3
    transaction.state.log.replication.factor=3
    transaction.state.log.min.isr=2
    # Group rebalance
    group.initial.rebalance.delay.ms=3000

---
apiVersion: policy/v1
kind: PodDisruptionBudget
metadata:
  name: kafka-pdb
  namespace: middleware
spec:
  minAvailable: 2
  selector:
    matchLabels:
      app: kafka

---
apiVersion: apps/v1
kind: StatefulSet
metadata:
  name: kafka
  namespace: middleware
  labels:
    app: kafka
spec:
  serviceName: kafka-headless
  replicas: 3
  podManagementPolicy: OrderedReady
  updateStrategy:
    type: RollingUpdate
    rollingUpdate:
      partition: 0
  selector:
    matchLabels:
      app: kafka
  template:
    metadata:
      labels:
        app: kafka
    spec:
      affinity:
        podAntiAffinity:
          requiredDuringSchedulingIgnoredDuringExecution:
            - labelSelector:
                matchLabels:
                  app: kafka
              topologyKey: kubernetes.io/hostname
      terminationGracePeriodSeconds: 180
      containers:
        - name: kafka
          image: bitnami/kafka:3.6
          ports:
            - containerPort: 9092
              name: client
            - containerPort: 9093
              name: controller
            - containerPort: 9999
              name: jmx
          env:
            - name: POD_NAME
              valueFrom:
                fieldRef:
                  fieldPath: metadata.name
            - name: POD_IP
              valueFrom:
                fieldRef:
                  fieldPath: status.podIP
            # 动态生成 broker.id 和 advertised.listeners
            - name: KAFKA_BROKER_ID
              valueFrom:
                fieldRef:
                  fieldPath: metadata.name
            - name: KAFKA_ADVERTISED_LISTENERS
              value: "PLAINTEXT://$(POD_NAME).kafka-headless.middleware.svc.cluster.local:9092"
          command:
            - sh
            - -c
            - |
              # 从 Pod 名提取序号作为 broker.id
              ORDINAL=${POD_NAME##*-}
              export KAFKA_BROKER_ID=${ORDINAL}
              # 动态设置 advertised.listeners
              export KAFKA_ADVERTISED_LISTENERS="PLAINTEXT://${POD_NAME}.kafka-headless.middleware.svc.cluster.local:9092"
              # JVM 参数
              export KAFKA_HEAP_OPTS="-Xms4g -Xmx4g -XX:MetaspaceSize=96m -XX:+UseG1GC -XX:MaxGCPauseMillis=20 -XX:InitiatingHeapOccupancyPercent=35 -XX:G1HeapRegionSize=16M"
              # JMX 监控
              export JMX_PORT=9999
              export KAFKA_JMX_OPTS="-Dcom.sun.management.jmxremote -Dcom.sun.management.jmxremote.authenticate=false -Dcom.sun.management.jmxremote.ssl=false -Dcom.sun.management.jmxremote.host=0.0.0.0 -Djava.rmi.server.hostname=${POD_IP}"
              exec kafka-server-start.sh /config/server.properties
          resources:
            requests:
              cpu: 1
              memory: 4Gi
            limits:
              cpu: 2
              memory: 6Gi
          readinessProbe:
            tcpSocket:
              port: 9092
            initialDelaySeconds: 15
            periodSeconds: 10
            timeoutSeconds: 5
          livenessProbe:
            tcpSocket:
              port: 9092
            initialDelaySeconds: 60
            periodSeconds: 30
            timeoutSeconds: 10
          startupProbe:
            tcpSocket:
              port: 9092
            failureThreshold: 30
            periodSeconds: 10
          lifecycle:
            preStop:                        # 优雅关闭:先停止接受新请求,等待副本同步
              exec:
                command:
                  - sh
                  - -c
                  - |
                    kafka-server-stop.sh
                    # 等待 Controller 将该 Broker 的分区 Leader 转移到其他 Broker
                    sleep 60
          volumeMounts:
            - name: kafka-data
              mountPath: /bitnami/kafka/data
            - name: config
              mountPath: /config/server.properties
              subPath: server.properties
      volumes:
        - name: config
          configMap:
            name: kafka-config
  volumeClaimTemplates:
    - metadata:
        name: kafka-data
      spec:
        accessModes: ["ReadWriteOnce"]
        storageClassName: fast-ssd
        resources:
          requests:
            storage: 100Gi

2.5 KRaft 模式部署(完整 YAML)

---
apiVersion: v1
kind: ConfigMap
metadata:
  name: kafka-kraft-config
  namespace: middleware
data:
  server.properties: |
    # KRaft 模式 - 节点角色
    process.roles=broker,controller
    
    # 节点 ID(启动脚本中动态设置)
    # node.id=1
    
    # Controller Quorum 投票者列表
    controller.quorum.voters=1@kafka-kraft-0.kafka-headless:9093,2@kafka-kraft-1.kafka-headless:9093,3@kafka-kraft-2.kafka-headless:9093
    
    # 监听器
    listeners=PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093
    # advertised.listeners 启动脚本中动态设置
    listener.security.protocol.map=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
    inter.broker.listener.name=PLAINTEXT
    controller.listener.names=CONTROLLER
    
    # 存储目录
    log.dirs=/bitnami/kafka/data
    
    # Topic 默认配置
    num.partitions=3
    default.replication.factor=3
    min.insync.replicas=2
    offsets.topic.replication.factor=3
    transaction.state.log.replication.factor=3
    transaction.state.log.min.isr=2
    
    # 日志保留
    log.retention.hours=168
    log.segment.bytes=1073741824
    
    # 性能
    num.network.threads=3
    num.io.threads=8
    
    # 安全
    auto.create.topics.enable=false
    unclean.leader.election.enable=false

---
apiVersion: apps/v1
kind: StatefulSet
metadata:
  name: kafka-kraft
  namespace: middleware
spec:
  serviceName: kafka-headless
  replicas: 3
  podManagementPolicy: OrderedReady
  selector:
    matchLabels:
      app: kafka-kraft
  template:
    metadata:
      labels:
        app: kafka-kraft
    spec:
      affinity:
        podAntiAffinity:
          requiredDuringSchedulingIgnoredDuringExecution:
            - labelSelector:
                matchLabels:
                  app: kafka-kraft
              topologyKey: kubernetes.io/hostname
      terminationGracePeriodSeconds: 180
      initContainers:
        - name: format-storage
          image: bitnami/kafka:3.6
          command:
            - sh
            - -c
            - |
              ORDINAL=${POD_NAME##*-}
              NODE_ID=$((ORDINAL + 1))
              # 检查是否已格式化(避免重复格式化丢数据)
              if [ -f /bitnami/kafka/data/.kafka-cluster-id ]; then
                EXISTING_ID=$(cat /bitnami/kafka/data/.kafka-cluster-id)
                echo "Storage already formatted with cluster ID: ${EXISTING_ID}"
              else
                # 生成集群 ID(所有节点必须用同一个)
                CLUSTER_ID="${KAFKA_CLUSTER_ID:-$(kafka-storage.sh random-uuid)}"
                echo "${CLUSTER_ID}" > /bitnami/kafka/data/.kafka-cluster-id
                # 临时写入 node.id 到配置
                sed -i "s/^# node.id=.*/node.id=${NODE_ID}/" /config/server.properties
                kafka-storage.sh format \
                  --cluster-id "${CLUSTER_ID}" \
                  --config /config/server.properties
                echo "Storage formatted with cluster ID: ${CLUSTER_ID}, node ID: ${NODE_ID}"
              fi
          env:
            - name: POD_NAME
              valueFrom:
                fieldRef:
                  fieldPath: metadata.name
            - name: KAFKA_CLUSTER_ID
              value: "xtz7N0gGTPaCgkUbk1p7bg"  # 所有 Pod 必须一致!
          volumeMounts:
            - name: kafka-data
              mountPath: /bitnami/kafka/data
            - name: config
              mountPath: /config
      containers:
        - name: kafka
          image: bitnami/kafka:3.6
          ports:
            - containerPort: 9092
              name: client
            - containerPort: 9093
              name: controller
            - containerPort: 9999
              name: jmx
          env:
            - name: POD_NAME
              valueFrom:
                fieldRef:
                  fieldPath: metadata.name
            - name: POD_IP
              valueFrom:
                fieldRef:
                  fieldPath: status.podIP
          command:
            - sh
            - -c
            - |
              ORDINAL=${POD_NAME##*-}
              NODE_ID=$((ORDINAL + 1))
              # 动态设置 node.id 和 advertised.listeners
              sed -i "s/^# node.id=.*/node.id=${NODE_ID}/" /config/server.properties
              export KAFKA_ADVERTISED_LISTENERS="PLAINTEXT://${POD_NAME}.kafka-headless.middleware.svc.cluster.local:9092"
              # JVM 参数
              export KAFKA_HEAP_OPTS="-Xms4g -Xmx4g -XX:MetaspaceSize=96m -XX:+UseG1GC -XX:MaxGCPauseMillis=20"
              # JMX
              export JMX_PORT=9999
              export KAFKA_JMX_OPTS="-Dcom.sun.management.jmxremote -Dcom.sun.management.jmxremote.authenticate=false -Dcom.sun.management.jmxremote.ssl=false -Djava.rmi.server.hostname=${POD_IP}"
              exec kafka-server-start.sh /config/server.properties
          resources:
            requests:
              cpu: 1
              memory: 4Gi
            limits:
              cpu: 2
              memory: 6Gi
          readinessProbe:
            tcpSocket:
              port: 9092
            initialDelaySeconds: 15
            periodSeconds: 10
          livenessProbe:
            tcpSocket:
              port: 9092
            initialDelaySeconds: 60
            periodSeconds: 30
          startupProbe:
            tcpSocket:
              port: 9092
            failureThreshold: 30
            periodSeconds: 10
          lifecycle:
            preStop:
              exec:
                command:
                  - sh
                  - -c
                  - "kafka-server-stop.sh && sleep 60"
          volumeMounts:
            - name: kafka-data
              mountPath: /bitnami/kafka/data
            - name: config
              mountPath: /config
      volumes:
        - name: config
          configMap:
            name: kafka-kraft-config
  volumeClaimTemplates:
    - metadata:
        name: kafka-data
      spec:
        accessModes: ["ReadWriteOnce"]
        storageClassName: fast-ssd
        resources:
          requests:
            storage: 100Gi

KRaft 关键差异总结

  • 不需要 zookeeper.connect,改用 controller.quorum.voters
  • 每个 Broker 同时是 broker,controller(或单独拆分角色)
  • 首次启动需 kafka-storage.sh format 格式化存储目录
  • process.roles 决定节点角色(broker,controllerbrokercontroller
  • node.id 替代 broker.id(KRaft 模式下统一用 node.id
  • 所有节点必须使用同一个 cluster.id

2.6 Topic 管理常用命令

# 创建 Topic
kafka-topics.sh --bootstrap-server kafka-0.kafka-headless:9092 \
  --create --topic order-events \
  --partitions 12 \
  --replication-factor 3 \
  --config min.insync.replicas=2 \
  --config retention.ms=604800000

# 查看 Topic 列表
kafka-topics.sh --bootstrap-server kafka-0.kafka-headless:9092 --list

# 查看 Topic 详情(分区/副本/ISR)
kafka-topics.sh --bootstrap-server kafka-0.kafka-headless:9092 \
  --describe --topic order-events

# 查看所有 Under-replicated Partitions
kafka-topics.sh --bootstrap-server kafka-0.kafka-headless:9092 \
  --describe --under-replicated-partitions

# 增加 Partition(只能增加不能减少)
kafka-topics.sh --bootstrap-server kafka-0.kafka-headless:9092 \
  --alter --topic order-events --partitions 24

# 修改 Topic 配置
kafka-configs.sh --bootstrap-server kafka-0.kafka-headless:9092 \
  --alter --topic order-events \
  --add-config retention.ms=2592000000   # 30 天

# 删除 Topic(需开启 delete.topic.enable=true)
kafka-topics.sh --bootstrap-server kafka-0.kafka-headless:9092 \
  --delete --topic order-events

# Producer 生产消息
kafka-console-producer.sh --bootstrap-server kafka-0.kafka-headless:9092 \
  --topic order-events \
  --producer-property acks=all \
  --producer-property compression.type=lz4

# Consumer 消费消息
kafka-console-consumer.sh --bootstrap-server kafka-0.kafka-headless:9092 \
  --topic order-events \
  --from-beginning \
  --group order-consumer-group

# 查看 Consumer Group 偏移量 & Lag
kafka-consumer-groups.sh --bootstrap-server kafka-0.kafka-headless:9092 \
  --describe --group order-consumer-group

# 重置 Consumer Group 偏移量到最早
kafka-consumer-groups.sh --bootstrap-server kafka-0.kafka-headless:9092 \
  --group order-consumer-group \
  --reset-offsets --to-earliest \
  --topic order-events --execute

# 副本重新分配(扩容后重新平衡)
kafka-reassign-partitions.sh --bootstrap-server kafka-0.kafka-headless:9092 \
  --reassignment-json-file reassignment.json \
  --execute --throttle 50000000   # 限速 50MB/s

# 首选 Leader 选举(让 ISR 中的首选副本成为 Leader)
kafka-leader-election.sh --bootstrap-server kafka-0.kafka-headless:9092 \
  --election-type preferred \
  --all-topic-partitions

三、常见部署问题排障

3.1 Kafka Broker 常见问题

问题典型现象根因解决方案
Broker 启动失败进程退出,日志报 FATALbroker.id 冲突;log.dirs 权限不对;端口占用检查 broker.id 唯一性;chown -R kafka:kafka /data/kafka-logs;检查端口
advertised.listeners 配置错误Producer/Consumer 连不上 Broker内外网 IP 混淆;用了 0.0.0.0明确设置 advertised.listeners=PLAINTEXT://<可达域名>:9092不能0.0.0.0
Under-replicated Partitionskafka-topics --describe 显示 ISR < 副本数副本 Broker 宕机;网络抖动;磁盘 IO 瓶颈;GC 停顿kubectl describe pod 排查 Broker 状态;检查 GC 日志;kafka-reassign-partitions 重新分配
Leader 选举风暴频繁 Leader 切换,消费 lag 飙升ZK session 超时;GC 停顿过长;网络分区调大 session.timeout.ms;优化 JVM GC;排查网络
磁盘空间不足Broker 拒写,报 Disk usage 警告log.retention.hours 未生效;日志清理线程卡住检查 log.retention.bytes + log.retention.hours;手动 kafka-delete-records;扩容 PV
Producer 发送超时TimeoutExceptionBroker 不可达;acks=all 且 ISR 不足;消息过大检查 advertised.listeners;确认 min.insync.replicas;检查 max.request.size
Consumer rebalance 频繁消费间歇性暂停session.timeout.ms 太短;max.poll.interval.ms 不够;消费者处理太慢调大 session.timeout.msmax.poll.interval.ms;优化消费逻辑
消息丢失消费端拿不到部分消息acks=1acks=0unclean.leader.election=true;Consumer auto.offset.reset=latest设置 acks=all + min.insync.replicas=2;禁止非 ISR 选举;Consumer 用 earliest
消息重复消费端收到重复消息Producer 重试;Consumer 未手动提交偏移量就崩溃Producer 开启 enable.idempotence=true;Consumer 关闭自动提交,处理完成后手动提交
Controller 故障元数据操作卡住,无法创建/删除 TopicController Broker 宕机;ZK 不可达检查 Controller 状态:kafka-metadata-quorum.sh (KRaft) 或 ZK /controller 节点

3.2 Zookeeper 常见问题

问题典型现象根因解决方案
无法组成集群日志反复 LEADER Election 但选不出 Leadermyid 文件缺失/重复/值不对;server.x 配置错误;DNS 解析失败确认 /data/zookeeper/data/myid 唯一且与 zoo.cfgserver.x 对应;检查 Headless Service + publishNotReadyAddresses
脑裂 (Split-Brain)出现多个 Leader网络分区;initLimit/syncLimit 过大确保奇数节点(3/5/7);调小 syncLimit;排查网络分区
磁盘满导致 snapshot 失败ZK 进程 OOM 或 hung事务日志和快照混放同一磁盘;autopurge 未配事务日志(dataLogDir) 和快照(dataDir) 分盘存放;启用 autopurge.snapRetainCount=3; purgeInterval=1
ZK 抖动 → Kafka ISR 缩小Kafka 日志频繁 Shrinking ISRZK GC 停顿;磁盘 IO 争抢;网络延迟ZK 独立部署,不与 Kafka 共享节点;JVM 堆建议 2-4G;SSD 独占
Watch 堆积ZK 响应变慢,Kafka 元数据操作延迟客户端 watch 未清理;session 过期后 watch 重连风暴升级 ZK 3.6+;监控 zk_watch_count 指标
standalone 模式误用单节点 ZK,无容错误部署单节点 ZK生产环境必须 3/5/7 节点
事务日志写入慢ZK 写操作延迟高事务日志目录用了 HDD 或网络盘dataLogDir 必须用本地 SSD,不能用 NFS/Ceph RBD
session 频繁过期Kafka Broker 日志报 Session expiredZK 心跳超时;GC 停顿调大 tickTimesyncLimit;优化 ZK JVM GC

3.3 K8s 环境特有问题

问题根因解决方案
Pod 重启后数据丢失用了 emptyDir 或未正确绑定 PVC必须用 StatefulSet + PVCvolumeClaimTemplates 持久化 log.dirs / dataDir
Headless Service DNS 解析慢/失败CoreDNS 缓存;Pod 尚未 Ready 就被解析配合 publishNotReadyAddresses: true;检查 CoreDNS ready 插件配置
myid 在 StatefulSet 中如何动态生成每个 ZK Pod 需要不同 myid用 initContainer 从 hostname 提取序号:echo $((${HOSTNAME##*-}+1)) > /data/zookeeper/myid
Graceful Shutdown 数据损坏terminationGracePeriodSeconds 太短,Broker 未完成日志刷盘设置 terminationGracePeriodSeconds: 180+;preStop hook 执行 kafka-server-stop.sh && sleep 60
PV/PVC 绑定错 PodStorageClass volumeBindingMode 配置不当WaitForFirstConsumer + 节点亲和性确保 PV 绑定到正确节点
ZK 集群滚动更新导致不可用滚更时多个 ZK 同时重启podManagementPolicy: OrderedReady + 间隔等待;用 PodDisruptionBudget 保证至少 (N-1)/2+1 可用
Kafka 滚更期间消费中断Leader 副本所在 Pod 被先删preStop hook 等待 Controller 完成 Leader 迁移;设置足够长的 terminationGracePeriodSeconds
JMX 监控端口不通JMX 绑定了 127.0.0.1 而非 0.0.0.0设置 -Dcom.sun.management.jmxremote.host=0.0.0.0-Djava.rmi.server.hostname=<PodIP>
Pod 无法调度podAntiAffinity 要求分散到不同节点,但节点数不足确保节点数 ≥ Kafka/ZK 副本数;或改用 preferredDuringScheduling 软策略
StatefulSet 滚更卡住OrderedReady 模式下某个 Pod 健康检查不通过先修复该 Pod 问题;或临时调大 failureThreshold;不要手动删除 Pod

3.4 网络问题排障

问题排查命令说明
Pod 间网络不通kubectl exec <pod> -- ping <target-pod-ip>检查 CNI 插件状态
DNS 解析失败kubectl exec <pod> -- nslookup kafka-0.kafka-headless检查 CoreDNS
Headless Service 不返回地址kubectl get svc zk-headless 确认 clusterIP: None确认 publishNotReadyAddresses: true
外部访问 Kafka 失败telnet <nodeport-ip> <port>advertised.listeners 必须用外部可达地址;考虑 NodePort / LoadBalancer / Ingress
跨命名空间访问使用完整 FQDN:kafka-0.kafka-headless.middleware.svc.cluster.local跨命名空间必须用全限定域名

3.5 性能问题排查工具

# 1. Kafka 性能测试 - 生产者吞吐量
kafka-producer-perf-test.sh --bootstrap-server kafka-0.kafka-headless:9092 \
  --topic perf-test --num-records 100000 --record-size 1024 \
  --throughput -1 --producer-props acks=all compression.type=lz4

# 2. Kafka 性能测试 - 消费者吞吐量
kafka-consumer-perf-test.sh --bootstrap-server kafka-0.kafka-headless:9092 \
  --topic perf-test --messages 100000 --threads 3

# 3. 查看 Broker 的请求队列积压
kubectl exec kafka-0 -- jcmd <pid> Thread.print | grep -A5 "RequestHandler"

# 4. GC 日志分析
# 确保 JVM 参数中包含:
# -Xlog:gc*:file=/var/log/kafka/gc.log:time,uptime,level,tags

# 5. 磁盘 IO 延迟检查
kubectl exec kafka-0 -- sh -c 'dd if=/dev/zero of=/bitnami/kafka/data/test bs=1M count=100 oflag=dsync'

# 6. Kafka Metrics JMX 查询
kubectl exec kafka-0 -- kafka-run-class.sh kafka.tools.JmxTool \
  --object-name kafka.server:type=BrokerTopicMetrics,name=MessagesInPerSec \
  --jmx-url service:jmx:rmi:///jndi/rmi://localhost:9999/jmxrmi \
  --one-time true

四、生产环境最佳实践

4.1 资源规划

按集群规模分级

集群规模ZK 资源Kafka 资源ZK 磁盘Kafka 磁盘副本数
小型 (< 1k msg/s)0.5C / 2G2C / 4G20G SSD × 2100G SSDZK=3, Kafka=3
中型 (1k-10k msg/s)1C / 4G4C / 8G50G SSD × 2500G SSDZK=3, Kafka=3-5
大型 (10k-100k msg/s)2C / 4G8C / 16G100G SSD × 21T NVMe × 2ZK=5, Kafka=5-7
超大型 (> 100k msg/s)4C / 8G16C / 32G200G SSD × 22T NVMe × 4ZK=5-7, Kafka=7+

分区数规划

分区数 = max(目标吞吐量 / 单分区吞吐量, 消费者并发数)

示例:
  目标吞吐量: 100 MB/s
  单分区吞吐量: 10 MB/s (生产端)
  消费者并发: 20

  分区数 = max(100/10, 20) = 20 个分区

注意:分区数只能增加不能减少。初始规划宁可多设也不要后续频繁扩容。每个分区都会消耗 Broker 的内存(用于索引)和 ZK 的 watch(ZK 模式下)。

4.2 JVM 调优详解

# ================= Kafka Broker JVM =================
KAFKA_HEAP_OPTS="
  -Xms4g                           # 初始堆 = 最大堆,避免动态扩缩
  -Xmx4g
  -XX:MetaspaceSize=96m            # 元空间初始大小
  -XX:MaxMetaspaceSize=256m        # 元空间上限
  -XX:+UseG1GC                     # G1 垃圾收集器(适合大堆)
  -XX:MaxGCPauseMillis=20          # GC 停顿目标 20ms
  -XX:InitiatingHeapOccupancyPercent=35  # 堆使用率 35% 时触发 GC
  -XX:G1HeapRegionSize=16M         # G1 区域大小(> 2G 堆建议 16M)
  -XX:+ExplicitGCInvokesConcurrent # System.gc() 走并发 GC
  -XX:+HeapDumpOnOutOfMemoryError  # OOM 时生成堆转储
  -XX:HeapDumpPath=/var/log/kafka/heapdump
  -Xlog:gc*:file=/var/log/kafka/gc.log:time,uptime,level,tags  # GC 日志
"

# ================= Zookeeper JVM =================
SERVER_JVMFLAGS="
  -Xms2g                           # ZK 堆不要超过 4G
  -Xmx2g
  -XX:MetaspaceSize=128m
  -XX:+UseG1GC
  -XX:MaxGCPauseMillis=50          # ZK 对 GC 停顿更敏感
  -XX:InitiatingHeapOccupancyPercent=40
  -XX:+ExplicitGCInvokesConcurrent
  -Djute.maxbuffer=1048576         # ZK 节点数据上限(1MB)
  -XX:+HeapDumpOnOutOfMemoryError
  -XX:HeapDumpPath=/data/zookeeper/heapdump
"

JVM 调优要点

要点说明
Xms = Xmx避免堆动态扩缩引起 GC 停顿
G1GC适合 4G+ 堆,可控停顿时间
MaxGCPauseMillisKafka 建议 20ms,ZK 建议 50ms
InitiatingHeapOccupancyPercent35% 就开始并发标记,避免 Full GC
GC 日志必须开启,用于排查 ISR 缩小问题
堆转储OOM 时自动生成,用于内存泄漏分析

4.3 磁盘 IO 调优

# 1. 挂载参数优化(ext4/xfs)
# /etc/fstab 中加上 noatime
/dev/sdb /data/kafka xfs noatime,nodiratime 0 0

# 2. 文件描述符限制
# /etc/security/limits.conf
kafka soft nofile 100000
kafka hard nofile 100000

# 3. 内核参数调优
sysctl -w vm.dirty_ratio=80              # 脏页占比 80% 才开始同步写入
sysctl -w vm.dirty_background_ratio=5    # 脏页占比 5% 开始后台写入
sysctl -w vm.swappiness=1                # 尽量不使用 swap
sysctl -w vm.max_map_count=262144

# 4. Kafka 落盘策略
# 默认依赖 OS page cache,不强制 fsync
log.flush.interval.messages=9223372036854775807  # 实际上不限制
log.flush.interval.ms=null                         # 不强制刷盘
# 由 OS 的 pdflush/flush 线程管理落盘(性能最优)
# 数据可靠性由多副本保证,而非单机 fsync

关键:Kafka 的性能依赖于 OS Page Cache。不要强制 fsync,那样会严重降低吞吐量。数据可靠性通过多副本(replication.factor=3 + min.insync.replicas=2)保证。

4.4 滚动更新策略

# StatefulSet 滚动更新配置
updateStrategy:
  type: RollingUpdate
  rollingUpdate:
    partition: 0    # 从最高 ordinal 开始替换

# 确保 PDB 限制
# ZK: minAvailable: 2 (3 节点集群)
# Kafka: minAvailable: 2 (3 节点集群)

# preStop hook 确保优雅退出
lifecycle:
  preStop:
    exec:
      command:
        - sh
        - -c
        - |
          # Kafka: 停止 Broker,等待 Controller 迁移 Leader
          kafka-server-stop.sh
          sleep 60
          # ZK: 停止 ZK,等待 Leader 重新选举
          # zkServer.sh stop
          # sleep 30

滚动更新注意事项

  1. podManagementPolicy: OrderedReady — 有序逐个更新
  2. 每次更新一个 Pod,等待其 Ready 后再更新下一个
  3. preStop hook 中 sleep 给 Controller 足够时间做 Leader 迁移
  4. terminationGracePeriodSeconds 必须 > preStop 耗时
  5. 大集群(多分区)更新时考虑分区 Leader 迁移的耗时,可能需要 5-10 分钟

4.5 StorageClass 配置建议

# 推荐:使用 local-storage 或带高 IOPS 的块存储
apiVersion: storage.k8s.io/v1
kind: StorageClass
metadata:
  name: fast-ssd
provisioner: kubernetes.io/no-provisioner    # 本地存储(性能最优)
# 或使用云盘:
# provisioner: disk.csi.alibabacloud.com      # 阿里云 ESSD
# provisioner: kubernetes.io/aws-ebs          # AWS gp3
volumeBindingMode: WaitForFirstConsumer       # 延迟绑定,确保 PV 和 Pod 在同节点
reclaimPolicy: Retain                         # 删除 PVC 后保留数据
allowVolumeExpansion: true                    # 允许在线扩容
parameters:
  type: essd_pl1                              # 阿里云 ESSD PL1(性能 + 成本平衡)
  performanceLevel: PL1

关键:Kafka 的 log.dirs 和 ZK 的 dataLogDir 不要用网络存储(NFS、CephFS 等),IO 延迟不稳定会导致 ISR 频繁抖动。优先用 local-storage 或云盘(块存储)。


五、安全加固:SASL + TLS + ACL

5.1 SASL/SCRAM 认证完整配置

步骤 1:创建 SCRAM 凭据

# 在 Broker 中创建用户(admin 和 app-user)
kafka-configs.sh --bootstrap-server kafka-0.kafka-headless:9092 \
  --alter --add-config 'SCRAM-SHA-512=[password=admin-secret],SCRAM-SHA-256=[password=admin-secret]' \
  --entity-type users --entity-name admin

kafka-configs.sh --bootstrap-server kafka-0.kafka-headless:9092 \
  --alter --add-config 'SCRAM-SHA-512=[password=app-user-secret]' \
  --entity-type users --entity-name app-user

# 验证
kafka-configs.sh --bootstrap-server kafka-0.kafka-headless:9092 \
  --describe --entity-type users --entity-name admin

步骤 2:Broker 端配置

# server.properties
listeners=SASL_SSL://0.0.0.0:9093,SASL_PLAINTEXT://0.0.0.0:9094
advertised.listeners=SASL_SSL://<host>:9093,SASL_PLAINTEXT://<host>:9094
sasl.enabled.mechanisms=SCRAM-SHA-512,SCRAM-SHA-256
sasl.mechanism.inter.broker.protocol=SCRAM-SHA-512
listener.security.protocol.map=SASL_SSL:SASL_SSL,SASL_PLAINTEXT:SASL_PLAINTEXT

# Broker 间通信的 JAAS 配置
listener.name.sasl_ssl.scram-sha-512.sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required username="admin" password="admin-secret";
listener.name.sasl_plaintext.scram-sha-512.sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required username="admin" password="admin-secret";

步骤 3:K8s Secret 存储 JAAS 配置

apiVersion: v1
kind: Secret
metadata:
  name: kafka-sasl-secret
  namespace: middleware
type: Opaque
stringData:
  kafka_client_jaas.conf: |
    KafkaClient {
      org.apache.kafka.common.security.scram.ScramLoginModule required
      username="admin"
      password="admin-secret";
    };

步骤 4:客户端配置

# client.properties (Producer/Consumer)
security.protocol=SASL_SSL
sasl.mechanism=SCRAM-SHA-512
sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required \
  username="app-user" \
  password="app-user-secret";

# 使用时指定
kafka-console-producer.sh --bootstrap-server kafka:9093 \
  --topic order-events \
  --producer.config client.properties

5.2 TLS 加密配置

证书生成流程

# 1. 生成 CA 密钥和证书
openssl req -new -x509 -keyout ca-key -out ca-cert -days 365 \
  -subj "/CN=Kafka-CA" -passout pass:ca-password

# 2. 生成 Broker 密钥库
keytool -keystore kafka.server.keystore.jks -alias kafka -validity 365 \
  -genkey -keyalg RSA -dname "CN=kafka-0.kafka-headless.middleware.svc.cluster.local" \
  -storepass keystore-password -keypass keystore-password

# 3. 生成证书签名请求 (CSR)
keytool -keystore kafka.server.keystore.jks -alias kafka -certreq -file cert-file \
  -storepass keystore-password

# 4. 用 CA 签名
openssl x509 -req -CA ca-cert -CAkey ca-key -in cert-file -out cert-signed \
  -days 365 -CAcreateserial -passin pass:ca-password

# 5. 导入 CA 和签名证书到密钥库
keytool -keystore kafka.server.keystore.jks -alias CARoot -import -file ca-cert \
  -storepass keystore-password -noprompt
keytool -keystore kafka.server.keystore.jks -alias kafka -import -file cert-signed \
  -storepass keystore-password -noprompt

# 6. 生成信任库(Truststore)
keytool -keystore kafka.server.truststore.jks -alias CARoot -import -file ca-cert \
  -storepass truststore-password -noprompt

Broker TLS 配置

# server.properties
ssl.keystore.location=/etc/kafka/secrets/kafka.server.keystore.jks
ssl.keystore.password=keystore-password
ssl.key.password=keystore-password
ssl.truststore.location=/etc/kafka/secrets/kafka.server.truststore.jks
ssl.truststore.password=truststore-password
ssl.client.auth=required               # 要求客户端也提供证书
ssl.endpoint.identification.algorithm=  # 禁用主机名验证(Pod 域名动态变化时需要)

5.3 ACL 权限管理

# 创建 Secret 存储 ACL 配置
# 确认 broker 配置:
# authorizer.class.name=kafka.security.authorizer.AclAuthorizer  (旧版)
# authorizer.class.name=org.apache.kafka.metadata.authorizer.StandardAuthorizer  (KRaft)
# allow.everyone.if.no.acl.found=false
# super.users=User:admin

# 给 app-user 对 order-events Topic 的读写权限
kafka-acls.sh --bootstrap-server kafka:9093 \
  --command-config admin.properties \
  --add \
  --allow-principal User:app-user \
  --operation Read --operation Write \
  --topic order-events

# 给 app-user 对 consumer group 的权限
kafka-acls.sh --bootstrap-server kafka:9093 \
  --command-config admin.properties \
  --add \
  --allow-principal User:app-user \
  --operation Read \
  --group order-consumer-group

# 查看所有 ACL
kafka-acls.sh --bootstrap-server kafka:9093 \
  --command-config admin.properties \
  --list

# 删除权限
kafka-acls.sh --bootstrap-server kafka:9093 \
  --command-config admin.properties \
  --remove \
  --allow-principal User:app-user \
  --operation Write \
  --topic order-events

六、监控集成:Prometheus + Grafana

6.1 JMX Exporter 部署

# ConfigMap - JMX Exporter 配置
apiVersion: v1
kind: ConfigMap
metadata:
  name: kafka-jmx-config
  namespace: middleware
data:
  kafka-jmx.yml: |
    lowercaseOutputName: true
    lowercaseOutputLabelNames: true
    rules:
      # Broker 级指标
      - pattern: 'kafka.server<type=BrokerTopicMetrics, name=(MessagesInPerSec|BytesInPerSec|BytesOutPerSec)><>(Count|OneMinuteRate)'
        name: kafka_server_broker_topic_metrics_$1
        type: GAUGE
      # Under-replicated partitions
      - pattern: 'kafka.server<type=ReplicaManager, name=UnderReplicatedPartitions><>Value'
        name: kafka_server_under_replicated_partitions
        type: GAUGE
      # Offline partitions
      - pattern: 'kafka.controller<type=KafkaController, name=OfflinePartitionsCount><>Value'
        name: kafka_controller_offline_partitions_count
        type: GAUGE
      # Active controller count
      - pattern: 'kafka.controller<type=KafkaController, name=ActiveControllerCount><>Value'
        name: kafka_controller_active_controller_count
        type: GAUGE
      # ISR
      - pattern: 'kafka.server<type=ReplicaManager, name=(UnderReplicatedPartitions|UnderMinIsrPartitionCount)><>Value'
        name: kafka_server_replica_manager_$1
        type: GAUGE
      # Consumer Group lag
      - pattern: 'kafka.consumer<type=ConsumerFetcherManager, name=MaxLag, clientId=(.+)><>Value'
        name: kafka_consumer_max_lag
        type: GAUGE

6.2 Prometheus ServiceMonitor

apiVersion: monitoring.coreos.com/v1
kind: ServiceMonitor
metadata:
  name: kafka-jmx-exporter
  namespace: middleware
  labels:
    release: prometheus          # 匹配 Prometheus Operator 的 label
spec:
  selector:
    matchLabels:
      app: kafka
  endpoints:
    - port: jmx-metrics          # JMX Exporter 暴露的端口
      interval: 30s
      path: /metrics
      scrapeTimeout: 10s
  jobLabel: app
---
# ZK ServiceMonitor
apiVersion: monitoring.coreos.com/v1
kind: ServiceMonitor
metadata:
  name: zk-jmx-exporter
  namespace: middleware
spec:
  selector:
    matchLabels:
      app: zookeeper
  endpoints:
    - port: jmx-metrics
      interval: 30s
      path: /metrics

6.3 关键告警规则

# PrometheusRule
apiVersion: monitoring.coreos.com/v1
kind: PrometheusRule
metadata:
  name: kafka-alerts
  namespace: middleware
spec:
  groups:
    - name: kafka-broker
      rules:
        - alert: KafkaUnderReplicatedPartitions
          expr: kafka_server_under_replicated_partitions > 0
          for: 5m
          labels:
            severity: warning
          annotations:
            summary: "Kafka under-replicated partitions ({{ $value }})"
            description: "Broker {{ $labels.instance }} has under-replicated partitions for 5 minutes"

        - alert: KafkaOfflinePartitions
          expr: kafka_controller_offline_partitions_count > 0
          for: 1m
          labels:
            severity: critical
          annotations:
            summary: "Kafka offline partitions ({{ $value }})"
            description: "There are offline partitions - immediate attention required"

        - alert: KafkaNoActiveController
          expr: sum(kafka_controller_active_controller_count) == 0
          for: 2m
          labels:
            severity: critical
          annotations:
            summary: "Kafka has no active controller"
            description: "No active controller in the cluster"

        - alert: KafkaBrokerDown
          expr: up{job="kafka"} == 0
          for: 2m
          labels:
            severity: critical
          annotations:
            summary: "Kafka broker {{ $labels.instance }} is down"

    - name: zookeeper
      rules:
        - alert: ZookeeperNodeDown
          expr: up{job="zookeeper"} == 0
          for: 1m
          labels:
            severity: critical
          annotations:
            summary: "Zookeeper node {{ $labels.instance }} is down"

        - alert: ZookeeperHighLatency
          expr: zookeeper_avg_latency > 10
          for: 5m
          labels:
            severity: warning
          annotations:
            summary: "Zookeeper average latency > 10ms ({{ $value }}ms)"

6.4 推荐 Grafana Dashboard

Dashboard ID名称说明
721Kafka OverviewKafka 集群总览(吞吐量、分区、ISR)
7589Kafka Exporter Overview基于 kafka-exporter 的监控面板
10465Kafka Lag ExporterConsumer Lag 监控
10991Strimzi KafkaStrimzi Operator 管理的 Kafka 面板
9236ZookeeperZK 集群监控(延迟、watch、连接数)

七、快速排查 Checklist

7.1 集群状态检查

# 1. Pod 状态 & 节点分布
kubectl get pods -n middleware -o wide

# 2. PVC 绑定状态
kubectl get pvc -n middleware

# 3. Service 状态
kubectl get svc -n middleware

# 4. 事件排查
kubectl get events -n middleware --sort-by='.lastTimestamp' | tail -20

7.2 Zookeeper 排查

# 5. ZK 进程状态
kubectl exec zk-0 -n middleware -- sh -c 'echo srvr | nc localhost 2181'
# 输出示例:
# Zookeeper version: 3.8.1
# Mode: leader / follower
# Zxid: 0x100000003
# Connections: 2
# Outstanding: 0
# Node count: 123

# 6. ZK 监控指标
kubectl exec zk-0 -n middleware -- sh -c 'echo mntr | nc localhost 2181'

# 7. ZK 健康检查
kubectl exec zk-0 -n middleware -- sh -c 'echo ruok | nc localhost 2181'
# 预期返回: imok

# 8. 查看 ZK 中的 Kafka 元数据
kubectl exec zk-0 -n middleware -- zkCli.sh -server localhost:2181 \
  ls /brokers/ids
kubectl exec zk-0 -n middleware -- zkCli.sh -server localhost:2181 \
  ls /brokers/topics
kubectl exec zk-0 -n middleware -- zkCli.sh -server localhost:2181 \
  get /controller

# 9. ZK 日志
kubectl logs zk-0 -n middleware --tail=100

7.3 Kafka 排查

# 10. Broker 日志
kubectl logs kafka-0 -n middleware --tail=100

# 11. Topic 列表
kubectl exec kafka-0 -n middleware -- \
  kafka-topics.sh --bootstrap-server kafka-0.kafka-headless:9092 --list

# 12. Topic 详情(分区/副本/ISR)
kubectl exec kafka-0 -n middleware -- \
  kafka-topics.sh --bootstrap-server kafka-0.kafka-headless:9092 --describe

# 13. 查看 Under-replicated Partitions
kubectl exec kafka-0 -n middleware -- \
  kafka-topics.sh --bootstrap-server kafka-0.kafka-headless:9092 \
  --describe --under-replicated-partitions

# 14. Consumer Group Lag
kubectl exec kafka-0 -n middleware -- \
  kafka-consumer-groups.sh --bootstrap-server kafka-0.kafka-headless:9092 \
  --describe --group <group-id>

# 15. 所有 Consumer Group
kubectl exec kafka-0 -n middleware -- \
  kafka-consumer-groups.sh --bootstrap-server kafka-0.kafka-headless:9092 --list

# 16. Broker 配置查看
kubectl exec kafka-0 -n middleware -- \
  kafka-configs.sh --bootstrap-server kafka-0.kafka-headless:9092 \
  --describe --entity-type brokers --entity-default

# 17. KRaft 模式 - 检查 controller quorum 状态
kubectl exec kafka-kraft-0 -n middleware -- \
  kafka-metadata-quorum.sh --bootstrap-server kafka-kraft-0.kafka-headless:9092 \
  describe --status

# 18. 检查 Broker 是否在集群中注册
kubectl exec kafka-0 -n middleware -- \
  kafka-broker-api-versions.sh --bootstrap-server kafka-0.kafka-headless:9092

7.4 磁盘 & 性能排查

# 19. 节点磁盘使用率
kubectl get nodes -o json | jq '.items[].status.allocatable'

# 20. Pod 磁盘使用
kubectl exec kafka-0 -n middleware -- df -h /bitnami/kafka/data

# 21. Kafka 日志目录大小
kubectl exec kafka-0 -n middleware -- du -sh /bitnami/kafka/data/*

# 22. IO 性能测试
kubectl exec kafka-0 -n middleware -- \
  sh -c 'dd if=/dev/zero of=/bitnami/kafka/data/test bs=1M count=100 oflag=dsync'

# 23. Kafka Producer 性能测试
kubectl exec kafka-0 -n middleware -- \
  kafka-producer-perf-test.sh \
  --bootstrap-server kafka-0.kafka-headless:9092 \
  --topic perf-test --num-records 100000 --record-size 1024 --throughput -1

八、ZK → KRaft 迁移路径

迁移流程

阶段 1: 准备                  阶段 2: 部署 KRaft Controller    阶段 3: Broker 切换 KRaft
┌──────────────────┐         ┌──────────────────┐             ┌──────────────────┐
│ ZK + Kafka       │  ──→    │ ZK + Kafka       │  ──→        │ Kafka KRaft      │
│ (ZK 模式)        │         │ + KRaft Controller│             │ (KRaft only)     │
│                  │         │ (双元数据源)      │             │                  │
│ - 升级到 3.4+    │         │ - Controller 从  │             │ - 逐个 Broker    │
│ - 备份 ZK 数据   │         │   ZK 同步元数据  │             │   重启为 KRaft   │
│ - 验证集群健康   │         │ - 验证元数据一致 │             │ - 下线 ZK 集群   │
└──────────────────┘         └──────────────────┘             └──────────────────┘

详细步骤

阶段操作命令/步骤风险回滚方案
1. 准备升级 Kafka 到 3.4+滚动升级 Kafka Broker回退到旧版本镜像
1. 准备备份 ZK 数据zkCli.sh save <path> + 快照备份-
1. 准备备份 Kafka log.dirskubectl exec kafka-N -- tar czf /backup/kafka-N.tar.gz /bitnami/kafka/data-
1. 准备验证集群健康kafka-topics --describe 检查无 under-replicated--
2. 部署 Controller部署 KRaft Controller 节点部署 3 个 process.roles=controller 的节点删除 Controller Pod
2. 元数据迁移Controller 从 ZK 同步元数据kafka-storage.sh format --cluster-id <id> + 自动同步从 ZK 备份恢复
2. 验证验证元数据一致对比 ZK 和 KRaft 中的 Topic/Partition 信息--
3. 切换 Broker逐个 Broker 重启为 KRaft修改 process.roles=broker,controller + 移除 zookeeper.connect恢复 ZK 模式配置
3. 验证每个 Broker 切换后验证kafka-topics --describe + 生产消费测试--
4. 下线 ZK删除 ZK StatefulSetkubectl delete statefulset zk -n middleware重新部署 ZK
4. 清理删除 ZK PVC(确认后)kubectl delete pvc zk-data-zk-0 ...不可逆从备份恢复

迁移前必须

  1. 完整备份 ZK 数据 + Kafka log.dirs
  2. 在测试环境完整走一遍迁移流程
  3. 准备回滚方案
  4. 选择低峰期操作
  5. 确保有 3 个以上 Controller 节点保证 Raft 多数派

KRaft 迁移验证命令

# 验证 KRaft Controller Quorum 状态
kafka-metadata-quorum.sh --bootstrap-server <broker>:9092 describe --status

# 输出示例:
# ClusterId:              xtz7N0gGTPaCgkUbk1p7bg
# LeaderId:               1
# LeaderEpoch:            15
# HighWatermark:          1234567
# MaxFollowerLag:         0
# MaxFollowerLagTimeMs:   0
# CurrentVoters:          1,2,3
# CurrentObservers:       []

# 验证 Topic 元数据完整
kafka-topics.sh --bootstrap-server <broker>:9092 --describe

# 验证生产消费正常
kafka-console-producer.sh --bootstrap-server <broker>:9092 --topic verify-test
kafka-console-consumer.sh --bootstrap-server <broker>:9092 --topic verify-test --from-beginning

九、Strimzi Operator 方案简介

如果不想手写 StatefulSet YAML,可以使用 Strimzi Kafka Operator 在 K8s 上管理 Kafka 集群。

核心优势

优势说明
声明式管理一个 Kafka CR 定义整个集群(Broker + ZK + 监控)
自动运维滚动升级、扩缩容、证书轮转自动化
内置监控自动部署 JMX Exporter + Prometheus ServiceMonitor
安全加固自动生成和轮转 TLS 证书
Topic 管理通过 KafkaTopic CR 声明式管理 Topic
User 管理通过 KafkaUser CR 管理用户和 ACL

部署示例

# Strimzi Kafka CR
apiVersion: kafka.strimzi.io/v1beta2
kind: Kafka
metadata:
  name: my-cluster
  namespace: middleware
spec:
  kafka:
    version: 3.6.0
    replicas: 3
    listeners:
      - name: plain
        port: 9092
        type: internal
        tls: false
      - name: tls
        port: 9093
        type: internal
        tls: true
    config:
      offsets.topic.replication.factor: 3
      transaction.state.log.replication.factor: 3
      transaction.state.log.min.isr: 2
      default.replication.factor: 3
      min.insync.replicas: 2
      auto.create.topics.enable: "false"
    storage:
      type: jbod
      volumes:
        - id: 0
          type: persistent-claim
          size: 100Gi
          class: fast-ssd
          deleteClaim: false
    template:
      pod:
        affinity:
          podAntiAffinity:
            requiredDuringSchedulingIgnoredDuringExecution:
              - labelSelector:
                  matchLabels:
                    app.kubernetes.io/name: kafka
                topologyKey: kubernetes.io/hostname
  zookeeper:
    replicas: 3
    storage:
      type: persistent-claim
      size: 20Gi
      class: fast-ssd
      deleteClaim: false
  entityOperator:
    topicOperator: {}
    userOperator: {}
  cruiseControl: {}              # 自动集群平衡

Topic 管理示例

apiVersion: kafka.strimzi.io/v1beta2
kind: KafkaTopic
metadata:
  name: order-events
  namespace: middleware
  labels:
    strimzi.io/cluster: my-cluster
spec:
  partitions: 12
  replicas: 3
  config:
    retention.ms: 604800000
    min.insync.replicas: 2
    compression.type: lz4

User 管理示例

apiVersion: kafka.strimzi.io/v1beta2
kind: KafkaUser
metadata:
  name: app-user
  namespace: middleware
  labels:
    strimzi.io/cluster: my-cluster
spec:
  authentication:
    type: scram-sha-512
  authorization:
    type: simple
    acls:
      - resource:
          type: topic
          name: order-events
        operations: [Read, Write]
      - resource:
          type: group
          name: order-consumer-group
        operations: [Read]

选型建议:如果团队 K8s 运维能力强、偏好 GitOps,Strimzi 是更好的选择。如果需要精细控制每个配置项,手写 StatefulSet 更灵活。两者也可以混用:先用 StatefulSet 部署,后续迁移到 Strimzi 管理。


相关链接

参考文档