这是本节的多页打印视图。 .
模块:KAFKA
Kafka 是一个分布式事件流平台。Pigsty 的 KAFKA 模块使用 RPM/DEB 软件包,在纳管节点上部署 Apache Kafka 4.1+ 动态 KRaft 集群,并统一管理安全、资源、生命周期与可观测性。
当前 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 - 快速上手
本教程从一个最小单节点集群开始,完成 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:
这个配置会得到:
- 一个随机 Cluster ID;
- 一个动态 KRaft combined 节点;
- 默认 RF=1、minISR=1;
- 一个名为
quickstart.events的单 Partition Topic; - JMX Exporter
:9404与一个协议 Exporter:9308。
plaintext 没有传输加密、认证和 ACL,只能用于本机开发或可信隔离网络。
2. 纳管节点
如果该主机尚未完成 NODE 初始化,先执行检查模式:
审阅结果并取得批准后再纳管节点:
已经由 Pigsty 纳管、软件仓库和时间同步均正常的节点可以跳过这一步。NODE 的完整准备与日常管理见 节点管理。
3. 部署 Kafka
先对完整集群执行检查:
确认目标确实只有 kf-dev 的完整成员,审阅数据路径、软件包、端口与配置变化后执行:
角色会安装 Java 与 kafka-stack、生成随机身份和 Bootstrap Manifest、格式化 KRaft 存储、启动服务、创建 Topic,并注册监控目标。
4. 验证服务与 Quorum
登录 Kafka 节点,检查服务:
使用角色自有健康检查:
返回 JSON 中应有 "healthy": true。继续检查动态 Quorum 与 Topic:
应看到有效 LeaderId、包含本节点的 CurrentVoters,以及 RF=1、ISR=1 的 quickstart.events。
5. 生产与消费消息
启动 Console Producer:
输入几行消息后按 Ctrl-D 结束。在另一个终端消费:
到这里,单节点部署、Topic 收敛和消息读写已经完成。进一步的状态检查见 日常管理。
二、部署三节点安全 HA 演示基线
三节点示例是一套全新的 kf-main 集群,使用三个 combined 节点。它可以容忍一个 Controller 故障;业务 Topic 使用 RF=3/minISR=2,并启用 scram 生产安全档位。
1. 定义安全集群与资源
将以下组加入现有 all.children:
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. 纳管并部署
如果节点尚未纳管:
部署 Kafka 时必须选择全部三个成员:
不能只 -l 10.10.10.11:每个被选中的集群必须完整,部分选择会被拒绝。同时选择多个完整集群(-l kf-dev,kf-main)或不加 -l 裸跑全部集群则是允许的。
3. 验证三节点健康
从管理节点检查三个 Kafka 服务:
在任一 Broker 上执行完整健康检查:
查询 quorum 和 Topic:
上线前应看到:一个 Active Controller、三个 Current Voters、三个可用 Broker;所有 quickstart.events Partition 均有三副本、ISR=3,没有 Offline、Under Replicated 或 Under Min ISR Partition。
三、接入应用客户端
1. 分发 CA 公钥证书
将管理节点上的公共 CA 证书安全复制到应用主机:
ca.crt 是可以分发的公钥证书。绝不要复制、暴露或分发 files/pki/ca/ca.key。 应用主机上的 CA 文件建议由 root 管理并设为只读。已被 Pigsty 纳管的应用主机无需复制:NODE 模块已把同一 CA 安装在 /etc/pki/ca.crt,客户端可直接引用。
2. 创建客户端配置
在应用主机创建 /etc/kafka-client/client.properties:
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 的应用主机上执行:
消费时使用 ACL 允许的 Group 前缀:
生产应用还应显式评审客户端语义:
| 客户端配置 | 建议起点 | 说明 |
|---|---|---|
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 小时:
不要照抄 6G/8/24;这些值必须由 CPU、内存、连接数、消息大小、Partition 数、磁盘和 Page Cache 压测决定。
示例:增加 Partition 并缩短 Topic 保留
把 quickstart.events 从 12 个 Partition 增加到 24,并把保留时间改成三天:
Partition 不能减少。replication_factor 与现场不一致时,角色会拒绝普通收敛并要求显式 Partition Reassignment;不会自动搬迁既有副本。
应用变更
无论修改静态参数还是动态资源,都运行完整状态机:
不要只运行 -t kafka_config。角色会自动判断:静态变化执行严格逐节点滚动;只修改 Topic/User 等动态资源时不重启 Kafka。
以下键属于角色自身,不能放入 kafka_parameters:
全部 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 模块使用 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,无法容忍节点故障,只适合开发、测试与功能验证:
角色会从初始 Broker 数量推导 RF=1、minISR=1。不要把单节点拓扑或默认 plaintext 安全模式直接用于生产。
三节点复合部署
三个节点都承担 Broker 与 Controller,是紧凑的生产起点。省略全部角色字段即可使用默认 combined:
初始三个 Broker 会自动得到 RF=3、minISR=2 的角色自有复制策略,无需也不允许在 kafka_parameters 中覆盖内部 Topic RF、default.replication.factor 或 min.insync.replicas。示例中的 4G Heap 只是写法示意;生产应通过压测平衡 JVM Heap、操作系统 Page Cache 与同机其他进程。
Controller 与 Broker 分离
关键或较大集群可以把控制面与数据面分离。因为存在显式角色,所有成员都必须声明角色:
纯 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:
| 参数 | 默认值 | 用途 |
|---|---|---|
kafka_port |
9092 |
Broker、客户端与 Broker 间通信 |
kafka_controller_port |
9093 |
KRaft Controller 仲裁 |
kafka_exporter_port |
9308 |
协议 Exporter HTTP 指标 |
kafka_jmx_exporter_port |
9404 |
JMX Exporter HTTP 指标 |
固定监听器约定如下:
- 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:
Broker-capable 节点必须全部设置或全部省略 Rack。修改 Rack 会触发安全滚动,但不会自动迁移既有副本。
复制策略
首次 Bootstrap 根据初始 Broker 数量派生:
初始的未来 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 等非角色自有键。
以下模式由角色拥有,禁止覆盖:
出现任一保留键时,身份预检会在写文件前直接失败。
安全与声明式资源
kafka_security: scram 是一个完整生产档位,而不是一组可任意组合的开关。它自动启用:
- Pigsty CA 签发的每节点证书;
- Controller listener 双向 TLS;
- Broker/client 与 Broker 间 SASL_SSL + SCRAM-SHA-512;
StandardAuthorizer、默认拒绝,以及角色自有管理/监控身份;- 在协议 Exporter 启动前收敛其最小监控 ACL。
应用资源由两个领域对象声明:
资源收敛语义: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 项持久参数。拓扑、Listener、安全实现、存储子目录、复制安全与 Exporter 放置等细节由角色统一推导,不能作为额外持久变量覆盖。
参数概览
| 参数 | 层级 | 默认值 | 说明 |
|---|---|---|---|
kafka_cluster |
集群 | 必填 | Kafka 集群身份 |
kafka_seq |
实例 | 必填 | 集群内唯一 KRaft node.id |
kafka_role |
实例 | combined |
combined、broker 或 controller |
kafka_cluster_id |
集群 | 未设置 | 接管/恢复断言;新集群随机生成 |
kafka_data |
实例 | /data/kafka |
角色自有数据根目录 |
kafka_heap_opts |
实例 | -Xms1G -Xmx1G |
Kafka JVM Heap |
kafka_port |
实例 | 9092 |
Broker/client 端口 |
kafka_controller_port |
实例 | 9093 |
KRaft Controller 端口 |
kafka_rack |
实例 | 未设置 | Broker 故障域标签 |
kafka_parameters |
集群/实例 | {} |
非角色自有 Broker 参数 |
kafka_jmx_exporter_port |
实例 | 9404 |
JMX Exporter HTTP 端口 |
kafka_exporter_port |
实例 | 9308 |
协议 Exporter HTTP 端口 |
kafka_security |
集群 | plaintext |
plaintext 或生产 scram 档位 |
kafka_users |
集群 | [] |
用户凭据、ACL 与 Quota |
kafka_topics |
集群 | [] |
声明式 Topic |
kafka_cluster 与 kafka_seq 必须定义;kafka_role 有真实默认值。集群角色要么全部省略,要么全部显式声明。
身份与拓扑
kafka_cluster
必填的集群身份。必须以字母或数字开头,只能包含字母、数字、下划线和连字符:
它用于发现完整集群成员、生成实例名和定位 Bootstrap Manifest。每次 kafka.yml 生命周期操作必须用精确 -l 选择该集群的全部成员。
kafka_seq
必填的非负整数,在同一 kafka_cluster 中唯一,直接成为 KRaft node.id:
实例名派生为 ${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:
普通新建集群不要设置。角色会随机生成 Cluster ID,并写入每个成员的 /etc/kafka/manifest.yml。该参数不会重新标记现有数据;与 Manifest 或 meta.properties 冲突时会失败关闭。
kafka_rack
可选的 Broker 故障域标签,渲染为 broker.rack:
所有 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,默认:
生产应根据负载与内存压测设置,通常保持 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_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;非对象条目或未知顶层字段会在资源收敛前失败:
约束:
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;非对象条目或未知顶层字段会在资源收敛前失败:
身份预检只校验 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) |
强制显式 -l;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 模块把 Kafka 安装在 /opt/kafka,使用 Systemd 管理服务,并把持久意图保存在 pigsty.yml。节点上的生成文件不应手工修改。
以下 Kafka CLI 示例都使用角色生成的 /etc/kafka/admin.properties。即使当前是 plaintext 也建议始终保留 --command-config:切换到 scram 管理通道时命令结构不变。将 <broker>:9092 替换为可达的 inventory_hostname 与端口。
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 节点检查服务与最近日志:
协议 Exporter 只在 kafka_seq 最小的至多两个 Broker-capable 节点运行。被选择的节点再检查:
检查监听器与指标端点:
kafka_up 与 kafka_exporter_up 是 VictoriaMetrics 侧的记录指标,不一定出现在原始端点。JMX 端点应包含 jmx_scrape_error 0.0、JVM 指标和与节点角色匹配的 kafka_ 指标。
健康检查
角色的生命周期门禁不依赖 JMX,而是通过同一管理通道检查动态 Quorum、不可用 Partition、副本不足与 Under Min ISR:
返回 JSON 中 healthy: true 才表示该门禁通过。它适合只读诊断,但不能替代业务端到端验证。
该脚本还内置解析回归自检(pigsty-kafka-health selftest),每次剧本运行都会在安装后自动执行;若自检失败说明健康谓词本身不可信,应停止变更并排查。
KRaft 仲裁状态
从任一可用 Broker 查询动态 Quorum:
重点检查:
LeaderId存在且对应预期 Controller;CurrentVoters与预期成员一致(加入中的新节点会先出现在CurrentObservers);MaxFollowerLag与MaxFollowerLagTimeMs没有持续增长;- Dashboard 中恰好有一个 Active Controller。
如需确认动态 Quorum(KIP-853)特性级别,可用 /opt/kafka/bin/kafka-features.sh ... describe 查看 kraft.version。
查看 Controller 复制状态:
如果没有 Leader、成员长期落后或 Voter 集合与预期不一致,应先停止其他变更,保留日志、Manifest 与 meta.properties 证据再分析。死掉的 Voter 用 缩容 或 替换故障节点 流程摘除;不要手工改写 quorum 状态。
管理 Topic
生产 Topic 应优先在 pigsty.yml 的 kafka_topics 中声明:
修改声明后运行剧本收敛:
角色会幂等创建 Topic、只增加 Partition,并只修改声明的配置键。RF 变化会失败并要求显式 Partition Reassignment;从清单中移除条目不会删除 Topic。
只读查看 Topic:
临时或外部管理的 Topic 可以使用 Kafka CLI 创建,但不会自动写回 pigsty.yml。不要让声明式与手工管理同时拥有同一个 Topic。Topic 删除是业务数据删除动作,必须走独立审批、精确名称确认和恢复方案,本文不提供通用删除命令。
管理用户与权限
kafka_security: scram 时,应用身份应通过 kafka_users 管理:
完整剧本会幂等收敛密码、该用户的 ACL 集合与显式给出的 Quota 字段。密码不要以明文提交到仓库或输出到日志。移除用户条目不会自动删除 Principal/凭据;删除或彻底撤权需要独立受审流程。
验证消息读写
使用测试 Topic 做端到端验证。Console Producer/Consumer 使用同一客户端配置文件:
在另一个终端消费:
生产验收应从真实客户端网络执行,覆盖 DNS/advertised.listeners、证书校验、ACL、生产者 ACK、消费提交与端到端延迟,而不只验证 Broker 本机路径。
管理 Consumer Group
列出和查看 Consumer Group:
Lag 要结合消费速率与业务 SLO 判断:短暂积压可能是批处理行为,持续增长且消费速率低于生产速率才表示无法追平。重置 Offset 可能造成重复消费或跳过消息,必须有独立审批、精确 Group/Topic 确认与回放方案。
配置集群
修改 pigsty.yml 后以完整集群为目标执行:
角色根据现场健康和静态指纹自动选择路径:
- 集群不健康或停止:只启动已停止的 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 纳管,然后仍以完整集群为目标:
角色按成员类型自动选择路径,每次只处理一个新节点:
- 纯 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 选择集群的 真子集 即为成员退役(选择整个集群则是 集群下线)。退役会通过一台幸存成员,自动从现场元数据中摘除该节点:
执行内容依次为:注销监控 Target → 停止服务 → remove-controller 摘除 KRaft Voter 条目(若该成员是 Voter;多成员退役时严格串行)→ kafka-cluster.sh unregister 注销 Broker → 清理本机配置与数据(受 kafka_rm_data 控制)。Broker 注销步骤容忍失败,以便重入与处理已失联成员;只有在核对现场 Quorum、Broker 注册、副本健康以及目标本机状态后,才从 pigsty.yml 删除该成员条目。
退役前请自行确认:剩余 Controller 仍构成多数派、保持奇数个 Controller、剩余 Broker 数不低于现有 Topic 的最大 RF。如果被退役 Broker 上仍有 Partition 副本,角色会打印警告:这些 Partition 将保持副本不足,直到同 kafka_seq 的替换节点重新加入(自动继承副本分配并补数据),或你显式执行 Reassignment 将副本迁走。计划内缩容应当先 Reassignment 排空、再退役。
替换故障节点
节点永久损坏(磁盘丢失、机器报废)时,保持其 IP 与 kafka_seq 不变,三步完成补换:
第 ① 步的所有元数据操作都委派给幸存成员执行,因此对已经无法连接的死节点同样有效;它还会一并清理监控 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 变更等数据搬迁操作仍使用独立评审的运行手册。
日志诊断
VictoriaLogs/Grafana 查询:
常见诊断顺序是:服务日志 → 监听端口 → 管理通道健康 → 动态 Quorum → Broker/Partition/ISR → 客户端地址与证书/ACL → Consumer Lag。详细面板与告警映射见 监控告警。
5 - 预置剧本
KAFKA 模块提供两个剧本:kafka.yml 用于部署 Apache Kafka 4.1+ 动态 KRaft 集群并收敛其安全、
资源与监控状态;kafka-rm.yml 用于下线集群或移除成员。
每个被选中的 kafka_cluster 必须包含其全部成员:部分选择会在写入前失败;选择一个集群、多个完整集群或不加 -l 裸跑全部集群都是允许的。先对完全相同的目标执行 --check;真实运行前仍需人工核验备份/重建意图、容量、业务窗口、回退方案与变更批准。
kafka.yml
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 采集任务下:
Target 文件每次完整运行按当前 Exporter 放置刷新;Target 的删除由 kafka-rm.yml 的注销步骤完成。
受保护轮换
轮换变量是一次性 extra-vars,不应写入 pigsty.yml。两种动作互斥,每次只能执行其一;前提是所有成员已格式化、集群健康、安全模式为 scram、角色自有 Secret 材料存在,且 kafka_rotate_confirm 与集群名完全一致。
内部凭据轮换
角色使用 active/standby 内部身份:先通过活管理通道更新非活动凭据,再原子切换本地受保护记录,并进入正常严格滚动。旧 active 保留为下一轮 standby,使中断后的重跑可恢复。
证书轮换
角色废弃共享 PKI 树中已签发的节点证书,用同一 Pigsty CA 为每个节点重新签发私钥与证书,更新节点上的 PEM 证书包并进入严格滚动。新旧证书由同一 CA 签发、彼此互信,因此不需要分阶段互换信任;健康预检失败时不会开始轮换,节点上的现有证书保持不变。
kafka-rm.yml
移除动作不在 kafka.yml 中,而是使用独立的 kafka-rm.yml 剧本。
该剧本 强制要求非空 -l/--limit,裸跑会在进入角色前失败;-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_data 是专用的安全绝对路径:不含 ./.. 路径段,且不能是 /、/data、/var、/etc、/opt、/usr、/home、/root 或 /pg。
防误删开关是 kafka_safeguard:设置为 true(命令行或清单中)时剧本直接中止,不删除任何东西。身份冲突、Exporter 异常或一般启动失败都不是删除数据的理由——先用 kafka.yml 收敛并读取失败原因。
集群下线
kafka_rm_data 默认为 true:一次默认参数的 kafka-rm.yml 就会删除所选节点的数据/KRaft 元数据与 /etc/kafka 恢复状态。剧本没有确认字符串等额外闸门,执行前必须人工核对 -l 目标、备份或明确重建意图,并评估生产者/消费者影响。
成员退役
部分退役要求 -l 之外至少保留一个 Broker-capable(combined/broker)成员和一个 Controller-capable(combined/controller)成员;
两者可以是同一台 Combined 节点。缺少任一幸存锚点时,剧本会在注销或停服前失败。
通过这些幸存成员,剧本尝试摘除目标的 KRaft Voter 条目(remove-controller,多成员时严格串行)并注销其 Broker 注册(unregister),再执行本机清理。
元数据操作委派给幸存成员,因此对已经死亡、无法连接的目标节点同样适用——这也是 替换故障节点 的第一步。
注销 Broker 的命令被设计为可重入并容忍失败;真实运行后必须检查现场 Quorum、Broker 注册和副本健康,不能只凭剧本返回状态判定退役完成。
退役自动化不等于免除规划:缩容后剩余 Controller 应保持奇数并构成多数派,剩余 Broker 数不能低于现有 Topic 的最大 RF;若被退役 Broker 仍持有 Partition 副本,剧本会打印警告——计划内缩容应当先完成 Reassignment 排空。
剧本边界
两个剧本都不会自动完成 Partition Reassignment 与数据均衡、Topic/用户删除、plaintext 到 scram 的在线迁移、版本升级与 Feature Level 终结、数据备份与灾难恢复,也不部署 Connect、Schema Registry、MirrorMaker、Cruise Control 等生态组件。完整清单见 模块边界;日常只读检查和资源管理见 日常管理。
6 - 监控告警
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 采集任务下:
单 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:
Kafka Overview
集群与全局总览。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。
Kafka Instance
以 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。
Kafka Topic
以 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。
Kafka Consumer
以 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
检查采集目标:
检查某集群复制健康:
检查 Consumer Lag:
检查请求饱和与延迟:
日志查询
Kafka 服务把标准输出与错误写入 Journald,节点 Vector 的 Journald Source 会转发到 VictoriaLogs,统一使用 job:syslog。
Kafka Instance Dashboard 的日志面板使用类似查询,并展示时间、级别、Systemd Unit 与消息。诊断时应把日志与同一时间窗口内的 KRaft、ISR、请求队列、GC、磁盘 I/O 和网络指标对齐。
验证监控链路
在 Kafka 节点验证原始端点:
在 Infra 节点检查文件发现(每实例一个文件,被选中节点的文件含 JMX 与协议 Exporter 两个目标):
然后在 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 模块使用两类指标源,都注册在同一 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 - 常见问题
当前 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 地址。
依次检查:
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 目标):
完整运行会按当前放置刷新每个实例的 Target 文件,不应只针对单节点运行注册标签。注意:若 Exporter 放置因拓扑变化而转移,曾被选中节点上的旧 kafka_exporter 服务不会被普通剧本自动停止,需要手工或通过 kafka-rm.yml 清理。
为什么 JMX 端点可访问,但 jmx_scrape_error=1?
HTTP 可访问只说明 Java Agent 已加载;jmx_scrape_error=1 表示本轮 MBean 采集失败:
检查 /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/网络异常。
再检查 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 目标、可恢复备份或明确重建意图与业务停用状态。成员退役中的 Broker 注销命令会容忍失败,真实运行后还必须核对 Quorum、Broker 注册与副本健康。完整语义见 预置剧本:kafka-rm.yml。