Skip to content

附录 A:Kafka 速查工具书

开机即用的 Kafka 工具书。命令、配置、JMX、内部 Topic、错误信息、SMT、版本里程碑全部以速查表形式收纳,Ctrl+F 搜关键词即可。

基于 Apache Kafka 3.8(KRaft 模式),对 ZK 时代的差异会单独标注 ⚠️。

目录

  1. 命令行工具速查(10 大 kafka-*.sh
  2. Broker 关键配置速查(30+ 项)
  3. Producer 关键配置速查(20+ 项)
  4. Consumer 关键配置速查(20+ 项)
  5. JMX 重要指标速查(30+ 项)
  6. 内部 Topic 速查
  7. 错误信息辞典(15+ 条)
  8. SMT(Single Message Transforms)速查
  9. Kafka 版本里程碑速查
  10. 一行流(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 管理

用法命令示例说明
Topickafka-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=compactTopic 级覆盖 Broker 默认
全部 Topickafka-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 --describekafka-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 —— 消费组管理

用法命令示例说明
列所有 groupkafka-consumer-groups.sh --bootstrap-server :9092 --list
看 Lag--describe --group g1输出 CURRENT-OFFSET / LOG-END-OFFSET / LAG
看 group 状态--describe --group g1 --stateStable / Empty / PreparingRebalance / CompletingRebalance / Dead
看成员--describe --group g1 --members --verboseclient-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 g1group 必须空
删除 单条 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=compactTopic 级覆盖
删除 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=20971521MB/s 写 / 2MB/s 读
改 User 配额--entity-type users --entity-name alice --alter --add-config request_percentage=200I/O 线程占用上限

1.6 kafka-acls.sh —— ACL 管理(启用 authorizer.class.name 后)

用法命令示例说明
列所有 ACLkafka-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 t1Deny 优先
给事务 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 --verifyIn progress / Completed / Failed
限流 迁移--reassignment-json-file plan.json --execute --throttle 5000000050MB/s,避免打满网络
取消迁移--cancel3.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 —— 日志文件解码(运维 / 排障 / 教学神器)

用法命令示例说明
.logkafka-dump-log.sh --files /var/lib/kafka/data/t1-0/00000000000000000000.log --print-data-log打印每条 Record
.index--files .../00000000000000000000.indexOffset → Position 索引
.timeindex--files .../00000000000000000000.timeindexTimestamp → 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-decoderKRaft 元数据日志
校验--verify-index-only只校验 .index 与 .log 对应关系

1.10 kafka-metadata-shell.sh —— KRaft 元数据交互(KRaft 专属)

用法命令示例说明
进入 shellkafka-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 t1Streams 应用 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 网络 / 监听器

配置默认推荐一句话作用
listenersPLAINTEXT://:9092显式声明Broker 监听的 listener 列表(IP:PORT)
advertised.listeners同 listeners必填,对外可达地址注册到元数据、客户端拿到的连接地址
listener.security.protocol.mapPLAINTEXT:PLAINTEXT,SSL:SSL,...多 listener 隔离给每个 listener 名映射协议
inter.broker.listener.namePLAINTEXT单独 INTERNALBroker 之间副本同步走的 listener
num.network.threads33 ~ 8Acceptor + Processor 数量,CPU 核数的 1/3
num.io.threads88 ~ 16磁盘 I/O 线程数,建议 ≈ 磁盘个数 × 2
socket.send.buffer.bytes / socket.receive.buffer.bytes1024001048576(1MB)高带宽网络下放大
socket.request.max.bytes104857600 (100MB)默认单次请求最大体积,超出 OversizedMessageException

2.2 日志 / 存储

配置默认推荐一句话作用
log.dirs/tmp/kafka-logs多盘并写:/data1/kafka,/data2/kafka数据目录,多个目录支持 JBOD
num.partitions13 / 6默认分区数(建 Topic 不指定时用)
default.replication.factor13默认副本因子
min.insync.replicas12(3 副本时)acks=all 至少要写多少副本才算成功
log.retention.hours168(7d)168 / 72按时间清理
log.retention.bytes-1按容量定按字节清理(与时间是 OR 关系)
log.segment.bytes1073741824 (1GB)512MB ~ 1GB单 segment 大小
log.roll.hours16824强制滚 segment 时间(与 bytes OR)
log.cleaner.enabletrue保持开(compact Topic 必需)Log Compaction 后台线程总开关
log.cleaner.threads12 ~ 4Compaction 线程数
log.cleaner.dedupe.buffer.size134217728 (128MB)256MB ~ 1GBCompaction 去重缓冲(影响一次 clean 的 segment 数)
log.flush.interval.messages / log.flush.interval.ms不主动刷保持默认Kafka 主动 fsync 配置;生产靠副本不靠 fsync,开了反而慢
auto.create.topics.enabletruefalse防止误产生奇怪 Topic
delete.topic.enabletrue(3.0+)true允许真删 Topic
unclean.leader.election.enablefalse(2.0+)false默认不允许非 ISR 选 Leader(保数据不丢)
auto.leader.rebalance.enabletruetrue自动周期把 Leader 切回 preferred replica
leader.imbalance.per.broker.percentage1010触发自动 rebalance 的不均衡阈值

2.3 副本 / 复制

配置默认推荐一句话作用
replica.lag.time.max.ms3000030000Follower 落后多久从 ISR 踢出
replica.fetch.max.bytes1048576 (1MB)message.max.bytes 对齐Follower 一次抓取的最大字节
replica.fetch.wait.max.ms500默认Follower fetch 阻塞超时
num.replica.fetchers14 ~ 8每 Broker 拉副本的线程数
replica.fetch.response.max.bytes10485760 (10MB)默认Follower 单次 fetch 响应上限

2.4 Group / 协议

配置默认推荐一句话作用
group.initial.rebalance.delay.ms30003000(教学 0)第一次 join 时延迟,攒一波客户端
offsets.topic.replication.factor33__consumer_offsets 副本因子
transaction.state.log.replication.factor33__transaction_state 副本因子
transaction.state.log.min.isr22事务元数据最小 ISR
group.max.session.timeout.ms1800000 (30min)默认session.timeout.ms 上限
group.min.session.timeout.ms6000默认session.timeout.ms 下限

2.5 KRaft 专属

配置默认推荐一句话作用
process.rolesbroker,controller(教学)/ 分离(生产)进程角色
node.id唯一 IntKRaft 替代 broker.id
controller.quorum.voters1@h1:9093,2@h2:9093,3@h3:9093Controller Quorum 列表
controller.listener.namesCONTROLLERKRaft 通信 listener 名
metadata.log.dirlog.dirs[0]单独高速盘__cluster_metadata 存放位置
metadata.log.segment.bytes1073741824默认元数据 segment 大小

2.6 安全

配置一句话作用
authorizer.class.nameACL 授权器,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.authnone / requested / required,是否校验客户端证书

3. Producer 关键配置速查

Java key 与 librdkafka (confluent-kafka) key 通常一致;下方括号内是 librdkafka 名(仅在不同时给出)。

配置默认推荐一句话作用
bootstrap.serversh1:9092,h2:9092,h3:9092集群引导地址
client.id""业务名_实例排障 / 配额必备
acksall(3.0+ 默认改为 allall0(不等)、1(Leader 写入)、all(ISR 全确认)
enable.idempotencetrue(3.0+ 默认)true开幂等,自动设 acks=all/retries=MAX/MaxInFlight≤5
retries2147483647默认失败重试次数
delivery.timeout.ms120000 (2min)120000 ~ 300000send → 成功 / 失败的总时长上限
request.timeout.ms3000030000单次请求超时
max.in.flight.requests.per.connection5≤5(开幂等时强制)同一连接未确认的请求上限
linger.ms05 ~ 50攒批等待时长,0 = 来一条发一条(吞吐差)
batch.size16384 (16KB)32KB ~ 256KB每个分区批次的字节上限
buffer.memory33554432 (32MB)64MB ~ 256MBProducer 内存缓冲区总量
max.block.ms6000060000buffer 满 / 元数据未就绪时 send 阻塞上限
compression.typenonezstdlz4压缩算法(端到端,Broker 不解压)
partitioner.classDefaultPartitioner(Sticky)默认自定义路由规则;3.0+ 默认 Sticky 而非纯 Murmur2
partitioner.ignore.keysfalsefalse配 Sticky 时是否忽略 Key
key.serializer / value.serializerStringSerializer / ByteArraySerializerJava 序列化器
transactional.idtx-<service>-<instance>必须稳定开事务必填,重启后重启 Producer 用同 ID 触发 PID Fencing
transaction.timeout.ms60000≤ Broker transaction.max.timeout.ms事务最大持续时间
metadata.max.age.ms300000 (5min)默认元数据强制刷新周期
connections.max.idle.ms540000 (9min)默认空闲连接关闭时间
interceptor.classes监控 / 染色Producer 拦截器链
send.buffer.bytes / receive.buffer.bytes131072 / 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实例稳定 IDStatic Membership,避免重启触发 Rebalance
client.id""业务名_实例排障 / 配额
key.deserializer / value.deserializerStringDeserializer / ByteArrayDeserializerJava 反序列化
auto.offset.resetlatestearliest(流处理)/ latest(实时业务)找不到 offset 时的行为,none 抛异常
enable.auto.committruefalse关掉自动提交,手动提交才准
auto.commit.interval.ms50005000自动提交周期
isolation.levelread_uncommittedread_committed(消费事务消息时)是否过滤未提交事务
max.poll.records500100 ~ 1000单次 poll 返回的最大记录数
max.poll.interval.ms300000 (5min)业务实际处理时间 × 1.5两次 poll 最大间隔,超时被踢出
session.timeout.ms45000(3.0+)45000 ~ 60000心跳超时,超过就踢出 group
heartbeat.interval.ms3000session/3心跳间隔
partition.assignment.strategy[RangeAssignor, CooperativeStickyAssignor](3.0+)CooperativeStickyAssignor分配策略,4 选 1
fetch.min.bytes11 ~ 50000Broker 至少攒多少字节才返回
fetch.max.bytes52428800 (50MB)默认单次 fetch 总字节上限
max.partition.fetch.bytes1048576 (1MB)message.max.bytes 对齐单分区单次 fetch 上限
fetch.max.wait.ms500500fetch.min.bytes 没攒够时最多等多久
request.timeout.ms3000030000请求超时
connections.max.idle.ms540000默认空闲连接关闭
interceptor.classes染色 / 监控Consumer 拦截器
check.crcstruetrue校验消息 CRC,关了快但有风险

5. JMX 重要指标速查

全部走 kafka.* MBean,配 Prometheus JMX Exporter 即可暴露。 下表分四类:Broker / Topic / Producer / Consumer

5.1 Broker 指标(30 个里挑了最关键的)

MBean类型一句话含义报警阈值建议
kafka.controller:type=KafkaController,name=ActiveControllerCountGauge当前 Broker 是否是 Controller集群和必须 = 1
kafka.controller:type=KafkaController,name=OfflinePartitionsCountGauge没有 Leader 的分区数> 0 立即告警
kafka.server:type=ReplicaManager,name=UnderReplicatedPartitionsGauge副本不同步的分区数> 0 持续 5min 告警
kafka.server:type=ReplicaManager,name=UnderMinIsrPartitionCountGaugeISR 数 < min.insync.replicas 的分区> 0 立即告警(影响写入)
kafka.server:type=ReplicaManager,name=AtMinIsrPartitionCountGaugeISR 数刚好 = min.isr 的分区持续较高需扩容
kafka.server:type=BrokerTopicMetrics,name=MessagesInPerSecMeter每秒入站消息数趋势
kafka.server:type=BrokerTopicMetrics,name=BytesInPerSec / BytesOutPerSecMeter每秒入 / 出字节趋势 + 容量规划
kafka.server:type=BrokerTopicMetrics,name=FailedProduceRequestsPerSec / FailedFetchRequestsPerSecMeter每秒失败请求数> 0 排障
kafka.network:type=RequestMetrics,name=RequestsPerSec,request={Produce,Fetch,Metadata}Meter各类请求 QPS趋势
kafka.network:type=RequestMetrics,name=TotalTimeMs,request=ProduceHistogramProduce 端到端耗时p99 > 1s 告警
kafka.network:type=RequestMetrics,name=LocalTimeMs / RequestQueueTimeMs / RemoteTimeMs / ResponseQueueTimeMs / ResponseSendTimeMsHistogram请求各阶段耗时排障定位瓶颈
kafka.server:type=KafkaRequestHandlerPool,name=RequestHandlerAvgIdlePercentMeterI/O 线程平均空闲率< 0.3 表示 IO 线程吃紧
kafka.network:type=SocketServer,name=NetworkProcessorAvgIdlePercentGauge网络线程平均空闲率< 0.3 表示网络线程吃紧
kafka.log:type=LogFlushStats,name=LogFlushRateAndTimeMsHistogramfsync 频率与耗时突增看磁盘
kafka.log:type=LogManager,name=LogDirectoryOfflineGauge离线日志目录数> 0 表示磁盘故障
kafka.log:type=Log,name=Size,topic=t1,partition=0Gauge单分区磁盘占用容量规划
kafka.log:type=Log,name=LogStartOffset / LogEndOffsetGauge单分区起始 / 结束 Offset排障
kafka.log.cleaner:type=LogCleanerManager,name=max-dirty-percentGaugeCompact Topic 最高 dirty 比例> 50% 看 cleaner 是否挂
kafka.log.cleaner:type=LogCleaner,name=cleaner-recopy-percentGaugeCleaner 工作量趋势
kafka.controller:type=ControllerStats,name=LeaderElectionRateAndTimeMsHistogramLeader 选举频率频繁选举 = 集群抖动
kafka.controller:type=ControllerStats,name=UncleanLeaderElectionsPerSecMeterunclean 选举次数> 0 排障(数据可能丢)
kafka.server:type=ReplicaFetcherManager,name=MaxLag,clientId=ReplicaGaugeFollower 最大 Lag> 10000 排障
kafka.server:type=ReplicaManager,name=IsrShrinksPerSec / IsrExpandsPerSecMeterISR 收缩 / 扩张速率频繁波动 = 抖动
kafka.server:type=DelayedOperationPurgatory,name=PurgatorySize,delayedOperation={Produce,Fetch}Gauge延迟操作队列长度突增表示 Broker 在等 ack / 数据
java.lang:type=GarbageCollector,name=*CounterJVM GC 时长 / 次数Full GC 频繁告警
java.lang:type=MemoryGauge堆 / 非堆使用> 80% 告警
java.lang:type=OperatingSystem,name=SystemCpuLoadGauge系统 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 / maxfetch 平均 / 最大耗时
records-lag-max最大 Lag(业务最关心的指标)
records-lead-min最小 Lead(消费位点距离 LogStartOffset)→ 太小说明数据要被清掉
commit-latency-avgoffset 提交耗时
commit-rate每秒 commit 次数
kafka.consumer:type=consumer-coordinator-metrics,client-id=*,join-rate / sync-rate / heartbeat-rateGroup 协议交互速率
last-rebalance-seconds-ago距离上次 Rebalance 时间,频繁掉低 = 抖动
assigned-partitions当前分配到的分区数

6. 内部 Topic 速查

Topic用途cleanup.policy副本因子建议查看方式
__consumer_offsets消费组提交的 offset / group 元数据compact3kafka-consumer-groups.sh --describe;裸读 kafka-console-consumer --topic __consumer_offsets --formatter "kafka.coordinator.group.GroupMetadataManager\$OffsetsMessageFormatter"
__transaction_state事务元数据(PIDEpoch、参与分区、状态)compact3kafka-dump-log.sh --transaction-log-decoder
__cluster_metadataKRaft 集群元数据日志(替代 ZK)delete(带快照)等于 Controller Quorum 大小kafka-metadata-shell.sh --snapshot <file>kafka-dump-log.sh --cluster-metadata-decoder
_schemasConfluent Schema Registry 存放 schemacompact3kafka-console-consumer --topic _schemas --from-beginning
connect-configsKafka Connect 配置compact3同上
connect-offsetsConnect Source connector 的位点compact25 partitions / 3 副本同上
connect-statusConnect 任务状态compact3同上
<app-id>-changelog-*Kafka Streams 状态 store 的 changelogcompact(KTable)/ compact,delete(带 retention 的 store)跟随 Streams 配置用 Streams Reset 工具清理
<app-id>-repartition-*Streams 中间重分区 Topicdelete + 短 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.sizemessage.max.bytesreplica.fetch.max.bytes(三件套)或拆分消息
BufferExhaustedException / Failed to allocate memory within timeoutbuffer.memory 满了buffer.memory / 加 linger.ms 让批次更大、消费下游加速
OutOfOrderSequenceException幂等 Producer 序列号乱序,多半是网络重排 + max.in.flight > 5max.in.flight ≤ 5,开 enable.idempotence=true
ProducerFencedExceptiontransactional.id 的新 Producer 起来了,旧的被挤掉正常,旧 Producer 直接退出
InvalidProducerEpochException事务超时被中止业务侧重新 beginTransaction()
NotEnoughReplicasException / NotEnoughReplicasAfterAppendExceptionISR 数 < min.insync.replicas,写不进去(acks=all 时)修复掉队 Broker;临时降级可调小 min.insync.replicas(不推荐)
UnknownTopicOrPartitionExceptionTopic 不存在或刚删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=noneearliest/latest;或人工 seek
InconsistentGroupProtocolException同 group 不同实例 partition.assignment.strategy 不一致全部统一
RebalanceInProgressException收到这个错时正在 Rebalance,提交无效onPartitionsRevoked 钩子里提交

7.3 Broker / 集群

错误 / 现象含义处理
Unable to write to meta.properties / cluster.id ... doesn't matchKRaft 节点 cluster ID 与磁盘上不一致删数据卷重建;不要手动改 meta.properties
Shutdown broker because all log dirs in /xxx have failed全部数据盘故障修磁盘或迁移到新节点
Cluster authorization failedACL 拒绝kafka-acls --list 确认;事务必须授 IdempotentWriteWrite
INVALID_FETCH_SESSION_EPOCHFetch 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 中某些字段提到 keyfields
HoistField$Key/Value把 schemaless 包一层 structfield
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$Valuevalue 字段抽到 Headerfields, headers
DropHeaders删指定 Headerheaders

典型链(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=false

9. Kafka 版本里程碑速查

版本发布年关键特性
0.72012LinkedIn 原始版本,Scala 写的雏形
0.82013副本机制首次出现(Replication);Producer / Consumer 协议成型
0.92015新 Consumer API(基于 Group Coordinator)、SSL / SASL 安全Quotas、Kafka Connect 雏形
0.102016Kafka Streams 发布;消息加 timestamp;rack-aware 副本分配
0.10.1 / 0.10.22016max.poll.interval.ms(解决长处理踢出);客户端 / Broker 解耦
0.112017幂等 Producer + 事务 + Exactly Once(PID + Epoch + 序列号 + Transaction Coordinator);V2 RecordBatch(更小、加 Header)
1.02017协议稳定 GA;JBOD 改进;Java 9 支持
1.12018单 Broker 支持百万级分区;Controller 启动加速
2.02018OAuthBearer SASL;前缀 ACL;Java 11 支持;unclean.leader.election.enable 默认改为 false
2.12018zstd 压缩;CreateTopics 协议改进
2.22019Default partitioner 改进(KIP-389 限制单连接请求);ACL 支持 prefix
2.32019Incremental Cooperative Rebalance(KIP-429);Static Membership(KIP-345,group.instance.id
2.42019KIP-392 follower fetching(最近副本读,跨机房友好);KIP-466 用 Java 11
2.52020KIP-447 事务 API 增强(每分区一个 Producer);TLS 1.3
2.62020client.id quota 改进;KIP-573 文档 metric 化
2.72020KRaft 早期预览(KIP-500);End-to-End Latency 指标
2.82021KRaft 模式 Preview 可用(无需 ZK 启动);KIP-516 topic ID
3.02021KRaft 进入 Production-ready (Preview);移除 0.10 之前的旧协议;默认 acks=all + enable.idempotence=true;StickyAssignor 默认
3.1 / 3.22022KRaft GA 准备;Tiered Storage 早期 KIP;ZK 模式开始进入「不推荐」
3.32022KRaft 模式 GA(KIP-833),生产可用
3.42023KRaft 元数据快照、ZK → KRaft 迁移工具早期版(KIP-866)
3.52023KRaft 迁移工具进一步完善;新协议 Consumer Group 设计开始(KIP-848)
3.62023Tiered Storage Early Access(KIP-405),冷数据落 S3
3.72024KIP-848 next-gen Consumer Group Protocol Preview;MirrorMaker 2 改进
3.82024本教程基线;KIP-1011 KRaft 多版本兼容;JBOD on KRaft
3.92024KRaft 模式 ZK 迁移功能稳定,最后一个支持 ZK 的版本系列
4.02025完全移除 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 全家桶速查脑图


📌 本附录的姐妹篇