楼层: 首页/ 软件技术/ 中间件全景/ 消息中间件:Kafka 与 RabbitMQ
04

消息中间件:Kafka 与 RabbitMQ

Message Queue · Kafka / RabbitMQ

消息中间件就是系统之间的邮局。以前 A 系统要通知 B 系统,直接打电话(同步调用),B 挂了 A 也卡住。现在 A 把信扔进邮局(消息队列),该干嘛干嘛,B 自己去邮局取。削峰填谷、系统解耦、异步处理,这就是消息队列干的三件事。Kafka 就是消息界的高铁——跑得快、运量大;RabbitMQ 就是快递小哥——路由灵活、可靠送达。

3.1 Kafka:分布式流处理平台

论Kafka 的脾气

分布式流处理平台,LinkedIn 出品,高吞吐/持久化/分布式,日志/事件流/大数据管道首选。它不是传统意义的"消息队列",更像一个"分布式提交日志"。写磁盘用顺序写(比随机写快 100 倍),零拷贝 sendfile,单机每秒能写几十万条消息。它的设计哲学就是不惜一切代价追求吞吐量。

3.2 Kafka 核心概念(必须背下来)

Kafka 术语表,面试必问
术语解释
BrokerKafka 服务器节点,一台机器就是一个 Broker
Topic消息的分类(像数据库的表),生产者写 Topic,消费者读 Topic
PartitionTopic 的分区。一个 Topic 可以分成多个 Partition,分布在不同 Broker 上,这是并行和扩展的关键
Producer生产者,往 Topic 写消息
Consumer消费者,从 Topic 读消息
Consumer Group消费者组。同组内每个 Partition 只能被一个消费者消费,组间互不影响
Offset消费者在 Partition 里的消费位置(序号),记录"读到哪了"
KRaft3.x 后内置的 Raft 协议,替代 ZooKeeper 管理集群元数据

3.3 安装:Docker 单机版

docker-compose.yml —— Kafka KRaft 单机版

version: '3' services: kafka: image: apache/kafka:3.8.0 ports: - "9092:9092" environment: KAFKA_NODE_ID: 1 KAFKA_PROCESS_ROLES: broker,controller KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092 KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT KAFKA_CONTROLLER_QUORUM_VOTERS: 1@localhost:9093 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
$ docker compose up -d $ docker exec -it kafka /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server localhost:9092 \ --create --topic test-events --partitions 3 --replication-factor 1 Created topic test-events.

3.4 生产者配置:怎么发消息最可靠

生产者关键参数
参数含义与取值
acks0=发了就不管(最快但可能丢);1=Leader 写了就回(默认);all/-1=所有 ISR 副本都写了才回(最可靠,生产推荐)
retries发送失败重试次数,设个 3-5 次
batch.size批量发送的大小(字节),攒一批一起发,提高吞吐
linger.ms等多久凑一批(毫秒),比如 5ms 内攒够 batch 就发
compression.type压缩格式(gzip/snappy/lz4),省网络带宽

3.5 消费者配置:怎么不丢消息

# consumer.properties group.id=order-consumers # 消费者组名,同组内分区负载均衡 auto.offset.reset=earliest # 没有 offset 时从头消费(latest=从最新开始) enable.auto.commit=false # 关闭自动提交,手动提交更安全 # 处理完消息后手动 commitSync(),确保处理完才提交 offset

论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 只读已提交的消息。

开启幂等生产者(最简单,单分区不跨系统推荐)

spring: kafka: producer: acks: all enable-idempotence: true # 开启幂等:生产者重试导致的重复写,Broker 自动去重 retries: 3 # 这一步只解决"生产者重复",不解决"消费者重复"。 # 绝大多数线上场景:开幂等 + 消费者做幂等键,就够了。

跨系统事务(真正端到端 Exactly-Once,代价大)

# 生产者侧:给一个事务 ID,把"读 A → 处理 → 写 B"包成事务 props.put("transactional.id", "tx-1"); ProducerFactory pf = new DefaultKafkaProducerFactory(props); KafkaTransactionManager ktm = new KafkaTransactionManager(pf); # 消费者侧:只读取已提交的事务消息 spring: kafka: consumer: properties: 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 依赖

<dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> </dependency>

application.yml

spring: kafka: bootstrap-servers: localhost:9092 producer: acks: all retries: 3 key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.springframework.kafka.support.serializer.JsonSerializer consumer: group-id: order-consumers auto-offset-reset: earliest enable-auto-commit: false

OrderEventProducer.java —— 生产者

