兄弟们,Kafka这块的知识点,说难不难,说简单也真不简单。我见过太多人,平时用着没问题,一上生产或者一面试就露馅,知识点全是散的。最近后台问Kafka的人特别多,从“Windows怎么装Kafka”到“Kafka消息延迟高咋排查”都有,今天我就把这些东西串起来,写成一份能直接落地、能照着敲的速记笔记。这份内容不是抄官方文档,是我自己在本地玩、在线上踩坑总结出来的,从部署到命令实操,从性能排查到面试题,一条线捋下来,你花一晚上过一遍,绝对比零散刷十篇博客有用。
先说好,本文默认你用Docker,因为现在真没必要自己下载tar包去配ZooKeeper了,KRaft模式才是主流。我不讲那些过时的配置,直接告诉你现在该怎么玩。
1. 先立坐标系:Kafka到底在解决什么问题
很多人学Kafka上来就背术语,结果越背越乱。我建议你先搞清楚一件事:Kafka在系统里到底扮演什么角色。
1.1 一个消息引擎的三个身份
Kafka本质上是一个“分布式提交日志”,但它在不同的场景下有三副面孔。
第一副面孔:消息队列。这是它最广为人知的身份。生产者往里扔消息,消费者从里面取消息。和RabbitMQ这类传统队列比,Kafka的消费模型更灵活:同一个topic,可以被多个消费组同时消费,互不干扰,这就是发布订阅;同一个消费组内,消息又会被分摊到不同消费者手里,这就是点对点。一套系统两种模式都占了。
第二副面孔:日志系统。很多人忽略了这一点。Kafka的存储模型是“追加写日志”,每条消息都有一个offset(偏移量),消息按照顺序追加到分区里。因为只能追加不能修改,它天生就是为日志这种海量顺序写场景设计的。这也是为什么ELK这套日志方案里,Kafka能作为日志缓冲层的核心。
第三副面孔:流处理平台。基于Kafka Streams或者KSQL,你可以直接对topic里的数据做实时计算,不依赖Spark Streaming或Flink也能处理部分流式需求。
你理解了这三副面孔,面试时被问到“Kafka和RabbitMQ的区别”就不会只答一个“吞吐量高”了。Kafka的吞吐量高是因为它把“队列”做成了“日志”,这在架构设计上是两种完全不同的思路。
1.2 五个核心概念一条线串起来
Topic、Partition、Offset、Consumer Group、ISR。这五个词是Kafka的全部。
我打个比方。Topic就是一个快递集散中心,Partition就是集散中心里的一条条传送带,Offset就是传送带上的包裹编号,Consumer Group就是负责卸货的工人班组,ISR就是“目前确认还在岗且能正常干活的班组名单”。
一条消息从生产到消费,完整路径是这样的:
- 生产者指定Topic,按key哈希或轮询方式选择写入哪个Partition。
- 消息追加到Partition的尾部,获得一个递增的Offset。
- Broker上该Partition的所有副本(Replica)同步这条消息,只要ISR列表里的副本同步成功,消息就算“提交”了。
- 消费者从自己记录的消费位置(也是Offset)开始拉取消息,处理完之后提交新的Offset。
- 如果消费者挂了,同一个Group里的其他消费者会接管它的Partition,从上次提交的Offset继续消费。
这里面最重要的一个认知是:Partition是Kafka并行度的上限。一个Partition只能被一个消费组里的一个消费者线程消费。所以一个topic的Partition数决定了它的最大消费并发度。你生产端再快,消费端Partition数不够,消息一样积压。
ISR这个概念后面排查问题会反复用到。记住,ISR不是一成不变的,Broker宕机、网络抖动、消费太慢导致副本同步超时,节点都会被踢出ISR列表。ISR列表收缩,意味着你的集群容错能力在下降,这是监控里必须盯死的指标。
2. 从单机到集群:Docker部署与Win11三节点实战
现在部署Kafka,我强烈建议直接用KRaft模式,也就是不依赖ZooKeeper的模式。从Kafka 3.x开始KRaft就已经生产可用了,Apache Kafka 4.0已经彻底移除了ZooKeeper。如果你现在还在看网上那些“先装ZooKeeper再装Kafka”的教程,可以关掉了,过时了。
2.1 单机快速起一个Kafka(KRaft模式)
本地开发、验证命令,单节点完全够用。我用的是Docker Compose方式,配置如下:
services: kafka: image: bitnami/kafka:3.7 ports: - "9092:9092" environment: - KAFKA_CFG_NODE_ID=1 - KAFKA_CFG_PROCESS_ROLES=controller,broker - KAFKA_CFG_CONTROLLER_QUORUM_VOTERS=1@kafka:9093 - KAFKA_CFG_LISTENERS=PLAINTEXT://:9092,CONTROLLER://:9093 - KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092 - KAFKA_CFG_CONTROLLER_LISTENER_NAMES=CONTROLLER - KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT - KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE=true volumes: - kafka_data:/bitnami/kafka volumes: kafka_data:启动命令很简单:
docker-compose up -d跑起来之后验证一下:
docker exec -it kafka-kafka-1 kafka-topics.sh --bootstrap-server localhost:9092 --list能返回空列表就说明服务正常。
这里有几个关键点必须说明白。
为什么用Bitnami镜像?因为它对新手最友好,把Kafka原生的复杂环境变量做了统一封装,统一用KAFKA_CFG_前缀,不用去记不同版本之间配置项的变化。你要用Apache官方镜像也不是不行,但官方镜像对KRaft模式的支持是后来才完善的,配置起来更费劲,不适合速记场景。
为什么ADVERTISED_LISTENERS要写localhost?这是最常踩的坑。ADVERTISED_LISTENERS是告诉客户端“你该往哪个地址连”。如果你在容器里跑生产者,写kafka:9092没问题;但如果你在宿主机上用命令行工具或者代码连,必须写成localhost:9092,否则客户端拿到Broker返回的内网地址,根本连不上。
用bitnami/kafka:3.7时,如果你创建topic时想指定副本数,单节点只能指定为1;用kafka-topics.sh --describe能看到分区和leader信息。
2.2 Win11上硬起三个节点组成集群
Win11上部署Kafka集群,本质就是起三个容器,让它们的controller节点组成一个quorum。这里直接给配置。
项目目录结构:
kafka-cluster/ ├── docker-compose.yml └── .envdocker-compose.yml:
services: kafka1: image: bitnami/kafka:3.7 container_name: kafka1 ports: - "9092:9092" environment: - KAFKA_CFG_NODE_ID=1 - KAFKA_CFG_PROCESS_ROLES=controller,broker - KAFKA_CFG_CONTROLLER_QUORUM_VOTERS=1@kafka1:9093,2@kafka2:9093,3@kafka3:9093 - KAFKA_CFG_LISTENERS=PLAINTEXT://:9092,CONTROLLER://:9093 - KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://192.168.x.x:9092 - KAFKA_CFG_CONTROLLER_LISTENER_NAMES=CONTROLLER - KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT - KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE=true volumes: - kafka1_data:/bitnami/kafka kafka2: image: bitnami/kafka:3.7 container_name: kafka2 ports: - "9093:9092" environment: - KAFKA_CFG_NODE_ID=2 - KAFKA_CFG_PROCESS_ROLES=controller,broker - KAFKA_CFG_CONTROLLER_QUORUM_VOTERS=1@kafka1:9093,2@kafka2:9093,3@kafka3:9093 - KAFKA_CFG_LISTENERS=PLAINTEXT://:9092,CONTROLLER://:9093 - KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://192.168.x.x:9093 - KAFKA_CFG_CONTROLLER_LISTENER_NAMES=CONTROLLER - KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT - KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE=true volumes: - kafka2_data:/bitnami/kafka kafka3: image: bitnami/kafka:3.7 container_name: kafka3 ports: - "9094:9092" environment: - KAFKA_CFG_NODE_ID=3 - KAFKA_CFG_PROCESS_ROLES=controller,broker - KAFKA_CFG_CONTROLLER_QUORUM_VOTERS=1@kafka1:9093,2@kafka2:9093,3@kafka3:9093 - KAFKA_CFG_LISTENERS=PLAINTEXT://:9092,CONTROLLER://:9093 - KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://192.168.x.x:9094 - KAFKA_CFG_CONTROLLER_LISTENER_NAMES=CONTROLLER - KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT - KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE=true volumes: - kafka3_data:/bitnami/kafka volumes: kafka1_data: kafka2_data: kafka3_data:注意几点:
CONTROLLER_QUORUM_VOTERS必须三个节点完全一致,这是它们互相发现、选举的基础。- 每个容器把内部的9092映射到宿主机的不同端口:9092、9093、9094。
ADVERTISED_LISTENERS不要写localhost,要写你的宿主机局域网IP,不然从宿主机访问第二个节点时会误连回第一个节点。这也是Win11部署集群最常见的坑。你拿不准IP的话,Win11里用ipconfig查一下。
启动:
docker-compose up -d docker exec -it kafka1 kafka-metadata.sh --bootstrap-server localhost:9092 describe能看到三个broker的信息,就说明元数据层面已经没问题了。再创建一个3副本的topic验证:
docker exec -it kafka1 kafka-topics.sh --bootstrap-server localhost:9092 --create --topic test-cluster --partitions 3 --replication-factor 3 docker exec -it kafka1 kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic test-cluster如果Leader分布在不同broker上,就没问题了。可以试试docker stop kafka2,再用--describe看,你会发现原来Leader在kafka2上的分区会重新选举,ISR列表缩成两个节点。这就是集群容灾的基础。
2.3 集群部署后的第一件自检
集群起来了不要急着扔消息,先做三件事:
第一件,验证集群controller的选举。KRaft模式下controller是集群的“大脑”,元数据变更全靠它。你用一个节点,它挂了整个集群元数据操作就断了。所以至少3个节点,容忍1个节点故障。
第二件,确认自动创建topic的行为。开发环境无所谓,生产环境建议把AUTO_CREATE_TOPICS_ENABLE设成false。否则有人拼错topic名字,Kafka会默默给你建一个单分区单副本的topic,后面数据全写进去才发现不对,还要费劲迁移。
第三件,检查监听器协议。如果你的客户端需要走SSL,必须在LISTENER_SECURITY_PROTOCOL_MAP里预先配好,容器启动后再改协议,客户端连接会被拒绝,报错往往还看不出来是协议不匹配。
3. 生产消费命令的“反直觉”细节:为什么不退出、看不到数据
热搜词里有个问题非常有代表性:“Kafka生产消费命令启动一次会一直运行吗?” 我可以明确告诉你:**默认情况下,两个命令都会一直运行,这不是bug,是设计如此。**但很多人在这里愣了好几分钟,以为卡死了。
3.1 生产命令为什么启动后一直挂着
你执行:
docker exec -it kafka1 kafka-console-producer.sh --bootstrap-server localhost:9092 --topic test-cluster回车之后,命令行进入等待输入状态,你输入什么,回车,就发什么。它本质是阻塞式的,不停地从stdin读数据。所以它当然会一直运行,直到你按Ctrl+C退出。
它可以接受多行输入,但每行是一条独立的Kafka消息。如果你想要带key的消息,需要指定参数:
docker exec -it kafka1 kafka-console-producer.sh --bootstrap-server localhost:9092 --topic test-cluster --property parse.key=true --property key.separator=:这样输入格式就变成了key:value。默认分隔符是Tab,我习惯配成冒号,方便敲。
3.2 消费命令为什么不是“跑一次出结果”
很多新手执行:
docker exec -it kafka1 kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic test-cluster然后发现,它启动后看不到任何数据,就一直挂在那里。他们以为命令执行完了没输出,其实不是,而是这个topic当前没有新消息进来。
Kafka消费者默认就是从“最新位置”开始消费,它只消费启动之后产生的新消息。所以你开两个终端,一个生产一个消费,在生产端敲一条,消费端立刻打印一条,这才能体会到它为什么是“一直挂着”的。
如果你想消费这个topic里的历史数据,必须加--from-beginning:
docker exec -it kafka1 kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic test-cluster --from-beginning但这里又引入一个新问题:如果你从头消费,消息非常多,它就会一直刷屏,没有自动停下来的机制。这时候你希望“看一定数量的消息后就退出”,就得两个参数配合:
docker exec -it kafka1 kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic test-cluster --from-beginning --max-messages 10这样它消费完10条就自动退出了,这是查看topic里有没有数据最常用的命令。另外一个参数--timeout-ms 5000也可以让它空闲5秒后自动退出,适合确认某个topic在某个时间段内有没有消息进来。
3.3 查看topic里到底有什么的三种姿势
命令行看数据,有三种姿势,适用场景完全不同,我都列出来。
姿势一:控制台消费,看业务数据长什么样。就是上面说的,用--max-messages限制数量。缺点是你看到的是序列化后的字节,如果消息用了Avro或Protobuf序列化,不装对应的Serializer是看不了真值的。
姿势二:查看堆积情况,用kafka-get-offsets.sh。
docker exec -it kafka1 kafka-get-offsets.sh --bootstrap-server localhost:9092 --topic test-cluster输出类似:
test-cluster:0:0 test-cluster:1:15 test-cluster:2:7第三列是每个分区当前最新的offset(也就是下一条消息的写入位置)。注意,它不代表“生产了多少条”,只代表“写到了哪个位置”。结合消费组的提交offset,就能算出积压量。
姿势三:直接看日志段文件。这个适合排查消息有没有写进去、有没有损坏。Kafka的数据文件是日志段(log segment),用自带的DumpLogSegments工具看:
docker exec -it kafka1 kafka-run-class.sh kafka.tools.DumpLogSegments --files /bitnami/kafka/data/test-cluster-0/00000000000000000000.log --print-data-log这条命令能看到每个日志项的offset、时间戳、key和value的字节内容。生产上排查“消息到底写没写进去”“写入的value是不是空的”非常管用。
3.4 消费组偏移量:最容易翻车的环节
消费者最核心的机制就是“提交Offset”。命令行的消费者工具默认会生成一个随机的消费组ID,你Ctrl+C退出后,组信息可能已经提交了;下次再跑,就是一个新的随机组,从头或者从最新开始取决于参数。
如果你想看某个消费组的消费位置和积压量,用:
docker exec -it kafka1 kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group my-group输出里有CURRENT-OFFSET(当前消费到哪)、LOG-END-OFFSET(最新消息位置)、LAG(积压量)三列。LAG不为0,说明消费速度跟不上生产速度,这是排查延迟的第一手信息。
生产上经常要“重放消息”,也就是把消费组的位置重置到过去某个时间点或offset。命令是:
docker exec -it kafka1 kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group my-group --reset-offsets --to-earliest --topic test-cluster --execute--to-earliest是重置到最开头,--to-latest是跳到最新,--to-datetime 2024-01-01T00:00:00.000可以重置到指定时间点,--shift-by -100可以回退100条。注意,重置offset必须要消费组处于非活跃状态,否则命令会报错“Assignments can only be reset if the group is inactive”。所以流程是:先停掉应用,再重置,再启动应用。
这一步是运维Kafka的人的日常操作,一定要练熟。
4. 两个高频事故排查:消息延迟高和OOM
热搜词里“kafka消息延迟高”和“kafka oom”排在前面,这是线上最痛的两个问题。我直接给出排查思路。
4.1 消息延迟高:先分端再下结论
消息从生产到消费,链路是:生产端 → Broker → 消费端。任何一个环节慢,最终表现的指标都是消费端“消息延迟高”或者“LAG持续上涨”。所以排查第一步,不是猜,是分端观察。
生产端排查。看客户端的发送耗时。如果发送耗时高,先看是不是acks=all且min.insync.replicas配置得不合理。acks=all要求所有ISR里的副本都确认写入才算成功,如果ISR只有一个副本,你还开着acks=all,性能不会差太多;但如果ISR里有3个副本,其中有一个磁盘性能很差,整个生产链路的耗时就会被这个慢副本拖住。其次是看batch.size和linger.ms。Kafka生产者默认是攒一批消息再发,如果你的batch设得太小,或者linger.ms=0,每条消息都单独发,网络往返次数剧增,吞吐自然上不去。
props.put("batch.size", 16384); props.put("linger.ms", 5);这个配置的意思是:攒够16KB或最多等5毫秒就发送。这是最简单的调优参数。
Broker端排查。用iostat -x 1看磁盘。Kafka号称“顺序写所以快”,这是有前提的——磁盘不能有其他随机IO干扰。如果你把Kafka的数据目录和操作系统、数据库放在同一块磁盘上,随机IO会把顺序写的优势完全抹掉。另外,Page Cache回收(Page Cache Reclaim)也会导致写入延迟毛刺,这种情况看vmstat的wa(IO Wait)会飙升。
再一个容易忽略的点:副本同步延迟。Broker之间同步副本走的是网络,如果机器之间的网络有抖动或者带宽被打满,ISR里的副本跟不上,leader会等待或剔除副本,对生产端的表现就是“发送超时”。这时候用kafka-topics.sh --describe --topic xxx --under-replicated-partitions看一下有没有分区处于under-replicated状态。
消费端排查。消费端延迟高,最常见的原因是消费逻辑里有耗时操作,比如每消费一条就去调用一次外部API。这时候max.poll.interval.ms默认5分钟,如果一条消息处理超过5分钟,消费者会被踢出消费组,触发rebalance,rebalance期间整个组停止消费,LAG瞬间暴涨。
我踩过的一个真实场景:某业务消费消息后调用一个第三方接口,接口偶尔延迟10秒,平时没事;某天大促,接口整体变慢,消费者处理一条消息超过5分钟,被踢出消费组,不停rebalance,LAG从几百涨到几百万。表面上看起来是Kafka的问题,实际是下游接口的问题。
所以排查链路应该是:
- 先看消费组LAG:
kafka-consumer-groups.sh --describe,确认是某个分区还是所有分区都积压。 - 看消费者日志:有没有rebalance日志,有没有处理超时异常。
- 看Broker指标:CPU、磁盘、网络、under-replicated分区数。
- 按上面三个方向定位,不要一上来就加分区。
4.2 Kafka OOM:堆内还是堆外
Kafka本身是JVM应用,但它和普通Java应用有个本质区别:Kafka大量使用Page Cache(页缓存)来读写数据,所以它的堆内存不需要设置得特别大,反而是堆外内存和操作系统缓存对性能的影响更大。
网上搜索“kafka oom”,很大一部分案例其实是这几种情况:
第一种,JVM堆内存溢出。Kafka的堆主要用来存放分区、副本、客户端的元数据信息。如果你topic数量非常多(比如几千个分区),堆内存会涨得很快。这种OOM的解决办法不是盲目加大-Xmx,而是先看堆里到底什么对象占据内存——用jmap -dump导出堆转储,用MAT分析。我见过一个案例,某团队为了性能把-Xmx设成16G,结果GC一次停顿好几秒,消息延迟抖动严重,这就是堆太大导致的,最后调回6G反而稳了。
Kafka的内存分配,我建议直接通过环境变量控制:
environment: - KAFKA_HEAP_OPTS=-Xmx4g -Xms4g -XX:+UseG1GC第二种,Java堆外内存溢出。消费者或生产者的receive.buffer.bytes和send.buffer.bytes设置过大,堆积大量缓冲时,进程的常驻内存(RSS)会超过容器限制,被OOM Killer干掉。这种情况表现是:进程直接消失,日志里没有Java堆栈,docker inspect可以看到OOMKilled=true。
第三种,最容易忽略的:宿主机的“可用内存”假象。Kafka会尽量利用Page Cache,文件没有写入磁盘而是留在缓存里,这会让free命令看到的内存所剩无几。很多人误以为Kafka内存泄漏了,其实不是——那些是缓存,可以随时回收。判断Kafka内存是否真的有问题,要看JVM的堆内存指标,而不是看free。
所以我的建议是:Kafka的容器内存限制不要卡得太死。比如-Xmx4g的JVM,容器内存至少要给6g-8g,要给Page Cache留空间。如果限制太紧,操作系统会在内存压力下频繁回收缓存,磁盘IO飙升,对Kafka的性能影响比OOM本身还大。
5. 接入可观测性:ELK与OTel的常见接法
提到Kafka,绕不开“ELK”这个词。很多公司的日志链路是:应用 → Kafka → Logstash → Elasticsearch → Kibana。另一个热词是“otel kafka”,也就是OpenTelemetry对Kafka的监控和链路追踪支持。我把这两块都过一遍。
5.1 日志接入:Filebeat加Kafka的经典组合
为什么日志链路中间要加一层Kafka?两个原因:削峰填谷和消费解耦。
业务高峰期,日志量可能是平均值的十倍。如果应用直接往Elasticsearch写,ES分分钟被打爆,集群直接雪崩。用Kafka缓冲之后,Logstash侧的消费速度就是固定的“匀速排水”,ES集群压力可控。
这个链路的标准配置是:
应用日志 -> Filebeat -> Kafka Topic(如 app-log) -> Logstash -> ElasticsearchFilebeat这边只需要配置Kafka output:
output.kafka: hosts: ["192.168.x.x:9092", "192.168.x.x:9093", "192.168.x.x:9094"] topic: "app-log" partition.hash: hash: ["fields.app"] # 用Kafka的key来保证同一应用实例的日志有序 key: "my-app"注意,这个topic建议在创建的时候就规划好分区数。日志场景分区数不要太多,3个到6个就够用了,太多反而增加Broker的元数据负担。副本数至少2个,保证一台机器挂了日志不丢。
Logstash这边用Kafka input:
input { kafka { bootstrap_servers => "192.168.x.x:9092,192.168.x.x:9093,192.168.x.x:9094" topics => ["app-log"] group_id => "logstash-es" auto_offset_reset => "latest" codec => "json" } } output { elasticsearch { hosts => ["http://192.168.x.x:9200"] index => "app-log-%{+YYYY.MM.dd}" } }这里group_id很重要。如果你起了多个Logstash实例并配置同一个group_id,它们在Kafka里属于同一个消费组,会分摊topic分区的消费任务,这就是日志消费水平的横向扩容方法。
5.2 OTel Collector与Kafka指标采集
OpenTelemetry对Kafka的接入,分两层:一是Kafka作为OTel数据管道的一部分(OTel Collector用Kafka做receiver/exporter),二是采集Kafka自身的运行指标。
先说Kafka自身的指标采集。Kafka暴露JMX指标,最常用的采集方案是JMX Exporter加Prometheus,但JMX Exporter的配置多少有点繁琐。用OTel Collector的jmx receiver可以一步到位:
receivers: jmx: jar_path: /opt/opentelemetry-java-contrib-jmx-metrics.jar endpoint: localhost:9099 target_system: kafka-jvm collection_interval: 10starget_system: kafka-jvm会自动抓取Kafka的关键指标,包括内存、GC、请求处理时间等,不需要手写一大堆MBean规则。
再说OTel Collector用Kafka传输span和metric。如果你的集群已经有一套Kafka,完全可以让OTel Collector把采集到的链路数据或指标先发到Kafka,再由下游的消费者决定如何导出。这是因为Kafka的缓冲能力比OTel Collector网络导出的方式更稳定,大规模采集时不容易丢数据。
exporter配置:
exporters: kafka: broker: - 192.168.x.x:9092 - 192.168.x.x:9093 - 192.168.x.x:9094 protocol_version: 2.0.0 topic: otlp-spans-topicreceiver配置(下游消费):
receivers: kafka: brokers: - 192.168.x.x:9092 topic: otlp-spans-topic这套组合适合那种“OTel Collector实例特别多,但不想每个实例都直连后端”的场景。日志、指标、Trace都进Kafka,统一做缓冲和分发。
5.3 运维监控要盯哪几个数
监控Kafka,不是把Grafana上所有面板都看一遍,而是盯着那几个“出事之前一定有前兆”的指标。
集群级,必须盯的:
| 指标 | 含义 | 告警阈值建议 |
|---|---|---|
| UnderReplicatedPartitions | 分区副本数低于配置值 | 持续0,一旦>0立即告警 |
| OfflinePartitionsCount | 分区leader下线 | 必须为0 |
| ActiveControllerCount | 活跃controller数量 | 正常情况为1 |
| ISR收缩/扩展次数 | ISR列表频繁变动 | 趋势性上升就要查 |
Broker级:CPU使用率、GC暂停时间(Kafka对GC很敏感,Full GC一次对应一次生产延迟尖峰)、磁盘读写延迟、网络吞吐。这些用普通的Node Exporter加JMX Exporter就能采集。
Topic级:每秒消息数、每秒字节数、分区Leader分布是否均匀。我在实际运维中养成一个习惯:每次发版前和发版后都看一眼这几个指标,很多时候消息量突增或突降都跟代码变更直接相关。
特别提醒一下,Kafka的JMX指标非常丰富,但新手容易在浩瀚的指标里迷失,看完一堆图表还是不知道集群健不健康。我的建议是:先只看上面表格里的几个集群级指标,其他指标等出了问题再对照查,效率反而更高。
6. 面试速记:Kafka高频问题怎么答
热搜词里“kafka面试题”也是一个独立的热门方向。我把最常被问的几个问题,整理成可以直接说出口的“口语化答案”。不求面面俱到,但求每道题都能答出核心逻辑。
6.1 为什么Kafka这么快
这道题5分钟讲不完,但核心就四点。
第一,顺序写磁盘。Kafka不修改已写入的数据,只追加写。顺序写磁盘的速度接近内存随机读的速度,大约每秒几百MB到1GB。你用普通的机械硬盘也能获得很高的吞吐,前提是分区内的写入是顺序的。这是Kafka高性能的根基。
第二,Page Cache(页缓存)。Kafka读写数据主要走操作系统的Page Cache,而不是JVM堆内存。这意味着数据写入时会先写到Page Cache就被视为成功,操作系统延时刷盘。同时,消费者读数据时命中Cache的概率很高,读不出盘。
第三,零拷贝(Zero Copy)。消费者从Kafka拉数据时,数据从磁盘到Page Cache,再到网卡,全程不经过JVM堆和应用缓冲区,省去了用户态和内核态之间的拷贝。这就是sendfile系统调用的效果。用kafka-run-class.sh看消费者服务端日志时,能看到这个机制在起作用。
第四,批量与压缩。生产者在内存里攒一批消息再发,减少了网络往返次数;Broker存储时用LZ4或ZSTD压缩,减少了磁盘IO和网络带宽占用。
面试官如果追问“那Topic越多吞吐越慢吗”,你要能接住:是的。每个Partition都有对应的文件句柄和内存映射,几百上千个Partition会让随机IO变多,顺序写的优势被削弱,所以Kafka集群的Partition总数不是越多越好。
6.2 消息不丢怎么保证
这是个老生常谈的题,但很多人答不全。我建议按三端来说。
生产端不丢:设置acks=all,表示ISR里所有副本都写入成功才返回成功;开启重试retries(记得配retry.backoff.ms),并开启幂等enable.idempotence=true防止重试时消息重复。
Broker端不丢:Topic副本数(replication.factor)至少3,最小同步副本数(min.insync.replicas)至少2,这样即使一个Broker挂了,还有两个副本有数据,消息不会丢。同时,leader选举时,Kafka只会从ISR里的副本选新leader,那些比leader落后太多的副本没资格成为leader。
消费端不丢:核心是“先处理后提交”。默认配置是enable.auto.commit=true,每5秒自动提交一次Offset,如果消费者在自动提交前崩溃,重启后会重复消费,不会丢,但如果消费者在处理消息前就提交了Offset,然后处理过程中挂了,那这批消息就丢了。所以生产环境建议改成:
props.put("enable.auto.commit", "false");处理完消息后再手动提交:
consumer.commitSync();用一句话总结就是:生产端把ack开满,Broker端把副本铺厚,消费端先干活再记账。
6.3 Rebalance什么时候发生
消费组里任何一个消费者退出或加入,都会触发Rebalance。Rebalance就是所有消费者重新分配分区,这个过程中整个消费组是暂停消费的,所以Rebalance越多,消费效率越低。
触发条件主要有:
- 消费者主动或被动退出(崩溃、网络超时、处理超时被踢)。
- 消费者数量变化(新增或删减)。
- Topic的分区数变化(扩容分区)。
- 订阅关系变化(代码里改订阅的topic)。
要避免频繁Rebalance,核心是调好两个参数:
max.poll.interval.ms:消费者两次poll之间的最大时间间隔。如果处理一条消息耗时逼近这个值,会被判定为“卡死”,被踢出消费组。建议根据消息处理最慢耗时,设置一个安全余量。session.timeout.ms:消费者和Broker之间的会话超时时间。心跳超时也会被踢出组。但这个值和心跳间隔heartbeat.interval.ms要成比例设置,通常心跳间隔是超时时间的三分之一左右。
如果发生Rebalance,Kafka提供ConsumerRebalanceListener机制,你可以在onPartitionsRevoked里做“提交offset并清理资源”,在onPartitionsAssigned里做“初始化连接”,尽量减少Rebalance造成的数据重复和消费中断。
我个人在实际排查和面试中的体会是,Rebalance是一个“表面现象”,它的背后几乎总是隐藏着“消费者处理太慢”或“网络不稳定”这几个真正的问题。所以面试时回答Rebalance,一定别忘了说最后一句:一旦发现集群频繁Rebalance,优先排查消费者线程是否卡死、GC是否长时间停顿、网络是否存在抖动,而不是急着调参。这样才能体现你有实战经验,而不是背得熟。
最后再分享一个工作习惯:我每接手一套Kafka集群,都会在本地备一个速记文档,记录这套集群的部署参数、关键Topic列表、常见异常排查步骤。排障时先看三件事——消费组LAG、UnderReplicatedPartitions、GC日志。这三个数正常,集群大概率没事;这三个数异常,基本就能顺藤摸瓜找到问题根源。这套速记方法,建议你也试试。