Kafka on Kubernetes 运维面试题
7 道题- 分类
- Kubernetes
- 子分类
- middleware-ops
- 题目数
- 7 道
1 Kafka on Kubernetes 的部署方案对比(Strimzi / Confluent Operator / Banzaicloud Koperator)
答案:
Kubernetes 上部署 Kafka 的主流方案包括 Strimzi、Confluent for Kubernetes(CFK)和 Koperator(原 Banzaicloud Koperator,2021 年 Banzai Cloud 被 Cisco 收购),三者均通过 Operator 模式实现声明式管理与自动化运维。
| 维度 | Strimzi | Confluent for Kubernetes | Banzaicloud Koperator |
|---|---|---|---|
| 开源协议 | Apache 2.0(CNCF Sandbox) | Confluent Community License | Apache 2.0 |
| Kafka 发行版 | Apache Kafka | Confluent Platform(含企业特性) | Apache Kafka |
| ZooKeeper 管理 | Operator 一并管理 ZK 集群 | Operator 管理 ZK | 需用户自行管理 |
| KRaft 支持 | Strimzi 0.34 引入实验性;0.40+ 正式 GA | 完整支持 | 不支持 |
| Topic 管理 | Topic Operator(KafkaTopic CRD) | 原生集成 | 通过 KafkaTopic CRD |
| 用户管理 | User Operator(KafkaUser CRD) | Confluent RBAC | 需手动管理 |
| MirrorMaker 2 | KafkaMirrorMaker2 CRD | Confluent Replicator | 不支持 |
| Connect 管理 | KafkaConnect CRD | Confluent Control Center 集成 | 不支持 |
| Cruise Control | 内置集成 | 内置集成 | 不支持 |
| Schema Registry | 需额外部署 | 原生集成 | 不支持 |
| 监控 | JMX Exporter + Prometheus annotations | Confluent Control Center | 需自行集成 |
| 网络 | TLS/SCRAM-SHA-512/OAuth2/mTLS | mTLS/SASL/SSO | TLS/SASL |
| 存储 | PVC/JBOD/Tiered Storage | PVC | PVC |
| 升级策略 | 支持滚动升级/手动控制 | 支持滚动升级 | 有限支持 |
| 适用场景 | 社区版 Kafka,CNCF 生态,中小规模 | 企业合规,Confluent 生态,大规模 | 简单部署场景 |
Strimzi 优势:CNCF 生态兼容性最佳,CRD 覆盖全生命周期,社区活跃度高。Confluent Operator 优势:企业级功能(RBAC、Control Center、Schema Registry、Tiered Storage)开箱即用,商业支持完善。Banzaicloud Koperator 优势:轻量,部署简单。
2 Strimzi Operator 的架构与核心 CRD
答案:
Strimzi Operator 采用分层架构,通过一组 Kubernetes CRD 描述 Kafka 集群、Topic、用户和数据集成组件,Operator 负责将 CRD 描述转换为实际的 StatefulSet、Deployment、Service 等 Kubernetes 原生资源。
核心 CRD 列表:
| CRD | 职责 | 对应 K8s 资源 |
|---|---|---|
Kafka | 定义 Kafka 集群拓扑,含 Broker、ZK/KRaft、存储、网络、认证配置 | StatefulSet、Service、ConfigMap、PVC |
KafkaTopic | 声明式管理 Topic 的分区数、副本因子、配置 | 无(通过 Kafka Admin API 创建) |
KafkaUser | 定义客户端用户及 ACL 权限 | Secret(存储凭证) |
KafkaConnect | 部署和管理 Kafka Connect 集群 | Deployment/StatefulSet、Service |
KafkaConnector | 管理单个 Connector 实例 | 无(通过 Connect REST API) |
KafkaMirrorMaker2 | 部署 Mirrormaker 2 跨集群复制 | Deployment、Service |
KafkaBridge | 部署 HTTP Bridge,支持 HTTP 协议访问 Kafka | Deployment、Service |
KafkaRebalance | 触发 Cruise Control 执行分区均衡 | 无(通过 Cruise Control API) |
KafkaNodePool | 定义 Broker 节点池(支持异构 Broker) | StatefulSet |
架构工作流:
Kafka CR 创建 → Cluster Operator Watch →
创建 ZK StatefulSet(非 KRaft 模式)→
创建 Broker StatefulSet →
Entity Operator(Topic Operator + User Operator)启动 →
监听 KafkaTopic / KafkaUser CRD →
通过 Admin API 操作 Kafka 集群
Strimzi Operator 控制平面部署在 kafka 命名空间,支持单命名空间或集群级别监听。Cluster Operator 通过 STRIMZI_NAMESPACE 环境变量控制监听范围。
3 Strimzi 的 Kafka 集群 StatefulSet 部署
答案:
Strimzi 将每个 Kafka Broker 和 ZooKeeper 节点映射为 StatefulSet 中的 Pod,利用 StatefulSet 的有序部署、稳定网络标识和持久化存储实现有状态服务的 Kubernetes 化管理。
StatefulSet 关键配置:
apiVersion: kafka.strimzi.io/v1beta2
kind: Kafka
metadata:
name: my-cluster
spec:
kafka:
replicas: 3
version: 3.7.0
listeners:
- name: plain
port: 9092
type: internal
tls: false
- name: tls
port: 9093
type: internal
tls: true
- name: external
port: 9094
type: nodeport
tls: true
storage:
type: jbod
volumes:
- id: 0
type: persistent-claim
size: 100Gi
deleteClaim: false
- id: 1
type: persistent-claim
size: 100Gi
deleteClaim: false
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
template:
pod:
affinity:
podAntiAffinity:
requiredDuringSchedulingIgnoredDuringExecution:
- labelSelector:
matchExpressions:
- key: strimzi.io/name
operator: In
values:
- my-cluster-kafka
topologyKey: kubernetes.io/hostname
resources:
requests:
memory: 8Gi
cpu: 2
limits:
memory: 16Gi
cpu: 4
StatefulSet 网络标识:每个 Broker 分配稳定的 DNS 名称 my-cluster-kafka-0.my-cluster-kafka-brokers.kafka.svc.cluster.local,Broker ID 与 Pod 序号对应。跨集群重启后 Broker ID 保持不变。
PodAntiAffinity 与 Rack Awareness:通过 podAntiAffinity 将不同 Broker Pod 调度到不同 K8s 节点,配合 rack 配置实现机架感知副本分布,确保分区副本跨物理故障域分布。
JBOD 存储:Strimzi 支持 JBOD(Just a Bunch of Disks)配置,每个 Broker 可挂载多个独立 PVC,Kafka 日志目录分布在不同卷上,提升 I/O 并行度。
4 Strimzi 的 Entity Operator(Topic Operator + User Operator)
答案:
Entity Operator 是 Strimzi 中负责 Topic 和 User 生命周期管理的子组件,以 Sidecar 或独立 Deployment 形式运行,通过 Kafka Admin API 与 Kafka 集群交互。
Entity Operator 架构:
| 子组件 | CRD | 职责 | 通信方式 |
|---|---|---|---|
| Topic Operator | KafkaTopic | 管理 Topic 创建、分区数、副本因子、配置更新 | Kafka Admin API |
| User Operator | KafkaUser | 管理用户凭证生成、ACL 绑定、配额设置 | Kafka Admin API |
| TLS Sidecar | — | 管理 TLS 证书轮换与推送 | 文件系统共享卷 |
Topic Operator 工作流程:
KafkaTopic CR 创建/更新 → Topic Operator Watch 事件 →
反序列化 CR → 校验分区数/副本因子/配置 →
调用 AdminClient.createTopics() / AdminClient.describeConfigs() →
对比期望状态与当前状态 → 执行变更 → 更新 KafkaTopic.status
User Operator 工作流程:
KafkaUser CR 创建/更新 → User Operator Watch 事件 →
生成凭证(SCRAM-SHA-512 密码 / TLS 证书)→
存储到 Secret → 调用 Admin API 创建/更新用户 →
绑定 ACL 规则 → 设置 Quota → 更新 KafkaUser.status
Topic Operator 双向同步策略:
| 策略 | 行为 |
|---|---|
UNIDIRECTIONAL(单向) | 仅从 KafkaTopic CR → Kafka 集群同步,集群内手工变更不被 CR 感知 |
BIDIRECTIONAL(双向) | KafkaTopic CR ↔ Kafka 集群双向同步,集群内手工变更会更新 CR 状态 |
常见故障处理:Topic Operator 与 Kafka 集群网络中断时,CR 状态进入 NotReady,连接恢复后自动重建 ZooWatcher 并重新同步。证书过期时 TLS Sidecar 自动触发证书轮换。
5 Kafka Broker 的 Rack Awareness 在 K8s 上的实现
答案:
Rack Awareness 确保 Kafka 分区的不同副本分布在不同的故障域(K8s 节点 / 可用区)上,避免单个节点或 AZ 故障导致数据不可用。
K8s 上实现 Rack Awareness 的三种方式:
| 方式 | 配置方法 | 适用场景 |
|---|---|---|
| Kubernetes Node Label | rack.topologyKey: topology.kubernetes.io/zone | 多 AZ 集群 |
| 手动指定 Rack ID | template.pod.topologySpreadConstraints | 精细化控制 |
| Pod Anti-Affinity | affinity.podAntiAffinity | 确保不同 Broker 在不同 Node |
Strimzi Rack Awareness 配置:
apiVersion: kafka.strimzi.io/v1beta2
kind: Kafka
metadata:
name: my-cluster
spec:
kafka:
replicas: 3
rack:
topologyKey: topology.kubernetes.io/zone
listeners:
- name: plain
port: 9092
type: internal
affinity:
podAntiAffinity:
requiredDuringSchedulingIgnoredDuringExecution:
- labelSelector:
matchExpressions:
- key: strimzi.io/name
operator: In
values:
- my-cluster-kafka
topologyKey: kubernetes.io/hostname
topologySpreadConstraints: # topologySpreadConstraints 与 podAntiAffinity 并列
- maxSkew: 1
topologyKey: topology.kubernetes.io/zone
whenUnsatisfiable: DoNotSchedule
labelSelector:
matchLabels:
strimzi.io/name: my-cluster-kafka
工作原理:Strimzi 读取每个 Broker Pod 所在 K8s 节点的 topology.kubernetes.io/zone Label,将其注入 Kafka 的 broker.rack 配置。Kafka Controller 在分区分配时确保同一分区的不同副本落在不同 AZ 的 Broker 上。
Pod Topology Spread Constraints 实现细粒度分布:
template:
pod:
topologySpreadConstraints:
- maxSkew: 1
topologyKey: topology.kubernetes.io/zone
whenUnsatisfiable: DoNotSchedule
labelSelector:
matchLabels:
strimzi.io/name: my-cluster-kafka
- maxSkew: 1
topologyKey: kubernetes.io/hostname
whenUnsatisfiable: DoNotSchedule
labelSelector:
matchLabels:
strimzi.io/name: my-cluster-kafka
Rack Awareness 与分区副本分布规则:
- 分区副本分布在
min(rackCount, replicationFactor)个机架上。 - ISO(In-Sync Replica)集合优先跨 Rack 分布。
- Rack 级故障时,剩余 Rack 上的 Follower 副本自动选举为新 Leader。
19 Kafka 的 MirrorMaker 2 跨集群复制
答案:
MirrorMaker 2(MM2)基于 Kafka Connect 框架实现,支持集群间的 Topic、Consumer Group Offset 和 ACL 复制。
KafkaMirrorMaker2 CRD:
apiVersion: kafka.strimzi.io/v1beta2
kind: KafkaMirrorMaker2
metadata:
name: dr-mirror
spec:
version: 3.7.0
replicas: 2
connectCluster: target-cluster
clusters:
- alias: source-cluster
bootstrapServers: source-kafka-bootstrap.source-ns.svc:9092
- alias: target-cluster
bootstrapServers: target-kafka-bootstrap.target-ns.svc:9092
mirrors:
- sourceCluster: source-cluster
targetCluster: target-cluster
sourceConnector:
config:
replication.factor: 3
offset-syncs.topic.replication.factor: 3
sync.topic.acls.enabled: true
refresh.topics.enabled: true
refresh.topics.interval.seconds: 60
tasks.max: 4
topicsPattern: "orders.*|payments.*|inventory.*"
groupsPattern: "order-.*|payment-.*"
checkpointConnector:
config:
checkpoints.topic.replication.factor: 3
refresh.groups.enabled: true
refresh.groups.interval.seconds: 60
sync.group.offsets.enabled: true
emit.checkpoints.enabled: true
tasks.max: 2
heartbeatConnector:
config:
heartbeats.topic.replication.factor: 3
tasks.max: 1
MM2 核心 Connector:
| Connector | 功能 |
|---|---|
| MirrorSourceConnector | 从源集群消费消息,写入目标集群(Topic 名加 source-cluster. 前缀) |
| MirrorCheckpointConnector | 同步 Consumer Group Offset(从源集群到目标集群) |
| MirrorHeartbeatConnector | 生成心跳消息,监控复制链路的健康状况 |
MM2 复制语义:
源集群 Topic: orders
目标集群 Topic: source-cluster.orders
复制逻辑:
MM2 从源集群消费 orders → 写入目标集群 source-cluster.orders →
依赖 Kafka Connect 的 Exactly-Once 支持保证不丟不重 →
心跳 Topic(source-cluster.heartbeats)持续写入监控延迟 →
Checkpoint Topic 记录 Consumer Group Offset 映射
MM2 部署架构:
[Source DC]
Kafka Cluster (source-cluster)
|
MM2 Workers (active-standby)
|
[Target DC]
Kafka Cluster (target-cluster)
- source-cluster.orders
- source-cluster.payments
- offsets synced
生产关注事项:
- MM2 Worker 建议部署在目标集群侧,减少跨 DC 延迟对写入的影响。
tasks.max按源集群的分区总数合理设置,单个 Task 处理的分区数不超过 20。- Offset 同步间隔(
emit.checkpoints.interval.seconds)按 RPO 目标设置,典型值 60s。 - MM2 对网络带宽要求高,需确保跨 DC 链路带宽足够(估算:源集群写入吞吐 x 压缩率)。
30 Kafka on Kubernetes 生产环境最佳实践
答案:
Kafka on Kubernetes 生产部署需覆盖集群规划、资源配置、安全加固、监控告警和运维自动化五个维度。
集群规划:
| 维度 | 建议 |
|---|---|
| K8s 版本 | >= 1.27(支持 Strimzi 0.39+),优先选择长期支持版本 |
| 节点隔离 | Kafka Broker 专用节点组(NodeSelector / Taint + Toleration) |
| 可用区分布 | Broker 跨 3 个 AZ 分布(Rack Awareness 配置) |
| 网络 CNI | Cilium / Calico,支持 NetworkPolicy |
| 存储 | 独立的 StorageClass(NVMe SSD 优先),禁用 PVC 自动删除 |
| 集群规模 | 单集群分区数 < 200K,单个 Broker 分区数 < 6K |
资源配置:
spec:
kafka:
replicas: 3
resources:
requests:
memory: 8Gi
cpu: 2
limits:
memory: 16Gi
cpu: 4
jvmOptions:
-Xms: 6144m
-Xmx: 6144m # heap = limit * 0.5(剩余给 Page Cache)
-XX:+UseG1GC
-XX:MaxGCPauseMillis: 20
-XX:G1HeapRegionSize: 32m
-XX:MetaspaceSize: 128m
-XX:MaxMetaspaceSize: 256m
-XX:+ExitOnOutOfMemoryError
storage:
type: jbod
volumes:
- id: 0
type: persistent-claim
size: 500Gi
class: nvme-ssd
template:
pod:
affinity:
podAntiAffinity:
requiredDuringSchedulingIgnoredDuringExecution:
- labelSelector:
matchExpressions:
- key: strimzi.io/name
operator: In
values:
- my-cluster-kafka
topologyKey: kubernetes.io/hostname
tolerations:
- key: dedicated
operator: Equal
value: kafka
effect: NoSchedule
serviceAccount:
metadata:
annotations:
iam.gke.io/gcp-service-account: [email protected]
安全加固清单:
| 安全维度 | 措施 |
|---|---|
| 网络 | TLS 加密所有 Listener,NetworkPolicy 仅允许必要端口入站 |
| 认证 | SASL SCRAM-SHA-512 或 mTLS,禁用 Plain Listener |
| 授权 | KafkaUser ACL 最小权限,Entity Operator 配置 superUsers |
| 存储 | PVC 加密(StorageClass 开启 Encryption),Secret 使用 External Secrets Operator |
| 审计 | 启用 Kafka Authorizer 日志,汇聚至集中日志平台 |
| 镜像 | 使用官方 Strimzi 镜像,定期扫描 CVE 并更新 |
运维自动化:
# Cruise Control 异常检测
spec:
cruiseControl:
config:
anomaly.detection.goals: >-
com.linkedin.kafka.cruisecontrol.analyzer.goals.RackAwareGoal,
com.linkedin.kafka.cruisecontrol.analyzer.goals.ReplicaCapacityGoal
failed.brokers.zk.session.timeout.ms: 15000
metric.anomaly.finder.class: >-
com.linkedin.kafka.cruisecontrol.detector.NoopMetricAnomalyFinder
goal.violation.detection.interval.ms: 300000
| 运维场景 | 自动化方案 |
|---|---|
| 滚动升级 | Strimzi Operator 自动滚动升级,可配置 inter.broker.protocol.version 逐步升级 |
| 集群扩缩容 | 修改 kafka.replicas 后 Operator 自动处理,配合 Cruise Control Rebalance |
| 证书轮换 | Strimzi CA 自动轮换,证书过期前 30 天自动更新 |
| 分区均衡 | KafkaRebalance CRD 或 Cruise Control Anomaly Detector 定时均衡 |
| 备份恢复 | Velero 定时备份,MM2 跨集群灾备 |
| 容量规划 | Prometheus 指标 + 线性回归预测磁盘/分区增长趋势 |
关键告警规则:
groups:
- name: kafka-critical
rules:
- alert: KafkaUnderReplicatedPartitions
expr: kafka_server_replica_manager_underreplicatedpartitions > 0
for: 5m
labels:
severity: critical
annotations:
summary: "Kafka Under Replicated Partitions = {{ $value }}"
- alert: KafkaUnderMinISR
expr: kafka_server_replica_manager_underminisrpartitioncount > 0
for: 2m
labels:
severity: critical
- alert: KafkaBrokerDown
expr: kafka_server_replicamanager_leadercount < 1
for: 2m
labels:
severity: critical
- alert: KafkaActiveControllerCount
expr: sum(kafka_controller_activecontroller) != 1
for: 1m
labels:
severity: critical
- alert: KafkaHighConsumerLag
expr: kafka_consumergroup_group_lag > 100000
for: 10m
labels:
severity: warning
- alert: KafkaDiskUsageHigh
expr: (kafka_log_log_size / kafka_log_log_size_limit) > 0.85
for: 5m
labels:
severity: warning
版本升级策略:
Strimzi 版本升级流程:
1. 阅读 Strimzi 版本升级说明(Upgrade Guide)
2. 备份 Kafka CRD 和配置
3. 升级 Strimzi Cluster Operator(Helm / OLM)
4. 升级 CRD 定义
5. 升级 Kafka CR 中的 version 字段(Kafka 版本)
6. Operator 自动触发 Broker 滚动升级
7. 验证集群状态和客户端兼容性
跨大版本升级原则:
- 不跳版本升级(如 3.5 → 3.6 → 3.7)
- 先升级 inter.broker.protocol.version,再升级 log.message.format.version
- 客户端版本与 Broker 版本的兼容性:**旧客户端连新 Broker 保证兼容(N-1)**;新客户端连旧 Broker 不保证