ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Kafka入门完全指南:从消息队列原理到集群部署与故障排查

Kafka入门完全指南:从消息队列原理到集群部署与故障排查 如果让我用一句话形容第一次接触Kafka的感受那就是每个术语都认识连起来不知道在说什么。topic、partition、offset、consumer group、broker、rebalance……我当初也是开了十几个标签页越看越晕。后来真正动手装了集群、写了生产者和消费者、被线上消息队列的重复消费和消息延迟折磨过几次之后才意识到Kafka入门这件事最缺的不是文档而是一条清晰的路线它是用来解决什么问题的、内部长什么样、怎么装起来、怎么用起来、出错怎么排查。这篇文章就是按这条路线写的主题就是一个完整的消息队列入门目标是让从没碰过Kafka的人用一下午把它跑起来并且能理解自己在干什么——不是背命令而是真的掌握。1. 为什么你的系统迟早会需要消息队列1.1 没有消息队列时两个服务是连坐的先说一个很常见的场景。订单服务创建订单之后要调用用户服务发短信、要调用积分服务加积分、还要调风控服务做校验。最早大家图省事直接在订单服务里同步调用这三个接口。结果是什么呢任何一个下游接口变慢订单接口就跟着慢任何一个下游服务重启订单服务就报错如果赶上活动大促瞬时流量打进来数据库先扛不住紧接着一堆服务跟着雪崩。这是典型的耦合问题。同步调用把上游的可用性、性能跟下游绑死在了一起三个下游只要有一个出问题整条链路就遭殃。消息队列在这里干的事就是解耦和削峰订单服务只管把“订单创建成功”这个事件写进队列用户服务、积分服务、风控服务各自去消费自己关心的消息互不等待互不拖累。流量高峰时队列先把消息存下来消费者根据自己的处理能力慢慢消化而不是把所有压力瞬间压到数据库头上。所以如果你在架构评审的时候听到“我们要做异步化”“要做削峰填谷”“服务之间要解耦”背后基本都是同一个需求引入一个中间件让消息的“生产者”和“消费者”不用同时在线、不用互相等待。这个中间件最常见的选择就是消息队列。Kafka、RabbitMQ、RocketMQ甚至Windows生态里的MSMQ都属于这一类产品。1.2 Kafka不是一个“普通消息队列”它是提交日志很多人第一反应是把Kafka和RabbitMQ当成同一类东西这能理解毕竟都叫消息队列。但Kafka的定位其实更特殊它本质上是一个分布式提交日志distributed commit log。什么意思就是消息进来之后不是消费完就删掉而是按照顺序一直追加写在磁盘上消费者用一个叫offset的东西记住“我读到哪里了”。你可以随时从头重新读一遍历史消息就像翻日志一样。这个设计带来两个明显好处一是吞吐量极高因为日志追加写是顺序IO比随机写入快好几个数量级二是不丢消息且可重放消费者挂了恢复之后从之前记录的offset继续读就行。相比之下RabbitMQ更偏向传统的消息代理消息被消费确认后就会删除强调灵活的 routing、延迟队列、死信队列这类企业级功能。两者没有谁绝对好但如果你要的是大数据管道、日志采集、高吞吐流处理Kafka几乎是首选。而如果只是几个微服务之间传一些任务消息、需要复杂的路由规则RabbitMQ往往更合适。下面是个简化的对比方便刚入门的人建立直觉特性KafkaRabbitMQ核心模型分布式提交日志消息代理消息是否删除按保留策略默认保留7天消费确认后删除吞吐能力非常高百万级/秒相对较低万级/秒消息顺序分区内有序队列内有序典型场景日志管道、流处理、削峰业务消息、任务分发、延迟消息1.3 哪些场景其实不需要上KafkaKafka虽好但真不是所有消息场景都该上。我见过不少团队把系统搞复杂就是因为明明只有几个业务队列硬要上Kafka集群结果运维成本上来了功能还别扭。Kafka在延迟队列、定时消息、按条件路由这几种能力上很弱这些恰恰是RabbitMQ的强项。比如“订单支付成功后30分钟关单”这种延迟任务用Kafka实现要自己拼轮询或者额外搭存储非常不划算。还有一种常见误用拿Kafka当RPC用。A服务发一条消息等着B服务处理完再返回结果。Kafka本身就是异步消息模型做这种同步请求应答非常别扭远不如直接HTTP或gRPC。我的建议是只有当你明确需要高吞吐、持久化、消息可重放、多消费者组独立消费这几种能力时才优先考虑Kafka。否则一个简单的内存队列、或者RabbitMQ可能反而更合适。这也是这篇文章的价值——先搞清楚“为什么选它”再谈“怎么用它”。2. 把Kafka的核心概念揉碎了讲2.1 Topic、Partition、Offset仓库、货架、位置编号Kafka里最基础的概念是Topic。你可以把Topic理解成一个仓库每个仓库存放同一类事件比如“订单事件”“用户行为日志”。但一个Topic在物理存储上不是一整块而是被拆成若干个Partition分区每个分区就是仓库里的一个货架消息按顺序放在货架的一层层格子里。每当一条新消息进来Kafka把它追加到某个分区的末尾对应一个单调递增的编号这个编号就是Offset。为什么要拆分区核心目的是并发。单条消息只能写在一个文件里的一个位置上如果所有消息都写一个分区写并发就被一个文件锁死了。拆成多个分区之后多个分区可以并行写入、并行读取吞吐量才能上去。代价是跨分区的消息不再有全局顺序只有“同一个分区内的消息是有序的”。这是一个非常重要的取舍后面所有关于顺序的讨论都建立在这条规则上。分区的数据也不是只有一份。为了不丢消息每个分区还有若干副本分布在不同的BrokerKafka节点上其中一个是Leader负责接收读写请求其余是Follower只负责同步Leader的数据。当Leader挂了Kafka会从副本里选出新的Leader继续服务。这也是Kafka高可用的基础。2.2 Consumer Group一家人怎么分工还是几家各干各的消费者的概念比生产者稍微绕一点。假设你有一个Topic叫“order-events”积分服务和用户服务都要消费它但它们关心的事情不一样处理逻辑也不同。Kafka用Consumer Group来解决这个问题每个服务一个消费组组与组之间消费同一批消息互不影响各记各的进度offset相当于几家公司各自派货车来同一个仓库拉货。在同一个消费组内部规则是同一个分区同一时刻只能被组内的一个消费者实例消费。假如某个Topic有4个分区一个消费组里有2个消费者那么每个消费者各分2个分区如果有4个消费者就一人分一个如果加到了6个消费者那就会有2个消费者闲置。这个分配动作发生在消费者加入、退出、或者Topic分区数变化的时候也就是常说的Rebalance再平衡。理解这个模型特别重要。很多人问“为什么我加消费者消费速度没提升”大概率就是分区数不够消费者加再多也只能闲着。同样的同一个组里如果两个消费者跑在不同的机器上、配置完全一样它们消费的是一半的消息而不是重复消费所有消息——这也是新手经常误解的一点。2.3 ACK与Offset提交所有丢消息和重复消息的源头消息从生产到消费要经过两次“确认”ACK这两次确认决定了可靠性的边界。第一次是生产者确认。生产者发消息给Broker时可以配置acks参数acks0表示发了就不管不管Broker收没收到acks1表示Leader分区写成功就算成功acksall表示所有同步副本都写成功才算成功。从“不丢消息”的角度生产端应该用acksall同时配合重试机制。只不过可靠性越高延迟也会越高这本身就是权衡。第二次是消费端的Offset提交。消费者拉取消息、处理完业务之后需要告诉Kafka“这批消息我处理完了下次从下一个位置继续”这个动作叫提交Offset。如果启用了自动提交Kafka会每隔几秒自动提交当前的消费位置。问题是如果业务逻辑还没处理完但自动提交时间到了此时进程突然挂了恢复后就会从已提交的offset继续消费——那么崩溃前没处理完的那批消息就丢了。反过来如果业务处理完了但offset还没提交进程就挂了恢复后就会再次消费到这批消息——重复消费。所以你发现没有丢消息和重复消费这两件事根源往往都在“处理”和“提交”之间没有一个原子保证。这不是Kafka独有的问题分布式系统里几乎都存在。理解了这一点后面看故障排查就会非常顺畅。3. 从下载到跑通单机和集群安装实录3.1 先选模式KRaft还是ZooKeeper老教程里装Kafka基本都要先装一套ZooKeeper因为Kafka早期的元数据管理、Broker选举、消费者组管理全都依赖ZooKeeper。有两个系统要维护部署和排查都很麻烦。从Kafka 2.8开始官方引入KRaft模式用自己的Raft协议来做元数据管理到3.3版本已经生产可用目前主流的3.x版本搭建新集群我建议直接用KRaft模式不需要再装ZooKeeper了。本篇以Kafka 3.7.0为例可以直接到Apache官网下载二进制包wget https://downloads.apache.org/kafka/3.7.0/kafka_2.13-3.7.0.tgz tar -xzf kafka_2.13-3.7.0.tgz cd kafka_2.13-3.7.0下载解压之后bin目录下有各组件脚本config里是各种配置文件。建议先不要动一堆参数默认配置足够跑通单机。3.2 Linux单机部署的完整命令KRaft模式单机部署只需要三步生成集群ID、格式化存储目录、启动服务。KAFKA_CLUSTER_ID$(bin/kafka-storage.sh random-uuid) bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c config/kraft/server.properties bin/kafka-server-start.sh config/kraft/server.properties第一条命令生成一个随机UUID作为这个Kafka集群的唯一标识第二条命令用它格式化存储目录相当于初始化数据目录第三条就是启动Broker服务。启动后看到“Kafka Server started”日志就说明成功了。有一点必须提醒kafka-storage.sh format这个命令会清空目录里的数据生产环境千万不要乱执行。格式化的实际作用是生成元数据、创建必要的内部Topic只在初始化时做一次。如果你误执行了格式化等于是把那台机器上的Kafka数据全部清空。如果只是想平时开发测试单节点完全够用。接下来我会讲Windows的坑再讲真正的集群怎么搭。3.3 Windows安装的三个坑Windows下安装Kafka其实不复杂bin/windows目录下提供了对应的bat脚本命令和Linux基本一样bin\windows\kafka-storage.bat random-uuid bin\windows\kafka-storage.bat format -t uuid -c config\kraft\server.properties bin\windows\kafka-server-start.bat config\kraft\server.properties但Windows上我踩过几个坑写出来你就能少走弯路。第一个是路径别带空格Kafka脚本对路径里的空格处理得不好最好解压到C:\kafka这样的短路径下我见过有人放在“C:\Program Files\Apache\kafka”下面启动时各种诡异报错。第二个是必须显式配置advertised.listenersWindows主机名经常带有特殊字符或者和局域网机器名不一致不配置的话Broker注册到元数据里的地址可能不对客户端根本连不上。建议在config/kraft/server.properties里写死listenersPLAINTEXT://0.0.0.0:9092 advertised.listenersPLAINTEXT://localhost:9092第三个坑是不要直接把Windows机器当成生产集群节点。Kafka依赖页缓存和顺序IOWindows下的文件句柄、网络栈调优和Linux不是一个量级的日常开发测试没问题但生产集群还是老老实实用Linux。3.4 扩展成3节点集群的配置要点单机跑通之后很多人会问集群怎么搞。其实和单机差别不大核心是每个Broker要有唯一的node.idcontroller的地址要互相配置一致advertised.listeners要配置成其他节点能访问到的地址。以3节点KRaft混合模式一个节点同时跑Broker和Controller为例每台机器的server.properties核心配置如下process.rolesbroker,controller node.id1 controller.quorum.voters1192.168.1.10:9093,2192.168.1.11:9093,3192.168.1.12:9093 listenersPLAINTEXT://192.168.1.10:9092,CONTROLLER://192.168.1.10:9093 advertised.listenersPLAINTEXT://192.168.1.10:9092 controller.listener.namesCONTROLLER log.dirs/data/kafka第二台机器node.id改成2listeners和advertised.listeners改成自己的IP第三台改成3。然后就像单机一样每台机器初始化一次格式化和启动。三台都起来之后集群就组好了。注意controller.quorum.voters里那三个地址必须是所有节点都能访问的否则Controller之间无法建立Raft选举。给生产集群的建议是Controller节点至少3个奇数是为了Raft投票。Broker数量可以根据数据量和吞吐需求横向扩容。默认情况下分区副本数replication-factor建议设置成至少2这样单个节点挂掉不会丢数据。3.5 部署完第一件事命令行验证服务起没起来最直接的验证方式是创建Topic并生产消费一条消息。先创建一个叫test-topic的Topic3个分区、1个副本bin/kafka-topics.sh --create --topic test-topic --partitions 3 --replication-factor 1 --bootstrap-server localhost:9092再查看Topic状态bin/kafka-topics.sh --describe --topic test-topic --bootstrap-server localhost:9092输出里能看到Partition的列表、每个分区的Leader和副本分布。然后打开一个终端运行生产者随便输入几行文字bin/kafka-console-producer.sh --topic test-topic --bootstrap-server localhost:9092另开一个终端运行消费者加上--from-beginning意思是从最早的消息开始读bin/kafka-console-consumer.sh --topic test-topic --from-beginning --bootstrap-server localhost:9092生产者那边输入“hello kafka”消费者这边马上就能看到。到这里环境就算真正跑通了。这个“从零到能收发”的阶段实际只需要十分钟比看文档快得多。4. 用命令行和Java客户端发第一条消息4.1 命令行工具走一遍生产消费全过程上面已经用过kafka-console-producer和console-consumer了这一节把它们的功能再挖深一点。命令行工具最大的价值不是日常使用而是快速验证环境、排查问题。比如你怀疑网络不通用console工具连一下就能确认怀疑Topic配置不对describe一眼就能看到分区和副本情况。开发阶段最常用的几个命令我再补充一个bin/kafka-consumer-groups.sh --describe --group test-group --bootstrap-server localhost:9092这条命令用来查看消费组test-group在每个分区上的消费进度里面有一列叫LAG表示“还剩多少条消息没有消费”。这个指标是排查消息积压的核心后面第5章还会反复提到。注意消费端有个隐藏细节如果消费组是第一次消费一个Topic而且之前没提交过offset那么earliest和latest这两个策略决定从哪开始读。earliest表示从最早的消息读起latest表示从最新的消息开始读。做数据补录通常用earliest做实时消费业务通常用latest。我见过有人因为默认配置是latest重启消费服务后“丢”了旧消息其实是没理解这个语义。4.2 Java客户端最小可运行代码命令行工具能跑通就说明集群没问题接下来该用客户端代码了。Java客户端是最主流的方式引入依赖dependency groupIdorg.apache.kafka/groupId artifactIdkafka-clients/artifactId version3.7.0/version /dependency生产者端的关键代码非常短Properties props new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, localhost:9092); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.ACKS_CONFIG, all); KafkaProducerString, String producer new KafkaProducer(props); producer.send(new ProducerRecord(test-topic, user-1001, {\orderId\:1001})); producer.close();消费者端稍微复杂一点核心是poll循环Properties props new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, localhost:9092); props.put(ConsumerConfig.GROUP_ID_CONFIG, order-consumer-group); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, earliest); KafkaConsumerString, String consumer new KafkaConsumer(props); consumer.subscribe(List.of(test-topic)); while (true) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecordString, String record : records) { // 在这里处理业务逻辑 System.out.printf(offset%d, key%s, value%s%n, record.offset(), record.key(), record.value()); } consumer.commitSync(); }这里我把ENABLE_AUTO_COMMIT_CONFIG设成false然后手动调用commitSync提交offset。如果你刚开始学建议就用这套“处理完再提交”的模板它能最大程度避免业务逻辑处理了一半就自动提交导致的丢消息问题。补充一句客户端不只有JavaC/C有librdkafkaPython有confluent-kafka。我甚至见过有人用Qt写桌面工具连Kafka本质上也是包装librdkafka。入门阶段用Java或者命令行建立心智模型就足够了语言真的不是关键。4.3 几个必须理解的Producer参数下面这些参数面试爱问线上排障也一定会碰到值得花几分钟彻底搞清楚。参数作用我的建议acks0/1/all控制写入成功的确认边界业务消息用all日志类可放宽retries发送失败重试次数大于0配合重试避免网络抖动丢消息batch.size生产者攒一批消息再发的字节数保持默认先测效果再调linger.ms等多久再发这一批攒批的时间窗口默认0吞吐优先可调大max.request.size单个请求最大字节数默认约1MB发送大消息时需调整enable.idempotence生产者幂等防止重试导致重复消息新版默认开启保持开启很多性能排查都跟这几个参数有关。比如吞吐上不去可以先看linger.ms是不是太小消息每次来一条就发一条每次都产生一次网络往返。把linger.ms调到5到10毫秒同时batch.size调大让Kafka把一批消息批量发送吞吐往往能明显改善。4.4 分区与Key顺序性和并发度的来源生产者发送消息时可以指定Key。Kafka对Key做哈希决定这条消息进哪个分区。于是有了一个非常实用的特性相同Key的消息一定进同一个分区也一定被同一个分区内的顺序消费。比如用户操作日志按用户ID做Key那么同一个用户的操作事件就是有序的这对状态恢复类业务极其重要。但Key的选择要当心。如果某个Key的消息量特别大它对应的分区就会变成热点分区其他分区却很空闲。最典型的例子是按城市ID做Key一线城市的数据远远多于小城市。这时候的解法是要么换一个更离散的Key比如用户ID要么引入一个随机前缀在需要分区内顺序的业务里前缀不要影响关键标识。分区数的设置同理分区太少并发上不去分区太多每个分区带来的文件句柄、内存、Rebalance开销都会增大。通用经验是分区数至少大于单台消费者的并发上限但不要盲目设到上千。5. 在线系统最常见的四个坑以及排查思路5.1 重复消费到底是谁重复了重复消费几乎是每个Kafka使用者都会遇到的问题。背后的原因链条其实不复杂前面2.3已经讲了基本原理消费端处理完数据但offset还没提交进程挂了或者Rebalance发生了重启后就从上次提交的位置重新消费。生产者的重试也会导致Broker收到重复消息即使开启了幂等也只能保证同一个Producer会话内的分区级去重跨会话还是可能重复。所以解决方案不能只在Kafka层面找必须在下游业务里做兜底。最通用也最可靠的办法是消费幂等在消费逻辑里把消息中的唯一业务主键订单号、用户ID写进一张有唯一索引的去重表只有插入成功才继续业务处理、再提交offset。重复消息来的时候唯一索引冲突直接忽略即可。我之前在订单消息里就是这么做消费之前先select一眼没有就insert ignore再写业务数据。这套方案弱化了Kafka自身的精确一次语义把可靠性下沉到业务侧实现成本很低但效果非常稳定。如果你硬要追求Kafka的精确一次exactly-once可以研究事务API和读已提交隔离级别但先问自己业务侧的幂等表真的实现不了吗5.2 消息延迟高从LAG表开始定位“Kafka消息延迟高”是热搜里的高频词遇到这个问题我的第一个动作永远是查消费组的LAG。用前面提到的kafka-consumer-groups.sh --describe看哪个分区LAG在增长增长得厉害说明消费速度跟不上生产速度这是消费侧问题LAG一直为0但整体链路还是慢说明延迟可能出在生产或者网络环节。消费侧最常见的原因有三个。第一单条消息处理太慢比如每条消息都查一次数据库、调用一次外部接口整个消费吞吐被压到每秒几十条。第二处理超时导致消费者被踢出消费组触发频繁Rebalance。Kafka消费者有个max.poll.interval.ms参数默认5分钟你一次poll拿到的消息如果处理了超过5分钟服务端会认为消费者已死亡把它踢出组并触发重平衡这个过程中谁都消费不了这个分区延迟会瞬间拉高。解法是把max.poll.records调小让一批消息控制在几秒内处理完或者调大max.poll.interval.ms但最根本的还是要优化处理速度。第三分区数小于消费者数有一部分消费者空转白白浪费资源。定位这些问题的顺序我建议是先看LAG分布再查消费端日志里有没有频繁Rebalance然后用监控看poll和单条处理耗时。别一上来就调一堆Kafka参数很多延迟问题根源不在Kafka而在你的业务代码。5.3 1MB大消息Kafka不是对象存储“kafka接收1m”这个热搜词本质上是很多人第一次往Kafka里发大消息时撞上了默认限制。Kafka的单条消息默认上限是1MB左右message.max.bytes默认1048588字节所以当你试图发送一条比这更大的消息时Broker会直接报RecordTooLargeException。要突破这个限制需要同时调整多方参数位置参数说明Brokermessage.max.bytes消息的最大字节数Topicmax.message.bytesTopic维度覆盖Broker默认值Producermax.request.size单次请求最大字节数必须大于消息体Consumerfetch.max.bytes拉取响应最大字节数参数能调通但我的真实建议是大消息根本不该进Kafka。一条几MB、几十MB的消息放进消息队列意味着大量网络传输、磁盘占用、内存消耗还会拖慢整个Topic的所有消费者。更合理的做法是把大文件传到对象存储或者文件系统Kafka里只传元数据和下载地址消费者拿到地址再去拉取。这也是业界处理大消息的通行模式不只是Kafka任何消息队列都适用。5.4 消费积压先扩容还是先做优化消息积压的排查先搞清楚积压属于“读不过来”还是“处理不过来”。如果是读不过来看消费者是不是太少、分区数是不是不够如果是处理不过来就要看每条消息在业务逻辑里花了多长时间、有没有慢SQL、有没有外部接口超时。很多人一看到积压第一反应是加消费者。但前面说过同组内单个分区只能被一个消费者消费如果你的Topic只有5个分区消费者加到10个也没用剩下5个完全空闲。加消费者之前必须先确认分区数大于当前消费者数量。另一种扩容思路是子Topic分发把积压的消息转移到另一个分区数更多的临时Topic里让更多消费者并行消费处理完再由程序重新写回原Topic。这个方案有点重但紧急积压时确实能解燃眉之急。如果积压不是几个小时内能消化的还得注意Kafka的日志保留时间默认retention是168小时7天。积压太久的旧消息会被清理掉这不是bug是默认策略。必要时可以临时调大log.retention.hours或者用compact策略按Key保留最新消息但这都属于进阶操作了日常先把消费性能优化好才是正路。6. 日常运维需要哪些工具来“看”Kafka6.1 免费可视化工具怎么选“kafka有没有ui界面”——有而且选择相当多。对于入门和中小规模集群我现在最推荐的是Kafka UIprovectus/kafka-ui它Web界面做得顺手能同时管理多个集群可以查看Topic列表、分区分布、消费组Lag还支持直接在界面里查看消息内容、发送测试消息。部署也简单一条Docker命令就能跑起来docker run -p 8080:8080 -e KAFKA_CLUSTERS_0_NAMElocal -e KAFKA_CLUSTERS_0_BOOTSTRAPSERVERSlocalhost:9092 provectuslabs/kafka-ui装好之后打开浏览器指向local集群就能看到集群里的Topic和消费组状态了。还有几个工具也值得认识Kafdrop更轻量适合快速查看消息但管理功能弱CMAK是老牌的Kafka Manager集群运维功能比较全只是界面偏旧Offset Explorer原Kafka Tool是桌面客户端适合Windows下本地连测试集群不依赖Docker。我的建议是本地学习阶段用命令行日常开发调试用Kafka UI或Offset Explorer生产集群监控则重点依赖指标采集系统而不是靠人肉看界面。6.2 上线后需要盯住的几个指标可视化工具只是看个大概真正要保证Kafka集群平稳运行需要有监控。先记住这几个指标它们能覆盖90%的集群故障。UnderReplicatedPartitions副本未同步的分区数长期大于0说明副本同步异常。OfflinePartitions离线分区数大于0说明有的分区没有Leader在服务这是严重故障。Consumer Lag消费积压量重点看是否持续增长而不是看绝对值。RequestHandlerAvgIdlePercent请求处理线程空闲率如果很低说明Broker的CPU/线程已经忙不过来了。磁盘使用率和网络IOKafka重度依赖磁盘磁盘满了集群会直接写不进去。采集层面上Kafka的Broker和客户端都有JMX指标可以配置JMX exporter由Prometheus抓取再用Grafana做可视化。对于小团队Kafka UI自带Lag看板已经够用等集群规模上来了再逐步完善基础设施。运维这件事没必要一步到位但该有的底线指标必须心里有数。7. 面试被问到底层原理怎么答得像个用过的人7.1 为什么Kafka能这么快顺序写、Page Cache、零拷贝Kafka吞吐量高的原因从数据落盘那一刻就开始了。普通数据库更新数据是随机写要先找到目标页再修改慢Kafka做的却是顺序追加写新消息永远写到日志文件的末尾磁盘顺序写的速度可以达到每秒几百MB甚至更高所以写入瓶颈基本不在磁盘。再加上Page Cache的加持Kafka写入的时候优先写到操作系统的页缓存里并不立刻刷盘消费者读的时候也优先命中的是同一份缓存读写完全绕过了应用层与内核层之间多次无谓拷贝。消费端的零拷贝技术让数据从磁盘到网卡的过程中跳过用户态内存拷贝又省掉一大块CPU开销。这些机制叠加起来才是Kafka高吞吐的真正底气。7.2 ISR、HW、LEO副本同步到底在同步什么副本同步机制是Kafka可靠性的基石。每个分区有一个Leader和多个FollowerFollower主动从Leader拉取数据并追加到自己的日志里。所有保持同步的副本集合叫ISR。LEOLog End Offset是每个副本当前日志的最后一条偏移量HWHigh Watermark是ISR里所有副本LEO最小的那个值消费者只能读到HW之前的数据。这么设计是为了避免读到“未确认”的数据。想象一下Leader写入了offset 10到15但只有自己保存了Follower还没同步完这时Leader挂了新的Leader从Follower中选举那么10到15这些消息就丢了。如果消费者之前已经读到了它们就会出现数据“倒退”。HW规定消费者不能读超过所有副本共识的位置从机制上避免了这个矛盾。acksall配合min.insync.replicas大于1就是要求至少几个副本同步才算成功代价是写入延迟更高但安全。7.3 消息文件和清理策略Kafka为什么不删除消息Kafka说“不删除消息”更准确的说法是“不因为有消费者消费了而删除”。消息在Topic里按分区存储在若干个segment日志文件里文件达到一定大小后滚动生成新文件历史文件留着等待清理策略。清理策略有两种。一种是delete按时间默认7天或文件大小删除旧segment适合日志类场景。另一种是compact不按时间删而是按Key保留每个Key的最新一条消息适合保存“实体最新状态”的场景。Kafka内部也有大量用compact的Topic比如消费者组的offset记录。理解这一点你就能回答“为什么重启消费者还能读到之前的历史消息”也能设计消息的保留策略。7.4 高频面试题速查表最后放一张速查表都是面试里出现频率极高的问题一句话答案先记住再去展开细节。问题一句话答案Kafka为什么快顺序追加写、Page Cache、批量发送、零拷贝怎么保证消息不丢生产端acksallretriesBroker端副本数≥2且ISR同步消费端手动提交offset怎么保证消息不重复Kafka幂等只能在单会话分区内去重跨会话要业务侧幂等表兜底怎么保证顺序相同Key进同一分区分区内有序全局有序只能用单分区但牺牲吞吐分区数怎么定大于消费者数结合吞吐评估不宜过大避免Rebalance和文件开销消费者组最佳数量等于分区数时最均衡超过分区数的消费者只会闲置消息积压怎么处理先看LAG和Rebalance再加消费者或扩容临时Topic同时优化单条处理耗时Kafka vs RabbitMQKafka是分布式日志重吞吐、可回放适合日志与流处理RabbitMQ重路由、延迟队列等适合业务消息最后说点个人体会。我带过不少新人发现一个规律能把这篇文章里的坑都踩一遍差不多就理解了Kafka在处理什么。如果从头看到尾还没动手我建议你现在就下载一个单机版本创建3个分区的Topic开两个消费者实例自己看一次Rebalance日志比什么都强。Kafka的技术细节还有很多但入门的门槛其实不高关键是把链路跑通、把概念安到真实场景里。之后无论是深挖源码还是做性能调优都不会再觉得“每个词都认识”了。
返回列表