这是本节的多页打印视图。
点击此处打印.
返回本页常规视图.
模块:Kafka
使用 Pigsty 部署、保护与监控 Apache Kafka 4.1+ 动态 KRaft 集群。
Kafka 是一个分布式事件流平台。Pigsty 的 KAFKA 模块使用 RPM/DEB 软件包,在纳管节点上部署 Apache Kafka 4.1+ 动态 KRaft 集群,并统一管理安全、资源、生命周期与可观测性。
当前状态:Beta 模块
当前 Kafka 模块处于 Beta 状态。用于严肃生产环境前请务必充分测试,确保满足业务需求。
包括动态 KRaft、严格滚动、TLS/SCRAM/ACL、声明式 Topic/User、凭据与证书轮换,以及完整监控链路。
模块能力
KAFKA 模块当前提供:
- 原生动态 KRaft:不安装 ZooKeeper,也不渲染静态
controller.quorum.voters - 三种原生角色
combined / broker / controller,支持复合与控制面/数据面分离拓扑 - 新集群随机生成 Cluster ID 与 Controller Directory ID,由最小 Bootstrap Manifest 冻结身份,冲突时失败关闭
- 按实时健康状态自动选路:冷启动/修复、Broker 串行准入、Controller 动态加入或严格单节点滚动
- 滚动前后检查 Controller 多数派与 Voter 追平、Offline Partition、Under Min ISR 与 ISR 追平
- 成员退役与故障节点替换由剧本编排:
kafka-rm.yml 真子集退役(含死节点),三条命令完成补换 - 两种安全档位:
plaintext 与生产 scram(TLS、SCRAM-SHA-512、Controller mTLS、ACL 与默认拒绝授权) - 声明式收敛 Topic、用户凭据、ACL 与 Quota,不隐式删除业务 Topic;内部凭据与证书支持保护性轮换
- 完整可观测性:JMX 与协议双 Exporter、19 条 Recording Rule、15 条告警规则、4 个 Grafana Dashboard、日志入 VictoriaLogs
模块架构
KAFKA 模块依赖 NODE 完成节点纳管、仓库与基础监控,依赖 INFRA 提供 VictoriaMetrics、VictoriaLogs、Grafana 与 Alertmanager。
flowchart LR
admin["Pigsty 管理节点"] -->|"kafka.yml / exact cluster"| kafka["Kafka 4.1+ / 动态 KRaft"]
kafka --> jmx["每个 Kafka JVM / JMX :9404"]
kafka --> exporter["最多两个 Broker / kafka_exporter :9308"]
kafka --> journal["Journald"]
jmx --> vm["VictoriaMetrics"]
exporter --> vm
journal --> vector["Vector"] --> vl["VictoriaLogs"]
vm --> grafana["Grafana"]
vl --> grafana
vm --> alert["Alertmanager"]
style kafka fill:#70C1B3,stroke:#4f968b,color:#fff
style vm fill:#E66B7A,stroke:#b84e5c,color:#fff
style vl fill:#C98367,stroke:#9e634e,color:#fff
style grafana fill:#F29C64,stroke:#c77845,color:#fff每个 Kafka JVM 都注入 JMX Exporter 并注册为 job=kafka。协议型 kafka_exporter 只在按 kafka_seq 排序后的前两个 Broker-capable 节点运行;单 Broker 集群只运行一个,纯 Controller 不运行。它们返回的是同一逻辑集群视图,Recording Rule 会先去重再聚合。
文档导航
| 文档 | 内容 |
|---|
| 快速上手 | 从单节点到三节点安全集群、客户端接入、参数修改与上线检查 |
| 集群配置 | 拓扑、动态 KRaft、网络、存储、安全与资源声明 |
| 参数参考 | 15 项持久公开参数及临时运维变量 |
| 日常管理 | 状态检查、Topic、消息、Consumer Group 与拓扑变更 |
| 预置剧本 | kafka.yml 生命周期、任务标签、轮换与清理保护 |
| 监控告警 | 指标链路、Dashboard、日志查询与告警规则 |
| 指标定义 | JMX、协议 Exporter 与 Recording Rule 指标字典 |
| 常见问题 | 角色、身份、安全、Exporter 与扩缩容答疑 |
第一次使用
快速上手 提供一条从零开始、由浅入深的完整路径:单节点开发集群 → 三节点 TLS/SCRAM/ACL 安全集群 → 应用客户端接入 → 参数与资源变更 → 上线检查。
如果您已经熟悉 Kafka 与 Pigsty,可以直接进入 集群配置 或 参数参考。
默认端口
| 端口 | 服务 | 部署范围 | plaintext | scram |
|---|
9092 | Kafka Broker | Broker-capable 节点 | PLAINTEXT | SASL_SSL + SCRAM-SHA-512 |
9093 | KRaft Controller | Controller-capable 节点 | PLAINTEXT | 双向 TLS |
9308 | kafka_exporter | 最多两个 Broker-capable 节点 | HTTP 指标 | HTTP 指标,后端使用 TLS/SCRAM |
9404 | JMX Exporter | 所有 Kafka 节点 | HTTP 指标 | HTTP 指标 |
四个端口必须彼此不同,均可通过参数调整。JMX 与协议 Exporter 的 HTTP 端口仍应通过防火墙限制在监控网络内。
当前边界
当前角色提供的是 Kafka 核心部署基线,不替代完整的流平台或托管服务。下列能力仍需显式运行手册或独立组件:
- Broker 扩容后的既有 Partition Reassignment 与副本再均衡(成员的加入/退役/替换已由剧本编排,数据搬迁仍需显式计划)
- 扩容后提升冻结的
default.replication.factor:Kafka 4.3 需要显式数据迁移与静态配置维护窗口 - 已有 Topic 的副本因子变更、Topic 删除与用户删除
- 已格式化集群从
plaintext 在线迁移到 scram - Kafka 版本升级、Feature Level 终结、数据备份、恢复与灾难演练
- 多 Listener、NAT/公网地址、同一 Broker 多客户端网络、Tiered Storage
- Kafka Connect、Schema Registry、MirrorMaker 2、Cruise Control 与 Web UI
这些边界应在生产方案、审批流程与演练中明确记录,不能用普通清单重跑代替。
1 - 快速上手
从零部署单节点与三节点 Kafka,完成安全接入、参数调整和上线检查。
本教程从一个最小单节点集群开始,完成 Topic 创建与消息读写;随后部署一套独立的三节点安全集群,配置应用用户、ACL、Quota 和生产 Topic;最后演示核心参数修改、客户端接入、监控验证与上线检查。
教程范围
这里的“从零开始”是指从尚未部署 Kafka 开始。您需要先有一套可用的 Pigsty 管理节点,并已部署基础 INFRA 服务;如果还没有,请先完成 Pigsty 快速安装。目标节点需要 SSH/Sudo 权限,并可被 NODE 模块纳管。
学习路径
| 阶段 | 目标 | 最终结果 |
|---|
| 1 | 部署单节点开发集群 | 1 个 combined 节点、PLAINTEXT、RF=1 Topic、CLI 读写 |
| 2 | 部署三节点安全 HA 演示基线 | 3 个 combined 节点、动态 KRaft、TLS/SCRAM/ACL、RF=3/minISR=2 |
| 3 | 接入应用客户端 | 使用应用 Principal、Pigsty CA 与 SASL_SSL 生产/消费 |
| 4 | 修改核心参数 | 演示 Heap、Broker 参数、Topic Partition/保留和安全滚动 |
| 5 | 上线验收 | 检查 Quorum、ISR、端到端读写、监控、容量与运行手册 |
两个示例是独立集群
下面的 kf-dev 与 kf-main 是两套独立新集群。如果确有需要,也可以给单节点 kf-dev 声明两个新的 combined 节点后重跑 ./kafka.yml -l kf-dev,角色会逐个完成格式化、Observer 追平与 add-controller 提升,把它原地扩成三 Controller 集群——但演示环境仍建议直接建新集群,扩容语义详见扩容集群。
开始前准备
以下命令默认在 Pigsty 管理节点的项目目录执行:
开始前确认:
pigsty.yml 是当前环境的配置源,先备份并审阅现有内容;- Kafka 节点的
inventory_hostname 可以被所有 Kafka 成员和客户端直接解析、路由; - 管理节点与 Kafka 节点时间同步;
9092、9093、9308、9404 没有端口冲突;/data/kafka 对应专用数据盘或专用目录,且没有混放其他数据;- 每次
kafka.yml 都使用 -l 精确选择同一 Kafka 集群的全部成员; - 真实变更前先执行
--check,审阅输出并取得变更批准。
配置清单必须保留 all.children 层级。下面的组应合并到现有 pigsty.yml,不要用示例覆盖已有的 all.vars、infra、etcd、pgsql 等配置。
一、部署单节点 Kafka
1. 定义集群
将以下 kf-dev 组加入 all.children。该节点省略 kafka_role,因此使用默认 combined,同时承担 Broker 与 Controller:
all:
children:
# 现有 infra、etcd、pgsql 等分组继续保留
kf-dev:
hosts:
10.10.10.10: { kafka_seq: 1 }
vars:
kafka_cluster: kf-dev
kafka_data: /data/kafka
kafka_security: plaintext
kafka_topics:
- name: quickstart.events
partitions: 1
replication_factor: 1
config:
retention.ms: 86400000 # 1 天,仅用于教程
这个配置会得到:
- 一个随机 Cluster ID;
- 一个动态 KRaft combined 节点;
- 默认 RF=1、minISR=1;
- 一个名为
quickstart.events 的单 Partition Topic; - JMX Exporter
:9404 与一个协议 Exporter :9308。
plaintext 没有传输加密、认证和 ACL,只能用于本机开发或可信隔离网络。
2. 纳管节点
如果该主机尚未完成 NODE 初始化,先执行检查模式:
./node.yml --check -l kf-dev
审阅结果并取得批准后再纳管节点:
已经由 Pigsty 纳管、软件仓库和时间同步均正常的节点可以跳过这一步。NODE 的完整准备与日常管理见 节点管理。
3. 部署 Kafka
先对完整集群执行检查:
./kafka.yml --check -l kf-dev
确认目标确实只有 kf-dev 的完整成员,审阅数据路径、软件包、端口与配置变化后执行:
角色会安装 Java 与 kafka-stack、生成随机身份和 Bootstrap Manifest、格式化 KRaft 存储、启动服务、创建 Topic,并注册监控目标。
4. 验证服务与 Quorum
登录 Kafka 节点,检查服务:
systemctl is-active kafka kafka_exporter
journalctl -u kafka --since '-10 min' --no-pager
使用角色自有健康检查:
sudo -u kafka /usr/local/bin/pigsty-kafka-health cluster \
--bootstrap-server 10.10.10.10:9092 \
--command-config /etc/kafka/admin.properties
返回 JSON 中应有 "healthy": true。继续检查动态 Quorum 与 Topic:
/opt/kafka/bin/kafka-metadata-quorum.sh \
--bootstrap-server 10.10.10.10:9092 \
--command-config /etc/kafka/admin.properties \
describe --status
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server 10.10.10.10:9092 \
--command-config /etc/kafka/admin.properties \
--describe --topic quickstart.events
应看到有效 LeaderId、包含本节点的 CurrentVoters,以及 RF=1、ISR=1 的 quickstart.events。
5. 生产与消费消息
启动 Console Producer:
/opt/kafka/bin/kafka-console-producer.sh \
--bootstrap-server 10.10.10.10:9092 \
--command-config /etc/kafka/admin.properties \
--topic quickstart.events
输入几行消息后按 Ctrl-D 结束。在另一个终端消费:
/opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server 10.10.10.10:9092 \
--command-config /etc/kafka/admin.properties \
--topic quickstart.events \
--group quickstart.demo \
--from-beginning
到这里,单节点部署、Topic 收敛和消息读写已经完成。进一步的状态检查见 日常管理。
二、部署三节点安全 HA 演示基线
三节点示例是一套全新的 kf-main 集群,使用三个 combined 节点。它可以容忍一个 Controller 故障;业务 Topic 使用 RF=3/minISR=2,并启用 scram 生产安全档位。
1. 定义安全集群与资源
将以下组加入现有 all.children:
all:
children:
# 现有分组继续保留
kf-main:
hosts:
10.10.10.11: { kafka_seq: 1 }
10.10.10.12: { kafka_seq: 2 }
10.10.10.13: { kafka_seq: 3 }
vars:
kafka_cluster: kf-main
kafka_data: /data/kafka
kafka_heap_opts: '-Xms4G -Xmx4G'
kafka_security: scram
kafka_parameters:
num.partitions: 12
num.network.threads: 6
num.io.threads: 16
log.retention.hours: 168
log.segment.bytes: 1073741824
kafka_users:
- name: quickstart-app
password: "{{ vault_kafka_quickstart_password }}"
acls:
- resource: topic
name: quickstart.
pattern: prefixed
operations: [Read, Write, Describe]
- resource: group
name: quickstart.
pattern: prefixed
operations: [Read]
- resource: cluster
name: kafka-cluster
operations: [Describe, IdempotentWrite]
quota:
producer_byte_rate: 10485760
consumer_byte_rate: 20971520
kafka_topics:
- name: quickstart.events
partitions: 12
replication_factor: 3
config:
min.insync.replicas: 2
cleanup.policy: delete
retention.ms: 604800000
vault_kafka_quickstart_password 必须由现有的 Ansible Vault、KMS 或其他秘密注入机制提供,至少 12 个字符。不要把真实密码直接提交到 Git、日志或工单。
这个配置的关键语义:
- 三个节点全部省略
kafka_role,因此一致使用 combined; - 新集群直接 Bootstrap 为动态 KRaft;
scram 同时启用节点 TLS、Controller mTLS、SCRAM-SHA-512、ACL 和默认拒绝;- 三 Broker 初始复制策略自动派生为 RF=3、minISR=2;
quickstart.events 显式创建 12 个 Partition、3 副本;quickstart-app 可读写 quickstart.* Topic、读取 quickstart.* Group,并可使用幂等 Producer;- 最多两个 Broker 运行
kafka_exporter,三个 Kafka JVM 都运行 JMX Exporter。
如果三个 Broker 确实位于不同故障域,可以在全部节点上分别增加 kafka_rack: az-a/az-b/az-c。不要用虚构 Rack 标签制造不存在的容灾保证,详细规则见 集群配置:Rack。
2. 纳管并部署
如果节点尚未纳管:
./node.yml --check -l kf-main
./node.yml -l kf-main
部署 Kafka 时必须选择全部三个成员:
./kafka.yml --check -l kf-main
./kafka.yml -l kf-main
不能只 -l 10.10.10.11:每个被选中的集群必须完整,部分选择会被拒绝。同时选择多个完整集群(-l kf-dev,kf-main)或不加 -l 裸跑全部集群则是允许的。
3. 验证三节点健康
从管理节点检查三个 Kafka 服务:
ansible kf-main -b -m command -a 'systemctl is-active kafka'
在任一 Broker 上执行完整健康检查:
sudo -u kafka /usr/local/bin/pigsty-kafka-health cluster \
--bootstrap-server 10.10.10.11:9092 \
--command-config /etc/kafka/admin.properties
查询 quorum 和 Topic:
/opt/kafka/bin/kafka-metadata-quorum.sh \
--bootstrap-server 10.10.10.11:9092 \
--command-config /etc/kafka/admin.properties \
describe --status
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server 10.10.10.11:9092 \
--command-config /etc/kafka/admin.properties \
--describe --topic quickstart.events
上线前应看到:一个 Active Controller、三个 Current Voters、三个可用 Broker;所有 quickstart.events Partition 均有三副本、ISR=3,没有 Offline、Under Replicated 或 Under Min ISR Partition。
三、接入应用客户端
1. 分发 CA 公钥证书
将管理节点上的公共 CA 证书安全复制到应用主机:
files/pki/ca/ca.crt -> /etc/kafka-client/pigsty-ca.crt
ca.crt 是可以分发的公钥证书。绝不要复制、暴露或分发 files/pki/ca/ca.key。 应用主机上的 CA 文件建议由 root 管理并设为只读。已被 Pigsty 纳管的应用主机无需复制:NODE 模块已把同一 CA 安装在 /etc/pki/ca.crt,客户端可直接引用。
2. 创建客户端配置
在应用主机创建 /etc/kafka-client/client.properties:
bootstrap.servers=10.10.10.11:9092,10.10.10.12:9092,10.10.10.13:9092
security.protocol=SASL_SSL
sasl.mechanism=SCRAM-SHA-512
sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required username="quickstart-app" password="<secret-from-vault>";
ssl.truststore.type=PEM
ssl.truststore.location=/etc/kafka-client/pigsty-ca.crt
ssl.endpoint.identification.algorithm=https
Kafka Java 客户端支持 SASL_SSL + SCRAM,并支持 PEM Truststore。实际应用应在运行时从 Secret Manager 注入密码,而不是把包含密码的文件提交到仓库。完整字段见 Kafka 4.3 SASL/SCRAM 与 Producer 配置。
3. 为什么应用应直连多个 Broker
Kafka 客户端本身就具备集群感知能力。bootstrap.servers 只用于取得初始元数据;连接成功后,客户端根据元数据直接连接各 Partition 的 Leader Broker,并在 Leader 变化后刷新路由。因此生产环境的常规做法是:
- 在
bootstrap.servers 中配置至少两个、通常三个位于不同故障域的 Broker 地址; - 放通应用到所有 Broker 的
9092,并保证 Broker 宣告的 inventory_hostname 可解析、可路由; - 让 Producer/Consumer 使用 Kafka 客户端自身的重试、元数据刷新、幂等与 Consumer Group 协议;
- 不把 HAProxy、Keepalived VIP、四层 LB 或七层反向代理放在 Kafka 数据面前方。
单个 VIP/LB 既不能替代元数据中的 Broker 地址,也不能把一个连接透明转发到正确的 Partition Leader,只会增加长连接状态、故障定位与容量规划的复杂度。若平台必须提供统一发现入口,DNS 名称或 TCP LB 可以只承担 bootstrap,但 advertised.listeners 仍必须返回客户端可直达的每个 Broker 地址,应用也不能只获准访问 LB。跨 NAT、公网、Kubernetes 或多网络场景需要为每个 Broker 设计独立的外部可达地址与额外 Listener;当前模块固定宣告清单地址,不支持这类映射。
4. 用应用身份验证读写
在安装了 Kafka 4.3 CLI 的应用主机上执行:
kafka-console-producer.sh \
--bootstrap-server 10.10.10.11:9092,10.10.10.12:9092,10.10.10.13:9092 \
--command-config /etc/kafka-client/client.properties \
--topic quickstart.events
消费时使用 ACL 允许的 Group 前缀:
kafka-console-consumer.sh \
--bootstrap-server 10.10.10.11:9092,10.10.10.12:9092,10.10.10.13:9092 \
--command-config /etc/kafka-client/client.properties \
--topic quickstart.events \
--group quickstart.demo \
--from-beginning
生产应用还应显式评审客户端语义:
| 客户端配置 | 建议起点 | 说明 |
|---|
acks | all | 与 RF=3/minISR=2 配合,避免只等待 Leader |
enable.idempotence | true | 降低重试导致重复写入的风险,需要 IdempotentWrite ACL |
group.id | 独立稳定名称 | 不同业务/消费语义不要复用 Group |
| Offset 提交 | 按业务选择 | 自动提交简单;手动提交更容易绑定业务处理结果 |
client.id | 可识别实例名 | 便于日志、Quota 与客户端诊断 |
客户端 acks、重试、幂等、批量、压缩和 Offset 策略属于应用配置,不应写入 Broker 的 kafka_parameters。
四、修改核心参数
Kafka 的持久意图始终修改 pigsty.yml,不要直接编辑 /etc/kafka/server.properties。常见意图对应关系:
| 目标 | 参数 | 行为 |
|---|
| 调整 JVM Heap | kafka_heap_opts | 静态变化,健康集群进入严格单节点滚动 |
| 调整线程、保留、Segment | kafka_parameters | 非角色自有 Broker 参数;静态变化需要滚动 |
| 调整 Topic Partition/保留 | kafka_topics | 在线资源收敛;Partition 只增不减 |
| 调整应用密码/ACL/Quota | kafka_users | 在线资源收敛;密码由秘密系统提供 |
| 声明故障域 | kafka_rack | 所有 Broker-capable 节点全有或全无;变化会滚动但不搬迁数据 |
| 选择安全档位 | kafka_security | 只能在新集群 Bootstrap 时决定,不能普通重跑在线切换 |
示例:调整 Heap 与 Broker 默认参数
假设压测后决定将 Heap 调整为 6G、提高线程数,并把新 Topic 的默认保留时间改成 72 小时:
kf-main:
vars:
kafka_cluster: kf-main
kafka_heap_opts: '-Xms6G -Xmx6G'
kafka_parameters:
num.partitions: 12
num.network.threads: 8
num.io.threads: 24
log.retention.hours: 72
log.segment.bytes: 1073741824
不要照抄 6G/8/24;这些值必须由 CPU、内存、连接数、消息大小、Partition 数、磁盘和 Page Cache 压测决定。
示例:增加 Partition 并缩短 Topic 保留
把 quickstart.events 从 12 个 Partition 增加到 24,并把保留时间改成三天:
kafka_topics:
- name: quickstart.events
partitions: 24
replication_factor: 3
config:
min.insync.replicas: 2
cleanup.policy: delete
retention.ms: 259200000
Partition 不能减少。replication_factor 与现场不一致时,角色会拒绝普通收敛并要求显式 Partition Reassignment;不会自动搬迁既有副本。
应用变更
无论修改静态参数还是动态资源,都运行完整状态机:
./kafka.yml --check -l kf-main
./kafka.yml -l kf-main
不要只运行 -t kafka_config。角色会自动判断:静态变化执行严格逐节点滚动;只修改 Topic/User 等动态资源时不重启 Kafka。
以下键属于角色自身,不能放入 kafka_parameters:
kafka_parameters:
min.insync.replicas: 2 # 错误:角色拥有
default.replication.factor: 3 # 错误:角色拥有
listeners: ... # 错误:角色拥有
全部 15 项公开参数、默认值和保留键见 参数参考。
五、上线前关键检查
拓扑与数据安全
- 生产至少使用三个 Broker,并使用奇数 Controller;关键/大型集群考虑 3 Controller + N Broker 分离拓扑;
- Topic RF、minISR 与生产者
acks 形成一致的故障模型; kafka_rack 只表达真实故障域,且副本放置已经核验;- 数据盘容量、吞吐、延迟、保留时间、峰值写入和恢复时间已经压测;
- 新 Broker 加入后有显式 Reassignment 计划,现有 Topic RF 不会自动提高;
- 已明确 Kafka 数据备份/重建与灾难恢复流程,并演练过故障节点三步替换与成员退役。
安全与网络
- 新生产集群从 Bootstrap 起就使用
kafka_security: scram; - 应用密码由 Vault/KMS/Secret Manager 注入,未进入 Git 或日志;
- 只向客户端分发 CA 公钥证书,不分发 CA 私钥;
- 客户端可以解析并直达所有 Broker 的
inventory_hostname; 9092/9093 只向必要主体开放,9308/9404 只向监控网络开放;- 已建立应用 Principal、Topic/Group/Cluster ACL 与 Quota 审核清单;
- 已安排内部凭据和证书的 受保护轮换。
运行与监控
/usr/local/bin/pigsty-kafka-health cluster 返回健康;- 动态 Quorum 只有一个 Leader,所有预期 Controller 都在 Current Voters;
- 没有 Offline、Under Replicated 或 Under Min ISR Partition;
- 使用真实应用网络、真实 Principal 完成生产与消费验证;
- Kafka Overview、Kafka Instance、Kafka Topic 与 Kafka Consumer 数据正常;
- 告警路由、日志检索、容量阈值、值班责任和回退条件已经确认;
- 升级、Feature Level、Topic 删除、用户删除与集群下线均有独立审批流程。
详细告警与 PromQL 见 监控告警,指标语义见 指标定义。
文档索引与下一步
建议按以下路径继续阅读:
| 您接下来要做什么 | 对应文档 |
|---|
| 规划 combined 或 Controller/Broker 分离拓扑、网络、Rack、存储与安全 | 集群配置 |
| 查找 15 项公开参数、默认值、Schema 和保留键 | 参数参考 |
| 查看 Quorum、Topic、用户、消息、Consumer Group 与扩缩容操作 | 日常管理 |
理解 kafka.yml 生命周期、严格滚动、轮换与集群下线 | 预置剧本 |
| 使用 Dashboard、告警、PromQL 和 VictoriaLogs | 监控告警 |
| 理解每一项 JMX/Exporter/Recording Rule 指标 | 指标定义 |
| 排查身份冲突、连接、SCRAM、Exporter、Lag 与扩缩容问题 | 常见问题 |
| 回到模块能力、默认端口与边界总览 | Kafka 模块首页 |
一条推荐阅读链路是:快速上手 → 集群配置 → 参数参考 → 日常管理 → 预置剧本 → 监控告警 → 常见问题。
2 - 集群配置
规划 Kafka 动态 KRaft 拓扑、身份、网络、存储、安全与声明式资源。
KAFKA 模块使用 15 项持久公开参数表达集群意图,其余拓扑、监听器、存储子目录、复制安全、授权与 Exporter 放置由角色统一推导。首次部署建议先完成 快速上手;完整字段见 参数参考。
先规划,后格式化
kafka_seq 会写入 KRaft node.id;新集群的随机 Cluster ID、初始 Controller Identity、安全模式与初始复制策略会写入 Bootstrap Manifest。存储格式化后,不要随意修改身份、安全模式或 Controller 集合。角色会验证现场与 Manifest 并在冲突时失败关闭,不会自动覆盖或重新格式化数据。
部署前检查
填写清单前至少确认:
- 目标主机已由
NODE 纳管,软件仓库可用,inventory_hostname 可被所有 Kafka 成员与客户端直接路由 - 一次操作将用
-l 精确选择同一 kafka_cluster 的全部成员,而不是单节点、部分成员或多个集群 kafka_seq 在集群内唯一,Controller 为奇数,Broker 数量、故障域与容量目标匹配9092、9093、9308、9404 互不冲突,Infra 节点可以访问两个指标端口kafka_data 对应专用文件系统,并已按保留时间、写入峰值、复制流量、恢复时间与增长余量规划- 生产使用
kafka_security: scram;节点与管理端时间同步,Pigsty CA 可用,应用密码来自 Vault 等秘密来源 - Topic 的 Partition、副本、
min.insync.replicas、保留策略,以及客户端 acks、重试与消费恢复策略已经评审 - 扩缩容、Partition Reassignment、升级、备份、恢复与 Controller 成员变更有独立运行手册
角色与拓扑
kafka_role 只接受三个值:
| 角色 | Kafka process.roles | Broker 端口 | Controller 端口 | JMX | kafka_exporter |
|---|
combined | broker,controller | ✓ | ✓ | ✓ | 可被选择 |
broker | broker | ✓ | - | ✓ | 可被选择 |
controller | controller | - | ✓ | ✓ | - |
kafka_role 是全有或全无的:集群成员要么全部省略(一致使用 combined),要么全部显式声明——混写会在身份预检阶段被拒绝。集群必须至少包含一个 Controller-capable 节点和一个 Broker-capable 节点;偶数 Controller 会给出警告,生产通常使用 3 个 Controller。
单节点开发集群
单节点同时承担 Broker 与 Controller,无法容忍节点故障,只适合开发、测试与功能验证:
kf-dev:
hosts:
10.10.10.10: { kafka_seq: 1 }
vars:
kafka_cluster: kf-dev
角色会从初始 Broker 数量推导 RF=1、minISR=1。不要把单节点拓扑或默认 plaintext 安全模式直接用于生产。
三节点复合部署
三个节点都承担 Broker 与 Controller,是紧凑的生产起点。省略全部角色字段即可使用默认 combined:
kf-main:
hosts:
10.10.10.11: { kafka_seq: 1 }
10.10.10.12: { kafka_seq: 2 }
10.10.10.13: { kafka_seq: 3 }
vars:
kafka_cluster: kf-main
kafka_heap_opts: '-Xms4G -Xmx4G'
kafka_security: scram
kafka_parameters:
num.partitions: 3
num.network.threads: 6
num.io.threads: 16
kafka_topics:
- name: order.events
partitions: 12
replication_factor: 3
config:
min.insync.replicas: 2
cleanup.policy: delete
初始三个 Broker 会自动得到 RF=3、minISR=2 的角色自有复制策略,无需也不允许在 kafka_parameters 中覆盖内部 Topic RF、default.replication.factor 或 min.insync.replicas。示例中的 4G Heap 只是写法示意;生产应通过压测平衡 JVM Heap、操作系统 Page Cache 与同机其他进程。
Controller 与 Broker 分离
关键或较大集群可以把控制面与数据面分离。因为存在显式角色,所有成员都必须声明角色:
kf-main:
hosts:
10.10.10.11: { kafka_seq: 1, kafka_role: controller }
10.10.10.12: { kafka_seq: 2, kafka_role: controller }
10.10.10.13: { kafka_seq: 3, kafka_role: controller }
10.10.10.21: { kafka_seq: 4, kafka_role: broker }
10.10.10.22: { kafka_seq: 5, kafka_role: broker }
10.10.10.23: { kafka_seq: 6, kafka_role: broker }
vars:
kafka_cluster: kf-main
kafka_security: scram
纯 Controller 不监听 9092,也不运行协议 Exporter;它仍通过 JMX 暴露 KRaft 与 JVM 状态。最多两个 kafka_exporter 会放在 kafka_seq 最小的 Broker-capable 节点上。
动态 KRaft 与 Bootstrap Manifest
新集群直接使用动态 Quorum:所有节点渲染 controller.quorum.bootstrap.servers,不会生成静态 controller.quorum.voters。首次格式化时:
- Cluster ID 随机生成,不由集群名哈希;
- 初始 Controller 的 Directory ID 随机生成并冻结;
- 每个节点显式使用
--initial-controllers 或 --no-initial-controllers 格式化模式; - 首次 Bootstrap 启动后,角色等待动态 Quorum 选出 Leader,并校验每个初始 Controller 的 Directory ID 都已进入现场 Quorum。
Bootstrap-only 事实保存在每个集群成员节点上:
scram 集群的每个成员还持有 /etc/kafka/secrets.yml。管理节点不保存任何 Kafka 状态:Manifest 与 Secret 在每次运行时从任一成员副本解析,签发的节点证书放在共享 PKI 树 files/pki/kafka/(CSR 在 files/pki/csr/),丢失时直接用 Pigsty CA 重签。Manifest 只记录集群身份、初始 Controller Identity、安全模式和初始 RF/minISR。活集群始终是运行事实权威:
- Manifest 与现场身份或安全模式冲突时,普通剧本失败关闭;
- 旧 Manifest 存在但全部数据盘为空时拒绝复活旧集群;
- 所有成员都找不到 Manifest 副本而存储已格式化时,失败关闭并提示先在任一成员上恢复该文件;
- 已格式化的
scram 集群在所有成员都没有 Secret 副本时同样失败关闭。
Manifest 是集群的"出生证明":首次 Commission 之后,成员关系以 Raft 现场状态为权威。此后在清单中新增的 Combined/Controller 节点会由剧本编排加入动态 Quorum(全新格式化 → Observer 追平 → add-controller 提升),退役则由 kafka-rm.yml 真子集选择完成(自动 remove-controller 与 Broker 注销),详见扩容集群与缩容集群。
身份参数
| 身份 | 来源 | 示例 | 约束 |
|---|
| 集群名 | kafka_cluster | kf-main | 字母或数字开头,只含字母、数字、下划线和连字符 |
| 节点号 | kafka_seq | 1 | 非负整数,同一集群内唯一 |
| 实例名 | 自动生成 | kf-main-1 | ${kafka_cluster}-${kafka_seq} |
| 节点角色 | kafka_role | combined | 三种原生角色之一 |
| KRaft Cluster ID | Bootstrap 随机生成 | 22 字符 Kafka UUID | kafka_cluster_id 仅作接管/恢复断言 |
已格式化节点会从 ${kafka_data}/metadata/meta.properties 读取 cluster.id 与 node.id,并与 Manifest 及清单交叉校验;初始 Controller 的 Directory ID 则在启动后与活 quorum 比对。身份不匹配是保护性失败,不应通过删除 meta.properties 或清空数据绕过。
网络与监听器
角色只公开端口,不公开 bind、advertised address 或 listener map:
固定监听器约定如下:
- Broker listener 绑定
0.0.0.0,Controller listener 绑定 inventory_hostname; - Broker 的
advertised.listeners 使用 inventory_hostname; - Controller bootstrap 地址也使用
inventory_hostname; plaintext:BROKER 与 CONTROLLER 都使用 PLAINTEXT;scram:BROKER 使用 SASL_SSL + SCRAM-SHA-512,CONTROLLER 使用双向 TLS。
因此客户端必须能够解析并直达每一个 Broker 的 inventory_hostname。当前 v1 不支持 NAT、公网映射、同一 Broker 多客户端网络或任意 raw listener 覆盖;这些场景不能通过 kafka_parameters 拼装绕过。
Kafka 的标准接入模型是智能客户端直连 Broker:bootstrap.servers 配置多个种子地址,客户端获取集群元数据后直接连接 Partition Leader。HAProxy、Keepalived VIP、云 LB 不应作为常规 Kafka 数据面入口,因为它们不了解 Kafka 元数据和 Partition Leader,且无法免除客户端访问所有 advertised.listeners 地址的要求。DNS 或 TCP LB 最多作为可选的 bootstrap 发现入口;即使如此,应用网络仍必须直达全部 Broker。详见快速上手:接入应用客户端。
最小网络流向:
| 来源 | 目标 | 端口 | 用途 |
|---|
| Kafka 客户端、其他 Broker | 所有 Broker | 9092 | Produce、Fetch、元数据与 Broker 间通信 |
| 所有 Kafka 成员 | 所有 Controller | 9093 | KRaft 元数据仲裁 |
| Infra/VictoriaMetrics | 所有 Kafka 节点 | 9404 | JVM/Kafka 指标 |
| Infra/VictoriaMetrics | 被选择的 Exporter 节点 | 9308 | 集群/Topic/Consumer 指标 |
指标端口为 HTTP,即使 Kafka 使用 scram,也应通过防火墙限制在监控网络内。
存储、Heap 与 Rack
用户只设置根目录:
角色固定派生 Topic 数据目录 ${kafka_data}/data 与 KRaft 元数据目录 ${kafka_data}/metadata。kafka_data 必须是专用绝对路径,不能是 /、/data、/var、/etc、/opt、/usr、/home、/root 或 /pg。
生产规划至少考虑保留时间、消息峰值、复制流量、Partition/Segment 数、磁盘延迟与吞吐、文件描述符、恢复时间、JVM Heap 与 Page Cache。当前角色只生成一个 log.dirs;多盘 JBOD、磁盘替换和自动数据迁移需要独立运行手册。
跨故障域部署可以在所有 Broker-capable 节点上一致声明 kafka_rack:
10.10.10.21: { kafka_seq: 4, kafka_role: broker, kafka_rack: az-a }
10.10.10.22: { kafka_seq: 5, kafka_role: broker, kafka_rack: az-b }
10.10.10.23: { kafka_seq: 6, kafka_role: broker, kafka_rack: az-c }
Broker-capable 节点必须全部设置或全部省略 Rack。修改 Rack 会触发安全滚动,但不会自动迁移既有副本。
复制策略
首次 Bootstrap 根据初始 Broker 数量派生:
replication_factor = min(3, broker_count)
min_insync_replicas = max(1, replication_factor - 1)
初始的未来 Topic 默认 RF、内部 Topic RF 与集群 minISR 都会写入 Manifest 并冻结。扩容后:
default.replication.factor 保持初建值;Kafka 4.3 不允许通过动态 Broker 配置在线修改它;- 已有内部/业务 Topic 的 RF 不会自动提高;
- 角色不会把“Broker 已加入”报告成“数据已均衡”;
- RF 变化必须使用经过评审的
kafka-reassign-partitions.sh 计划;提升静态默认值还需要
Controller 高可用或明确维护窗口,并通过完整集群安全滚动生效。
生产者 acks、幂等、重试、批量和压缩属于客户端策略,不是 Kafka Broker 角色参数。
kafka_parameters
kafka_parameters 是唯一的 Broker 参数逃生舱,默认 {},只渲染到 Broker-capable 节点。它适合 num.partitions、线程数、Buffer、保留与 Segment 等非角色自有键。
以下模式由角色拥有,禁止覆盖:
process.roles
node.id
controller.quorum.*
listeners
advertised.listeners
listener.security.protocol.map
inter.broker.listener.name
controller.listener.names
log.dirs
metadata.log.dir
min.insync.replicas
default.replication.factor
offsets.topic.replication.factor
transaction.state.log.replication.factor
transaction.state.log.min.isr
share.coordinator.state.topic.replication.factor
share.coordinator.state.topic.min.isr
broker.rack
authorizer.class.name
super.users
allow.everyone.if.no.acl.found
sasl.*
ssl.*
listener.*
出现任一保留键时,身份预检会在写文件前直接失败。
安全与声明式资源
kafka_security: scram 是一个完整生产档位,而不是一组可任意组合的开关。它自动启用:
- Pigsty CA 签发的每节点证书;
- Controller listener 双向 TLS;
- Broker/client 与 Broker 间 SASL_SSL + SCRAM-SHA-512;
StandardAuthorizer、默认拒绝,以及角色自有管理/监控身份;- 在协议 Exporter 启动前收敛其最小监控 ACL。
应用资源由两个领域对象声明:
kafka_security: scram
kafka_users:
- name: order-service
password: "{{ vault_kafka_order_password }}"
acls:
- resource: topic
name: order.
pattern: prefixed
operations: [Read, Write, Describe]
- resource: group
name: order.
pattern: prefixed
operations: [Read]
quota:
producer_byte_rate: 10485760
consumer_byte_rate: 20971520
kafka_topics:
- name: order.events
partitions: 12
replication_factor: 3
config:
min.insync.replicas: 2
cleanup.policy: delete
资源收敛语义:Topic 创建幂等、Partition 只增加、只更新显式声明的配置;RF 变化会拒绝并提示 Reassignment。声明用户的密码、ACL 与给出的 Quota 字段会幂等收敛。移除 Topic/User 条目不会作为隐式删除流程。
安全模式在 Bootstrap 后不能通过普通剧本切换。内部凭据与证书可以使用 受保护轮换,但 plaintext 到 scram 的在线迁移仍需未来的显式状态机。
软件包与文件布局
角色通过平台映射安装 java-runtime 与 kafka-stack。2026-07-16 验证的载荷为 Kafka 4.3.1、kafka_exporter 1.9.0、JMX Exporter 1.6.0;实际版本仍以目标平台仓库与已安装包为准。
| 路径 | 用途 |
|---|
/opt/kafka/ | Kafka 程序与 CLI |
/etc/kafka/server.properties | 角色生成的服务配置 |
/etc/kafka/admin.properties | 角色生成的 Broker 管理通道;CLI 应始终使用 |
/etc/kafka/controller.properties | 角色生成的 Controller 管理通道 |
/etc/kafka/log4j2.yaml | Journald 日志配置 |
/etc/kafka/jmx_exporter.yml | 有界 JMX 指标规则 |
/etc/kafka/manifest.yml | 节点上的 Bootstrap Manifest 权威副本 |
/etc/kafka/secrets.yml | scram 节点上的内部 Secret 副本 |
/etc/kafka/.pigsty-applied-static.sha256 | 已证明生效的静态配置指纹,滚动重启的判定依据 |
/etc/kafka/pki/kafka.pem | scram 节点 PEM 私钥与证书;信任锚使用系统 /etc/pki/ca.crt |
${kafka_data}/data/ | Topic 日志数据 |
${kafka_data}/metadata/ | KRaft 元数据与 meta.properties |
files/pki/kafka/ | 管理节点上签发的节点证书(<cluster>-<seq>.key/.crt,CSR 在 files/pki/csr/) |
这些文件由角色管理。持久意图应写入 pigsty.yml,不要在节点上直接编辑生成文件,也不要把密码、私钥或角色自有 Secret 内容复制到清单、日志或工单。
3 - 参数参考
KAFKA 模块 15 项持久公开参数与临时受保护运维变量。
KAFKA 角色刻意只公开 15 项持久参数。拓扑、Listener、安全实现、存储子目录、复制安全与 Exporter 放置等细节由角色统一推导,不能作为额外持久变量覆盖。
参数概览
kafka_cluster 与 kafka_seq 必须定义;kafka_role 有真实默认值。集群角色要么全部省略,要么全部显式声明。
身份与拓扑
kafka_cluster
必填的集群身份。必须以字母或数字开头,只能包含字母、数字、下划线和连字符:
它用于发现完整集群成员、生成实例名和定位 Bootstrap Manifest。每次 kafka.yml 生命周期操作必须用精确 -l 选择该集群的全部成员。
kafka_seq
必填的非负整数,在同一 kafka_cluster 中唯一,直接成为 KRaft node.id:
10.10.10.11: { kafka_seq: 1 }
实例名派生为 ${kafka_cluster}-${kafka_seq}。节点格式化后不要修改或复用仍有关联数据的序号。
kafka_role
默认 combined,只接受:
| 值 | Kafka process.roles | 语义 |
|---|
combined | broker,controller | Broker 与 Controller 合设 |
broker | broker | 纯 Broker |
controller | controller | 纯 Controller |
集群所有成员都省略时一致使用 combined;只要任一成员显式设置,所有成员都必须显式设置。不提供旧角色别名。
kafka_cluster_id
默认未设置,仅用于接管或恢复时断言现有集群身份,必须是 22 字符 Kafka UUID:
kafka_cluster_id: MkU3OEVBNTcwNTJENDM2Qk
普通新建集群不要设置。角色会随机生成 Cluster ID,并写入每个成员的 /etc/kafka/manifest.yml。该参数不会重新标记现有数据;与 Manifest 或 meta.properties 冲突时会失败关闭。
kafka_rack
可选的 Broker 故障域标签,渲染为 broker.rack:
10.10.10.21: { kafka_seq: 4, kafka_role: broker, kafka_rack: az-a }
所有 Broker-capable 节点必须全部声明或全部省略。纯 Controller 不使用该值。修改 Rack 属于静态变化,会进入严格滚动,但不会重新分配既有副本。
存储、JVM 与网络
kafka_data
数据根目录,默认 /data/kafka:
角色固定派生 ${kafka_data}/data 与 ${kafka_data}/metadata。该路径必须是专用绝对路径,不能是 /、/data、/var、/etc、/opt、/usr、/home、/root 或 /pg。kafka-rm.yml 默认会删除整个根目录,因此不要混放其他服务或业务文件。
kafka_heap_opts
Kafka JVM Heap,默认:
kafka_heap_opts: '-Xms1G -Xmx1G'
生产应根据负载与内存压测设置,通常保持 Xms 与 Xmx 相同,并为操作系统 Page Cache 与其他进程留出足够内存。
kafka_port
Broker/client 监听端口,默认 9092,只在 Broker-capable 节点监听。plaintext 模式使用 PLAINTEXT;scram 模式使用 SASL_SSL + SCRAM-SHA-512。
kafka_controller_port
KRaft Controller 监听端口,默认 9093(Kafka KRaft 惯例端口),只在 Controller-capable 节点监听。与其他服务共用节点时请自行确认端口无冲突,角色不会自动检测跨服务端口占用。
四个公开端口必须彼此不同。Broker listener 绑定 0.0.0.0,Controller listener、Broker advertised address 与 Controller bootstrap address 固定使用 inventory_hostname,不另设地址参数。
kafka_parameters
默认 {},是唯一的 Kafka Broker 参数逃生舱,只渲染到 Broker-capable 节点:
kafka_parameters:
num.partitions: 12
num.network.threads: 6
num.io.threads: 16
log.retention.hours: 168
log.segment.bytes: 1073741824
以下键或模式由角色拥有,不能通过该映射覆盖:
process.roles
node.id
controller.quorum.*
listeners
advertised.listeners
listener.security.protocol.map
inter.broker.listener.name
controller.listener.names
log.dirs
metadata.log.dir
min.insync.replicas
default.replication.factor
offsets.topic.replication.factor
transaction.state.log.replication.factor
transaction.state.log.min.isr
share.coordinator.state.topic.replication.factor
share.coordinator.state.topic.min.isr
broker.rack
authorizer.class.name
super.users
allow.everyone.if.no.acl.found
sasl.*
ssl.*
listener.*
身份、监听器、安全、存储与复制策略必须保持单一权威;包含保留键时预检会直接失败。
可观测性
kafka_jmx_exporter_port
JMX Exporter HTTP 端口,默认 9404。角色为每个 Kafka JVM 无条件注入 JMX Exporter Java Agent,并注册为 job=kafka;没有单独的开关参数。生命周期健康门禁使用角色自有 Kafka CLI/metadata 通道,不依赖 JMX。Infra 监控节点必须可以访问该端口;端点不因 kafka_security: scram 自动启用 HTTPS,应通过监控网络和防火墙保护。
kafka_exporter_port
协议型 kafka_exporter HTTP 端口,默认 9308。角色只在按 kafka_seq 排序后的前两个 Broker-capable 节点配置、启动与注册;单 Broker 集群只运行一个。监控 Target 文件每次完整运行都会按当前放置刷新,但曾经被选中节点上的旧 Exporter 服务不会被普通剧本自动停止。
Exporter 使用的 Kafka 协议版本、TLS/SCRAM 参数和副本放置均为角色内部约定,没有额外公开开关或 options 参数。
安全与资源
kafka_security
默认 plaintext,只接受:
| 值 | Broker/client | Controller | 授权 | 用途 |
|---|
plaintext | PLAINTEXT | PLAINTEXT | 无 | 开发或可信隔离网络 |
scram | SASL_SSL + SCRAM-SHA-512 | 双向 TLS | StandardAuthorizer,默认拒绝 | 生产安全基线 |
scram 同时配置 Pigsty CA 签发的节点证书、角色自有管理/监控/内部身份、TLS/SCRAM 与 ACL 启用顺序。安全模式写入 Bootstrap Manifest;集群格式化后,普通重跑不能把 plaintext 切换成 scram,也不能反向切换。
节点证书的有效期沿用 Pigsty 共享的 CA 参数 cert_validity(默认 7300d),KAFKA 模块不提供独立的证书有效期参数。
kafka_users
默认 [],仅允许在 scram 模式声明。每个对象只接受 name、password、acls、quota:
kafka_users:
- name: order-service
password: "{{ vault_kafka_order_password }}"
acls:
- resource: topic
name: order.
pattern: prefixed
operations: [Read, Write, Describe]
- resource: group
name: order.
pattern: prefixed
operations: [Read]
- resource: transactional_id
name: order.
pattern: prefixed
operations: [Write, Describe]
quota:
producer_byte_rate: 10485760
consumer_byte_rate: 20971520
约束:
name 在列表中唯一;password 必填且至少 12 个字符,应引用秘密管理系统;- ACL
resource 为 topic、group、transactional_id、cluster; pattern 为 literal(默认)或 prefixed;- 操作为
Read、Write、Create、Delete、Alter、Describe、ClusterAction、DescribeConfigs、AlterConfigs、IdempotentWrite; - Quota 键为
producer_byte_rate、consumer_byte_rate、request_percentage、controller_mutation_rate。
角色为声明用户收敛 SCRAM 密码、完整 ACL 集合与显式给出的 Quota 字段。移除用户条目不会隐式删除 Principal 或凭据;删除/撤权需要独立受审操作。
kafka_topics
默认 []。每个对象只接受 name、partitions、replication_factor、config:
kafka_topics:
- name: order.events
partitions: 12
replication_factor: 3
config:
min.insync.replicas: 2
cleanup.policy: delete
retention.ms: 604800000
身份预检只校验 name 在列表中唯一;Partition 数与 RF 的合法性(至少为 1、RF 不超过当前 Broker 数)由 Kafka 在创建时判定,因此这类错误会在资源收敛阶段暴露,而不是在 --check 阶段。收敛语义是:
- Topic 不存在时幂等创建;
- Partition 只允许增加,减少会失败;
- RF 与现场不同时拒绝普通收敛,并要求显式 Reassignment;
- 只更新
config 中声明的键; - 从列表中移除 Topic 永远不会删除 Topic。
临时受保护运维变量
以下变量只通过命令行 -e 用于一次性运维动作,不属于 15 项持久 API,也不应写入 pigsty.yml:
| 动作 | 剧本 | 临时变量 | 保护条件 |
|---|
| 轮换内部凭据 | kafka.yml | kafka_rotate_credentials=true、kafka_rotate_confirm=<cluster> | 健康、全员已格式化的 scram 集群 |
| 轮换证书 | kafka.yml | kafka_rotate_certificates=true、kafka_rotate_confirm=<cluster> | 健康、全员已格式化的 scram 集群 |
| 下线集群 | kafka-rm.yml | kafka_rm_data(默认 true)、kafka_rm_pkg(默认 false)、kafka_safeguard(默认 false) | kafka_safeguard=true 时中止一切删除 |
两种轮换动作互斥,且必须以精确完整集群为目标。kafka-rm.yml 默认删除数据目录与节点上的 /etc/kafka 恢复状态;kafka_rm_data=false 会同时保留二者。执行前必须显式确认目标集群与备份/重建意图,命令与完整语义见 预置剧本。
kafka_safeguard
仅供 kafka-rm.yml 使用,默认 false。设为 true 时,移除角色会在注销、退群、停服和删除之前直接中止;这是布尔保护开关,不会探测集群是否存活。
kafka_rm_data
仅供 kafka-rm.yml 使用,默认 true。启用时删除整个 kafka_data 和 /etc/kafka;后者包含 Manifest、凭据副本及重新接管保留存储所需的恢复状态。设为 false 会同时保留这两处,但仍会注销监控目标、停止服务并删除运行时集成配置。
kafka_rm_pkg
仅供 kafka-rm.yml 使用,默认 false。设为 true 时卸载平台映射中的 kafka-stack 软件包(Kafka、Kafka Exporter 与 JMX Exporter 载荷);共享的 Java Runtime 不会被卸载。
4 - 日常管理
Kafka 集群的状态检查、Topic 与用户管理、配置变更、扩容缩容、故障节点替换与安全轮换。
KAFKA 模块把 Kafka 安装在 /opt/kafka,使用 Systemd 管理服务,并把持久意图保存在 pigsty.yml。节点上的生成文件不应手工修改。
以下 Kafka CLI 示例都使用角色生成的 /etc/kafka/admin.properties。即使当前是 plaintext 也建议始终保留 --command-config:切换到 scram 管理通道时命令结构不变。将 <broker>:9092 替换为可达的 inventory_hostname 与端口。
Console 工具的 --command-config 需要 Kafka 4.2+ CLI
KIP-1147 从 Kafka 4.2 起把所有 CLI 的配置文件参数统一为 --command-config、键值参数统一为 --command-property。节点上 /opt/kafka/bin 的 CLI 由 Pigsty 仓库提供(当前载荷 4.3.x),可直接使用;若从 4.1 或更早的外部 CLI 执行,Console Producer/Consumer 仍须使用旧名 --producer.config / --consumer.config。管理类工具(kafka-topics.sh、kafka-configs.sh、kafka-acls.sh、kafka-consumer-groups.sh、kafka-metadata-quorum.sh 等)一直使用 --command-config,不受影响。
速查手册
| 操作 | 命令 | 说明 |
|---|
| 创建集群 | ./kafka.yml -l <cls> | 创建或收敛 Kafka 集群,裸跑处理全部集群 |
| 扩容集群 | ./kafka.yml -l <cls> | 声明新成员后收敛:Broker 准入,Controller 加入 |
| 缩容集群 | ./kafka-rm.yml -l <ip> | 退役成员:摘除 Voter 条目与 Broker 注册 |
| 销毁集群 | ./kafka-rm.yml -l <cls> | 下线整个集群,默认删除数据 |
| 替换故障节点 | 退役 → 纳管 → 重入 | 三条命令补换死节点,自动继承副本分配 |
| 配置集群 | ./kafka.yml -l <cls> | 修改清单后在门禁保护下滚动生效 |
| 管理 Topic | ./kafka.yml -l <cls> | 声明式创建 Topic、扩分区、改配置 |
| 管理用户 | ./kafka.yml -l <cls> | 声明式收敛用户、ACL 与 Quota |
| 轮换密钥证书 | ./kafka.yml -e kafka_rotate_... | 受保护的内部凭据 / 证书轮换 |
集群定义与参数详见 集群配置,剧本语义详见 预置剧本,监控排障详见 监控告警。
状态检查
在任意 Kafka 节点检查服务与最近日志:
systemctl status kafka
systemctl is-enabled kafka
journalctl -u kafka --since '-30 min' --no-pager
协议 Exporter 只在 kafka_seq 最小的至多两个 Broker-capable 节点运行。被选择的节点再检查:
systemctl status kafka_exporter
journalctl -u kafka_exporter --since '-30 min' --no-pager
检查监听器与指标端点:
ss -lntp | grep -E ':9092|:9093|:9308|:9404'
curl -fsS http://<kafka-ip>:9404/metrics | grep -E '^(jmx_scrape_error|kafka_server_raft_state|kafka_server_broker_messages_in_total)'
curl -fsS http://<exporter-ip>:9308/metrics | grep -E '^(kafka_brokers|kafka_topic_partitions)'
kafka_up 与 kafka_exporter_up 是 VictoriaMetrics 侧的记录指标,不一定出现在原始端点。JMX 端点应包含 jmx_scrape_error 0.0、JVM 指标和与节点角色匹配的 kafka_ 指标。
健康检查
角色的生命周期门禁不依赖 JMX,而是通过同一管理通道检查动态 Quorum、不可用 Partition、副本不足与 Under Min ISR:
sudo -u kafka /usr/local/bin/pigsty-kafka-health cluster \
--bootstrap-server <broker>:9092 \
--command-config /etc/kafka/admin.properties
返回 JSON 中 healthy: true 才表示该门禁通过。它适合只读诊断,但不能替代业务端到端验证。
该脚本还内置解析回归自检(pigsty-kafka-health selftest),每次剧本运行都会在安装后自动执行;若自检失败说明健康谓词本身不可信,应停止变更并排查。
KRaft 仲裁状态
从任一可用 Broker 查询动态 Quorum:
/opt/kafka/bin/kafka-metadata-quorum.sh \
--bootstrap-server <broker>:9092 \
--command-config /etc/kafka/admin.properties \
describe --status
重点检查:
LeaderId 存在且对应预期 Controller;CurrentVoters 与预期成员一致(加入中的新节点会先出现在 CurrentObservers);MaxFollowerLag 与 MaxFollowerLagTimeMs 没有持续增长;- Dashboard 中恰好有一个 Active Controller。
如需确认动态 Quorum(KIP-853)特性级别,可用 /opt/kafka/bin/kafka-features.sh ... describe 查看 kraft.version。
查看 Controller 复制状态:
/opt/kafka/bin/kafka-metadata-quorum.sh \
--bootstrap-server <broker>:9092 \
--command-config /etc/kafka/admin.properties \
describe --replication
如果没有 Leader、成员长期落后或 Voter 集合与预期不一致,应先停止其他变更,保留日志、Manifest 与 meta.properties 证据再分析。死掉的 Voter 用缩容或替换故障节点流程摘除;不要手工改写 quorum 状态。
管理 Topic
生产 Topic 应优先在 pigsty.yml 的 kafka_topics 中声明:
kafka_topics:
- name: orders
partitions: 12
replication_factor: 3
config:
min.insync.replicas: 2
retention.ms: 604800000
修改声明后运行剧本收敛:
./kafka.yml --check -l kf-main
./kafka.yml -l kf-main
角色会幂等创建 Topic、只增加 Partition,并只修改声明的配置键。RF 变化会失败并要求显式 Partition Reassignment;从清单中移除条目不会删除 Topic。
只读查看 Topic:
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server <broker>:9092 \
--command-config /etc/kafka/admin.properties \
--list
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server <broker>:9092 \
--command-config /etc/kafka/admin.properties \
--describe --topic orders
临时或外部管理的 Topic 可以使用 Kafka CLI 创建,但不会自动写回 pigsty.yml。不要让声明式与手工管理同时拥有同一个 Topic。Topic 删除是业务数据删除动作,必须走独立审批、精确名称确认和恢复方案,本文不提供通用删除命令。
管理用户与权限
kafka_security: scram 时,应用身份应通过 kafka_users 管理:
kafka_users:
- name: order-service
password: "{{ vault_kafka_order_password }}"
acls:
- resource: topic
name: orders
operations: [Read, Write, Describe]
- resource: group
name: order-worker
operations: [Read]
quota:
producer_byte_rate: 10485760
consumer_byte_rate: 20971520
完整剧本会幂等收敛密码、该用户的 ACL 集合与显式给出的 Quota 字段。密码不要以明文提交到仓库或输出到日志。移除用户条目不会自动删除 Principal/凭据;删除或彻底撤权需要独立受审流程。
验证消息读写
使用测试 Topic 做端到端验证。Console Producer/Consumer 使用同一客户端配置文件:
/opt/kafka/bin/kafka-console-producer.sh \
--bootstrap-server <broker>:9092 \
--command-config /etc/kafka/admin.properties \
--topic ops-smoke
在另一个终端消费:
/opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server <broker>:9092 \
--command-config /etc/kafka/admin.properties \
--topic ops-smoke \
--from-beginning \
--group ops-smoke-check
生产验收应从真实客户端网络执行,覆盖 DNS/advertised.listeners、证书校验、ACL、生产者 ACK、消费提交与端到端延迟,而不只验证 Broker 本机路径。
管理 Consumer Group
列出和查看 Consumer Group:
/opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server <broker>:9092 \
--command-config /etc/kafka/admin.properties \
--list
/opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server <broker>:9092 \
--command-config /etc/kafka/admin.properties \
--describe --group order-worker
Lag 要结合消费速率与业务 SLO 判断:短暂积压可能是批处理行为,持续增长且消费速率低于生产速率才表示无法追平。重置 Offset 可能造成重复消费或跳过消息,必须有独立审批、精确 Group/Topic 确认与回放方案。
配置集群
修改 pigsty.yml 后以完整集群为目标执行:
./kafka.yml --check -l kf-main
./kafka.yml -l kf-main
角色根据现场健康和静态指纹自动选择路径:
- 集群不健康或停止:只启动已停止的 Controller,恢复并追平 Quorum 后再启动 Broker;若同时存在静态变化,仍在线成员随后进入严格滚动;
- 存在待加入的 Controller-capable 节点:逐个以 Observer 追平后
add-controller 提升为 Voter; - 健康集群新增纯 Broker:逐个格式化、启动并确认注册;
- 健康集群存在静态变化:严格逐节点滚动,每节点重启前后执行 Controller 零 Lag/最近追平、Quorum、Offline Partition、Under Min ISR 与 ISR 追平门禁;
- 没有静态变化:不重启 Kafka。
不要用 -t kafka_config 绕过完整状态机。动态 Topic/User/ACL/Quota 收敛位于 kafka_provision 资源收敛阶段,静态变化是否重启由角色决定。
扩容集群
健康集群可以直接在清单中声明新成员:kafka_role: broker、combined 或 controller 都可以。为新节点分配从未使用过的 kafka_seq(一台主机同一时间只能属于一个 Kafka 集群),确保节点已被 Pigsty 纳管,然后仍以完整集群为目标:
./node.yml --check -l 10.10.10.14 # 纳管新节点
./kafka.yml --check -l kf-main # 先空跑
./kafka.yml -l kf-main # 逐个准入 / 加入新成员
角色按成员类型自动选择路径,每次只处理一个新节点:
- 纯 Broker:格式化、启动,并验证 Broker 已注册且未 Fenced(
admit); - Combined / Controller:以
--no-initial-controllers 全新格式化、以 Observer 身份启动并追平元数据,再通过 add-controller 提升为 Voter,最后验证其已进入 Voter 集合且集群完整健康(join)。
运行结束时的 quorum-join-hosts / broker-admission-hosts 摘要会列出本次实际处理的节点。两点提醒:
- 新增 Controller-capable 节点会改变所有成员的
controller.quorum.bootstrap.servers,因此存量节点会随之执行一轮门禁保护下的严格滚动,属于预期行为; - 扩出偶数个 Controller 时角色会打印警告:偶数 Quorum 不提升容错能力,请尽量保持奇数。
新 Broker 加入不会迁移已有 Partition。必须另外生成、评审并监控 kafka-reassign-partitions.sh 计划,控制磁盘/网络负载并准备回退。“服务已注册"不等于"扩容完成”。
复制策略也不会随 Broker 数自动放大。尤其是 Kafka 4.3 的
default.replication.factor 不能动态修改:由 1 Broker 扩到 3 Broker 后,它仍为初建的
RF=1,未来未显式指定 RF 的 Topic 也仍按 RF=1 创建。应先完成既有 Partition
Reassignment,再规划 Controller 高可用或维护窗口,最后让新的静态默认值通过完整集群
安全滚动生效;不能为了改默认值绕过停机门禁。
缩容集群
用 kafka-rm.yml 选择集群的真子集即为成员退役(选择整个集群则是集群下线)。退役会通过一台幸存成员,自动从现场元数据中摘除该节点:
./kafka-rm.yml -l 10.10.10.13 # 退役单个成员:摘除 Voter 条目、注销 Broker、清理本机
执行内容依次为:注销监控 Target → 停止服务 → remove-controller 摘除 KRaft Voter 条目(若该成员是 Voter;多成员退役时严格串行)→ kafka-cluster.sh unregister 注销 Broker → 清理本机配置与数据(受 kafka_rm_data 控制)。完成后从 pigsty.yml 中删除该成员条目。
退役前请自行确认:剩余 Controller 仍构成多数派、保持奇数个 Controller、剩余 Broker 数不低于现有 Topic 的最大 RF。如果被退役 Broker 上仍有 Partition 副本,角色会打印警告:这些 Partition 将保持副本不足,直到同 kafka_seq 的替换节点重新加入(自动继承副本分配并补数据),或你显式执行 Reassignment 将副本迁走。计划内缩容应当先 Reassignment 排空、再退役。
替换故障节点
节点永久损坏(磁盘丢失、机器报废)时,保持其 IP 与 kafka_seq 不变,三步完成补换:
./kafka-rm.yml -l 10.10.10.13 # ① 退役死者:摘除 Voter 条目与 Broker 注册(节点不可达也能执行)
./node.yml -l 10.10.10.13 # ② 纳管替换机器(修复或换新,保持 IP)
./kafka.yml -l kf-main # ③ 重新加入:格式化、追平、准入/提升,自动继承原副本分配并补数据
第 ① 步的所有元数据操作都委派给幸存成员执行,因此对已经无法连接的死节点同样有效;它还会一并清理监控 Target,避免死节点持续触发 KafkaDown 告警。第 ③ 步中,同 kafka_seq 的 Broker 会自动继承原 Partition 分配并从副本重新同步数据,无需手工 Reassignment。
如果跳过第 ① 步直接重装节点并重跑 kafka.yml,角色会在配置阶段快速失败,并在报错中给出残留 Voter 条目的 Directory ID 与确切的 kafka-rm.yml 命令——按提示执行后重跑即可。加入流程可安全重入:任一步骤被中断后,重跑 kafka.yml 会从现场状态继续。
变更地址与端口
角色固定使用 inventory_hostname 作为 Broker advertised address 与 Controller bootstrap address。修改清单地址、kafka_port 或 kafka_controller_port 会影响客户端元数据、Broker 通信或 Quorum,属于静态高风险变更;必须同步检查 DNS、证书 SAN、路由、防火墙、Bootstrap 地址、监控 Target 与所有成员。
轮换密钥与证书
已格式化且健康的 scram 集群支持两种互斥的受保护动作:内部凭据轮换和证书轮换。两者都要求精确完整集群、匹配的 kafka_rotate_confirm 确认字符串,并且建议先执行 --check。证书由同一 Pigsty CA 重新签发,新旧证书互信,轮换通过严格滚动逐节点生效。
具体命令和失败语义见 预置剧本:受保护轮换。安全模式本身是 Bootstrap-only 属性;这些动作不等于支持 plaintext 到 scram 的在线迁移。
数据保护与恢复
Kafka 的数据保护依赖跨故障域副本、正确的 minISR、生产者 ACK 和经过演练的恢复流程。当前角色不提供 Kafka 数据备份、自动 Broker Drain(计划内缩容需先手工 Reassignment)或跨地域灾难恢复。
发生磁盘或节点故障时:
- 先查看 Kafka Overview/Instance、Quorum、ISR、Offline Partition 与 Under Min ISR;
- 保存
journalctl -u kafka、节点指标、Manifest、server.properties 与 meta.properties 证据; - 确认节点角色、
node.id、Cluster ID、Directory ID 与剩余副本可用性; - 节点确认无法恢复时,按替换故障节点三步走:
kafka-rm.yml 退役 → node.yml 纳管 → kafka.yml 重入;磁盘尚存、仅服务异常时不要急于退役或删除 meta.properties,先尝试普通收敛拉起; - 对 Reassignment、RF 变更等数据搬迁操作仍使用独立评审的运行手册。
日志诊断
journalctl -u kafka -f
journalctl -u kafka_exporter -f
journalctl SYSLOG_IDENTIFIER=kafka --since today
journalctl SYSLOG_IDENTIFIER=kafka_exporter --since today
VictoriaLogs/Grafana 查询:
job:syslog unit:kafka
job:syslog app:kafka
job:syslog unit:kafka_exporter
常见诊断顺序是:服务日志 → 监听端口 → 管理通道健康 → 动态 Quorum → Broker/Partition/ISR → 客户端地址与证书/ACL → Consumer Lag。详细面板与告警映射见 监控告警。
5 - 预置剧本
使用 kafka.yml 与 kafka-rm.yml 执行动态 KRaft 生命周期、严格滚动、资源收敛、轮换与下线。
KAFKA 模块提供两个剧本:kafka.yml 用于部署 Apache Kafka 4.1+ 动态 KRaft 集群并收敛其安全、
资源与监控状态;kafka-rm.yml 用于下线集群或移除成员。
集群完整性约束
每个被选中的 kafka_cluster 必须包含其全部成员:部分选择会在写入前失败;选择一个集群、多个完整集群或不加 -l 裸跑全部集群都是允许的。先对完全相同的目标执行 --check;真实运行前仍需人工核验备份/重建意图、容量、业务窗口、回退方案与变更批准。
kafka.yml
./kafka.yml --check -l kf-main # 先空跑
./kafka.yml -l kf-main # 创建或收敛单个集群
./kafka.yml # 裸跑:一次创建/收敛清单中的所有 Kafka 集群
Limit 规则是:每个被选中的集群必须完整。可以选择一个集群、多个集群,或不加 -l 对全部集群裸跑(集群内严格串行、集群间并发推进);但部分选择某个集群的成员会被直接拒绝。
检查模式验证公开 API、完整集群、角色、Rack、端口、Manifest 与可检查的文件变化,但会跳过格式化、服务启动和实时健康验收。因此 --check 成功不等于运行时一定成功。
执行阶段
kafka.yml 本身是一个薄封装:单一 Play 依次执行 node_id 与 kafka 两个角色,与 pgsql.yml 的结构一致。角色内部把生命周期拆成六个任务阶段;所有跨节点排序(并行 Bootstrap、逐个 Controller 加入、逐个 Broker 准入、严格逐节点滚动)由启动阶段统一负责:
| 阶段 | 标签 | 作用 |
|---|
| 身份预检 | kafka-id | 派生并断言身份、集群完整性、角色、Rack、端口与保留键 |
| 安装 | kafka_install | 创建 kafka 系统用户,安装 java-runtime 与 kafka-stack 软件包 |
| 配置 | kafka_config | 读取/恢复/创建 Manifest,签发安全材料,渲染配置,计算静态指纹,格式化空存储,判定生命周期路径 |
| 启动 | kafka_launch | 收敛不健康集群、逐个加入 Controller 与准入 Broker、严格滚动,确认 Manifest 与已生效静态状态 |
| 资源收敛 | kafka_provision | 收敛动态 minISR、用户凭据、ACL、Quota 与声明式 Topic,报告内部 Topic RF 漂移 |
| 监控 | kafka_monitor | 配置协议 Exporter 并注册 VictoriaMetrics Target |
Play 使用 any_errors_fatal: true。某个阶段失败时,后续危险推进会停止;修正原因后可以重跑完整集群,角色会从现场状态和持久指纹恢复,而不是盲目重复格式化。
生命周期路径
配置阶段使用角色自有管理通道判断集群健康,并选择唯一后续路径:
冷启动、首次部署或修复
当集群停止或健康谓词不通过时,进入 Converge:
- 启动所有 Controller-capable 节点;
- 等待 Controller listener 与动态 Quorum Leader;
- 首次 Bootstrap 时验证初始 Controller Directory ID 已进入现场 Quorum;
- 启动纯 Broker;
- 等待 Broker listener 并要求完整集群健康;
- 只有配置已证明成功运行后,才持久化静态指纹。
JMX 不参与生命周期门禁:启动、准入与滚动的判定完全基于角色自有的 Kafka CLI/metadata 管理通道。
健康集群新增 Broker 或 Controller
新格式化的 kafka_role: broker 逐个准入(admit):启动后要求它已经注册且未 Fenced 才继续下一个。
新的 Combined/Controller 节点则逐个加入动态 Quorum(join):已 Commission 的集群以 --no-initial-controllers 全新格式化该节点,它以 Observer 身份启动并追平元数据,随后角色执行 add-controller 将其提升为 Voter,并用健康后置检查确认它进入 Voter 集合且集群完整健康。加入流程可重入:中断后重跑会从现场状态继续;若其 node.id 在 Quorum 中残留着死去前任的 Voter 条目,配置阶段会快速失败并给出先行 kafka-rm.yml 退役的确切命令。
准入/加入只证明服务成为成员;已有 Partition 不会自动迁移到新 Broker,必须另行执行显式 Reassignment。
健康集群静态变化
当渲染后的静态指纹变化时,严格滚动每次只处理一个节点:
- 重启前检查 Controller 多数派、全部 Voter 零 Lag 且最近完成追平、Offline Partition、Under Replicated、Under Min ISR,以及移除目标后每个 Partition 的有效 ISR;
- 重启后要求目标 Controller 回到 Voter 且重新追平、目标 Broker 注册且未 Fenced、其副本重新进入 ISR;
- 任一门禁失败立即停止后续节点。
如果故障修复与静态变化同时存在,Converge 只启动已停止的成员,不并行重启仍在线成员;Quorum 恢复并追平后,尚未加载的静态变化继续进入严格滚动。
如果静态指纹没有变化,Kafka 不重启。动态资源变化仍会在资源收敛阶段在线生效。
任务标签
| 标签 | 阶段/作用 |
|---|
kafka-id | 始终执行的身份、完整集群与拓扑派生断言 |
kafka_install | 安装阶段总入口 |
kafka_user | 创建 kafka 系统用户与用户组 |
kafka_pkg | 按平台映射安装 java-runtime 与 kafka-stack 软件包 |
kafka_config | Manifest、安全材料、配置渲染、静态指纹、存储格式化与路径判定 |
kafka_launch | Converge、Controller 串行加入、Broker 串行准入、严格滚动与 Manifest Commission |
kafka_provision | 动态 minISR、Topic、User、ACL 与 Quota 收敛 |
kafka_monitor / monitor | 协议 Exporter 配置与监控注册总入口 |
kafka_register / register / add_metrics | 仅刷新 VictoriaMetrics 文件发现 Target |
正常配置变更应运行完整 kafka.yml,让角色自行选择生命周期路径。阶段标签主要用于开发、诊断和受控修复;不能用 -t kafka_config 或只限制单节点来绕过完整状态机。
身份、格式化与 Manifest
角色在写配置前校验:
- 每个被选中的集群包含其全部成员;
kafka_seq 唯一,角色全部省略或全部显式;- 至少一个 Controller 和一个 Broker;
- Rack 在所有 Broker-capable 节点上全有或全无;
- 端口有效、互不冲突,角色自有键未被
kafka_parameters 覆盖; - Manifest、安全模式、
meta.properties 与现场集群身份一致。
新集群随机生成 Cluster ID 和初始 Controller Directory ID,并以显式动态 Quorum 模式格式化每个节点。已有 ${kafka_data}/metadata/meta.properties 时在本地验证 Cluster ID 与 Node ID;初始 Controller Directory ID 只在首次 Bootstrap 启动后与现场 Quorum 比对,Commission 之后成员关系以 Raft 现场状态为准。角色不会自动重新格式化已有存储。
Bootstrap Manifest 的权威副本位于每个集群成员上:
scram 集群的每个成员另有 /etc/kafka/secrets.yml;管理节点不保存任何 Kafka 状态,每次运行时从任一成员副本解析。活集群是运行事实权威,但普通剧本不会在冲突时擅自改写任何一方:
- 所有成员都没有 Manifest 副本而存储已格式化时,失败关闭并提示先在任一成员上恢复该文件;
- Manifest 存在而所有数据盘为空时失败关闭;
- Cluster ID、安全模式或 Controller Identity 冲突时失败关闭;
- 新节点的
node.id 在 Quorum 中残留前任 Voter 条目时快速失败,要求先用 kafka-rm.yml 退役。
不要删除 meta.properties、Manifest 或 Secret 来绕过保护。
静态指纹与可恢复重跑
角色对影响 Kafka 进程的静态文件计算期望指纹,并只在以下条件之一成立后写入 /etc/kafka/.pigsty-applied-static.sha256:
- Converge 已经成功启动并通过全局健康检查;
- 严格滚动已经让该节点重启、追平并通过后置门禁。
如果执行中断,未被证明生效的变化不会被记成“已应用”。下一次完整重跑仍能识别待处理的静态重启。
资源收敛与监控注册
完整健康后,资源收敛与监控阶段依次:
- 收敛角色拥有的动态 cluster minISR;
- 幂等处理
kafka_users 的凭据、ACL 与声明 Quota; - 幂等处理
kafka_topics 的创建、Partition 增长与显式配置; - 检查内部 Topic RF 漂移,但不自动 Reassignment;
- 在按
kafka_seq 排序后的前两个 Broker-capable 节点配置并启动协议 Exporter; - 在全部 Infra 节点刷新文件发现 Target。
每个实例对应一个 Target 文件,JMX 目标与(被选中节点的)协议 Exporter 目标都在同一 kafka 采集任务下:
/infra/targets/kafka/<kafka_instance>.yml
Target 文件每次完整运行按当前 Exporter 放置刷新;Target 的删除由 kafka-rm.yml 的注销步骤完成。
受保护轮换
轮换变量是一次性 extra-vars,不应写入 pigsty.yml。两种动作互斥,每次只能执行其一;前提是所有成员已格式化、集群健康、安全模式为 scram、角色自有 Secret 材料存在,且 kafka_rotate_confirm 与集群名完全一致。
内部凭据轮换
./kafka.yml --check -l kf-main \
-e kafka_rotate_credentials=true \
-e kafka_rotate_confirm=kf-main
./kafka.yml -l kf-main \
-e kafka_rotate_credentials=true \
-e kafka_rotate_confirm=kf-main
角色使用 active/standby 内部身份:先通过活管理通道更新非活动凭据,再原子切换本地受保护记录,并进入正常严格滚动。旧 active 保留为下一轮 standby,使中断后的重跑可恢复。
证书轮换
./kafka.yml --check -l kf-main \
-e kafka_rotate_certificates=true \
-e kafka_rotate_confirm=kf-main
./kafka.yml -l kf-main \
-e kafka_rotate_certificates=true \
-e kafka_rotate_confirm=kf-main
角色废弃共享 PKI 树中已签发的节点证书,用同一 Pigsty CA 为每个节点重新签发私钥与证书,更新节点上的 PEM 证书包并进入严格滚动。新旧证书由同一 CA 签发、彼此互信,因此不需要分阶段互换信任;健康预检失败时不会开始轮换,节点上的现有证书保持不变。
kafka-rm.yml
移除动作不在 kafka.yml 中,而是使用独立的 kafka-rm.yml 剧本。-l 选中一个集群的全部成员即为集群下线,选中真子集即为成员退役,两者共用同一执行顺序:
注销 VictoriaMetrics Target(kafka_deregister)→ 停止并禁用 kafka/kafka_exporter 服务(kafka)→ 经幸存成员摘除 KRaft Voter 条目与 Broker 注册(kafka_retire,仅在选中真子集时有幸存成员可用)→ 删除 Exporter 配置、Systemd 环境/Unit 与辅助脚本(kafka_config)→ 删除数据目录与节点上的 /etc/kafka 恢复状态(kafka_data,受 kafka_rm_data 控制)→ 可选卸载软件包(kafka_pkg,受 kafka_rm_pkg 控制)。
防误删开关是 kafka_safeguard:设置为 true(命令行或清单中)时剧本直接中止,不删除任何东西。身份冲突、Exporter 异常或一般启动失败都不是删除数据的理由——先用 kafka.yml 收敛并读取失败原因。
集群下线
./kafka-rm.yml -l kf-main # 移除集群:注销监控、停服务,默认删除数据与 /etc/kafka 恢复状态
./kafka-rm.yml -l kf-main -e kafka_rm_data=false # 保留磁盘数据与 /etc/kafka 恢复状态,只移除服务集成
./kafka-rm.yml -l kf-main -e kafka_rm_pkg=true # 同时卸载 kafka-stack 软件包(共享的 Java 运行时不会卸载)
永久删除
kafka_rm_data 默认为 true:一次默认参数的 kafka-rm.yml 就会删除所选节点的数据/KRaft 元数据与 /etc/kafka 恢复状态。剧本没有确认字符串等额外闸门,执行前必须人工核对 -l 目标、备份或明确重建意图,并评估生产者/消费者影响。
成员退役
./kafka-rm.yml -l 10.10.10.13 # 退役单个成员:摘除 Voter 条目与 Broker 注册,再清理本机
剧本通过一台幸存成员摘除该节点的 KRaft Voter 条目(remove-controller,多成员时严格串行)并注销其 Broker 注册(unregister),再执行本机清理。所有元数据操作都委派给幸存成员,因此对已经死亡、无法连接的节点同样适用——这也是替换故障节点的第一步。
退役自动化不等于免除规划:缩容后剩余 Controller 应保持奇数并构成多数派,剩余 Broker 数不能低于现有 Topic 的最大 RF;若被退役 Broker 仍持有 Partition 副本,剧本会打印警告——计划内缩容应当先完成 Reassignment 排空。
剧本边界
两个剧本都不会自动完成 Partition Reassignment 与数据均衡、Topic/用户删除、plaintext 到 scram 的在线迁移、版本升级与 Feature Level 终结、数据备份与灾难恢复,也不部署 Connect、Schema Registry、MirrorMaker、Cruise Control 等生态组件。完整清单见 模块边界;日常只读检查和资源管理见 日常管理。
6 - 监控告警
Kafka 指标采集、Grafana Dashboard、日志查询与告警规则。
Pigsty 为 KAFKA 模块提供指标、日志、Dashboard 与告警一体化的可观测能力。监控同时覆盖 Kafka JVM 内部状态与 Kafka 协议视角,避免只看到进程存活而看不到 Partition、ISR 与 Consumer Lag,也避免只看到集群元数据而看不到 JVM、请求队列与 KRaft Controller 健康。
采集架构
KAFKA 模块使用两个互补的 Exporter:
| 采集面 | 服务/方式 | Job | 节点范围 | 主要内容 |
|---|
| JVM 与 Kafka 内部 | JMX Exporter Java Agent :9404 | kafka(带 role 标签) | 所有 Kafka 节点 | JVM、Broker 吞吐、复制、请求路径、KRaft、Controller |
| Kafka 协议视角 | kafka_exporter :9308 | kafka(无 role 标签) | kafka_seq 最小的至多两个 Broker-capable 节点 | Broker、Topic、Partition、Offset、Consumer Group、Lag |
| 主机资源 | node_exporter | node | 纳管节点 | CPU、内存、磁盘、网络、文件系统 |
| 日志 | Journald → Vector → VictoriaLogs | syslog | 所有 Kafka 节点 | Kafka 与 Exporter 结构化检索日志 |
角色在每一个 Infra 节点为每个实例生成一个文件发现目标,JMX 目标与(被选中节点的)协议 Exporter 目标都在同一文件、同一 kafka 采集任务下:
/infra/targets/kafka/<kafka_instance>.yml
单 Broker 集群只运行一个协议 Exporter;多 Broker 集群最多运行两个。纯 Controller 只注册 JMX 目标;未被选择的 Broker 与纯 Controller 都没有协议 Exporter 目标,这是预期行为。Target 文件每次完整运行按当前放置刷新;实例 Target 的删除由 kafka-rm.yml 的注销步骤完成。
标签模型
两类目标都注册在同一 job=kafka 采集任务下,通过有无 role 标签区分。
JMX 目标
| 标签 | 含义 | 示例 |
|---|
job | 采集任务 | kafka |
cls | Kafka 集群名 | kf-main |
ins | Kafka 实例名 | kf-main-1 |
ip | 清单主机地址 | 10.10.10.11 |
instance | JMX 抓取端点 | 10.10.10.11:9404 |
role | Pigsty Kafka 角色 | combined、broker 或 controller |
node_id | KRaft 节点号 | 1 |
协议 Exporter 目标
协议 Exporter 目标只包含 cls、ins、ip 与 instance(10.10.10.11:9308),没有 role/node_id 标签。vmagent 端的记录规则据此区分两类可用性:kafka_up 为 up{job="kafka",role=~".+"},kafka_exporter_up 为 up{job="kafka",role=""}。
Exporter 从 Broker 查询整个 Kafka 集群,因此同一集群的两个 Exporter 可能返回相同 Topic/Partition/Consumer Group 视图。集群级 Recording Rule 会先在 Exporter 实例间去重,再汇总逻辑集群速率。scram 模式下,Exporter 连接 Kafka 所需的 TLS/SCRAM 参数由角色自有监控身份自动生成。
Grafana Dashboard
Pigsty 提供四个互补 Dashboard:
集群与全局总览。cls=All 是全部 Kafka 集群的 Overview;选择具体 cls 后,同一 Dashboard 就成为该 Kafka Cluster 的总览,而不是另一套独立面板。
主要内容:
- 集群、Broker、Topic、Partition 与 Consumer Group 清单
- Broker 可用性、Exporter 健康与集群工作负载
- Leaderless、Under Replicated、ISR Deficit、Non-Preferred Replica
- Topic Offset 进展、Consumer Commit 进展与总 Lag
- Consumer Group 成员、Lag 排名和 Topic/Group 下钻
- Kafka/Exporter 日志量、Firing Alerts 与日志明细
常用变量:cls、members、topic、group、topk。
以 ins 变量选择任意 Kafka Broker/Controller JVM,包括纯 Controller,并联动宿主机资源。
主要内容:
- 实例身份、角色、JMX 可用性与抓取质量
- JVM Heap、GC、Thread、Buffer Pool、CPU、FD 与 Uptime
- Broker 吞吐、复制状态、请求错误/延迟/队列和 Handler/Network Idle
- KRaft Member State、Metadata Log、Controller 健康与事件延迟
- 节点 CPU/内存、磁盘 I/O、网络、文件系统与 Kafka 日志
常用变量:cls、ins、ip。
以 cls 与 topic 选择逻辑 Topic,查看 Topic/Partition 的协议状态。
主要内容:
- Topic 与 Partition 清单、Leader、副本、ISR 和 Preferred Leader
- Current Offset、保留跨度与消息追加速率
- Leaderless、ISR Deficit 和 Non-Preferred Replica
- 关联 Consumer Group、提交进度与 Lag
常用变量:cls、topic、topk。
以 cls 与 group 选择 Consumer Group,查看成员、提交 Offset、消费进展与积压。
主要内容:
- Consumer Group 清单与成员数量
- Group/Topic/Partition 的已提交 Offset
- Commit Rate、总 Lag、最大 Partition Lag 与积压趋势
- Group 到 Topic/Partition 的下钻
常用变量:cls、group、topic、topk。
Dashboard 选择
| 问题 | 首选 Dashboard | 下钻方向 |
|---|
| 哪个集群或 Topic 出现异常? | Kafka Overview | 选择 cls、topic、group |
| 某个 Consumer Group 为什么积压? | Kafka Consumer | Group → Topic → Partition Offset |
| 某个 Topic/Partition 是否异常? | Kafka Topic | Topic → Partition → Consumer |
| 某个 Broker 是否过载? | Kafka Instance | 请求路径 → JVM → Node 资源 |
| KRaft Controller 是否健康? | Kafka Instance | KRaft Metadata Plane → Controller Health |
| 是否存在 Leaderless/URP/ISR 问题? | Kafka Overview | Cluster → Kafka Instance / Topic |
| Exporter 缺数还是 Kafka 本身异常? | Overview + Instance | 对比 kafka_exporter_up 与 kafka_up |
Recording Rule
Kafka 规则文件位于 /infra/rules/kafka.yml。主要记录指标如下:
| 指标 | 含义 |
|---|
kafka:topic:msg_rate1m/5m | Topic 当前 Offset 的 1/5 分钟正向变化速率 |
kafka:cls:msg_rate1m/5m | 去重后的集群消息追加速率 |
kafka:csg_topic:commit_rate5m | Consumer Group/Topic 的 5 分钟提交进展速率 |
kafka:csg_topic:lag | Consumer Group/Topic 的总 Lag |
kafka:csg:lag | Consumer Group 跨 Topic 的总 Lag |
kafka:cls:lag | Kafka 集群全部 Consumer Group 的总 Lag |
kafka:ins:jvm_heap_used_ratio | Kafka JVM Heap 使用率 |
kafka:ins:jvm_cpu_cores | Kafka JVM 消耗的 CPU Core 数 |
kafka:ins:load / kafka:cls:load | 实例最忙请求线程池与集群平均负载 |
kafka:ins:jvm_gc_time_rate5m | 5 分钟 GC 时间速率 |
kafka:ins:messages_in_rate5m | Broker 5 分钟消息接收速率 |
kafka:ins:bytes_in_rate5m | Broker 5 分钟客户端入站字节速率 |
kafka:ins:bytes_out_rate5m | Broker 5 分钟客户端出站字节速率 |
kafka:ins:request_error_rate5m | Broker 5 分钟请求错误速率 |
kafka:cls:under_replicated_partitions | 集群 Under Replicated Partition 总数 |
kafka:cls:offline_partitions | 集群 Offline Partition 数 |
基于 Offset 变化得到的是进展速率,不是客户端请求数。日志截断、Offset 回退或 Exporter 重启可能造成瞬时负变化;规则使用 clamp_min(..., 0) 只保留正向进展。
告警规则
| 告警 | 条件 | 持续时间 | 级别 | 首选下钻 |
|---|
KafkaDown | up{job="kafka",role=~".+"} < 1 | 1m | CRIT | Kafka Instance / ins |
KafkaExporterDown | up{job="kafka",role=""} < 1 | 1m | CRIT | Kafka Instance / ins |
KafkaJmxScrapeError | jmx_scrape_error{job="kafka"} > 0 | 3m | WARN | Kafka Instance / JMX Collector |
KafkaJvmHeapHigh | Heap 使用率 > 90% | 15m | WARN | Kafka Instance / JVM Memory |
KafkaJvmDeadlock | JVM Deadlocked Thread > 0 | 1m | CRIT | Kafka Instance / JVM Threads |
KafkaRequestHandlerSaturated | Handler Idle < 10% | 10m | WARN | Kafka Instance / Request Path |
KafkaNetworkProcessorSaturated | Network Processor Idle < 10% | 10m | WARN | Kafka Instance / Request Path |
KafkaUnderReplicatedPartitions | URP > 0 | 5m | WARN | Kafka Instance / Replication |
KafkaUnderMinISR | Under Min ISR > 0 | 1m | CRIT | Kafka Instance / Replication |
KafkaOfflineLogDirectory | Offline Log Directory > 0 | 1m | CRIT | Kafka Instance / Disk Pressure |
KafkaOfflinePartitions | Controller Offline Partition > 0 | 1m | CRIT | Kafka Overview / cls |
KafkaControllerCountMismatch | Active Controller 数不等于 1 | 1m | CRIT | Kafka Overview / cls |
KafkaFencedBrokers | Fenced Broker > 0 | 5m | WARN | Kafka Overview / cls |
KafkaUncleanLeaderElection | 5 分钟出现不干净 Leader 选举 | 立即 | CRIT | Kafka Overview / cls |
KafkaConsumerLagGrowing | Group Lag > 100000 且 30 分钟仍增长 | 30m | WARN | Kafka Consumer / group |
不干净 Leader 选举可能意味着数据丢失,应立即保留 Controller/Broker 日志,确认受影响 Topic 与副本,再决定恢复动作。
常用 PromQL
检查采集目标:
kafka_up
kafka_exporter_up
up{job="kafka"}
检查某集群复制健康:
sum by (cls) (kafka_server_replica_manager_under_replicated_partitions{job="kafka"})
sum by (cls) (kafka_server_replica_manager_under_min_isr_partitions{job="kafka"})
max by (cls) (kafka_controller_offline_partition_count{job="kafka"})
检查 Consumer Lag:
topk(20, kafka_consumergroup_lag_sum{cls="kf-main"})
检查请求饱和与延迟:
kafka_server_request_handler_idle_ratio{job="kafka",cls="kf-main"}
max by (ins,request,quantile) (
kafka_network_request_total_time_seconds{job="kafka",cls="kf-main",quantile=~"0.95|0.99"}
)
日志查询
Kafka 服务把标准输出与错误写入 Journald,节点 Vector 的 Journald Source 会转发到 VictoriaLogs,统一使用 job:syslog。
job:syslog unit:kafka
job:syslog app:kafka
job:syslog unit:kafka_exporter
ip:10.10.10.11 job:syslog (unit:kafka OR app:kafka)
Kafka Instance Dashboard 的日志面板使用类似查询,并展示时间、级别、Systemd Unit 与消息。诊断时应把日志与同一时间窗口内的 KRaft、ISR、请求队列、GC、磁盘 I/O 和网络指标对齐。
验证监控链路
在 Kafka 节点验证原始端点:
curl -fsS http://<kafka-ip>:9404/metrics | grep '^jmx_scrape_error'
curl -fsS http://127.0.0.1:9308/metrics | grep '^kafka_brokers'
在 Infra 节点检查文件发现(每实例一个文件,被选中节点的文件含 JMX 与协议 Exporter 两个目标):
ls -l /infra/targets/kafka/
cat /infra/targets/kafka/kf-main-1.yml
然后在 VictoriaMetrics 查询 up{job="kafka"}(或记录指标 kafka_up 与 kafka_exporter_up)。自定义 exporter 指标在抓取失败后可能短暂保留旧样本,端点存活应以 Prometheus 原生 up 为准。若原始端点正常但记录指标缺失,依次检查文件发现、VictoriaMetrics Target、网络可达性、规则加载与标签;若 JMX HTTP 正常但 jmx_scrape_error 为 1,检查 Kafka 日志和 /etc/kafka/jmx_exporter.yml 的 MBean 匹配情况。
完整指标语义参阅 指标定义。
7 - 指标定义
Kafka JMX、协议 Exporter 与 Recording Rule 指标字典。
KAFKA 模块使用两类指标源,都注册在同一 job=kafka 采集任务下:JMX 目标(带 role 标签)采集每个 JVM 的内部状态;协议 Exporter 目标(无 role 标签)通过 Kafka 协议采集逻辑集群、Topic、Partition 与 Consumer Group 状态。协议 Exporter 只放在 kafka_seq 最小的至多两个 Broker-capable 节点上,单 Broker 集群只运行一个。
JMX 配置采用白名单,只导出 JVM 基线和有界的 Broker、复制、请求路径与 KRaft 指标;高基数的 per-client 与 per-partition JMX MBean 被有意排除,Partition 详情由协议 Exporter 提供。
公共标签
| 指标源 | 公共标签 |
|---|
JMX 目标(:9404) | job, cls, ins, ip, instance, role, node_id |
协议 Exporter 目标(:9308) | job, cls, ins, ip, instance |
两类目标的 job 都是 kafka;是否携带 role 标签是区分两类序列的依据。
部分指标还有 topic、partition、broker、consumergroup、request、version、error、quantile、state 或 operation 等维度。
可用性与抓取指标
| 指标 | 类型 | 含义 |
|---|
kafka_up | Gauge/Recording | JMX 目标抓取可用性:up{job="kafka",role=~".+"} |
kafka_exporter_up | Gauge/Recording | 协议 Exporter 目标抓取可用性:up{job="kafka",role=""} |
up | Gauge | VictoriaMetrics 对原始 Target 的抓取状态 |
jmx_scrape_error | Gauge | JMX Exporter 最近一次抓取是否出错,健康值为 0 |
jmx_scrape_duration_seconds | Gauge | JMX 抓取耗时 |
jmx_scrape_cached_beans | Gauge | JMX Exporter 缓存的 MBean 数量 |
scrape_duration_seconds | Gauge | VictoriaMetrics 抓取 Exporter 的耗时 |
scrape_samples_scraped | Gauge | 本次抓取的样本数量 |
协议 Exporter 指标
以下指标来自协议 Exporter 目标。同一集群的多个 Exporter 会看到相同的逻辑集群状态,直接做集群聚合时必须按语义去重,不能简单把所有 ins 相加。
Broker 与 Topic
| 指标 | 类型 | 关键维度 | 含义 |
|---|
kafka_brokers | Gauge | 集群 | Exporter 发现的 Broker 数量 |
kafka_broker_info | Gauge | id, address 等 | Broker 信息,以值 1 携带标签 |
kafka_topic_partitions | Gauge | topic | Topic 的 Partition 数量 |
kafka_topic_partition_current_offset | Gauge | topic, partition | Partition 当前 Log End Offset |
kafka_topic_partition_oldest_offset | Gauge | topic, partition | Partition 当前最早可读 Offset |
kafka_topic_partition_leader | Gauge | topic, partition | 当前 Leader Broker ID;无 Leader 时用于识别异常 |
kafka_topic_partition_replicas | Gauge | topic, partition, broker | 分配给 Partition 的副本集合 |
kafka_topic_partition_in_sync_replica | Gauge | topic, partition, broker | 当前 ISR 成员 |
kafka_topic_partition_under_replicated_partition | Gauge | topic, partition | Partition 是否处于副本不足状态 |
kafka_topic_partition_leader_is_preferred | Gauge | topic, partition | 当前 Leader 是否为 Preferred Replica |
current_offset - oldest_offset 可以估计当前可保留的 Offset Span,但 Offset 数量不等于字节数,Compact Topic 也不等于精确消息条数。
Consumer Group
| 指标 | 类型 | 关键维度 | 含义 |
|---|
kafka_consumergroup_members | Gauge | consumergroup | Group 当前成员数 |
kafka_consumergroup_current_offset | Gauge | consumergroup, topic, partition | Group 已提交 Offset |
kafka_consumergroup_current_offset_sum | Gauge | consumergroup, topic | 已提交 Offset 汇总 |
kafka_consumergroup_lag | Gauge | consumergroup, topic, partition | Partition 级消费滞后 |
kafka_consumergroup_lag_sum | Gauge | consumergroup, topic | Group/Topic 消费滞后汇总 |
没有提交 Offset 的临时消费者、使用外部 Offset 存储的客户端,或尚未消费某 Topic 的 Group,不一定产生这些时间序列。
Exporter 自身
| 指标 | 类型 | 含义 |
|---|
kafka_exporter_build_info | Gauge | Exporter 版本、Revision 与构建信息 |
process_* | Gauge/Counter | Exporter 进程 CPU、内存、FD、启动时间等 |
go_* | Gauge/Counter | Exporter Go Runtime、GC、Goroutine 与内存状态 |
promhttp_metric_handler_* | Counter | /metrics 请求处理状态 |
JMX:JVM 基线
excludeJvmMetrics: false 使 JMX Exporter 暴露标准 JVM/进程指标。Kafka Instance Dashboard 主要使用:
| 指标 | 含义 |
|---|
jvm_memory_used_bytes | 按 Heap/Non-Heap 与 Memory Pool 划分的已用内存 |
jvm_memory_committed_bytes | JVM 已提交内存 |
jvm_memory_max_bytes | JVM 可用最大内存 |
jvm_gc_collection_seconds_count | GC 次数 |
jvm_gc_collection_seconds_sum | GC 累计耗时 |
jvm_threads_state | 按线程状态统计的线程数 |
jvm_threads_deadlocked | 检测到的死锁线程循环数 |
jvm_buffer_pool_used_bytes | Direct/Mapped Buffer Pool 使用量 |
process_cpu_seconds_total | Kafka JVM 累计 CPU 时间 |
process_open_fds / process_max_fds | 已打开与最大文件描述符 |
process_start_time_seconds | Kafka JVM 启动时间 |
JMX:Broker 流量
| 指标 | 类型 | 含义 |
|---|
kafka_server_broker_messages_in_total | Counter | Broker 接收的消息总数 |
kafka_server_broker_bytes_in_total | Counter | Broker 接收的客户端字节总数 |
kafka_server_broker_bytes_out_total | Counter | Broker 发送的客户端字节总数 |
kafka_server_broker_replication_bytes_in_total | Counter | Broker 接收的复制字节总数 |
kafka_server_broker_replication_bytes_out_total | Counter | Broker 发送的复制字节总数 |
kafka_server_broker_produce_requests_total | Counter | Produce 请求总数 |
kafka_server_broker_failed_produce_requests_total | Counter | 失败 Produce 请求总数 |
kafka_server_broker_fetch_requests_total | Counter | Fetch 请求总数 |
kafka_server_broker_failed_fetch_requests_total | Counter | 失败 Fetch 请求总数 |
这些是 Broker 总量,不包含 Topic 维度,避免 JMX Series 随 Topic 数膨胀。Topic 级 Offset 与进展来自协议 Exporter。
JMX:复制与存储
| 指标 | 类型 | 含义 |
|---|
kafka_server_replica_manager_under_replicated_partitions | Gauge | ISR 少于已分配副本的 Partition 数 |
kafka_server_replica_manager_under_min_isr_partitions | Gauge | ISR 低于 min.insync.replicas 的 Partition 数 |
kafka_server_replica_manager_at_min_isr_partitions | Gauge | ISR 恰好等于 min.insync.replicas 的 Partition 数 |
kafka_server_replica_manager_offline_replicas | Gauge | 当前 Broker 上离线副本数 |
kafka_server_replica_manager_partitions | Gauge | 当前 Broker 承载的副本数 |
kafka_server_replica_manager_leaders | Gauge | 当前 Broker 领导的 Partition 数 |
kafka_server_replica_manager_isr_shrinks_total | Counter | ISR 收缩事件总数 |
kafka_server_replica_manager_isr_expands_total | Counter | ISR 扩张事件总数 |
kafka_server_replica_manager_failed_isr_updates_total | Counter | ISR 更新失败总数 |
kafka_server_replica_manager_reassigning_partitions | Gauge | 正在进行 Reassignment 的 Leader Partition 数 |
kafka_server_delayed_operation_purgatory_size | Gauge | 按 operation 划分的延迟操作等待数 |
kafka_log_manager_offline_log_directories | Gauge | Kafka 标记为离线的日志目录数 |
Under Replicated 表示副本没有全部同步;Under Min ISR 更严重,表示写入可用性或持久性条件已经低于设置的最小 ISR。At Min ISR 虽未越线,但已经没有额外副本余量。
JMX:请求路径
| 指标 | 类型 | 额外标签 | 含义 |
|---|
kafka_network_request_total | Counter | request, version | 各 Kafka API 请求总数 |
kafka_network_request_errors_total | Counter | request, error | 各 API/错误码响应错误总数 |
kafka_network_request_total_time_seconds | Gauge | request, version, quantile | API 总耗时 P50/P95/P99 |
kafka_network_request_queue_size | Gauge | - | 等待 Request Handler 的请求数 |
kafka_network_response_queue_size | Gauge | - | 等待 Network Processor 的响应数 |
kafka_server_request_handler_idle_ratio | Gauge | - | Request Handler 平均空闲比例 |
kafka_network_processor_idle_ratio | Gauge | - | Network Processor 平均空闲比例 |
排查高延迟时,应同时查看请求量、错误码、P95/P99、两个队列、Handler/Processor Idle、GC、CPU、磁盘 I/O 与网络。单独看到低 Idle 不足以判断瓶颈位置。
JMX:KRaft 与 Broker 元数据
| 指标 | 类型 | 含义 |
|---|
kafka_server_raft_state | Gauge | 当前成员的 KRaft 状态,以 state 标签表示 |
kafka_server_raft_current_leader | Gauge | 当前 KRaft Leader Node ID,-1 表示未知 |
kafka_server_raft_current_epoch | Gauge | 当前 KRaft Epoch |
kafka_server_raft_high_watermark | Gauge | 元数据日志 High Watermark |
kafka_server_raft_log_end_offset | Gauge | 元数据日志 Log End Offset |
kafka_server_broker_metadata_last_applied_record_lag_seconds | Gauge | Broker 应用元数据记录的时间滞后 |
kafka_server_broker_metadata_load_errors_total | Counter | Broker 加载元数据错误总数 |
kafka_server_broker_metadata_apply_errors_total | Counter | Broker 应用元数据镜像错误总数 |
kafka_server_metadata_snapshot_bytes | Gauge | 最近生成或加载的元数据 Snapshot 大小 |
kafka_server_metadata_snapshot_age_seconds | Gauge | 最近元数据 Snapshot 的年龄 |
log_end_offset - high_watermark 可辅助判断元数据提交滞后;还应结合成员角色、当前 Leader、Epoch 和 Controller 事件延迟判断。
JMX:Controller
这些 MBean 只存在于带 Controller 角色的 Kafka 进程中:
| 指标 | 类型 | 含义 |
|---|
kafka_controller_active_controller_count | Gauge | Active Controller 上为 1,其他 Controller 为 0 |
kafka_controller_fenced_broker_count | Gauge | Active Controller 观察到的 Fenced Broker 数 |
kafka_controller_active_broker_count | Gauge | Active Broker 数 |
kafka_controller_global_topic_count | Gauge | Controller 观察到的 Topic 数 |
kafka_controller_global_partition_count | Gauge | Controller 观察到的 Partition 数 |
kafka_controller_offline_partition_count | Gauge | 离线的非内部 Partition 数 |
kafka_controller_preferred_replica_imbalance_count | Gauge | Leader 不是 Preferred Replica 的 Partition 数 |
kafka_controller_metadata_errors_total | Counter | Controller 元数据处理错误总数 |
kafka_controller_last_applied_record_lag_seconds | Gauge | Controller 应用元数据记录的时间滞后 |
kafka_controller_timed_out_broker_heartbeats_total | Counter | Broker Heartbeat 超时总数 |
kafka_controller_elections_total | Counter | 本节点观察到的新 Active Controller 选举总数 |
kafka_controller_unclean_leader_elections_total | Counter | 不干净 Leader 选举总数 |
kafka_controller_event_queue_time_seconds | Gauge | Controller 事件排队 P50/P95/P99 |
kafka_controller_event_processing_time_seconds | Gauge | Controller 事件处理 P50/P95/P99 |
健康集群应恰好存在一个 Active Controller。offline_partition_count、metadata_errors_total 与 unclean_leader_elections_total 的增加都应优先处理。
Recording Rule 指标
Offset 进展
| 指标 | 聚合层级 | 窗口 | 含义 |
|---|
kafka:topic:msg_rate1m | Topic | 1m | Exporter 间去重后的 Current Offset 正向增长速率 |
kafka:topic:msg_rate5m | Topic | 5m | Exporter 间去重后的 Current Offset 正向增长速率 |
kafka:cls:msg_rate1m | 逻辑集群 | 1m | Exporter 间去重后的消息追加速率 |
kafka:cls:msg_rate5m | 逻辑集群 | 5m | Exporter 间去重后的消息追加速率 |
kafka:csg_topic:commit_rate5m | Group/Topic | 5m | Commit Offset 正向增长速率 |
kafka:csg_topic:lag | Group/Topic | 当前值 | Partition Lag 去重后汇总 |
kafka:csg:lag | Consumer Group | 当前值 | Group 跨 Topic 总 Lag |
kafka:cls:lag | 逻辑集群 | 当前值 | 集群跨 Consumer Group 总 Lag |
JVM 与 Broker
| 指标 | 含义 |
|---|
kafka:ins:jvm_heap_used_ratio | Heap Used / Heap Max |
kafka:ins:jvm_cpu_cores | 5 分钟 JVM CPU Core 消耗 |
kafka:ins:load | 实例最忙请求线程池的饱和度 |
kafka:cls:load | 集群实例平均负载 |
kafka:ins:jvm_gc_time_rate5m | 5 分钟 GC 时间速率 |
kafka:ins:messages_in_rate5m | 5 分钟 Broker 消息接收速率 |
kafka:ins:bytes_in_rate5m | 5 分钟 Broker 客户端入站字节速率 |
kafka:ins:bytes_out_rate5m | 5 分钟 Broker 客户端出站字节速率 |
kafka:ins:request_error_rate5m | 5 分钟非 NONE 请求错误速率 |
kafka:cls:under_replicated_partitions | 集群 Under Replicated Partition 总数 |
kafka:cls:offline_partitions | 集群 Offline Partition 数 |
基数与解释注意事项
- 不要把同一
cls 的多个 kafka_exporter 结果直接求和;它们可能是同一集群视图的副本。 kafka_topic_partition_current_offset 是 Offset,不是精确字节、请求或业务事件数量。- Consumer Lag 只覆盖 Kafka 中可见且已提交 Offset 的 Group。
- 纯 Controller 缺少 Broker 指标和协议 Exporter 指标属于正常角色差异;未被选择的 Broker 没有协议 Exporter 指标也属于正常放置结果。
- 某个 MBean 在具体 Kafka 版本/角色中不存在时,对应 JMX Series 也不会出现;应结合
role 判断。 - per-client/per-partition JMX 指标被白名单排除,以避免不可预测的时间序列基数。
Dashboard 与告警使用方式参阅 监控告警。
8 - 常见问题
Pigsty Kafka 4.1+ 动态 KRaft 模块常见问题与故障排查。
当前 KAFKA 模块是什么成熟度?
当前角色已实现生产级 v1 基线:动态 KRaft、完整集群护栏、冷启动/修复、Broker 串行准入与 Controller 动态加入、成员退役(含死节点)、故障节点三步替换、严格滚动、TLS/SCRAM/ACL、Topic/User 声明式收敛、内部凭据/证书轮换以及完整监控链路。
它不是托管 Kafka 产品。生产仍需使用 kafka_security: scram、奇数 Controller、足够 Broker/RF/minISR,并补充容量规划、Reassignment/数据均衡、升级、备份、恢复与故障演练。默认 plaintext 只适合开发或可信隔离网络。
为什么没有 ZooKeeper,也没有 controller.quorum.voters?
本模块面向 Kafka 4.1+,使用原生动态 KRaft,不安装 ZooKeeper,也不创建静态 Quorum。所有成员渲染 controller.quorum.bootstrap.servers;新集群显式使用 --initial-controllers/--no-initial-controllers 格式化,启动后角色会校验初始 Controller 的 Directory ID 已进入现场 Quorum。
初始 Controller Identity 写入 Bootstrap Manifest,但它只是"出生证明":集群首次 Commission 之后,现场 Quorum 的成员关系以 Raft 自身为准。后续 Controller 的增删由剧本编排完成——新增走 kafka.yml 的 Observer 追平 + add-controller 加入流程,删除走 kafka-rm.yml 真子集退役(自动 remove-controller)——你只需要编辑 inventory 并运行对应剧本。
combined、broker、controller 有什么区别?
combined:同时承担 Broker 与 Controller,监听 9092 和 9093,是默认值;broker:纯数据面,只监听 9092;controller:纯控制面,只监听 9093。
集群角色要么全部省略并一致使用 combined,要么全部显式声明。不再提供旧角色别名。
Controller 端口 9093 会和 Alertmanager 冲突吗?
不冲突。Pigsty 的 Alertmanager 监听 alertmanager_port 9059,集群端口为 9094,与 KRaft Controller 的惯例端口 9093 错开。若你改动过这些端口而发生碰撞,为该集群调整 kafka_controller_port 即可——角色只强制 9092、9093、9308、9404 四者互不相同,不会检测与其他服务的端口占用。
服务已启动,但远程客户端连不上?
Broker 的 advertised.listeners 固定使用 inventory_hostname。客户端连接 Bootstrap Server 后,还必须解析并访问元数据返回的每一个 Broker 地址。
依次检查:
grep '^advertised.listeners' /etc/kafka/server.properties
ss -lntp | grep ':9092'
getent hosts <inventory-hostname>
scram 客户端还要检查 CA、SASL mechanism、用户名/密码与 ACL。当前 v1 不提供自定义 advertised address、多 Listener 或 NAT/公网映射;如果客户端不能直接路由 inventory_hostname,该网络模型不在当前核心契约内,不能用 kafka_parameters 覆盖 raw listener 绕过。
为什么提示 Cluster ID、Node ID 或 Directory ID 不匹配?
角色会交叉校验 Bootstrap Manifest、${kafka_data}/metadata/meta.properties、inventory 与现场动态 Quorum。常见原因包括:
- 修改了
kafka_cluster 或 kafka_seq; - 把其他集群的数据盘挂载到当前节点;
- 恢复/接管时给出了错误的
kafka_cluster_id; - Controller 数据目录或 Directory ID 与现场 Voter 记录不一致;
- 选错了目标集群或使用了过期 Manifest。
这是保护性失败。不要删除 meta.properties、Manifest 或直接执行 kafka-rm.yml。先确认数据归属、剩余副本、真实 Cluster/Node/Directory Identity 与恢复目标。
Manifest 丢失或只剩旧 Manifest 会怎样?
每个集群成员都保留一份 Manifest 权威副本 /etc/kafka/manifest.yml(scram 集群另有 /etc/kafka/secrets.yml),管理节点不保存任何 Kafka 状态,每次运行时从任一成员副本解析,因此换管理节点或丢失本地检出都不影响集群管理。只有当所有成员的副本都丢失、而存储已经格式化时,角色才失败关闭并提示先在任一成员上恢复该文件;已格式化的 scram 集群在所有成员都找不到 Secret 副本时同样失败关闭。签发的节点证书缓存在 files/pki/kafka/,丢失时直接由 Pigsty CA 重签。
反过来,如果 Manifest 存在而全部 Kafka 数据盘为空,角色会失败关闭,避免用旧身份意外复活已消失的集群。确实要重建时必须先执行 kafka-rm.yml 和明确的重建流程。
为什么 kafka_parameters 中的某些键被拒绝?
身份、动态 Quorum、Listener、存储、复制、Rack 与安全必须保持单一权威,因此这些键由角色拥有:出现任意一个,身份预检都会在写文件前失败。完整保留列表见 kafka_parameters。
请改用对应的公开参数。角色不提供地址、路径子目录、Listener Map 或 Exporter options 变量。
如何启用 TLS、SCRAM 与 ACL?
新集群设置:
这会一次启用 Pigsty CA 节点证书、Controller mTLS、Broker/client SASL_SSL + SCRAM-SHA-512、StandardAuthorizer 与默认拒绝。应用用户通过 kafka_users 声明密码、ACL 和可选 Quota。
安全模式是 Bootstrap-only 属性。已格式化集群不能通过普通剧本从 plaintext 在线切换到 scram;这需要独立迁移状态机。健康 scram 集群可以使用受保护动作轮换内部凭据或证书。
kafka_topics 与 kafka_users 会删除资源吗?
不会因为从清单移除条目而隐式删除 Topic 或用户。
Topic 会幂等创建、Partition 只增加、只更新声明的配置;RF 变化要求显式 Reassignment。声明用户会收敛密码、完整 ACL 集合与给出的 Quota 字段。Topic 删除、用户删除或彻底撤权都是独立受审操作。
JMX Exporter 与 kafka_exporter 有什么区别?
JMX Exporter 注入每个 Kafka JVM,采集 JVM、Broker、复制、请求路径与 KRaft 内部指标,注册为带 role 标签的 job=kafka 目标。
kafka_exporter 通过 Kafka 协议查询逻辑集群、Topic、Partition、Offset、Consumer Group 与 Lag,注册为同一 job=kafka 下不带 role 标签的目标。角色只在按 kafka_seq 排序后的前两个 Broker-capable 节点运行;单 Broker 集群运行一个,纯 Controller 不运行。
两者互补。生命周期健康门禁使用角色自有 Kafka CLI/metadata 通道,不依赖任一 Exporter。
为什么某个 Broker 或纯 Controller 没有 kafka_exporter?
这是预期的派生放置。协议 Exporter 返回的是整个逻辑集群视图,不是节点指标;最多两个副本可以避免监控单点,同时控制重复采集成本。
检查当前目标(每实例一个文件,被选中节点的文件里含 :9308 的协议 Exporter 目标):
ls -l /infra/targets/kafka/
grep 9308 /infra/targets/kafka/*.yml
完整运行会按当前放置刷新每个实例的 Target 文件,不应只针对单节点运行注册标签。注意:若 Exporter 放置因拓扑变化而转移,曾被选中节点上的旧 kafka_exporter 服务不会被普通剧本自动停止,需要手工或通过 kafka-rm.yml 清理。
为什么 JMX 端点可访问,但 jmx_scrape_error=1?
HTTP 可访问只说明 Java Agent 已加载;jmx_scrape_error=1 表示本轮 MBean 采集失败:
journalctl -u kafka --since '-30 min' --no-pager
curl -fsS http://<kafka-ip>:9404/metrics | head -n 40
检查 /etc/kafka/jmx_exporter.yml 与当前 Kafka/JMX Exporter 包是否匹配,以及 JVM 是否已经过 startDelaySeconds。真实启动验收要求 jmx_scrape_error 0.0、JVM 指标和至少一项与角色匹配的 kafka_ 指标。
为什么 Consumer Lag 没有数据?
常见原因:Consumer 没使用 Group、未向 Kafka 提交 Offset、把 Offset 存在外部系统、Group 尚未消费目标 Topic,或协议 Exporter 的 TLS/SCRAM/ACL/网络异常。
/opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server <broker>:9092 \
--command-config /etc/kafka/admin.properties \
--describe --group <group>
再检查 kafka_exporter_up、Exporter 日志、Dashboard 变量和原始 kafka_consumergroup_* 指标。端点存活以 Prometheus 原生 up 为准,不要用抓取失败后可能短暂保留的自定义指标代替。
为什么两个 kafka_exporter 的集群指标不能相加?
两个 Exporter 查询同一逻辑集群,可能返回相同 Topic/Partition/Consumer Group 状态;直接求和会重复计算。Pigsty 的 kafka:cls:* Recording Rule 会先跨 Exporter 副本去重,再聚合到集群。
应用要经过 HAProxy、Keepalived VIP 或 LB 吗?
不要。Kafka Producer/Consumer 是集群感知的智能客户端:连上 bootstrap.servers 中任一种子取得元数据后,它直接连接各 Partition Leader。VIP 或通用 TCP LB 既不理解 Partition Leader,也不会改写元数据中的 Broker 地址,放在数据面只会增加长连接状态、故障点与排障复杂度。
若平台强制要求统一发现入口,DNS 或 TCP LB 可以只承担 bootstrap,但 advertised.listeners 仍返回每个 Broker 的可达地址,应用网络必须直达全部 Broker。跨 NAT、公网、多网络或 Kubernetes 暴露需要为每个 Broker 设计独立外部地址与额外 Listener,当前模块固定宣告清单地址,不支持这类映射。
详见快速上手:为什么应用应直连多个 Broker与集群配置:网络与监听器。
可以直接增删 Broker 或 Controller 吗?
可以。编辑 inventory 后由剧本编排完成 KRaft 成员变更的全部步骤:
- 增加:在 inventory 中声明新成员(
broker、combined、controller 均可),以完整集群为目标运行 ./kafka.yml -l <cls>(不能只 -l 新节点)。纯 Broker 逐个格式化、启动并验证注册;Combined/Controller 以 --no-initial-controllers 格式化,Observer 追平后 add-controller 提升为 Voter。全程逐节点、全程健康门禁。 - 移除:
./kafka-rm.yml -l <ip>(集群真子集)经幸存成员执行 remove-controller 与 Broker 注销,节点不可达也能完成,随后从 inventory 删除该成员。
仍需自行保证:变更后 Controller 保持奇数且多数派存活;一次只做一个方向的成员变更;被移除 Broker 上的 Partition 副本先行排空(或由同 kafka_seq 的替换节点接管)。加入后既有 Partition 不会自动迁移,需独立执行并监控 Reassignment——“Broker 已注册”不等于“容量已均衡”。
软件包版本由哪个参数控制?
角色使用 package_map['java-runtime'] 与 package_map['kafka-stack'],不提供 kafka_version、scala_version 或 Exporter 版本参数。实际版本由目标平台的 Pigsty 仓库和已安装包决定。
2026-07-16 验证的载荷为 Kafka 4.3.1、kafka_exporter 1.9.0、JMX Exporter 1.6.0。升级仍需单独评审兼容性、备份/回退、滚动顺序与 Feature Level,不能只替换包。
如何安全清空 Kafka 数据?
kafka.yml 永远不执行清理,删除动作只在独立的 kafka-rm.yml 中:-l 选中整个集群(或裸跑选中全部集群)即为集群下线,选中真子集则是成员退役。默认 kafka_rm_data=true 会永久删除数据/KRaft 元数据、节点上的 /etc/kafka 恢复状态与监控 Target;kafka_rm_data=false 保留数据与恢复状态,kafka_safeguard=true 中止一切删除。
该剧本没有确认字符串等额外闸门,执行前必须人工确认精确 -l 目标、可恢复备份或明确重建意图与业务停用状态。完整语义见 预置剧本:kafka-rm.yml。