
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拉数据尽量追平。ISRIn-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_LISTENERSPLAINTEXT://:9092 KAFKA_CFG_ADVERTISED_LISTENERSPLAINTEXT://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: 0kafka: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.rolesbroker,controller node.id1 controller.quorum.voters110.0.1.11:9093,210.0.1.12:9093,310.0.1.13:9093 listenersPLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093 advertised.listenersPLAINTEXT://10.0.1.11:9092 controller.listener.namesCONTROLLER listener.security.protocol.mapCONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT log.dirs/data/kafka-logs num.partitions3 default.replication.factor3 min.insync.replicas2节点2和节点3只需把node.id改成2、3advertised.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 CNbroker1 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服务端三个关键配置是listenersSSL://0.0.0.0:9094,CONTROLLER://0.0.0.0:9093 advertised.listenersSSL://10.0.1.11:9094 listener.security.protocol.mapCONTROLLER:PLAINTEXT,SSL:SSL,PLAINTEXT:PLAINTEXT ssl.keystore.location/path/to/server.keystore.jks ssl.keystore.passwordchangeit ssl.key.passwordchangeit ssl.truststore.location/path/to/server.truststore.jks ssl.truststore.passwordchangeit client.authnone客户端这边只要把 keytool 生成的客户端信任库指到CA并配置安全协议为SSLkafka-console-producer.sh --bootstrap-server 10.0.1.11:9094 --topic test \ --producer-property security.protocolSSL \ --producer-property ssl.truststore.location/path/to/client.truststore.jks \ --producer-property ssl.truststore.passwordchangeit最常见的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.enabletrue删除命令不会真正生效。即使开启了删除操作也是异步的不是立即可见kafka-topics.sh --bootstrap-server localhost:9092 --delete --topic order-events4.2 生产消费命令启动一次会一直运行吗这个问题在热搜里出现不是没道理的。很多新手执行了 kafka-console-producer.sh敲了一条消息回车后发现进程没有任何退出的意思还以为自己操作错了。实际上控制台生产者默认的行为就是不断从标准输入读数据并发送直到你按CtrlC或者输入流结束。这是一种“持续运行”的交互式工具不是执行一次发完就退出的命令。如果只想发固定数量的消息然后结束可以给生产者脚本配置 max-messages 参数但没有这个参数的场景下直接维护输入流或者CtrlC是标准做法。控制台消费者更让人困惑。kafka-console-consumer.sh 启动后会持续轮询Broker并等待新消息这一行为由消费组和位点提交机制决定。默认情况下它是“实时”的也就是说启动之前Topic里已有的历史数据它一条都不会消费只消费启动之后新到的消息。所以你会看到消费者启动后屏幕上没有任何输出那不是卡死是它在安静地等待。想让它启动后就把历史数据也读出来需要加 --from-beginningkafka-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.keytrue --property print.timestamptrue指定分区和偏移量从第100条开始读20条既能看到数据又不会把整个Topic读完。不要忘了默认的控制台消费者格式化的是消息体。想输出key和时间戳就加print.key和print.timestamp。想输出完整的头部信息可以加 --property print.headerstrue。如果Broker上的日志文件已经落盘你想直接检查Segment文件而不是通过网络拉数据可以用 kafka-dump-logkafka-run-class.sh kafka.tools.DumpLogSegments \ --files /data/kafka-logs/order-events-0/00000000000000000100.log \ --print-data-log这个命令常用于排查“数据写进了磁盘但消费不到”“消息内容异常”这类问题因为它绕开了Consumer API直接看到存储层的二进制内容解码结果。排查日志类问题时很有效但注意要生产环境谨慎使用大文件会输出巨量信息。5. 高频故障排查清单这一节全部是我自己踩过或者帮别人处理过的真实问题不是从文档里抄出来的。每个问题给出直接的排查路径和参数位置。5.1 消息延迟高先查哪里“消息延迟高”是个很笼统的说法必须先拆分场景。以我的经验至少分成三类生产者发不出去、Broker写不动、消费者消费不过来。生产者发不出去时先看请求是否超时。如果设置了 acksall生产者会等待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 describe5.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把业务日志写入KafkaLogstash从Kafka消费日志再写入Elasticsearch最后用Kibana做可视化。这就是最典型的Elastic路径Kafka在这里的角色是削峰填谷和异步解耦下游Logstash挂了也不影响业务日志的采集。监控Kafka自身时下面几个指标必须时刻盯住UnderReplicatedPartitions如果长期大于0说明有副本没跟上集群处于高可用降级状态OfflinePartitions必须为0只要有就说明有分区完全不可用ActiveControllerCount正常情况下集群里只有一个出现多个是脑裂的严重信号Consumer Lag每个消费组的累积延迟是判断业务是否受影响的最直接指标OTelOpenTelemetry这边很多人关心的是Kafka消息链路追踪。如果你用的是OTel Collector可以通过Kafka Receiver/Exporter把Trace数据和日志数据作为消息流接入Kafka。比如把OTel Collector作为生产者把业务服务打包成Trace数据发到Kafka下游再由Collector消费完成导出。这种方式的好处是架构统一Traces、Metrics、Logs都走同一条管道。实际落地时可以把Collector部署在Kafka集群附近避免网络跳数过大影响消息写入延迟。6. 高频面试题速记本这份速记的最后一节给准备面试的同学。我看过不少面试题整理很多是抄来抄去的概念解释太泛。这里我只挑真正能区分“用过”和“没用过”的问题。6.1 基础题必须能脱口而出先说说面试官几乎必问的几个基础问题。Kafka为什么快。这个前面讲原理时覆盖过回答时把顺序写、分区并行、页缓存零拷贝、批量发送和压缩这几点串起来讲会比单独背一两个名词更有说服力。ack机制怎么选。ack0表示不等待Broker确认可能丢消息ack1表示Leader写入成功就返回Leader挂了可能丢ackall表示ISR全部确认最可靠但延迟最高。关键是结合min.insync.replicas一起讲才不会显得只会背参数。消费者重平衡是个什么过程。当消费者组的成员变化、订阅的Topic分区变化时组协调器会触发重平衡把分区重新分配给消费者。重平衡期间消费会暂停所以“频繁重平衡”本身就是一个性能问题。回答说清楚“如何触发-协调器怎么选-分区分配合约”这三个环节基本就过关了。6.2 进阶题考察有没有真正踩过坑真正拉分的是这一部分。为什么分区不是越多越好。分区太多会导致文件句柄过多、副本同步压力大、Controller的元数据管理开销上升而且如果单分区消息量不大分区数增加并不会提升消费吞吐。大家常说的“分区瓶颈最大化”是个理想模型实际得权衡Broker数量和Topic数据量。如何保证消息不重复消费。严格意义上Kafka的at-least-once语义下重复消费是可能发生的比如消费者处理完消息但没来得及提交偏移量就挂了。解决方式一般是让消费逻辑支持幂等或者把偏移量提交和业务操作做成原子步骤。回答里如果能提到“用外部存储记录消费位点”的方案会比只背“enable.auto.commitfalse”更出彩。顺序消费怎么保证。单分区内有序跨分区没有全局顺序。如果业务要求严格有序要么一个业务只用一个分区要么把key哈希到同一个分区。但减少分区数会影响吞吐所以“顺序和并行度”其实是互相取舍的关系面试官就是想看你能不能说出这个权衡。这些问题往往没有唯一标准答案重点是体现出你对参数选型背后的代价和场景有实实在在的理解。比起背答案用一两句自己经历过的踩坑案例来讲效果会好得多。我自己在实际使用中还有一个习惯把Kafka的常用命令、常见错误码、参数默认值都集中放在一个本地备忘录里每次处理完问题就追加一条现在回头看已经攒了几百条。这份速记算是我那个备忘录里比较成体系的一部分。真遇到文档查不到的问题多看看Broker日志和客户端日志里的warn级别信息通常能找到比网上教程更准确的线索。希望这份速记能帮你少熬几个夜。