
清晨五点我盯着监控面板上堆积到几百万的日志数据第一次意识到原来消息队列不是一道面试题而是每天都要面对的现实。那会儿公司日志系统还是服务之间直接 HTTP 调用一到流量高峰整个调用链就卡成幻灯片。后来引入 Kafka把日志采集、指标上报、异步通知全部切开系统才真正喘过气来。这篇文章就是想把 Kafka 从零讲清楚它是什么、能解决什么问题、怎么部署、怎么用、哪些坑必须绕开。适合刚接触消息队列的开发者也适合那些已经在用但说不清原理的运维和架构师。1. 消息队列到底解决了什么问题很多人一开始接触消息队列都是被高并发削峰填谷这些词砸晕的。其实消息队列干的事情特别朴素让消息的发送方和接收方不直接见面。发送方把消息扔进队列就走接收方按自己的节奏来取双方各自独立谁也不拖累谁。1.1 从同步调用到异步解耦我打个比方。没有消息队列的系统就像你去餐厅吃饭每道菜都必须等厨师炒完端到你面前你才能点下一道。厨师忙不过来你就得一直等着你等得不耐烦厨师还得停下来安抚你。这个场景里顾客和服务员、厨师是强耦合的任何一方出问题整条链路就瘫了。用上消息队列之后点餐变成这样你写好菜单消息放到传菜窗口队列服务员拿走菜单去下单后厨做完菜叫号你凭号取餐。你和后厨之间隔了一个传菜窗口谁快谁慢都互不干扰。这就是解耦。实际业务里订单服务只需要把订单创建消息发给 Kafka库存服务、积分服务、短信服务各自订阅这个消息各自处理。订单服务根本不需要知道下游有几个系统、它们处理得慢不慢、会不会挂掉。下游出问题了消息还在队列里躺着等服务恢复后继续消费业务不受影响。1.2 削峰填谷把瞬间洪峰变成平缓水流削峰填谷是消息队列最出名的能力。我有个做电商的朋友大促那天的支付峰值是平时的二十倍。如果支付系统直接扛所有请求扩容成本高得吓人还容易在峰值瞬间被打垮。用 Kafka 做缓冲后所有支付请求先落进队列支付服务按照自己最大的稳定速度消费处理不过来就慢慢排队等峰值过去了积压的消息也能在低谷期逐步消化掉。这里有个容易被误解的点削峰填谷不是让消息变少而是让处理过程的时间分布变均匀。它不丢消息只是把突发流量拉平。就像水库暴雨来了先蓄水而不是让洪水直接冲向下游。这也是为什么 Kafka 这类消息队列在大促、秒杀、日志采集场景里几乎是标配。1.3 点对点、发布订阅与 Kafka 的选择初学者要厘清三类模式。点对点队列一条消息只有一个消费者能拿到拿完消息就没了发布订阅模型一条消息可以被多个订阅者各拿一份日志追加模型消息被持久化存储消费者可以反复读取历史数据。Kafka 属于日志追加 发布订阅的混合体消息写入分区后不是立刻删除而是按保留策略存一段时间消费者组各自记录自己的消费位置offset。这个设计思路让它比传统队列适用范围广很多——既可以做普通消息中间件也可以做事件流平台、日志管道。你得先明白自己要的是哪种模式再决定用不用 Kafka。2. Kafka 核心概念与整体架构我第一次看 Kafka 文档时被 Broker、Topic、Partition、Consumer Group 这些词搞得头大。其实把这些概念串成一个故事就特别清晰了。2.1 Broker、Topic、Partition 到底是谁Broker就是 Kafka 服务节点一台机器上跑一个 Kafka 进程就是一个 Broker。多个 Broker 组成集群集群里所有消息分散存在各个 Broker 上。Topic是消息的逻辑分类相当于数据库里的表。你发订单消息就发到 order-topic发日志消息就发到 log-topic不同业务互不干扰。Partition分区是 Topic 的物理分片。一个 Topic 可以拆成多个分区每个分区是一个有序的消息日志。为什么要有分区因为单台机器存不下所有消息也扛不住所有读写。分区把压力和容量分摊到多台机器上又能并行读写。分区数越多吞吐上限越高但也不是越多越好——分区太多意味着文件句柄、选举开销、客户端连接开销都会上涨。消息在分区内部是有序的靠偏移量offset标定位置。你可以把分区理解成一本书消息是一行行文字offset 就是行号。分区之间不保证全局有序所以全局严格有序在 Kafka 里是要靠设计来凑的后面我会专门讲。2.2 消费者组与 offset 的进阶理解消费者组Consumer Group是 Kafka 的一个精髓设计。同一个组内一个分区同一时刻只分配给一个消费者消费不同组之间互不影响都能完整消费一遍所有消息。这个机制带来的直接好处是你想提升消费吞吐就往组里加消费者Kafka 会自动重平衡Rebalance分区分配。比如一个 Topic 有 6 个分区组里 3 个消费者每人分 2 个加到 6 个消费者每人分 1 个。但如果你加到 7 个消费者那第 7 个人会闲着没分区可分因为它抢不到分区了。这也是判断分区数设置是否合理的经验之一消费者数超过分区数时超出部分消费力就白搭了。offset 是消费者在分区上的书签。Kafka 0.9 之前offset 存在 ZooKeeper 里后来改成存在内部主题 __consumer_offsets 中。消费者每消费完一批消息就提交一次 offset记录我读到哪了。下次再启动时从上次提交的位置继续读。这里面有个坑如果消费完业务逻辑还没处理完就提交了 offset进程一挂消息就丢了被跳过如果处理完了但没提交 offset重启后会重复消费。这是重复消费和消息丢失两大问题的根源之一第 6 节我会展开排坑。2.3 Kafka 高吞吐的底层逻辑Kafka 能做到单机每秒几十万条写入、消息堆积几亿条不崩溃靠的是一套组合拳。第一是顺序写磁盘。普通数据库随机写磁盘很慢但 Kafka 是追加写入消息挨着挨着往文件尾部写把随机 IO 变成了顺序 IO。机械硬盘顺序写也能跑到一两百兆每秒SSD 更快所以它的写入速度上限远高于直觉。第二是页缓存与零拷贝。Kafka 读写充分利用操作系统页缓存生产者写入的数据优先落在页缓存里消费者读取如果命中页缓存就直接从内存返回完全不经过应用层复制发送给用户时再用 sendfile 零拷贝减少两次内存拷贝和系统调用。所以 Kafka 进程本身占用内存不高因为它把内存交给 OS 管了。第三是批量与压缩。生产者把多条消息攒成一个批次再发送消费者也是拉取一批再处理减少网络往返次数。开启压缩后如 lz4、zstd网络传输量还能再降。我用一个生活类比来收尾这部分普通数据库像图书馆里到处插书每插一本都要查半天索引Kafka 像流水线上的传送带货物一个接一个往传送带上放放完就往下游推秩序决定了效率。3. 从零安装到跑通第一条消息扯了这么多原理不如上手跑一遍。这一节我带你把 Kafka 单机环境从下载到命令行收发消息完整跑通新手照抄即可注意我标注的版本坑。3.1 下载解压与基础配置Kafka 有两个主流模式老牌的 ZooKeeper 模式和 2.8 之后引入、3.x 之后逐渐成熟的 KRaft 模式去 ZooKeeper。单机入门我建议直接用 KRaft 模式少装一个 ZooKeeper省心很多。我现在用的是 Kafka 3.6 系列的二进制包。下载解压之后进入 config 目录核心文件是 server.properties。你要关注这几个配置项# 每个 Broker 的唯一 ID broker.id0 # 监听地址不要用默认的 localhost不然外部客户端连不上 listenersPLAINTEXT://0.0.0.0:9092 # 对外通告地址填你机器的实际 IP advertised.listenersPLAINTEXT://192.168.1.10:9092 # 日志存储目录一定放数据盘别放系统盘 log.dirs/data/kafka-logs # 日志保留时长默认 7 天按需调整 log.retention.hours168这里必须强调一个配置advertised.listeners。很多新手第 1 次从外部机器连不上 Kafka就是因为它没写。Broker 启动后会把这份地址告诉客户端客户端拿着它去连 Broker。如果配的是 localhost你在另一台机器上怎么都连不上。踩过一次这个坑你就永远记住了。3.2 格式化存储目录并启动KRaft 模式需要先生成集群 ID然后格式化存储目录。这个步骤千万别省否则启动直接报错。# 生成一个唯一集群 ID bin/kafka-storage.sh random-uuid # 用输出的 UUID 格式化存储目录假设 UUID 是 xxxxxxxx-xxxx-xxxx bin/kafka-storage.sh format -t xxxxxxxx-xxxx-xxxx -c config/server.properties格式化完成后启动bin/kafka-server-start.sh -daemon config/server.properties启动后看日志有没有报错确认进程在跑jps # 或 ps -ef | grep kafka看到 Kafka 进程就说明起来了。如果启动失败多数情况是端口被占、log.dirs目录没有写权限或者没格式化。把错误日志打开一行行看比乱猜管用。3.3 命令行创建 Topic 与收发消息接下来用自带脚本验证消息通路。创建一个双分区、双副本的测试主题bin/kafka-topics.sh --bootstrap-server 192.168.1.10:9092 \ --create --topic quickstart-events \ --partitions 2 --replication-factor 1注意单机模式下副本因子只能写 1因为副本需要跨 Broker 同步一个节点写 2 就是自找麻烦。创建完查询主题信息确认分区正常bin/kafka-topics.sh --bootstrap-server 192.168.1.10:9092 \ --describe --topic quickstart-events开一个终端启动消费者再开一个终端启动生产者。我习惯先启动消费者再发消息这样能立刻看到结果# 生产者终端 bin/kafka-console-producer.sh --bootstrap-server 192.168.1.10:9092 \ --topic quickstart-events # 消费者终端 bin/kafka-console-consumer.sh --bootstrap-server 192.168.1.10:9092 \ --topic quickstart-events --from-beginning生产端输入一行按回车消费端马上就能看到。整个流程和用 curl 测接口一样直接这就是 Kafka 的命令行自测方法。4. 生产环境部署与关键参数调优开发环境跑通只是开胃菜。真正上了生产你要面对的是集群怎么搭、副本怎么同步、参数怎么调、硬件怎么选。这一节我把生产部署里真正影响稳定性的东西拎出来讲。4.1 三节点集群搭建与脑裂预防生产环境至少三台 Broker这是行业共识。为什么是三台配合副本因子 3允许坏掉一台机器而不丢数据、仍能选举出新的领导者。两节点加副本因子 2坏一台照样缺副本而且容易出尴尬数的选举问题。搭建三节点集群关键步骤是三步三台机器各自安装 Kafka、每台的server.properties配置不同的broker.id、使用同一个cluster.id格式化存储目录。KRaft 模式还要多配一个controller.quorum.voters把三台节点的 controller 功能都列进去# 三台节点分别设 0、1、2 broker.id0 # controller 选举参与者三台都参与 controller.quorum.voters0192.168.1.10:9093,1192.168.1.11:9093,2192.168.1.12:9093有个细节我要提醒默认情况下一挂掉两台剩下的单节点会产生大多数吗不一定。KRaft 的 controller 选举依赖多数派三节点集群挂两台就只剩一个节点不够多数整个集群会拒绝写入。这是防止脑裂的必要代价——宁可不可用也不能出现两个大脑同时写数据然后数据分叉。你要知道这个行为上线前给业务方讲清楚别等故障了才解释。4.2 副本机制与 ISR 的取舍逻辑Kafka 的副本分两种角色Leader 和 Follower。读写都走 LeaderFollower 只负责同步数据。Leader 挂了控制器从 Follower 里挑一个新的 Leader 出来。ISRIn-Sync Replicas是和 Leader 保持同步的副本集合。不是所有 Follower 都有资格当备胎只有滞后程度在规定范围内的副本才在 ISR 里。这里有个关键参数min.insync.replicas它规定了最少要有几个副本确认写入成功消息才算写入成功。举一个生产里最典型的配置组合配置项推荐值说明replication.factor3每个分区保留 3 份数据min.insync.replicas2至少 2 个副本同步成功才算写入成功acksall生产者要求所有 ISR 副本都确认这三个配置配套使用才能同时做到不丢消息和容忍单节点故障。你可能会问为什么min.insync.replicas不设成 3因为设成 3 时只要坏掉一台机器所有写入就全被拒绝系统直接不可用。设成 2 的话坏一台还能继续写。数据可靠性和可用性之间的平衡就体现在这个 2 上。我在实际项目里基本都用这套组合只有对数据极度敏感、且能接受停机风险的业务才会用更激进的配置。4.3 内存、磁盘与存储保留策略Kafka 的性能和硬件关系非常直接尤其是磁盘和内存。先说磁盘Kafka 写入是顺序写但日志数据积累量大而且大量读操作靠页缓存所以优先选 SSD读多写少的日志场景 NVMe 的收益尤其明显。磁盘容量按消息体大小、日生产量、保留天数三个数算单日消息量GB × 保留天数 × 副本数 × 1.2预留余量 所需磁盘容量举个例子每天生产 50GB 消息保留 3 天副本 3 份那就是 50×3×3×1.2 约等于 540GB。别把计算想复杂它就是一道乘法题但很多人上线前根本没算过结果一周就把磁盘打满。内存方面我建议 Broker 所在机器的内存至少 32GB 起步。Kafka 堆内存不需要给太多4GB 到 6GB 通常足够剩下的内存都给页缓存。JVM 参数里有个新手必踩的坑-Xmx和-Xms要设成一样防止堆动态伸缩引发 Full GC。如果是容器部署务必用KAFKA_HEAP_OPTS明确设置默认值在某些发行版里不会生效。5. 客户端开发生产者与消费者实战命令行验证完就该写代码了。下面我用 Java 客户端演示生产和消费的核心用法同时把 ack、提交 offset 这些决定消息命运的细节讲透。你用 Python、Go 的时候概念完全一致只是 API 长得不同。5.1 生产者 API 与 ack 机制生产者核心就五个对象配置项、主题、序列化器、消息对象和发送回调。直接看代码Properties props new Properties(); props.put(bootstrap.servers, 192.168.1.10:9092,192.168.1.11:9092); // 三个关键参数acks、retries、batch.size props.put(acks, all); props.put(retries, 3); props.put(batch.size, 16384); props.put(linger.ms, 5); props.put(key.serializer, org.apache.kafka.common.serialization.StringSerializer); props.put(value.serializer, org.apache.kafka.common.serialization.StringSerializer); ProducerString, String producer new KafkaProducer(props); producer.send(new ProducerRecord(order-topic, key-001, order-created), (metadata, exception) - { if (exception null) { // 发送成功metadata 里有分区号和 offset } else { // 发送失败记录并补偿 } });acks参数决定了要等几个副本确认。我建议生产环境一律acksall配合上一节说的min.insync.replicas2。retries是发送失败的重试次数注意重试可能带来消息乱序如果业务严格要求顺序且没有其它方案可以在需要保持顺序的消息上压低重试值或用同步发送。batch.size和linger.ms是吞吐和延迟的平衡杠杆想把吞吐拉高就调大这两个值想降低单条延迟就调小linger.ms。5.2 消费者 API 与提交 offset 的两种姿势消费者比生产者复杂重点全在 offset 提交时机上。我见过太多生产事故都是 offset 提交姿势不对导致的。先看基本代码Properties props new Properties(); props.put(bootstrap.servers, 192.168.1.10:9092); props.put(group.id, order-consumer-group); props.put(enable.auto.commit, false); // 我习惯关掉自动提交 props.put(key.deserializer, org.apache.kafka.common.serialization.StringDeserializer); props.put(value.deserializer, org.apache.kafka.common.serialization.StringDeserializer); KafkaConsumerString, String consumer new KafkaConsumer(props); consumer.subscribe(Collections.singletonList(order-topic)); while (true) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecordString, String record : records) { // 处理业务逻辑 } consumer.commitSync(); // 拉取一批处理完再提交 offset }自动提交enable.auto.committrue是默认行为每隔一段时间自动提交当前消费的位置。看起来省事实际很危险如果业务处理比较慢拉取了一批消息但还没处理完自动提交的 offset 可能已经推进了这时进程崩溃重启没处理完的消息就再也消费不到了这就是消息丢失的一种来源。我推荐的做法是关掉自动提交手动提交并且先处理业务逻辑再提交 offset。这样能最大化避免丢消息。代价是可能重复消费——处理完了、提交之前挂了重启后同一批消息会再消费一遍。所以你的业务逻辑要么做成幂等要么能接受至少一次at-least-once语义。消息队列领域没有完美的恰好一次Kafka 提供的幂等生产者解决的是生产者重复发送问题而消费端幂等还得自己在业务层解决。这是选型时必须认清的现实。5.3 消息顺序性一个分区一个消费者Kafka 只保证分区内有序不保证跨分区有序。所以解决顺序问题就一句话让需要有序的消息进同一个分区。做法是给消息设置一个业务 keyKafka 用 key 做哈希同一个 key 的消息一定进同一个分区。比如订单状态变更用订单号做 key同一订单的所有状态流转就在一个分区里排队消费者串行处理顺序就保住了。另一个要求是这个分区只能被一个消费者线程处理。如果组里有多个消费者订阅了同一个分区同一组内不会发生或者你在消费端自行多线程处理同一个分区的消息顺序也会被打破。我见过一个项目用 key 保证同一订单进同一分区但消费端图快搞了个线程池并行处理结果订单状态回退和推进的顺序乱了最后返工改成单线程消费。所以规则是分区有序 该分区单消费者 消费处理不引入并行竞争三者缺一不可。6. 高频故障复盘重复消费、延迟高、消息丢失这一节是我最想写的部分因为网上教程大多停留在怎么用但真实世界全是怎么炸。我把这些年排过的高频故障整理成速查表并展开讲三个老大难。6.1 重复消费的根治思路重复消费的表现是消费端明明处理过这条消息又收到了一遍。常见原因有三个offset 未提交处理完业务但提交 offset 前进程挂了重启后从旧 offset 重新消费。消费端重平衡消费者加入或退出触发 Rebalance分区被重新分配新的消费者从上次提交的 offset 开始读而上次提交可能落后于实际处理位置。生产者重试生产者发送超时后重试Broker 其实已经写入了但生产者以为失败了又发一次就产生了两条相同业务内容的消息。方案也是三个层面。消费端做幂等把业务唯一键订单号、流水号查出或写入数据库时加唯一约束重复消息直接忽略。生产者端在需要时可开启幂等生产者enable.idempotencetrue防止重试导致的多写。消费端手动提交时如果对准确性要求高可以使用处理完再提交 提交失败则退出进程触发重平衡策略宁可重复也不可丢失。记住一个原则消息队列的语义天然是至少一次业务系统必须默认重试会发生。6.2 消息延迟高的排查路径消息延迟高是另一个高频故障表现形式是消费端 lag积压持续上涨。排查路径我按优先级列看消费端是否在干活Java 应用线程卡死、数据库慢查询、下游接口超时都可能让单条消息处理时间暴涨。先用kafka-consumer-groups.sh查看 group 的 lagbin/kafka-consumer-groups.sh --bootstrap-server 192.168.1.10:9092 \ --describe --group order-consumer-group输出里 CURRENT-OFFSET 和 LOG-END-OFFSET 的差距就是积压量。看分区分布如果某个分区 lag 特别高、其它分区低大概率是部分消息的 key 集中导致某个分区数据量失衡或者消费该分区的消费者处理能力弱。这时候要检查分区分配是否均匀必要时增加分区数并重新设计 key。看 Broker 瓶颈磁盘 IO 利用率、网络带宽、页缓存命中率。Kafka 的指标里kafka.server:typeBrokerTopicMetrics的 BytesInPerSec、BytesOutPerSec 能看出流量是否打满网卡。看消费线程数单消费者线程处理能力到顶时先加消费者或者用多线程消费模型。但注意前面说的顺序性约束多线程只适用于对顺序不敏感的消息。延迟问题 90% 是消费端慢不是 Kafka 慢。先怀疑自己的业务代码再怀疑中间件。6.3 消息丢失的所有可能路径消息丢失是最严重的事故分散在各环节。我把它拆成三段讲生产端acks0或acks1时Broker 写入主副本失败或未同步就返回成功消息就丢了。解决acksallmin.insync.replicas2。Broker 端unclean.leader.election.enabletrue时Leader 挂了会允许落后副本当选新 Leader那么 ISR 副本里没有的新消息就全丢了。解决把这个参数设成false宁可短暂不可用也不能丢数据。还有一个场景是磁盘损坏——所以副本因子至少 3并且定期做数据备份。消费端自动提交 offset 且处理失败不 catch或者手动提交时机提前都会造成消息读了但不处理。解决关自动提交 处理完再提交 抛异常就走重试。我自己有个习惯关键业务上线前会做一次故障演练分别杀掉一个 Broker、停掉消费者进程、模拟磁盘满看监控和告警能不能及时反映消息到底丢没丢。与其等生产炸了再复盘不如提前把事故预演一遍。Kafka 没有绝对不丢的魔法只有每个环节都用正确姿势堵住缺口之后的尽量不丢。6.4 可视化工具与日常运维建议命令行写多了眼睛疼我平常用两个工具。一个是 Kafka 官方生态里的 Kafka UI老项目名 Kafdrop后来几个项目合并过打开浏览器就能看到 Topic 列表、分区分布、消费组 lag 和消息内容做日常巡检非常直观。另一个是开源的 Kafka Eagle现在叫 EFak功能更重一些支持告警和监控面板。轻量看图用前者完整监控用后者。运维上补充两个容易忽略的日常动作一是给消费组的 lag 配告警超过阈值就报警很多延迟事故其实是在积压发生两小时后才被发现的二是定期检查磁盘占用和历史 Topic 清理Kafka 的日志保留策略是到时间删旧段但如果你创建了一堆不用的 Topic 还在持续写入磁盘照样会满。我见过太多团队 Topic 建了不清理最后磁盘告警响了一排还找不到是哪个业务在写。7. 选型对比Kafka、RabbitMQ 与 RocketMQ被问得最多的问题就是这三个到底怎么选我把它们放在一张表里讲清楚再给你我的选型逻辑。维度KafkaRabbitMQRocketMQ吞吐量极高百万级/秒中等万级/秒高十万级/秒延迟毫秒级但更侧重吞吐微秒级低延迟优先毫秒级消息可靠性高配合 acksall高支持多种确认机制高支持事务消息消息顺序分区内有序单队列有序队列内有序功能丰富度较基础重吞吐和流处理路由灵活Exchange、死信队列事务消息、定时消息、消息轨迹社区与云厂商支持极广几乎所有云都有广国内广泛使用我的选型经验是日志采集、埋点上报、大数据管道、需要超高吞吐的流式场景选 Kafka 没有悬念。它的设计目标就是海量数据、高吞吐、追加式消费拿它做业务消息中间件不是不行但你会发现很多高级功能要自己造轮子比如延迟消息、事务消息做起来都很费劲。内部系统间复杂的业务消息、需要灵活路由和精细控制选 RabbitMQ。它的 Exchange 路由模型能把一条消息按规则分发给多个队列配上死信队列做异常处理非常契合业务系统。缺点是吞吐上限低扛不住超大流量。有大厂加持、既要高吞吐又要事务和定时消息等高级特性选 RocketMQ。它在吞吐上逼近 Kafka功能上又比 Kafka 完善特别适合国内电商、支付类业务。如果团队已经熟悉 Kafka 生态而且业务主要是事件流就别为了功能多迁到 RocketMQ迁移成本远高于那点功能收益。另外提醒一个被忽略的点选型不只是选中间件还要看团队的维护能力和上下游生态。没人会调 JVM、没人懂分区扩缩容的团队别贸然上 Kafka 集群运维能力一般但业务需要可靠消息托管云服务或者 RabbitMQ 反而更稳。技术选型终究是平衡题没有标准答案。最后说点实际的我自己在这些年踩坑后总结了几条土办法第一Kafka 的配置参数动任何一个之前先想清楚它对丢消息、重复消费、延迟、吞吐这四个维度的影响不要只盯着某一个指标调第二任何消费处理逻辑都默认会重复执行幂等设计从第一天就做第三上线前一定配好 lag 和磁盘告警这两个是 Kafka 最闷声酿成大祸的指标第四多读生产集群的日志和监控Kafka 的官网文档写得不算通俗但它的故障日志和 JMX 指标比任何教程都诚实。如果你是从零开始建议先把这篇里的单机环境跑通亲手创建一个 Topic、收发一次消息、看一次分区描述再感受一下 kill 消费者进程后的重复消费现象。把基础动作变成肌肉记忆后面学集群、学调优、学源码都会顺很多。Kafka 不是几个月就能吃透的东西但搞清楚这条主线它就不再是黑盒了。