主题
附录 A:Kafka 速查工具书
开机即用的 Kafka 工具书。命令、配置、JMX、内部 Topic、错误信息、SMT、版本里程碑全部以速查表形式收纳,Ctrl+F 搜关键词即可。
基于 Apache Kafka 3.8(KRaft 模式),对 ZK 时代的差异会单独标注 ⚠️。
目录
- 命令行工具速查(10 大
kafka-*.sh) - Broker 关键配置速查(30+ 项)
- Producer 关键配置速查(20+ 项)
- Consumer 关键配置速查(20+ 项)
- JMX 重要指标速查(30+ 项)
- 内部 Topic 速查
- 错误信息辞典(15+ 条)
- SMT(Single Message Transforms)速查
- Kafka 版本里程碑速查
- 一行流(One-liner)实用清单
1. 命令行工具速查
所有命令均假设
KAFKA_HOME=/opt/kafka、Bootstrap 为localhost:9092。 在 docker 环境下:docker exec -it kafka1 /opt/kafka/bin/kafka-xxx.sh ...。
1.1 kafka-topics.sh —— Topic 管理
| 用法 | 命令示例 | 说明 |
|---|---|---|
| 建 Topic | kafka-topics.sh --bootstrap-server :9092 --create --topic learn.04.bench --partitions 6 --replication-factor 3 | 显式指定分区与副本,避免落到默认值 |
| 建 带覆盖配置 | --create --topic t1 --partitions 3 --replication-factor 3 --config retention.ms=86400000 --config cleanup.policy=compact | Topic 级覆盖 Broker 默认 |
| 列 全部 Topic | kafka-topics.sh --bootstrap-server :9092 --list | 含内部 Topic 加 --exclude-internal 反向 |
| 看 Topic 详情 | --describe --topic learn.04.bench | 输出分区、Leader、ISR、Replicas、Configs |
| 看 异常分区 | --describe --under-replicated-partitions / --unavailable-partitions | 排障必用 |
| 改 分区数(只能加不能减) | --alter --topic t1 --partitions 12 | ⚠️ 加分区会破坏 Key Hash 分布! |
| 删 Topic | --delete --topic t1(需 delete.topic.enable=true,3.0+ 默认 true) | 异步删除,看磁盘等几分钟 |
| 看 Topic 配置 | kafka-configs.sh --bootstrap-server :9092 --entity-type topics --entity-name t1 --describe | kafka-topics --describe 也带配置 |
1.2 kafka-console-producer.sh —— 控制台发消息
| 用法 | 命令示例 | 说明 |
|---|---|---|
| 简单发 | kafka-console-producer.sh --bootstrap-server :9092 --topic t1 | 一行一条消息 |
| 带 Key | --topic t1 --property "parse.key=true" --property "key.separator=:" | 输入 k1:v1 |
| 指定分区器 | --producer-property partitioner.class=org.apache.kafka.clients.producer.RoundRobinPartitioner | 默认 Sticky |
| 改 acks | --producer-property acks=all --producer-property enable.idempotence=true | 命令行试可靠性 |
| 从文件灌 | kafka-console-producer.sh ... < big.txt | 行模式 |
| 压测发 | kafka-producer-perf-test.sh --topic t1 --num-records 1000000 --record-size 1024 --throughput -1 --producer-props bootstrap.servers=:9092 acks=1 | 性能基准 |
1.3 kafka-console-consumer.sh —— 控制台收消息
| 用法 | 命令示例 | 说明 |
|---|---|---|
| 从头读 | kafka-console-consumer.sh --bootstrap-server :9092 --topic t1 --from-beginning | 不带 group 是临时 group |
| 指定 group | --topic t1 --group g1 | 提交 offset 到 __consumer_offsets |
| 显示 Key + 元信息 | --property print.key=true --property print.timestamp=true --property print.partition=true --property print.offset=true --property key.separator=' | ' | 调试神器 |
| 只读 1 个分区 | --topic t1 --partition 0 --offset earliest | --offset 支持 earliest/latest/<n> |
read_committed | --topic t1 --isolation-level read_committed | 事务读 |
| 压测读 | kafka-consumer-perf-test.sh --topic t1 --bootstrap-server :9092 --messages 1000000 --threads 1 | 性能基准 |
1.4 kafka-consumer-groups.sh —— 消费组管理
| 用法 | 命令示例 | 说明 |
|---|---|---|
| 列所有 group | kafka-consumer-groups.sh --bootstrap-server :9092 --list | |
| 看 Lag | --describe --group g1 | 输出 CURRENT-OFFSET / LOG-END-OFFSET / LAG |
| 看 group 状态 | --describe --group g1 --state | Stable / Empty / PreparingRebalance / CompletingRebalance / Dead |
| 看成员 | --describe --group g1 --members --verbose | 含 client-id / host / assignment |
| 重置 offset 到最早 | --reset-offsets --group g1 --to-earliest --topic t1 --execute | 必须先 --dry-run 看效果 |
| 重置 到时间戳 | --reset-offsets --group g1 --to-datetime 2024-01-01T00:00:00.000 --topic t1 --execute | 故障回放神器 |
| 重置 偏移 N | --reset-offsets --group g1 --shift-by -1000 --topic t1 --execute | 回退 1000 条 |
| 删除 group | --delete --group g1 | group 必须空 |
| 删除 单条 offset | --delete-offsets --group g1 --topic t1 | 不删 group,只删 offset |
1.5 kafka-configs.sh —— 动态配置
| 用法 | 命令示例 | 说明 |
|---|---|---|
| 看 Broker 配置 | kafka-configs.sh --bootstrap-server :9092 --entity-type brokers --entity-name 1 --describe | --all 看所有,含默认值 |
| 改 Broker 动态配置 | --entity-type brokers --entity-name 1 --alter --add-config log.cleaner.threads=4 | 不重启生效 |
| 改 Topic 配置 | --entity-type topics --entity-name t1 --alter --add-config retention.ms=3600000,cleanup.policy=compact | Topic 级覆盖 |
| 删除 Topic 配置 | --entity-type topics --entity-name t1 --alter --delete-config retention.ms | 回到 Broker 默认 |
| 改 Client 配额 | --entity-type clients --entity-name app1 --alter --add-config producer_byte_rate=1048576,consumer_byte_rate=2097152 | 1MB/s 写 / 2MB/s 读 |
| 改 User 配额 | --entity-type users --entity-name alice --alter --add-config request_percentage=200 | I/O 线程占用上限 |
1.6 kafka-acls.sh —— ACL 管理(启用 authorizer.class.name 后)
| 用法 | 命令示例 | 说明 |
|---|---|---|
| 列所有 ACL | kafka-acls.sh --bootstrap-server :9092 --list | |
| 给 Producer 授权 | --add --allow-principal User:alice --producer --topic t1 | 自动授 Write / Describe |
| 给 Consumer 授权 | --add --allow-principal User:bob --consumer --topic t1 --group g1 | 自动授 Read / Describe |
| 拒绝某 IP | --add --deny-principal User:* --deny-host 10.0.0.1 --operation All --topic t1 | Deny 优先 |
| 给事务 ID 授权 | --add --allow-principal User:alice --operation Write --transactional-id tx-* --resource-pattern-type prefixed | 事务必备 |
| 删 ACL | --remove --allow-principal User:alice --producer --topic t1 --force |
1.7 kafka-reassign-partitions.sh —— 分区迁移
| 用法 | 命令示例 | 说明 |
|---|---|---|
| 生成方案 | kafka-reassign-partitions.sh --bootstrap-server :9092 --topics-to-move-json-file topics.json --broker-list 1,2,3 --generate | 输出 current / proposed 两段 JSON |
| 执行迁移 | --reassignment-json-file plan.json --execute | 异步触发 |
| 查进度 | --reassignment-json-file plan.json --verify | In progress / Completed / Failed |
| 限流 迁移 | --reassignment-json-file plan.json --execute --throttle 50000000 | 50MB/s,避免打满网络 |
| 取消迁移 | --cancel | 3.0+ 才支持 |
1.8 kafka-leader-election.sh —— Leader 切换
| 用法 | 命令示例 | 说明 |
|---|---|---|
| 优先副本选举 | kafka-leader-election.sh --bootstrap-server :9092 --election-type preferred --all-topic-partitions | 把 Leader 切回 Replicas 列表第一个 |
| 指定分区 | --election-type preferred --topic t1 --partition 0 | |
| unclean 选举(非 ISR) | --election-type unclean --topic t1 --partition 0 | ⚠️ 数据丢失风险,Broker 端 unclean.leader.election.enable=true 才能用 |
1.9 kafka-dump-log.sh —— 日志文件解码(运维 / 排障 / 教学神器)
| 用法 | 命令示例 | 说明 |
|---|---|---|
解 .log | kafka-dump-log.sh --files /var/lib/kafka/data/t1-0/00000000000000000000.log --print-data-log | 打印每条 Record |
解 .index | --files .../00000000000000000000.index | Offset → Position 索引 |
解 .timeindex | --files .../00000000000000000000.timeindex | Timestamp → Offset |
解 __consumer_offsets | --files .../__consumer_offsets-12/00000000000000000000.log --offsets-decoder | 解码 OffsetCommit |
解 __transaction_state | --files .../__transaction_state-7/00000000000000000000.log --transaction-log-decoder | 解码事务元数据 |
解 __cluster_metadata | --files .../__cluster_metadata-0/*.log --cluster-metadata-decoder | KRaft 元数据日志 |
| 校验 | --verify-index-only | 只校验 .index 与 .log 对应关系 |
1.10 kafka-metadata-shell.sh —— KRaft 元数据交互(KRaft 专属)
| 用法 | 命令示例 | 说明 |
|---|---|---|
| 进入 shell | kafka-metadata-shell.sh --snapshot /var/lib/kafka/data/__cluster_metadata-0/00000000000000000000-0000000000.checkpoint | 类 unix 文件系统命令 |
| 内部 | ls / / ls /topics / cat /brokers/1 / find / | 浏览元数据 |
| 用途 | 离线诊断元数据快照、对比新旧 Controller 状态 | 替代 ZK 的 zookeeper-shell.sh |
1.11 其它常用脚本
| 脚本 | 一句话作用 |
|---|---|
kafka-storage.sh format -t <UUID> -c <config> | KRaft 模式首次格式化数据目录(生成 meta.properties) |
kafka-cluster.sh cluster-id --bootstrap-server :9092 | 查集群 UUID |
kafka-broker-api-versions.sh --bootstrap-server :9092 | 看 Broker 支持的 API 版本(排障客户端兼容性) |
kafka-log-dirs.sh --describe --bootstrap-server :9092 --json | 查每个 Broker 日志目录大小(容量规划) |
kafka-delete-records.sh --bootstrap-server :9092 --offset-json-file delete.json | 删指定 offset 之前的数据(GDPR / 误写补救) |
kafka-streams-application-reset.sh --application-id app1 --bootstrap-servers :9092 --input-topics t1 | Streams 应用 offset / 内部 Topic 一键重置 |
kafka-mirror-maker.sh / MirrorMaker 2(推荐 connect-mirror-maker.sh) | 跨集群同步 |
kcat -b :9092 -t t1 -C -o beginning | 第三方 CLI(前身 kafkacat),比官方脚本快 / 易用 |
2. Broker 关键配置速查
全集见官方文档 Broker Configs。下表是生产高频 30+ 项,「推荐值」按中等规模集群(3 ~ 10 Broker、TB 级日数据)经验给出,请按业务实际调。
2.1 网络 / 监听器
| 配置 | 默认 | 推荐 | 一句话作用 |
|---|---|---|---|
listeners | PLAINTEXT://:9092 | 显式声明 | Broker 监听的 listener 列表(IP:PORT) |
advertised.listeners | 同 listeners | 必填,对外可达地址 | 注册到元数据、客户端拿到的连接地址 |
listener.security.protocol.map | PLAINTEXT:PLAINTEXT,SSL:SSL,... | 多 listener 隔离 | 给每个 listener 名映射协议 |
inter.broker.listener.name | PLAINTEXT | 单独 INTERNAL | Broker 之间副本同步走的 listener |
num.network.threads | 3 | 3 ~ 8 | Acceptor + Processor 数量,CPU 核数的 1/3 |
num.io.threads | 8 | 8 ~ 16 | 磁盘 I/O 线程数,建议 ≈ 磁盘个数 × 2 |
socket.send.buffer.bytes / socket.receive.buffer.bytes | 102400 | 1048576(1MB) | 高带宽网络下放大 |
socket.request.max.bytes | 104857600 (100MB) | 默认 | 单次请求最大体积,超出 OversizedMessageException |
2.2 日志 / 存储
| 配置 | 默认 | 推荐 | 一句话作用 |
|---|---|---|---|
log.dirs | /tmp/kafka-logs | 多盘并写:/data1/kafka,/data2/kafka | 数据目录,多个目录支持 JBOD |
num.partitions | 1 | 3 / 6 | 默认分区数(建 Topic 不指定时用) |
default.replication.factor | 1 | 3 | 默认副本因子 |
min.insync.replicas | 1 | 2(3 副本时) | acks=all 至少要写多少副本才算成功 |
log.retention.hours | 168(7d) | 168 / 72 | 按时间清理 |
log.retention.bytes | -1 | 按容量定 | 按字节清理(与时间是 OR 关系) |
log.segment.bytes | 1073741824 (1GB) | 512MB ~ 1GB | 单 segment 大小 |
log.roll.hours | 168 | 24 | 强制滚 segment 时间(与 bytes OR) |
log.cleaner.enable | true | 保持开(compact Topic 必需) | Log Compaction 后台线程总开关 |
log.cleaner.threads | 1 | 2 ~ 4 | Compaction 线程数 |
log.cleaner.dedupe.buffer.size | 134217728 (128MB) | 256MB ~ 1GB | Compaction 去重缓冲(影响一次 clean 的 segment 数) |
log.flush.interval.messages / log.flush.interval.ms | 不主动刷 | 保持默认 | Kafka 主动 fsync 配置;生产靠副本不靠 fsync,开了反而慢 |
auto.create.topics.enable | true | false | 防止误产生奇怪 Topic |
delete.topic.enable | true(3.0+) | true | 允许真删 Topic |
unclean.leader.election.enable | false(2.0+) | false | 默认不允许非 ISR 选 Leader(保数据不丢) |
auto.leader.rebalance.enable | true | true | 自动周期把 Leader 切回 preferred replica |
leader.imbalance.per.broker.percentage | 10 | 10 | 触发自动 rebalance 的不均衡阈值 |
2.3 副本 / 复制
| 配置 | 默认 | 推荐 | 一句话作用 |
|---|---|---|---|
replica.lag.time.max.ms | 30000 | 30000 | Follower 落后多久从 ISR 踢出 |
replica.fetch.max.bytes | 1048576 (1MB) | 与 message.max.bytes 对齐 | Follower 一次抓取的最大字节 |
replica.fetch.wait.max.ms | 500 | 默认 | Follower fetch 阻塞超时 |
num.replica.fetchers | 1 | 4 ~ 8 | 每 Broker 拉副本的线程数 |
replica.fetch.response.max.bytes | 10485760 (10MB) | 默认 | Follower 单次 fetch 响应上限 |
2.4 Group / 协议
| 配置 | 默认 | 推荐 | 一句话作用 |
|---|---|---|---|
group.initial.rebalance.delay.ms | 3000 | 3000(教学 0) | 第一次 join 时延迟,攒一波客户端 |
offsets.topic.replication.factor | 3 | 3 | __consumer_offsets 副本因子 |
transaction.state.log.replication.factor | 3 | 3 | __transaction_state 副本因子 |
transaction.state.log.min.isr | 2 | 2 | 事务元数据最小 ISR |
group.max.session.timeout.ms | 1800000 (30min) | 默认 | session.timeout.ms 上限 |
group.min.session.timeout.ms | 6000 | 默认 | session.timeout.ms 下限 |
2.5 KRaft 专属
| 配置 | 默认 | 推荐 | 一句话作用 |
|---|---|---|---|
process.roles | 无 | broker,controller(教学)/ 分离(生产) | 进程角色 |
node.id | 无 | 唯一 Int | KRaft 替代 broker.id |
controller.quorum.voters | 无 | 1@h1:9093,2@h2:9093,3@h3:9093 | Controller Quorum 列表 |
controller.listener.names | 无 | CONTROLLER | KRaft 通信 listener 名 |
metadata.log.dir | 同 log.dirs[0] | 单独高速盘 | __cluster_metadata 存放位置 |
metadata.log.segment.bytes | 1073741824 | 默认 | 元数据 segment 大小 |
2.6 安全
| 配置 | 一句话作用 |
|---|---|
authorizer.class.name | ACL 授权器,KRaft 用 org.apache.kafka.metadata.authorizer.StandardAuthorizer,ZK 用 kafka.security.authorizer.AclAuthorizer |
super.users | 超级用户列表,绕过 ACL(如 User:admin) |
allow.everyone.if.no.acl.found | 找不到 ACL 时是否放行(默认 false,生产必须保 false) |
sasl.enabled.mechanisms | 启用的 SASL 机制(PLAIN,SCRAM-SHA-512,GSSAPI,OAUTHBEARER) |
ssl.client.auth | none / requested / required,是否校验客户端证书 |
3. Producer 关键配置速查
Java key 与 librdkafka (confluent-kafka) key 通常一致;下方括号内是 librdkafka 名(仅在不同时给出)。
| 配置 | 默认 | 推荐 | 一句话作用 |
|---|---|---|---|
bootstrap.servers | 无 | h1:9092,h2:9092,h3:9092 | 集群引导地址 |
client.id | "" | 业务名_实例 | 排障 / 配额必备 |
acks | all(3.0+ 默认改为 all) | all | 0(不等)、1(Leader 写入)、all(ISR 全确认) |
enable.idempotence | true(3.0+ 默认) | true | 开幂等,自动设 acks=all/retries=MAX/MaxInFlight≤5 |
retries | 2147483647 | 默认 | 失败重试次数 |
delivery.timeout.ms | 120000 (2min) | 120000 ~ 300000 | send → 成功 / 失败的总时长上限 |
request.timeout.ms | 30000 | 30000 | 单次请求超时 |
max.in.flight.requests.per.connection | 5 | ≤5(开幂等时强制) | 同一连接未确认的请求上限 |
linger.ms | 0 | 5 ~ 50 | 攒批等待时长,0 = 来一条发一条(吞吐差) |
batch.size | 16384 (16KB) | 32KB ~ 256KB | 每个分区批次的字节上限 |
buffer.memory | 33554432 (32MB) | 64MB ~ 256MB | Producer 内存缓冲区总量 |
max.block.ms | 60000 | 60000 | buffer 满 / 元数据未就绪时 send 阻塞上限 |
compression.type | none | zstd 或 lz4 | 压缩算法(端到端,Broker 不解压) |
partitioner.class | DefaultPartitioner(Sticky) | 默认 | 自定义路由规则;3.0+ 默认 Sticky 而非纯 Murmur2 |
partitioner.ignore.keys | false | false | 配 Sticky 时是否忽略 Key |
key.serializer / value.serializer | 无 | StringSerializer / ByteArraySerializer | Java 序列化器 |
transactional.id | 无 | tx-<service>-<instance>,必须稳定 | 开事务必填,重启后重启 Producer 用同 ID 触发 PID Fencing |
transaction.timeout.ms | 60000 | ≤ Broker transaction.max.timeout.ms | 事务最大持续时间 |
metadata.max.age.ms | 300000 (5min) | 默认 | 元数据强制刷新周期 |
connections.max.idle.ms | 540000 (9min) | 默认 | 空闲连接关闭时间 |
interceptor.classes | 空 | 监控 / 染色 | Producer 拦截器链 |
send.buffer.bytes / receive.buffer.bytes | 131072 / 32768 | -1(OS 默认) | TCP 缓冲 |
📌 「吞吐 vs 延迟」黄金口诀:吞吐优先 → 大
linger.ms+ 大batch.size+ zstd;延迟优先 →linger.ms=0+ 小 batch + lz4。
4. Consumer 关键配置速查
| 配置 | 默认 | 推荐 | 一句话作用 |
|---|---|---|---|
bootstrap.servers | 无 | 同上 | 集群引导地址 |
group.id | 无(手动 assign 可为空) | 业务名 | 消费组标识 |
group.instance.id | 无 | 实例稳定 ID | 开 Static Membership,避免重启触发 Rebalance |
client.id | "" | 业务名_实例 | 排障 / 配额 |
key.deserializer / value.deserializer | 无 | StringDeserializer / ByteArrayDeserializer | Java 反序列化 |
auto.offset.reset | latest | earliest(流处理)/ latest(实时业务) | 找不到 offset 时的行为,none 抛异常 |
enable.auto.commit | true | false | 关掉自动提交,手动提交才准 |
auto.commit.interval.ms | 5000 | 5000 | 自动提交周期 |
isolation.level | read_uncommitted | read_committed(消费事务消息时) | 是否过滤未提交事务 |
max.poll.records | 500 | 100 ~ 1000 | 单次 poll 返回的最大记录数 |
max.poll.interval.ms | 300000 (5min) | 业务实际处理时间 × 1.5 | 两次 poll 最大间隔,超时被踢出 |
session.timeout.ms | 45000(3.0+) | 45000 ~ 60000 | 心跳超时,超过就踢出 group |
heartbeat.interval.ms | 3000 | session/3 | 心跳间隔 |
partition.assignment.strategy | [RangeAssignor, CooperativeStickyAssignor](3.0+) | CooperativeStickyAssignor | 分配策略,4 选 1 |
fetch.min.bytes | 1 | 1 ~ 50000 | Broker 至少攒多少字节才返回 |
fetch.max.bytes | 52428800 (50MB) | 默认 | 单次 fetch 总字节上限 |
max.partition.fetch.bytes | 1048576 (1MB) | 与 message.max.bytes 对齐 | 单分区单次 fetch 上限 |
fetch.max.wait.ms | 500 | 500 | fetch.min.bytes 没攒够时最多等多久 |
request.timeout.ms | 30000 | 30000 | 请求超时 |
connections.max.idle.ms | 540000 | 默认 | 空闲连接关闭 |
interceptor.classes | 空 | 染色 / 监控 | Consumer 拦截器 |
check.crcs | true | true | 校验消息 CRC,关了快但有风险 |
5. JMX 重要指标速查
全部走
kafka.*MBean,配 Prometheus JMX Exporter 即可暴露。 下表分四类:Broker / Topic / Producer / Consumer。
5.1 Broker 指标(30 个里挑了最关键的)
| MBean | 类型 | 一句话含义 | 报警阈值建议 |
|---|---|---|---|
kafka.controller:type=KafkaController,name=ActiveControllerCount | Gauge | 当前 Broker 是否是 Controller | 集群和必须 = 1 |
kafka.controller:type=KafkaController,name=OfflinePartitionsCount | Gauge | 没有 Leader 的分区数 | > 0 立即告警 |
kafka.server:type=ReplicaManager,name=UnderReplicatedPartitions | Gauge | 副本不同步的分区数 | > 0 持续 5min 告警 |
kafka.server:type=ReplicaManager,name=UnderMinIsrPartitionCount | Gauge | ISR 数 < min.insync.replicas 的分区 | > 0 立即告警(影响写入) |
kafka.server:type=ReplicaManager,name=AtMinIsrPartitionCount | Gauge | ISR 数刚好 = min.isr 的分区 | 持续较高需扩容 |
kafka.server:type=BrokerTopicMetrics,name=MessagesInPerSec | Meter | 每秒入站消息数 | 趋势 |
kafka.server:type=BrokerTopicMetrics,name=BytesInPerSec / BytesOutPerSec | Meter | 每秒入 / 出字节 | 趋势 + 容量规划 |
kafka.server:type=BrokerTopicMetrics,name=FailedProduceRequestsPerSec / FailedFetchRequestsPerSec | Meter | 每秒失败请求数 | > 0 排障 |
kafka.network:type=RequestMetrics,name=RequestsPerSec,request={Produce,Fetch,Metadata} | Meter | 各类请求 QPS | 趋势 |
kafka.network:type=RequestMetrics,name=TotalTimeMs,request=Produce | Histogram | Produce 端到端耗时 | p99 > 1s 告警 |
kafka.network:type=RequestMetrics,name=LocalTimeMs / RequestQueueTimeMs / RemoteTimeMs / ResponseQueueTimeMs / ResponseSendTimeMs | Histogram | 请求各阶段耗时 | 排障定位瓶颈 |
kafka.server:type=KafkaRequestHandlerPool,name=RequestHandlerAvgIdlePercent | Meter | I/O 线程平均空闲率 | < 0.3 表示 IO 线程吃紧 |
kafka.network:type=SocketServer,name=NetworkProcessorAvgIdlePercent | Gauge | 网络线程平均空闲率 | < 0.3 表示网络线程吃紧 |
kafka.log:type=LogFlushStats,name=LogFlushRateAndTimeMs | Histogram | fsync 频率与耗时 | 突增看磁盘 |
kafka.log:type=LogManager,name=LogDirectoryOffline | Gauge | 离线日志目录数 | > 0 表示磁盘故障 |
kafka.log:type=Log,name=Size,topic=t1,partition=0 | Gauge | 单分区磁盘占用 | 容量规划 |
kafka.log:type=Log,name=LogStartOffset / LogEndOffset | Gauge | 单分区起始 / 结束 Offset | 排障 |
kafka.log.cleaner:type=LogCleanerManager,name=max-dirty-percent | Gauge | Compact Topic 最高 dirty 比例 | > 50% 看 cleaner 是否挂 |
kafka.log.cleaner:type=LogCleaner,name=cleaner-recopy-percent | Gauge | Cleaner 工作量 | 趋势 |
kafka.controller:type=ControllerStats,name=LeaderElectionRateAndTimeMs | Histogram | Leader 选举频率 | 频繁选举 = 集群抖动 |
kafka.controller:type=ControllerStats,name=UncleanLeaderElectionsPerSec | Meter | unclean 选举次数 | > 0 排障(数据可能丢) |
kafka.server:type=ReplicaFetcherManager,name=MaxLag,clientId=Replica | Gauge | Follower 最大 Lag | > 10000 排障 |
kafka.server:type=ReplicaManager,name=IsrShrinksPerSec / IsrExpandsPerSec | Meter | ISR 收缩 / 扩张速率 | 频繁波动 = 抖动 |
kafka.server:type=DelayedOperationPurgatory,name=PurgatorySize,delayedOperation={Produce,Fetch} | Gauge | 延迟操作队列长度 | 突增表示 Broker 在等 ack / 数据 |
java.lang:type=GarbageCollector,name=* | Counter | JVM GC 时长 / 次数 | Full GC 频繁告警 |
java.lang:type=Memory | Gauge | 堆 / 非堆使用 | > 80% 告警 |
java.lang:type=OperatingSystem,name=SystemCpuLoad | Gauge | 系统 CPU | > 80% 告警 |
5.2 Topic 维度指标
| MBean | 含义 |
|---|---|
kafka.server:type=BrokerTopicMetrics,name=MessagesInPerSec,topic=t1 | 单 Topic 入站速率 |
kafka.server:type=BrokerTopicMetrics,name=BytesInPerSec,topic=t1 | 单 Topic 入站字节 |
kafka.server:type=BrokerTopicMetrics,name=BytesOutPerSec,topic=t1 | 单 Topic 出站字节 |
kafka.server:type=BrokerTopicMetrics,name=BytesRejectedPerSec,topic=t1 | 因配额 / 错误被拒字节 |
kafka.server:type=BrokerTopicMetrics,name=ReplicationBytesInPerSec / OutPerSec | 副本同步字节 |
5.3 Producer 客户端指标(暴露在 Producer JVM)
| Metric | 含义 |
|---|---|
kafka.producer:type=producer-metrics,client-id=*,record-send-rate | 每秒发送条数 |
record-error-rate | 每秒发送失败条数(> 0 报警) |
record-retry-rate | 重试速率 |
record-size-avg / record-size-max | 消息平均 / 最大字节 |
request-latency-avg / request-latency-max | 请求平均 / 最大耗时 |
batch-size-avg / batch-size-max | 批次大小(用于评估 batch.size 是否合理) |
compression-rate-avg | 压缩率 |
buffer-available-bytes | 缓冲区剩余 |
buffer-exhausted-rate | 缓冲区满次数(> 0 加 buffer.memory) |
record-queue-time-avg | 消息在客户端排队时间 |
5.4 Consumer 客户端指标
| Metric | 含义 |
|---|---|
kafka.consumer:type=consumer-fetch-manager-metrics,client-id=*,records-consumed-rate | 每秒消费条数 |
bytes-consumed-rate | 每秒消费字节 |
fetch-latency-avg / max | fetch 平均 / 最大耗时 |
records-lag-max | 最大 Lag(业务最关心的指标) |
records-lead-min | 最小 Lead(消费位点距离 LogStartOffset)→ 太小说明数据要被清掉 |
commit-latency-avg | offset 提交耗时 |
commit-rate | 每秒 commit 次数 |
kafka.consumer:type=consumer-coordinator-metrics,client-id=*,join-rate / sync-rate / heartbeat-rate | Group 协议交互速率 |
last-rebalance-seconds-ago | 距离上次 Rebalance 时间,频繁掉低 = 抖动 |
assigned-partitions | 当前分配到的分区数 |
6. 内部 Topic 速查
| Topic | 用途 | cleanup.policy | 副本因子建议 | 查看方式 |
|---|---|---|---|---|
__consumer_offsets | 消费组提交的 offset / group 元数据 | compact | 3 | kafka-consumer-groups.sh --describe;裸读 kafka-console-consumer --topic __consumer_offsets --formatter "kafka.coordinator.group.GroupMetadataManager\$OffsetsMessageFormatter" |
__transaction_state | 事务元数据(PID、Epoch、参与分区、状态) | compact | 3 | kafka-dump-log.sh --transaction-log-decoder |
__cluster_metadata | KRaft 集群元数据日志(替代 ZK) | delete(带快照) | 等于 Controller Quorum 大小 | kafka-metadata-shell.sh --snapshot <file> 或 kafka-dump-log.sh --cluster-metadata-decoder |
_schemas | Confluent Schema Registry 存放 schema | compact | 3 | kafka-console-consumer --topic _schemas --from-beginning |
connect-configs | Kafka Connect 配置 | compact | 3 | 同上 |
connect-offsets | Connect Source connector 的位点 | compact | 25 partitions / 3 副本 | 同上 |
connect-status | Connect 任务状态 | compact | 3 | 同上 |
<app-id>-changelog-* | Kafka Streams 状态 store 的 changelog | compact(KTable)/ compact,delete(带 retention 的 store) | 跟随 Streams 配置 | 用 Streams Reset 工具清理 |
<app-id>-repartition-* | Streams 中间重分区 Topic | delete + 短 retention | 跟随 Streams 配置 | 同上 |
⚠️ 这些 Topic 绝对不要手动删!
__consumer_offsets删了所有消费组 Lag 全丢;__transaction_state删了所有进行中的事务永远卡住。
7. 错误信息辞典
7.1 Producer 端
| 错误 / 异常 | 含义 | 处理 |
|---|---|---|
org.apache.kafka.common.errors.TimeoutException: Topic xxx not present in metadata after 60000 ms | 元数据拉不到,多半是 bootstrap.servers 不通 / Topic 不存在 / 没权限 | 检查网络、ACL、auto.create.topics.enable |
RecordTooLargeException: The message is X bytes, larger than max.request.size | 单条消息超过客户端 max.request.size(默认 1MB) | 加大 max.request.size、message.max.bytes、replica.fetch.max.bytes(三件套)或拆分消息 |
BufferExhaustedException / Failed to allocate memory within timeout | buffer.memory 满了 | 加 buffer.memory / 加 linger.ms 让批次更大、消费下游加速 |
OutOfOrderSequenceException | 幂等 Producer 序列号乱序,多半是网络重排 + max.in.flight > 5 | 保 max.in.flight ≤ 5,开 enable.idempotence=true |
ProducerFencedException | 同 transactional.id 的新 Producer 起来了,旧的被挤掉 | 正常,旧 Producer 直接退出 |
InvalidProducerEpochException | 事务超时被中止 | 业务侧重新 beginTransaction() |
NotEnoughReplicasException / NotEnoughReplicasAfterAppendException | ISR 数 < min.insync.replicas,写不进去(acks=all 时) | 修复掉队 Broker;临时降级可调小 min.insync.replicas(不推荐) |
UnknownTopicOrPartitionException | Topic 不存在或刚删 | kafka-topics --list 确认 |
LeaderNotAvailableException | 分区 Leader 切换中 | 客户端会自动重试,多发几条就好 |
7.2 Consumer 端
| 错误 | 含义 | 处理 |
|---|---|---|
CommitFailedException: ... cannot be completed because the consumer is no longer part of the group | 处理太慢 → 超过 max.poll.interval.ms → 被踢 → 提交失败 | 加大 max.poll.interval.ms,或拆小 max.poll.records,或异步处理 |
OffsetOutOfRangeException | 提交的 offset 不在 [LogStartOffset, LogEndOffset](数据被清了 / 重置) | 用 auto.offset.reset 自动重置;或 seek 到合理位置 |
WakeupException | 主动调用 consumer.wakeup()(多发生在优雅关闭) | catch 后 consumer.close() |
NoOffsetForPartitionException | 第一次消费且 auto.offset.reset=none | 改 earliest/latest;或人工 seek |
InconsistentGroupProtocolException | 同 group 不同实例 partition.assignment.strategy 不一致 | 全部统一 |
RebalanceInProgressException | 收到这个错时正在 Rebalance,提交无效 | 在 onPartitionsRevoked 钩子里提交 |
7.3 Broker / 集群
| 错误 / 现象 | 含义 | 处理 |
|---|---|---|
Unable to write to meta.properties / cluster.id ... doesn't match | KRaft 节点 cluster ID 与磁盘上不一致 | 删数据卷重建;不要手动改 meta.properties |
Shutdown broker because all log dirs in /xxx have failed | 全部数据盘故障 | 修磁盘或迁移到新节点 |
Cluster authorization failed | ACL 拒绝 | kafka-acls --list 确认;事务必须授 IdempotentWrite 或 Write |
INVALID_FETCH_SESSION_EPOCH | Fetch session 失效 | 客户端会自动重建 session,不用管 |
8. SMT 速查
SMT = Single Message Transforms,Kafka Connect 在
transforms配置里挂的轻量级转换器,作用于单条消息。常用于改字段名、加时间戳、过滤路由、屏蔽敏感字段。
SMT 类(org.apache.kafka.connect.transforms.*) | 一句话作用 | 关键参数 |
|---|---|---|
InsertField$Key / InsertField$Value | 注入静态字段(offset/partition/timestamp/topic/static) | static.field, offset.field, timestamp.field |
ReplaceField$Key / ReplaceField$Value | 重命名 / 删除字段 | renames, exclude, include |
MaskField$Value | 屏蔽敏感字段值(脱敏) | fields, replacement |
ValueToKey | 把 value 中某些字段提到 key | fields |
HoistField$Key/Value | 把 schemaless 包一层 struct | field |
ExtractField$Key/Value | 抽出 struct 中的某字段做新值 | field |
Cast$Key/Value | 字段类型转换(int↔long↔string) | spec |
TimestampConverter$Value | 时间戳与字符串 / Unix 互转 | target.type, format, unix.precision |
RegexRouter | 用正则改写 Topic 名(动态路由) | regex, replacement |
Filter(Confluent 商业版 io.confluent.connect.transforms.Filter) | 按谓词过滤掉消息 | predicate, predicates.* |
Flatten$Value | 嵌套 struct 拍平成单层 | delimiter |
HeaderFrom$Value | value 字段抽到 Header | fields, headers |
DropHeaders | 删指定 Header | headers |
典型链(Debezium CDC → 路由 + 改名):
properties
transforms=route,unwrap
transforms.route.type=org.apache.kafka.connect.transforms.RegexRouter
transforms.route.regex=([^.]+)\\.([^.]+)\\.([^.]+)
transforms.route.replacement=cdc_$3
transforms.unwrap.type=io.debezium.transforms.ExtractNewRecordState
transforms.unwrap.drop.tombstones=false9. Kafka 版本里程碑速查
| 版本 | 发布年 | 关键特性 |
|---|---|---|
| 0.7 | 2012 | LinkedIn 原始版本,Scala 写的雏形 |
| 0.8 | 2013 | 副本机制首次出现(Replication);Producer / Consumer 协议成型 |
| 0.9 | 2015 | 新 Consumer API(基于 Group Coordinator)、SSL / SASL 安全、Quotas、Kafka Connect 雏形 |
| 0.10 | 2016 | Kafka Streams 发布;消息加 timestamp;rack-aware 副本分配 |
| 0.10.1 / 0.10.2 | 2016 | max.poll.interval.ms(解决长处理踢出);客户端 / Broker 解耦 |
| 0.11 | 2017 | 幂等 Producer + 事务 + Exactly Once(PID + Epoch + 序列号 + Transaction Coordinator);V2 RecordBatch(更小、加 Header) |
| 1.0 | 2017 | 协议稳定 GA;JBOD 改进;Java 9 支持 |
| 1.1 | 2018 | 单 Broker 支持百万级分区;Controller 启动加速 |
| 2.0 | 2018 | OAuthBearer SASL;前缀 ACL;Java 11 支持;unclean.leader.election.enable 默认改为 false |
| 2.1 | 2018 | zstd 压缩;CreateTopics 协议改进 |
| 2.2 | 2019 | Default partitioner 改进(KIP-389 限制单连接请求);ACL 支持 prefix |
| 2.3 | 2019 | Incremental Cooperative Rebalance(KIP-429);Static Membership(KIP-345,group.instance.id) |
| 2.4 | 2019 | KIP-392 follower fetching(最近副本读,跨机房友好);KIP-466 用 Java 11 |
| 2.5 | 2020 | KIP-447 事务 API 增强(每分区一个 Producer);TLS 1.3 |
| 2.6 | 2020 | client.id quota 改进;KIP-573 文档 metric 化 |
| 2.7 | 2020 | KRaft 早期预览(KIP-500);End-to-End Latency 指标 |
| 2.8 | 2021 | KRaft 模式 Preview 可用(无需 ZK 启动);KIP-516 topic ID |
| 3.0 | 2021 | KRaft 进入 Production-ready (Preview);移除 0.10 之前的旧协议;默认 acks=all + enable.idempotence=true;StickyAssignor 默认 |
| 3.1 / 3.2 | 2022 | KRaft GA 准备;Tiered Storage 早期 KIP;ZK 模式开始进入「不推荐」 |
| 3.3 | 2022 | KRaft 模式 GA(KIP-833),生产可用 |
| 3.4 | 2023 | KRaft 元数据快照、ZK → KRaft 迁移工具早期版(KIP-866) |
| 3.5 | 2023 | KRaft 迁移工具进一步完善;新协议 Consumer Group 设计开始(KIP-848) |
| 3.6 | 2023 | Tiered Storage Early Access(KIP-405),冷数据落 S3 |
| 3.7 | 2024 | KIP-848 next-gen Consumer Group Protocol Preview;MirrorMaker 2 改进 |
| 3.8 | 2024 | 本教程基线;KIP-1011 KRaft 多版本兼容;JBOD on KRaft |
| 3.9 | 2024 | KRaft 模式 ZK 迁移功能稳定,最后一个支持 ZK 的版本系列 |
| 4.0 | 2025 | 完全移除 ZooKeeper;KIP-848 新 Consumer Group Protocol GA;新 Queue 概念(KIP-932);Tiered Storage GA |
10. 一行流实用清单
收 100% 真实可复制的「一行能干一件事」命令。
bash
# 1) 生成 KRaft Cluster ID(22 字节 base64)
KAFKA_CLUSTER_ID=$(/opt/kafka/bin/kafka-storage.sh random-uuid)
# 2) 首次格式化数据目录(KRaft 必须)
/opt/kafka/bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c /opt/kafka/config/kraft/server.properties
# 3) 一行拿集群 UUID
/opt/kafka/bin/kafka-cluster.sh cluster-id --bootstrap-server :9092
# 4) 一行查所有 Topic 的分区与副本
/opt/kafka/bin/kafka-topics.sh --bootstrap-server :9092 --describe | awk '/^Topic:/'
# 5) 一行算所有 Topic 总磁盘占用(GB)
/opt/kafka/bin/kafka-log-dirs.sh --bootstrap-server :9092 --describe --json \
| python3 -c "import sys,json;d=json.load(sys.stdin);print(sum(p['size'] for b in d['brokers'] for l in b['logDirs'] for p in l['partitions'])/2**30, 'GB')"
# 6) 一行删 7 天前的某 Topic 数据(保留位点)
echo '{"partitions":[{"topic":"t1","partition":0,"offset":-1}],"version":1}' > delete.json
/opt/kafka/bin/kafka-delete-records.sh --bootstrap-server :9092 --offset-json-file delete.json
# 7) 一行回放某消费组到昨天 0 点
/opt/kafka/bin/kafka-consumer-groups.sh --bootstrap-server :9092 --group g1 --topic t1 \
--reset-offsets --to-datetime "$(date -d 'yesterday 00:00' '+%Y-%m-%dT%H:%M:%S.000')" --execute
# 8) 一行解一个 segment 文件
/opt/kafka/bin/kafka-dump-log.sh --files /var/lib/kafka/data/t1-0/00000000000000000000.log \
--print-data-log --deep-iteration | head -n 30
# 9) 一行起个 perf 压测(写)
/opt/kafka/bin/kafka-producer-perf-test.sh --topic t1 \
--num-records 1000000 --record-size 1024 --throughput -1 \
--producer-props bootstrap.servers=:9092 acks=all linger.ms=20 batch.size=131072 compression.type=zstd
# 10) 一行起个 perf 压测(读)
/opt/kafka/bin/kafka-consumer-perf-test.sh --bootstrap-server :9092 --topic t1 --messages 1000000 --threads 1
# 11) kcat 一行查最近 10 条
kcat -b :9092 -t t1 -C -o -10 -e -f '%p %o %k -> %s\n'
# 12) 一行批量给 Topic 加 retention 短期
for t in $(kafka-topics.sh --bootstrap-server :9092 --list | grep '^tmp\.'); do
kafka-configs.sh --bootstrap-server :9092 --entity-type topics --entity-name "$t" --alter --add-config retention.ms=3600000;
done
# 13) 一行查所有「ISR < 副本数」的分区
kafka-topics.sh --bootstrap-server :9092 --describe --under-replicated-partitions
# 14) 一行查所有「未同步且无 Leader」的分区(最严重故障)
kafka-topics.sh --bootstrap-server :9092 --describe --unavailable-partitions
# 15) 一行查 Broker JMX(前提 broker 启用 -Dcom.sun.management.jmxremote)
echo "get -b kafka.server:type=ReplicaManager,name=UnderReplicatedPartitions Value" | java -jar jmxterm.jar -l :9099 -n一图流:Kafka 全家桶速查脑图
📌 本附录的姐妹篇:
appendix_kafka_vs_others.md— Kafka vs 其它 MQ 横向对比appendix_pitfalls.md— 25+ 真实踩坑案例interview.md— 面试题总索引(210+ 题)