1. 这份速记为谁整理
上周帮运维同事处理一个Kafka集群的问题,前前后后折腾了一天,最后发现是消费者端的参数配错了。那天晚上我坐在工位上想,接触Kafka这几年,踩过的坑、翻过的文档、答过的面试题其实都能串成一条线,但每次遇到具体问题还是要现查。第二天刚好要做内部培训,我就干脆整理了一份速记,把原理、部署、命令、排障、面试题这五个方向全部写进去,培训完直接发给团队当手册用。
这份速记的目标读者很明确:正在做Kafka安装配置的运维、准备Kafka面试题的后端开发、以及刚开始用Kafka但对细节不太有把握的测试和SRE同学。它不替代官方文档,但把官方文档里散落的关键点,按实操顺序重新组织了一遍。你花半小时读完,遇到实际问题时能少走很多弯路。
先说好,这篇速记不是从零讲“Kafka是什么”的入门教材,而是把“怎么用好、怎么排障、怎么回答面试官”这些最常出现的问题一次性讲透。所以对Kafka有基本了解的朋友读起来最舒服,完全零基础的同学可以先补一点背景知识再回来看。
2. 先把Kafka的原理串一遍
很多人上来就配参数、装集群,结果出了问题完全不知道从哪下手,根本原因是对底层机制只有一个模糊的印象。这里我先把最核心的原理串一遍,后面所有操作和排障思路都建立在它上面。
2.1 消息是怎么落进磁盘的
Kafka在本质上不是一个传统意义上的消息队列,而是一个分布式的日志提交系统。每个Topic在Broker上对应一个目录,目录内部按分区组织,每个分区就是一段只能追加写入的日志文件。写入消息时,Broker把数据顺序追加到磁盘上的Segment文件里,同时更新索引。关键点在于“顺序追加”——这是Kafka在高吞吐场景下依然表现良好的基石,因为机械硬盘顺序写和随机写的性能差距可以达到几个数量级,即使是普通SSD,顺序IO也比随机IO稳定得多。
另一个容易被忽略的设计是“消费不删除消息”。普通消息队列的经典模型是消息被消费后就移除,但Kafka不是这样。Broker只根据保留策略删除过期消息,比如按时间(retention.ms,默认168小时)或按总大小(retention.bytes)清理,消费进度完全由消费者自己维护偏移量(offset)。正因为这样,你才能看到 Topic 里的历史数据,“从头消费”“重复消费”这类需求也才有实现的可能。
还有一个提得很多但很容易被说错的概念是零拷贝。Kafka消费消息时,数据从磁盘读到页缓存就可以直接通过sendfile系统调用发送到网卡,不需要在用户态内存里复制。这也是消费者能拉取大吞吐数据而不拖垮Broker的原因之一。面试官问“Kafka为什么快”的时候,顺序写、页缓存、零拷贝、批量操作这四个点能答出来,基本就能过关。
2.2 分区、副本和消费组
分区是Kafka并行度的核心单位。一个Topic的数据分散在多个分区里,每个分区内部有序,不同分区之间没有全局顺序。生产者发消息时,可以通过key的哈希决定往哪个分区写;没有key的话走round-robin或随机策略。消费者组(Consumer Group)里的每个消费者会被分配若干个分区,同一分区只会被同一个组里的一个消费者消费——这就是Kafka实现水平扩展的机制:增加分区数,然后给消费者组加消费者,消费吞吐就往上涨。
副本机制解决的是高可用问题。每个分区有多个副本,一个Leader和若干个Follower。生产者只往Leader写入,Follower去Leader拉数据,尽量追平。ISR(In-Sync Replicas)是“跟得上进度”的副本集合,Leader挂了以后,会从ISR里选一个新Leader出来。如果Follower落后太多,超过 replica.lag.time.max.ms,就会被踢出ISR。这个机制决定了消息的可靠性和可用性之间的平衡,后面讲ack和延迟排查时还会反复遇到。
2.3 从ZooKeeper到KRaft
如果是老版本的Kafka,都会有一个ZooKeeper集群负责元数据和Leader选举,这也是很多人部署Kafka时最烦的一件事:明明Kafka本身不复杂,却要先维护一套ZK,版本兼容还要对得上。Kafka 2.8开始引入KRaft模式,用内部的事件日志和Raft共识协议接管了ZK的职责。到3.x版本,KRaft已经生产可用,4.0之后ZK就被彻底移除了。
现在的集群部署,我建议直接用KRaft模式而不是再搭ZK,原因很简单:少一个依赖就少一类故障。KRaft给每个Broker分配节点ID,通过 controller.quorum.voters 参数指定Controller列表,启动前要用 kafka-storage random-uuid 生成集群ID,然后 kafka-storage format 格式化存储。这套流程比ZK时代简洁很多,出错概率也低。下面部署部分我会直接基于KRaft来写。
3. 从单机到集群的部署实操
部署这一块,网上的教程质量参差不齐,尤其是Windows下用Docker、WSL、Compose混着来的场景,坑特别多。我按从易到难的顺序给出三种情况:单机快速验证、三节点集群、加SSL认证。
3.1 在Windows上用Docker搭一个能跑的Kafka
很多人在Windows上装了Docker Desktop,拉一个Kafka镜像起来,结果发现生产者连不上、消费者连不上、看到日志在跑但什么数据都收不到。绝大多数问题出在 listeners 和 advertised.listeners 这两个参数上。Kafka的协议允许Broker同时监听多个地址,外部客户端连进来的地址却需要单独声明。容器里监听的是0.0.0.0:9092,但你物理机上的客户端得能通过某个地址访问到它。Docker Desktop下最简单的做法是把宿主机的地址映射为 host.docker.internal,所以环境变量可以这样配:
KAFKA_CFG_LISTENERS=PLAINTEXT://:9092 KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://host.docker.internal:9092我习惯用bitnami/apache-kafka镜像,理由有两个:一是镜像更新频率高,二是环境变量命名直接对应配置文件,不容易踩兼容性坑。Compose配置长这样:
services: kafka: image: bitnami/apache-kafka:3.7 ports: - "9092:9092" environment: KAFKA_CFG_NODE_ID: 0 KAFKA_CFG_PROCESS_ROLES: controller,broker KAFKA_CFG_CONTROLLER_QUORUM_VOTERS: 0@kafka:9093 KAFKA_CFG_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093 KAFKA_CFG_ADVERTISED_LISTENERS: PLAINTEXT://host.docker.internal:9092 KAFKA_CFG_CONTROLLER_LISTENER_NAMES: CONTROLLER KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT volumes: - kafka_data:/bitnami/kafka volumes: kafka_data:启动之后先在容器里测试,避免物理机和Docker网络差异干扰判断:
docker compose up -d docker exec -it kafka-kafka-1 /bin/bash kafka-topics.sh --bootstrap-server localhost:9092 --create --topic test --partitions 1 --replication-factor 1容器内能创建Topic,物理机上执行同样的命令却失败,就检查端口映射和防火墙。Windows上如果连不上 localhost,先用 127.0.0.1 试试,个别版本的WSL2网络转发偶尔会有些怪问题。
注意:Docker Desktop默认给WSL2分配的内存只有2GB左右,Kafka Broker默认的堆内存就比较吃紧,容器很容易被OOM杀掉。建议在Docker Desktop的Resources设置里把内存调到4GB以上,或者给容器额外加上 KAFKA_HEAP_OPTS=-Xmx512m -Xms512m。
3.2 三节点集群怎么配
单机验证通过以后,集群的架构其实没有本质变化,只是Controller和Broker角色分到多个节点上,每个节点都要有独立的node.id和日志目录。以三台机器为例,假设IP分别是10.0.1.11、10.0.1.12、10.0.1.13,通用的配置文件可以这样写:
process.roles=broker,controller node.id=1 controller.quorum.voters=1@10.0.1.11:9093,2@10.0.1.12:9093,3@10.0.1.13:9093 listeners=PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093 advertised.listeners=PLAINTEXT://10.0.1.11:9092 controller.listener.names=CONTROLLER listener.security.protocol.map=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT log.dirs=/data/kafka-logs num.partitions=3 default.replication.factor=3 min.insync.replicas=2节点2和节点3只需把node.id改成2、3,advertised.listeners改成对应IP就行。启动之前先做两件事,一是生成集群ID,二是格式化存储目录:
KAFKA_CLUSTER_ID="$(kafka-storage.sh random-uuid)" kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c /path/to/server.properties注意,多节点集群必须在每台机器上使用同一个KAFKA_CLUSTER_ID。如果三台机器拿到的集群ID不一致,Controller之间无法选主,表现就是进程启动后过一会自动退出,或Topic创建一直超时。格式化完,直接启动Broker进程即可。验证集群是否健康:
kafka-topics.sh --bootstrap-server 10.0.1.11:9092 --describe kafka-metadata.sh --bootstrap-server 10.0.1.11:9092 --describe把 create topic 的 replication-factor 设为3,查看描述时看到每个分区都有三个副本且ISR数量为3,就说明集群状态正常。
3.3 接入SSL的几个关键配置
SSL是我遇到最多“配置完一头雾水”的部分,因为涉及证书生成、listener命名、客户端参数三层,任何一层对不上都连不上。先用keytool生成三样东西:每个Broker的密钥库、信任库,以及一份CA证书。生产环境的证书通常由统一CA签发,这里只演示自签证书的思路。
为每个Broker生成密钥库并导出证书签名请求:
keytool -keystore server.keystore.jks -alias broker -validity 3650 -genkeypair -keyalg RSA -storepass changeit -dname "CN=broker1" keytool -keystore server.keystore.jks -alias broker -certreq -file broker1.csr -storepass changeit openssl x509 -req -CA ca.crt -CAkey ca.key -in broker1.csr -out broker1.crt -days 3650 -CAcreateserial keytool -keystore server.keystore.jks -alias broker -importcert -file ca.crt -storepass changeit keytool -keystore server.keystore.jks -alias broker -importcert -file broker1.crt -storepass changeit然后把CA证书导入信任库:
keytool -keystore server.truststore.jks -alias CARoot -importcert -file ca.crt -storepass changeit服务端三个关键配置是:
listeners=SSL://0.0.0.0:9094,CONTROLLER://0.0.0.0:9093 advertised.listeners=SSL://10.0.1.11:9094 listener.security.protocol.map=CONTROLLER:PLAINTEXT,SSL:SSL,PLAINTEXT:PLAINTEXT ssl.keystore.location=/path/to/server.keystore.jks ssl.keystore.password=changeit ssl.key.password=changeit ssl.truststore.location=/path/to/server.truststore.jks ssl.truststore.password=changeit client.auth=none客户端这边,只要把 keytool 生成的客户端信任库指到CA,并配置安全协议为SSL:
kafka-console-producer.sh --bootstrap-server 10.0.1.11:9094 --topic test \ --producer-property security.protocol=SSL \ --producer-property ssl.truststore.location=/path/to/client.truststore.jks \ --producer-property ssl.truststore.password=changeit最常见的SSL问题有两类:一类是证书里的CN或SAN和advertised.listeners的主机名对不上,客户端校验主机名失败;另一类是信任库没导入CA,或者客户端只导入了服务端证书却没有导入CA。遇到SSL握手失败先把自定义的 hostname.verification 功能放一边,确认“证书链完整”“主机名匹配”这两件事,八成问题就解决了。
4. 日常运维的命令行速查
命令这一块,我觉得比GUI工具更值得掌握,因为很多排查场景下你根本来不及打开工具,直接ssh到机器上敲命令是最快的路径。
4.1 主题的增删查改
创建Topic时的两个核心参数是 partitions 和 replication-factor。生产环境至少把副本数设为2或3,单副本在Broker宕机时整个分区直接不可用:
kafka-topics.sh --bootstrap-server localhost:9092 --create \ --topic order-events --partitions 6 --replication-factor 3查看现有Topic列表和详情:
kafka-topics.sh --bootstrap-server localhost:9092 --list kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic order-events修改分区数量时只能增加不能减少,因为减少分区涉及数据迁移和顺序重定义,Kafka设计上不支持:
kafka-topics.sh --bootstrap-server localhost:9092 --alter \ --topic order-events --partitions 12删除Topic时要注意,如果Broker没有开启 delete.topic.enable=true,删除命令不会真正生效。即使开启了,删除操作也是异步的,不是立即可见:
kafka-topics.sh --bootstrap-server localhost:9092 --delete --topic order-events4.2 生产消费命令启动一次会一直运行吗
这个问题在热搜里出现不是没道理的。很多新手执行了 kafka-console-producer.sh,敲了一条消息,回车后发现进程没有任何退出的意思,还以为自己操作错了。实际上,控制台生产者默认的行为就是不断从标准输入读数据并发送,直到你按Ctrl+C或者输入流结束。这是一种“持续运行”的交互式工具,不是执行一次发完就退出的命令。如果只想发固定数量的消息然后结束,可以给生产者脚本配置 max-messages 参数,但没有这个参数的场景下,直接维护输入流或者Ctrl+C是标准做法。
控制台消费者更让人困惑。kafka-console-consumer.sh 启动后,会持续轮询Broker并等待新消息,这一行为由消费组和位点提交机制决定。默认情况下它是“实时”的,也就是说启动之前Topic里已有的历史数据,它一条都不会消费,只消费启动之后新到的消息。所以你会看到消费者启动后屏幕上没有任何输出,那不是卡死,是它在安静地等待。想让它启动后就把历史数据也读出来,需要加 --from-beginning:
kafka-console-consumer.sh --bootstrap-server localhost:9092 \ --topic order-events --from-beginning --max-messages 100加了 --max-messages 100,消费者读满100条就会自动退出。如果想按时间退出,可以用 --timeout-ms 指定最长等待时间,时间到自动退出,适合在脚本里做数据抽样。
4.3 查看Topic里的历史数据
“查看topic中的数据”是运维里最常遇到的需求。最简单的方式就是用上面说的消费者加 --from-beginning 参数。但生产环境的Topic数据量动辄几千万条,直接从头消费会把控制台刷到爆炸,还会占用消费组位点资源(实际上控制台消费者不提交位点,但网络和内存开销仍然很大)。所以我一般会这样组合:
kafka-console-consumer.sh --bootstrap-server localhost:9092 \ --topic order-events \ --partition 0 --offset 100 \ --max-messages 20 \ --property print.key=true --property print.timestamp=true指定分区和偏移量,从第100条开始读20条,既能看到数据,又不会把整个Topic读完。不要忘了,默认的控制台消费者格式化的是消息体。想输出key和时间戳,就加print.key和print.timestamp。想输出完整的头部信息,可以加 --property print.headers=true。
如果Broker上的日志文件已经落盘,你想直接检查Segment文件而不是通过网络拉数据,可以用 kafka-dump-log:
kafka-run-class.sh kafka.tools.DumpLogSegments \ --files /data/kafka-logs/order-events-0/00000000000000000100.log \ --print-data-log这个命令常用于排查“数据写进了磁盘但消费不到”“消息内容异常”这类问题,因为它绕开了Consumer API,直接看到存储层的二进制内容解码结果。排查日志类问题时很有效,但注意要生产环境谨慎使用,大文件会输出巨量信息。
5. 高频故障排查清单
这一节全部是我自己踩过或者帮别人处理过的真实问题,不是从文档里抄出来的。每个问题给出直接的排查路径和参数位置。
5.1 消息延迟高,先查哪里
“消息延迟高”是个很笼统的说法,必须先拆分场景。以我的经验,至少分成三类:生产者发不出去、Broker写不动、消费者消费不过来。
生产者发不出去时,先看请求是否超时。如果设置了 acks=all,生产者会等待ISR里所有副本确认。ISR里如果有一个副本落后太多,写请求就会卡住。这个场景我遇到过好多次,客户端日志里全是 NotEnoughReplicas 或 TimeoutException。解决方式是把 min.insync.replicas 调低或者把 acks 降到1,但如果要求不丢数据,还是得从副本恢复下手,而不是降低可靠性。
Broker写不动,常见原因是磁盘IO瓶颈。Kafka的理想状态是依赖页缓存做写缓冲,但如果操作系统内存不足,或者磁盘本身是共享的云盘,写性能会直线下降。建议先看 iostat,特别是 avgqu-sz 和 %util 两个指标,如果持续接近100%,说明磁盘确实顶不住。另一个容易被忽略的点是分区数量过少导致单分区写入压力过大,合理增加分区数能立刻改善吞吐。
消费者滞后是最常见的“延迟高”来源。生产者发送速度正常,消费者消费不过来,Lag指标持续上涨。这时要看消费者的 max.poll.records、max.poll.interval.ms 和处理逻辑耗时。如果单条消息处理耗时超过 max.poll.interval.ms,消费者会被判定为死亡并触发重平衡,重平衡期间又暂停消费,造成“处理不过来-重平衡-继续处理不过来”的恶性循环。解决办法是调大 max.poll.records 或者开启异步处理。
下面这张表可以作为排查入口快速对照:
| 现象 | 重点排查项 | 常用命令/指标 |
|---|---|---|
| 生产超时 | acks设置、ISR状态、网络 | kafka-topics describe、客户端日志 |
| 磁盘写满/IO高 | 磁盘容量、IO Util、页缓存 | iostat、df -h |
| 消费者不消费 | 消费组状态、再平衡日志、偏移量 | kafka-consumer-groups describe |
| 单分区热点 | 分区数、key分布 | kafka-producer-perf-test、topic describe |
5.2 OOM和容器被杀
Kafka的OOM分两类,一类是Broker进程堆内存溢出,另一类是容器本身因为超内存被杀。先看JVM堆。Kafka的默认堆内存是1GB,对于生产业务来说往往不够,但也不是越大越好。Broker的核心是页缓存,堆内存只是用来处理网络连接和内部状态,堆设得过大反而挤压操作系统页缓存空间,得不偿失。我一般推荐4GB到6GB,根据Topic数量和连接数调整。通过 KAFKA_HEAP_OPTS 环境变量设置:
KAFKA_HEAP_OPTS="-Xmx4g -Xms4g"堆溢出排查时先看GC日志和堆转储。生产环境建议把GC日志打开,加上 -XX:+HeapDumpOnOutOfMemoryError,这样OOM时能拿到heap dump做分析。需要特别注意的是,不要只盯着Broker,消费者端也会OOM。有的消费者会把大量消息拉到本地内存做批量处理,如果 max.poll.records 设置过大、单条消息体又很大,JVM堆很容易被打爆。
容器被杀则是另一种情况。Docker会强制限制容器内存,如果JVM堆大小加上堆外内存、页缓存、网络缓冲区超出了容器限制,进程会被内核杀掉,表现就是容器直接退出,docker logs 里看不到明显的Java异常。排查命令是 docker inspect 看ExitCode,如果是137,基本就是OOM Kill。处理方法是在容器编排层面给Kafka单独设置内存上限,同时确保JVM堆设置小于容器内存上限,留出足够余量。
5.3 监控接入:ELK与OTel的姿势
Kafka本身的监控指标通过JMX暴露,生产环境一般会用JMX Exporter把指标拉到Prometheus,再用Grafana画板子。但和ELK结合是另一种常见的运维姿势:Kafka经常作为日志管道里的核心缓冲层,日志采集器(比如Filebeat)把业务日志写入Kafka,Logstash从Kafka消费日志,再写入Elasticsearch,最后用Kibana做可视化。这就是最典型的Elastic路径,Kafka在这里的角色是削峰填谷和异步解耦,下游Logstash挂了也不影响业务日志的采集。
监控Kafka自身时,下面几个指标必须时刻盯住:
- UnderReplicatedPartitions:如果长期大于0,说明有副本没跟上,集群处于高可用降级状态
- OfflinePartitions:必须为0,只要有就说明有分区完全不可用
- ActiveControllerCount:正常情况下集群里只有一个,出现多个是脑裂的严重信号
- Consumer Lag:每个消费组的累积延迟,是判断业务是否受影响的最直接指标
OTel(OpenTelemetry)这边,很多人关心的是Kafka消息链路追踪。如果你用的是OTel Collector,可以通过Kafka Receiver/Exporter把Trace数据和日志数据作为消息流接入Kafka。比如把OTel Collector作为生产者,把业务服务打包成Trace数据发到Kafka,下游再由Collector消费完成导出。这种方式的好处是架构统一,Traces、Metrics、Logs都走同一条管道。实际落地时可以把Collector部署在Kafka集群附近,避免网络跳数过大影响消息写入延迟。
6. 高频面试题速记本
这份速记的最后一节给准备面试的同学。我看过不少面试题整理,很多是抄来抄去的概念解释,太泛。这里我只挑真正能区分“用过”和“没用过”的问题。
6.1 基础题:必须能脱口而出
先说说面试官几乎必问的几个基础问题。
Kafka为什么快。这个前面讲原理时覆盖过,回答时把顺序写、分区并行、页缓存+零拷贝、批量发送和压缩这几点串起来讲,会比单独背一两个名词更有说服力。
ack机制怎么选。ack=0表示不等待Broker确认,可能丢消息;ack=1表示Leader写入成功就返回,Leader挂了可能丢;ack=all表示ISR全部确认,最可靠但延迟最高。关键是结合min.insync.replicas一起讲,才不会显得只会背参数。
消费者重平衡是个什么过程。当消费者组的成员变化、订阅的Topic分区变化时,组协调器会触发重平衡,把分区重新分配给消费者。重平衡期间消费会暂停,所以“频繁重平衡”本身就是一个性能问题。回答说清楚“如何触发-协调器怎么选-分区分配合约”这三个环节,基本就过关了。
6.2 进阶题:考察有没有真正踩过坑
真正拉分的是这一部分。
为什么分区不是越多越好。分区太多会导致文件句柄过多、副本同步压力大、Controller的元数据管理开销上升,而且如果单分区消息量不大,分区数增加并不会提升消费吞吐。大家常说的“分区瓶颈最大化”是个理想模型,实际得权衡Broker数量和Topic数据量。
如何保证消息不重复消费。严格意义上,Kafka的at-least-once语义下重复消费是可能发生的,比如消费者处理完消息但没来得及提交偏移量就挂了。解决方式一般是让消费逻辑支持幂等,或者把偏移量提交和业务操作做成原子步骤。回答里如果能提到“用外部存储记录消费位点”的方案,会比只背“enable.auto.commit=false”更出彩。
顺序消费怎么保证。单分区内有序,跨分区没有全局顺序。如果业务要求严格有序,要么一个业务只用一个分区,要么把key哈希到同一个分区。但减少分区数会影响吞吐,所以“顺序和并行度”其实是互相取舍的关系,面试官就是想看你能不能说出这个权衡。
这些问题往往没有唯一标准答案,重点是体现出你对参数选型背后的代价和场景有实实在在的理解。比起背答案,用一两句自己经历过的踩坑案例来讲,效果会好得多。
我自己在实际使用中还有一个习惯:把Kafka的常用命令、常见错误码、参数默认值都集中放在一个本地备忘录里,每次处理完问题就追加一条,现在回头看已经攒了几百条。这份速记算是我那个备忘录里比较成体系的一部分。真遇到文档查不到的问题,多看看Broker日志和客户端日志里的warn级别信息,通常能找到比网上教程更准确的线索。希望这份速记能帮你少熬几个夜。