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 心跳 - 集群启动时初始化选举
选举过程:
- 每个节点提出自己为 Leader 候选(投票格式:
(myid, ZXID)) - 收到其他节点投票后比较:先比 ZXID(大的优先),再比 myid(大的优先)
- 过半节点同意 → 选出 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 核心配置项详解
| 参数 | 推荐值 | 说明 | 调优场景 |
|---|---|---|---|
tickTime | 2000 | 基本时间单元(ms),心跳 = 1 tick | 网络延迟高的环境可调大到 3000 |
initLimit | 10 | follower 初始连接 leader 的最大 tick 数(= 20s) | 首次启动慢或数据量大时调大 |
syncLimit | 5 | follower 与 leader 同步的最大 tick 数(= 10s) | 网络抖动频繁时适当调大 |
maxClientCnxns | 60 | 单个客户端 IP 最大并发连接数 | 防止单客户端耗尽连接 |
autopurge.snapRetainCount | 3 | 保留的快照数量 | 最少保留 3 个用于恢复 |
autopurge.purgeInterval | 1 | 自动清理间隔(小时) | 设为 0 则禁用自动清理(不推荐) |
dataDir | /data/zookeeper/data | 快照目录(可放 HDD) | - |
dataLogDir | /data/zookeeper/log | 事务日志目录(必须 SSD) | IO 敏感,决定 ZK 写入性能 |
snapCount | 100000 | 每多少次事务触发一次快照 | 事务频繁时调大减少快照开销 |
maxSessionTimeout | 40000 | 最大 session 超时(ms) | 限制客户端 session 上限 |
4lw.commands.whitelist | srvr,stat,mntr | 允许的 4 字母命令白名单 | 监控需要 mntr,排查需要 stat |
JVM 堆 (-Xmx) | 2-4G | 不要超过 4G | GC 停顿会阻塞 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 即可,需要:
- 更新
zoo.cfg:添加新的server.N配置 - 逐个重启:所有现有节点需要加载新的
server.N列表 - 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 是否存活(返回 imok) | echo 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 | 是否只读模式(返回 RO 或 RW) | echo 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.ms | 30000 | Follower 落后超过此时间则移出 ISR |
min.insync.replicas | 1 | 配合 acks=all,ISR 少于此值时 Broker 拒绝写入 |
unclean.leader.election.enable | false | 是否允许非 ISR 副本成为 Leader(数据丢失风险) |
default.replication.factor | 1 | 默认副本数(生产建议 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 |
| compact | cleanup.policy=compact | KV 状态 Topic(如 __consumer_offsets) | 每个 key 只保留最新值 |
| delete+compact | cleanup.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.id | 0, 1, 2… | 集群内唯一,StatefulSet 从 Pod 序号生成 | - |
log.dirs | /data/kafka/logs | 日志目录(可逗号分隔多目录,跨磁盘) | 多磁盘环境分散 IO |
num.partitions | 3+ | 新建 Topic 的默认分区数 | 按吞吐量需求调整 |
default.replication.factor | 3 | 默认副本数 | 生产环境最少 3 |
min.insync.replicas | 2 | ISR 最小副本数 | 配合 acks=all |
log.retention.hours | 168 | 日志保留时间(7 天) | 按业务需求调整 |
log.retention.bytes | -1 | 日志保留大小(-1 不限制) | 磁盘有限时设置上限 |
log.segment.bytes | 1073741824 | 单个 segment 大小(1GB) | 小 segment 更快清理 |
zookeeper.connect | zk-0.zk-headless:2181,… | ZK 地址(ZK 模式) | KRaft 模式不需要 |
advertised.listeners | PLAINTEXT://<Pod域名>:9092 | 对外暴露地址 | 不能用 0.0.0.0 |
auto.create.topics.enable | false | 生产环境关闭自动创建 | 防止误创建 Topic |
unclean.leader.election.enable | false | 禁止非 ISR 副本成为 Leader | 防止数据丢失 |
replica.lag.time.max.ms | 30000 | Follower 落后超时移出 ISR | 网络差时适当调大 |
num.network.threads | 3 | 网络请求处理线程数 | 高并发时调大 |
num.io.threads | 8 | 磁盘 IO 线程数 | 多磁盘时调大 |
socket.send.buffer.bytes | 102400 | TCP 发送缓冲区 | 高吞吐场景调大 |
socket.receive.buffer.bytes | 102400 | TCP 接收缓冲区 | 高吞吐场景调大 |
log.flush.interval.messages | Long.MAX | 每多少条消息强制刷盘 | 默认由 OS 管理,不建议改 |
log.flush.interval.ms | null | 每多少 ms 强制刷盘 | 强制刷盘会降低性能 |
group.initial.rebalance.delay.ms | 3000 | Consumer Group 首次 rebalance 延迟 | 减少 rebalance 期间消息延迟 |
Producer 关键配置
| 参数 | 推荐值 | 说明 |
|---|---|---|
acks | all | 等待所有 ISR 副本确认(最安全) |
retries | 2147483647 | 无限重试(配合 delivery.timeout.ms) |
delivery.timeout.ms | 120000 | 消息发送的总超时(含重试) |
max.in.flight.requests.per.connection | 5 | 单连接未确认请求上限(保持顺序需 ≤ 5) |
enable.idempotence | true | 幂等 Producer(防重复消息) |
compression.type | lz4 | 压缩算法(lz4 性能最好) |
batch.size | 16384 | 批量发送大小(字节) |
linger.ms | 10 | 等待批量发送的时间(ms) |
buffer.memory | 33554432 | Producer 缓冲区大小(32MB) |
max.request.size | 1048576 | 单条消息最大大小(1MB) |
Consumer 关键配置
| 参数 | 推荐值 | 说明 |
|---|---|---|
enable.auto.commit | false | 关闭自动提交偏移量(手动提交更安全) |
auto.offset.reset | earliest | 无偏移量时从最早开始消费 |
session.timeout.ms | 45000 | Consumer 心跳超时 |
heartbeat.interval.ms | 15000 | 心跳间隔(建议 = session.timeout / 3) |
max.poll.records | 500 | 单次 poll 最大消息数 |
max.poll.interval.ms | 300000 | 两次 poll 之间的最大间隔(处理慢时调大) |
fetch.min.bytes | 1 | 最小拉取字节数 |
fetch.max.wait.ms | 500 | 最大等待时间 |
isolation.level | read_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,controller或broker或controller)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 启动失败 | 进程退出,日志报 FATAL | broker.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 Partitions | kafka-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 发送超时 | TimeoutException | Broker 不可达;acks=all 且 ISR 不足;消息过大 | 检查 advertised.listeners;确认 min.insync.replicas;检查 max.request.size |
| Consumer rebalance 频繁 | 消费间歇性暂停 | session.timeout.ms 太短;max.poll.interval.ms 不够;消费者处理太慢 | 调大 session.timeout.ms 和 max.poll.interval.ms;优化消费逻辑 |
| 消息丢失 | 消费端拿不到部分消息 | acks=1 或 acks=0;unclean.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 故障 | 元数据操作卡住,无法创建/删除 Topic | Controller Broker 宕机;ZK 不可达 | 检查 Controller 状态:kafka-metadata-quorum.sh (KRaft) 或 ZK /controller 节点 |
3.2 Zookeeper 常见问题
| 问题 | 典型现象 | 根因 | 解决方案 |
|---|---|---|---|
| 无法组成集群 | 日志反复 LEADER Election 但选不出 Leader | myid 文件缺失/重复/值不对;server.x 配置错误;DNS 解析失败 | 确认 /data/zookeeper/data/myid 唯一且与 zoo.cfg 中 server.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 ISR | ZK 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 expired | ZK 心跳超时;GC 停顿 | 调大 tickTime 或 syncLimit;优化 ZK JVM GC |
3.3 K8s 环境特有问题
| 问题 | 根因 | 解决方案 |
|---|---|---|
| Pod 重启后数据丢失 | 用了 emptyDir 或未正确绑定 PVC | 必须用 StatefulSet + PVC,volumeClaimTemplates 持久化 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 绑定错 Pod | StorageClass 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 / 2G | 2C / 4G | 20G SSD × 2 | 100G SSD | ZK=3, Kafka=3 |
| 中型 (1k-10k msg/s) | 1C / 4G | 4C / 8G | 50G SSD × 2 | 500G SSD | ZK=3, Kafka=3-5 |
| 大型 (10k-100k msg/s) | 2C / 4G | 8C / 16G | 100G SSD × 2 | 1T NVMe × 2 | ZK=5, Kafka=5-7 |
| 超大型 (> 100k msg/s) | 4C / 8G | 16C / 32G | 200G SSD × 2 | 2T NVMe × 4 | ZK=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+ 堆,可控停顿时间 |
MaxGCPauseMillis | Kafka 建议 20ms,ZK 建议 50ms |
InitiatingHeapOccupancyPercent | 35% 就开始并发标记,避免 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
滚动更新注意事项:
podManagementPolicy: OrderedReady— 有序逐个更新- 每次更新一个 Pod,等待其 Ready 后再更新下一个
- preStop hook 中
sleep给 Controller 足够时间做 Leader 迁移 terminationGracePeriodSeconds必须 > preStop 耗时- 大集群(多分区)更新时考虑分区 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 | 名称 | 说明 |
|---|---|---|
| 721 | Kafka Overview | Kafka 集群总览(吞吐量、分区、ISR) |
| 7589 | Kafka Exporter Overview | 基于 kafka-exporter 的监控面板 |
| 10465 | Kafka Lag Exporter | Consumer Lag 监控 |
| 10991 | Strimzi Kafka | Strimzi Operator 管理的 Kafka 面板 |
| 9236 | Zookeeper | ZK 集群监控(延迟、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.dirs | kubectl 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 StatefulSet | kubectl delete statefulset zk -n middleware | 低 | 重新部署 ZK |
| 4. 清理 | 删除 ZK PVC(确认后) | kubectl delete pvc zk-data-zk-0 ... | 不可逆 | 从备份恢复 |
迁移前必须:
- 完整备份 ZK 数据 + Kafka
log.dirs- 在测试环境完整走一遍迁移流程
- 准备回滚方案
- 选择低峰期操作
- 确保有 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 管理。
相关链接
- K8s 1.28-1.36 版本更新总结 — K8s 版本特性
- Redis 高可用与集群方案 — 同类中间件高可用参考
- MySQL 主从复制与高可用 — 数据层面高可用
- Linux 内核调优总览 — 底层性能调优
- 网络内核参数调优 — 网络参数调优