消息中间件:Kafka 与 RabbitMQ
消息中间件就是系统之间的邮局。以前 A 系统要通知 B 系统,直接打电话(同步调用),B 挂了 A 也卡住。现在 A 把信扔进邮局(消息队列),该干嘛干嘛,B 自己去邮局取。削峰填谷、系统解耦、异步处理,这就是消息队列干的三件事。Kafka 就是消息界的高铁——跑得快、运量大;RabbitMQ 就是快递小哥——路由灵活、可靠送达。
3.1 Kafka:分布式流处理平台
论Kafka 的脾气
分布式流处理平台,LinkedIn 出品,高吞吐/持久化/分布式,日志/事件流/大数据管道首选。它不是传统意义的"消息队列",更像一个"分布式提交日志"。写磁盘用顺序写(比随机写快 100 倍),零拷贝 sendfile,单机每秒能写几十万条消息。它的设计哲学就是不惜一切代价追求吞吐量。
3.2 Kafka 核心概念(必须背下来)
| 术语 | 解释 |
|---|---|
| Broker | Kafka 服务器节点,一台机器就是一个 Broker |
| Topic | 消息的分类(像数据库的表),生产者写 Topic,消费者读 Topic |
| Partition | Topic 的分区。一个 Topic 可以分成多个 Partition,分布在不同 Broker 上,这是并行和扩展的关键 |
| Producer | 生产者,往 Topic 写消息 |
| Consumer | 消费者,从 Topic 读消息 |
| Consumer Group | 消费者组。同组内每个 Partition 只能被一个消费者消费,组间互不影响 |
| Offset | 消费者在 Partition 里的消费位置(序号),记录"读到哪了" |
| KRaft | 3.x 后内置的 Raft 协议,替代 ZooKeeper 管理集群元数据 |
3.3 安装:Docker 单机版
docker-compose.yml —— Kafka KRaft 单机版
3.4 生产者配置:怎么发消息最可靠
| 参数 | 含义与取值 |
|---|---|
acks | 0=发了就不管(最快但可能丢);1=Leader 写了就回(默认);all/-1=所有 ISR 副本都写了才回(最可靠,生产推荐) |
retries | 发送失败重试次数,设个 3-5 次 |
batch.size | 批量发送的大小(字节),攒一批一起发,提高吞吐 |
linger.ms | 等多久凑一批(毫秒),比如 5ms 内攒够 batch 就发 |
compression.type | 压缩格式(gzip/snappy/lz4),省网络带宽 |
3.5 消费者配置:怎么不丢消息
论Kafka 三种投递语义
① 至多一次(At most once):offset 先提交再处理。如果处理崩了,这条消息就丢了。适合日志采集这种丢了无所谓的场景。
② 至少一次(At least once):先处理再提交 offset。如果处理完但提交崩了,下次重启会重复消费。绝大多数场景用这个,消费者要做幂等。
③ 恰好一次(Exactly once):Kafka 0.11+ 提供事务机制,保证消息只处理一次。性能开销大,只有金融等严格场景才用。
3.5b Exactly-Once 到底怎么实现:幂等 + 事务
引子:"至少一次"语义下,消费者重启就会重复消费。大多数业务靠"消费者自己做幂等"扛过去了,但金融场景(转账、计费)连"重复扣一次款"都不能接受。这时才上 Kafka 的 Exactly-Once,靠两把钥匙:幂等生产者和事务。
论两把钥匙各管一段
① 幂等生产者(enable.idempotence=true):解决"生产者自己重试导致重复写"。Kafka 给每个生产者一个 PID,每条消息带序号,Broker 发现同一个序号又来了就自动去重。开了它,acks=all、retries=Integer.MAX_VALUE 是强制的。
② 事务(transactional.id):解决"消费-处理-生产"这条链路里"读了又写"的重复。典型是流处理:从 A 读消息、处理后写进 B,要么这条"读+写"原子完成,要么压根不发生。生产者把一批写操作包进一个事务,消费者用 isolation.level=read_committed 只读已提交的消息。
开启幂等生产者(最简单,单分区不跨系统推荐)
跨系统事务(真正端到端 Exactly-Once,代价大)
别动不动就上事务。跨系统事务要维护事务协调、加锁、写事务日志,吞吐量掉一截、延迟涨一截。99% 的业务用"至少一次 + 消费者幂等(比如用订单号去重)"就够,既简单又快。判断标准只有一个:重复一次会不会出真金白银的事故?会,才上事务;不会,老老实实用幂等。
3.6 消息顺序性:Partition 内有序
Kafka 只保证同一个 Partition 内的消息是有序的。不同 Partition 之间的消息谁先谁后不保证。想让同一类业务消息有序(比如同一个订单的事件),就让它们按 key(如 order_id)哈希到同一个 Partition。
3.7 持久化:顺序写磁盘 + 日志分段
Kafka 消息不是放内存,而是顺序写磁盘。磁盘的顺序写比随机写快 100 倍,接近内存速度。消息按 segment 分段存储,老的 segment 按时间或大小自动删除(retention.ms/retention.bytes)。这就是 Kafka 能存海量消息又高性能的秘密。
3.8 Spring Boot 整合 Kafka
pom.xml 依赖
application.yml
OrderEventProducer.java —— 生产者
OrderEventConsumer.java —— 消费者(手动提交)
3.9 Kafka 集群部署与监控
生产环境至少 3 个 Broker,每个 Partition 配 2-3 个副本(replication-factor)。Leader 负责读写,Follower 只同步。Leader 挂了从 ISR(同步副本集合)里选新 Leader。监控用 JMX 暴露指标,配合 Kafka Eagle 或 kafka-manager 看图形。
Q1:Kafka 为什么这么快?
查看答案
① 顺序写磁盘(比随机写快 100 倍);② 零拷贝 sendfile;③ 批量发送;④ 分区并行;⑤ 页缓存。
Q2:Kafka 怎么保证消息不丢?
查看答案
生产者 acks=all + retries;Broker 副本因子>=2 + min.insync.replicas=2;消费者手动提交 offset(先处理再提交)。
Q3:Kafka 怎么保证消息顺序?
查看答案
同一 key 的消息哈希到同一 Partition,Partition 内有序。全局有序只能单 Partition(牺牲并行度)。
3.11 RabbitMQ:企业级消息代理
论RabbitMQ 的脾气
Erlang 写的 AMQP 消息代理,可靠性高/路由灵活/企业级首选。和 Kafka 的设计哲学完全不同:Kafka 追求吞吐量,RabbitMQ 追求可靠送达 + 灵活路由。它有四种交换机类型(Direct/Fanout/Topic/Headers),能实现复杂的路由规则。电商订单、支付通知这种"不能丢、要精准路由"的场景,RabbitMQ 更合适。
3.12 RabbitMQ 核心概念
| 术语 | 解释 |
|---|---|
| Broker | RabbitMQ 服务器 |
| VHost | 虚拟主机,逻辑隔离,像 MySQL 的 database |
| Exchange | 交换机,接收生产者消息,按规则路由到 Queue |
| Queue | 队列,消息实际存放的地方,消费者从这里取 |
| Binding | Exchange 和 Queue 之间的绑定关系,带 routing key |
3.13 四种 Exchange 类型
| 类型 | 路由规则 | 场景 |
|---|---|---|
| Direct | routing key 完全匹配才路由 | 点对点,如"订单创建"事件 |
| Fanout | 广播到所有绑定的 Queue(不管 key) | 系统通知、配置刷新 |
| Topic | routing key 模式匹配(* 匹配一个词,# 匹配多个) | 灵活路由,如 log.info.* 匹配 log.info.order |
| Headers | 按消息头属性匹配(不用 routing key) | 复杂多条件路由 |
3.14 安装:Docker + 管理界面
3.15 消息可靠性:生产者确认 + 持久化 + 手动 ack
RabbitMQ 的可靠性靠三道防线:① 生产者 Confirm 模式确认消息到了 Broker;② 队列和消息都设持久化(durable + delivery_mode=2);③ 消费者手动 ack,处理完才告诉 Broker"我收到了"。三道全上,消息基本不会丢。
3.16 延迟队列:订单超时取消
电商场景:用户下单后 30 分钟没付款,自动取消订单。RabbitMQ 没有现成的延迟队列,用TTL + 死信队列(DLX)实现:消息进一个"没人消费的队列",TTL 到了变成死信,路由到真正处理取消逻辑的队列。
3.17 Kafka vs RabbitMQ 选型
| 对比点 | Kafka | RabbitMQ |
|---|---|---|
| 吞吐量 | 极高(几十万/秒) | 高(几万/秒) |
| 延迟 | ms 级 | us 级(更低) |
| 可靠性 | 高(副本+持久化) | 极高(Confirm+ack+持久化) |
| 路由 | 弱(按 Partition 哈希) | 强(四种 Exchange 灵活路由) |
| 协议 | 自定义(二进制) | AMQP(标准协议) |
| 消息顺序 | Partition 内有序 | 单 Queue 有序 |
| 适用场景 | 日志采集、大数据管道、事件流 | 订单通知、业务解耦、RPC 异步 |
① Kafka = 高铁,高吞吐、顺序写磁盘、Partition 并行。acks=all + 手动提交 offset = 不丢消息。
② RabbitMQ = 快递小哥,路由灵活(四种 Exchange)、可靠性高(Confirm + 手动 ack)。延迟队列用 TTL+DLX。
③ 选型:大数据/日志/事件流选 Kafka;业务消息/复杂路由/低延迟选 RabbitMQ。