@Service public class OrderEventProducer { @Autowired private KafkaTemplate<String, OrderEvent> kafkaTemplate; public void sendOrderCreated(OrderEvent event) { # 按 orderId 做 key,保证同一订单的事件进同一 Partition kafkaTemplate.send("order-events", event.getOrderId(), event); System.out.println("发送成功:" + event.getOrderId()); } }

OrderEventConsumer.java —— 消费者(手动提交)

@Component public class OrderEventConsumer { @KafkaListener(topics = "order-events", groupId = "order-consumers") public void onMessage(OrderEvent event, Acknowledgment ack) { try { # 处理订单事件:扣库存、加积分、发通知... processOrder(event); ack.acknowledge(); # 处理完才提交 offset } catch (Exception e) { # 处理失败不提交,下次还会消费到(至少一次语义) log.error("处理失败,稍后重试", e); } } }

3.9 Kafka 集群部署与监控

生产环境至少 3 个 Broker,每个 Partition 配 2-3 个副本(replication-factor)。Leader 负责读写,Follower 只同步。Leader 挂了从 ISR(同步副本集合)里选新 Leader。监控用 JMX 暴露指标,配合 Kafka Eagle 或 kafka-manager 看图形。

3.10 Kafka 面试重点

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 核心概念

RabbitMQ 术语表
术语解释
BrokerRabbitMQ 服务器
VHost虚拟主机,逻辑隔离,像 MySQL 的 database
Exchange交换机,接收生产者消息,按规则路由到 Queue
Queue队列,消息实际存放的地方,消费者从这里取
BindingExchange 和 Queue 之间的绑定关系,带 routing key

3.13 四种 Exchange 类型

RabbitMQ 交换机路由方式
类型路由规则场景
Directrouting key 完全匹配才路由点对点,如"订单创建"事件
Fanout广播到所有绑定的 Queue(不管 key)系统通知、配置刷新
Topicrouting key 模式匹配(* 匹配一个词,# 匹配多个)灵活路由,如 log.info.* 匹配 log.info.order
Headers按消息头属性匹配(不用 routing key)复杂多条件路由

3.14 安装:Docker + 管理界面

docker run -d --name rabbitmq \ -p 5672:5672 \ # AMQP 协议端口 -p 15672:15672 \ # 管理界面端口 rabbitmq:3.13-management # 打开浏览器访问 http://localhost:15672 # 默认账号密码:guest / guest
$ curl -u guest:guest http://localhost:15672/api/overview {"rabbitmq_version":"3.13.x","management_version":"..."}

3.15 消息可靠性:生产者确认 + 持久化 + 手动 ack

RabbitMQ 的可靠性靠三道防线:① 生产者 Confirm 模式确认消息到了 Broker;② 队列和消息都设持久化(durable + delivery_mode=2);③ 消费者手动 ack,处理完才告诉 Broker"我收到了"。三道全上,消息基本不会丢。

# Spring Boot application.yml spring: rabbitmq: host: localhost port: 5672 username: guest password: guest publisher-confirm-type: correlated # 生产者确认 publisher-returns: true # 消息不可路由时返回 listener: simple: acknowledge-mode: manual # 手动 ack prefetch: 10 # QoS:一次最多推 10 条

3.16 延迟队列:订单超时取消

电商场景:用户下单后 30 分钟没付款,自动取消订单。RabbitMQ 没有现成的延迟队列,用TTL + 死信队列(DLX)实现:消息进一个"没人消费的队列",TTL 到了变成死信,路由到真正处理取消逻辑的队列。

下单消息
TTL=30min
→
TTL 队列(不消费)
等 30 分钟变成死信
→
死信交换机 DLX
路由到取消队列
order-cancel-queue → 消费者:取消订单、释放库存

3.17 Kafka vs RabbitMQ 选型

Kafka vs RabbitMQ 终极对比
对比点KafkaRabbitMQ
吞吐量极高(几十万/秒)高(几万/秒)
延迟ms 级us 级(更低)
可靠性高(副本+持久化)极高(Confirm+ack+持久化)
路由弱(按 Partition 哈希)强(四种 Exchange 灵活路由)
协议自定义(二进制)AMQP(标准协议)
消息顺序Partition 内有序单 Queue 有序
适用场景日志采集、大数据管道、事件流订单通知、业务解耦、RPC 异步
记
第三章小结

① Kafka = 高铁,高吞吐、顺序写磁盘、Partition 并行。acks=all + 手动提交 offset = 不丢消息。

② RabbitMQ = 快递小哥,路由灵活(四种 Exchange)、可靠性高(Confirm + 手动 ack)。延迟队列用 TTL+DLX。

③ 选型:大数据/日志/事件流选 Kafka;业务消息/复杂路由/低延迟选 RabbitMQ。