Spring AMQP/RocketMQ 与消息可靠性

共 18 题
📑 题目列表 18 题
#
★★★

1. Spring AMQP 的 @RabbitListener 与 RabbitMQ Streams 在 Spring Boot 3.5+ 的集成

Spring AMQP 的 @RabbitListener 与 RabbitMQ Streams 在 Spring Boot 3.5+ 中如何集成?

  • @RabbitListener 注解消费
  • RabbitMQ Streams 协议
  • Stream 与 AMQP 的差异

@RabbitListener(queue = "...") 声明式消费 AMQP 队列,Spring 的 SimpleMessageListenerContainer 管理连接与 ack。RabbitMQ Streams 是 RabbitMQ 3.9+ 的流式协议(rabbit-stream),基于日志追加、支持重放与消费者组,与经典 AMQP 队列不同(AMQP 是消费即删除,Stream 是持久日志)。Spring Boot 3.5+ 集成:Spring AMQP 提供 @RabbitListener 支持 Stream 容器(StreamListenerContainer),通过 StreamListener 消费 Stream;配置 spring.rabbitmq.listener 或使用 StreamMessageListenerContainer。Stream 适合高吞吐、重放、事件溯源场景,AMQP 队列适合传统任务队列。

@RabbitListener 抽象统一 AMQP 与 Stream 消费,底层用不同容器。Stream 是日志型的持久队列,支持重放与多个消费者共享。

#
★★★

2. Spring Boot 3.5+ 中 @TransactionalEventListener 与本地消息表(Outbox Pattern)的最终一致

Spring Boot 3.5+ 中 @TransactionalEventListener 与本地消息表(Outbox Pattern)如何实现最终一致?

  • Outbox Pattern 本地消息表
  • 事务边界的原子写
  • 异步投递与最终一致

Outbox Pattern(本地消息表):业务事务内同时更新业务表与"消息表"(outbox 表),事务提交后两者原子一致;随后一个后台任务(轮询/监听)读取 outbox 表中的未发送消息,投递给 MQ 并标记已发送。由于业务操作与消息表写入在同一数据库事务,保证"业务成功"与"消息待发送"原子,避免"业务成功但消息丢失"。@TransactionalEventListener(AFTER_COMMIT) 可在事务提交后触发投递逻辑(监听 outbox 变化),保证投递在事务提交后发生。最终一致性:outbox 表记录待投递消息,投递成功标记;若投递失败可重试。相比"事务内直接发 MQ",Outbox 避免了两阶段不一致。

Outbox 用"本地 DB 事务写 outbox 表 + 异步投递"实现业务与消息的最终一致,@TransactionalEventListener 保证投递在事务提交后触发。

#
★★★

3. RocketMQ 的 transactionalMessage 与本地事务回滚

RocketMQ 的 transactionalMessage(事务消息)与本地事务回滚如何实现?

  • 半消息(prepare 消息)
  • 本地事务执行
  • 提交/回滚与回查

RocketMQ 事务消息流程:先发送"半消息"(prepare,对消费者不可见)到 broker 并返回 UNKNOWN;回调本地事务方法(如写 DB),本地事务成功返回 COMMIT(半消息变为可见),失败返回 ROLLBACK(删除半消息)。若本地事务执行中 broker 未收到 COMMIT/ROLLBACK(如崩溃),broker 通过回查(checkLocalTransactionState)再次询问本地事务结果决定提交/回滚。这样"本地事务"与"消息发送"通过半消息 + 回查实现最终一致:本地事务成功则消息必投递,本地事务回滚则消息不投递,避免"业务成功但消息丢失或相反"。

RocketMQ 事务消息用"半消息 + 本地事务 + 回查"实现 DB 与消息的最终一致。回查是保证"本地事务状态未知时"能最终确定提交/回滚的关键。

#
★★★

4. RocketMQ 的事务回查(checkLocalTransactionState)

RocketMQ 的事务回查(checkLocalTransactionState)机制如何工作?

  • 回查触发条件
  • 回查回调实现
  • 保证最终一致

事务回查发生在 broker 收到半消息但长时间未收到 COMMIT/ROLLBACK 时(如应用崩溃或本地事务响应丢失),broker 主动向生产者发起回查(checkLocalTransactionState),询问该事务的本地业务是否成功。生产者实现回查回调:查询本地事务(如查 DB 状态),根据结果返回 COMMIT/ROLLBACK,broker 据此决定半消息是否可见。回查避免"本地事务已提交但消息未确认"的中间态,保证最终一致。回查次数有限(transactionCheckInterval 与次数),超限后消息进入死信。实现要点:回查需能确定本地事务结果(如事务表记录状态)。

回查是事务消息"状态未知"的兜底机制,通过重新查询本地事务结果决定提交/回滚,保证最终一致。需本地事务可被回查确认。

#
★★★

5. 幂等性实现(数据库唯一索引/Redis SET/状态机)

消息幂等性如何用数据库唯一索引、Redis SET、状态机实现?

  • 幂等性概念
  • 唯一索引、Redis SETNX、状态机
  • 选型

消息幂等性保证重复消费不产生副作用。实现:一是数据库唯一索引——给业务表加唯一约束(如订单号),重复插入报错被忽略;二是 Redis SET(SETNX)——用 SETNX key value 标记已处理,已处理则跳过;三是状态机——业务状态流转(如 NEW->PROCESSING->DONE),重复消息在错误状态被拒绝。选型:强一致、需要数据库原子性用唯一索引;高并发、轻量去重用 Redis SETNX;有明确状态流转的业务用状态机。三者可结合(如唯一索引 + 状态机)。幂等需"去重键"(消息 id/业务唯一键)稳定。

幂等是"同一消息重复消费只生效一次",靠"去重键 + 存储层原子判断"。唯一索引/Redis SETNX/状态机分别对应 DB 原子、Redis 原子、业务状态约束。

#
★★★

6. 消息幂等性的实现方案,唯一主键约束、Redis 去重表、数据库唯一索引

消息幂等性的实现方案有哪些(唯一主键约束、Redis 去重表、数据库唯一索引)?

  • 唯一主键/唯一索引
  • Redis 去重表
  • 数据库唯一索引

消息幂等方案:一是数据库唯一主键/唯一索引——业务表用消息 id 或业务唯一键做主键,重复插入因唯一约束失败被忽略,DB 层保证原子;二是 Redis 去重表——用 SETNX key 或记录已处理的消息 id,已存在则跳过,Redis 原子保证高并发下去重;三是数据库唯一索引——与唯一主键类似,某字段加唯一索引,重复记录插入报错。方案对比:DB 唯一索引可靠、强一致,适合需要落库幂等;Redis 去重表快、适合高并发,但需注意 Redis 与 DB 的一致性(去重表过期)。生产常用"DB 唯一索引为主 + Redis 去重为辅"。

幂等核心是"存储层原子去重"。DB 唯一索引可靠,Redis 去重表高并发,组合使用兼顾可靠与性能。

#
★★★

7. 消息重复消费的根本原因与幂等设计

消息重复消费的根本原因是什么,幂等设计如何应对?

  • 重复消费原因(at-least-once)
  • 消费端 ack/重试
  • 幂等设计

消息重复消费的根本原因是"at-least-once"投递语义:MQ 在消费端未确认、超时重试、broker 端重试、网络抖动等情况下,会重复投递同一消息。具体场景:消费端处理成功但 ack 丢失导致重投;消费者处理中崩溃未提交偏移,重启重新消费;broker 重试发送。幂等设计应对:消费端保证"同一消息重复处理结果一致"——用消息/业务唯一键做去重(DB 唯一索引、Redis SETNX、状态机),或业务本身幂等(如状态幂等更新)。幂等设计需"去重键稳定 + 原子判断 + 失败可重试"。核心是接受"可能重复"并在消费端去重。

MQ 大多提供 at-least-once,重复消费是常态,因此消费端必须幂等。幂等 = 去重键 + 原子生效,是消息可靠性的最后一环。

#
★★★

8. RabbitMQ 的 Publisher Confirms 与事务

RabbitMQ 的 Publisher Confirms 与事务(channel.tx)有何区别?

  • Publisher Confirms 确认机制
  • channel 事务
  • 可靠性与性能

RabbitMQ 的 Publisher Confirms(confirmSelect)让 broker 在消息落到队列后向生产者发送确认(basicAck),生产者据此确认消息已持久化,支持异步(waitForConfirms/批量 confirm),性能好、是推荐做法。事务(channel.txSelect + txCommit)让消息发送与事务提交绑定,commit 时批量发送,实现了"消息发送的原子性",但事务会显著降低吞吐(每条/每批 commit 开销大),且事务模式下无法配合 confirm。选择:现代 RabbitMQ 推荐 Publisher Confirms(可靠性 + 性能),事务已过时。两者都用于"生产者确认消息不丢失",但 confirm 更高效。

Publisher Confirms 是异步确认,性能好且可靠;channel 事务吞吐低。生产用 confirm 而非事务,配合 mandatory 与交换机确认。

#
★★

9. RabbitMQ 的 Quorum Queue(仲裁队列)

RabbitMQ 的 Quorum Queue(仲裁队列)是什么?

  • Quorum Queue 的特性
  • 数据本地化与 Raft
  • 与经典队列对比

Quorum Queue 是 RabbitMQ 3.8+ 的队列类型,基于 Raft 共识算法在多副本间复制,保证高可用与数据一致性(多数派副本确认)。特性:数据复制到多个节点(通常 3 副本),支持故障转移(leader 故障自动选新 leader);消息被多数派接受才算成功,避免经典队列(镜像队列)的脑裂与数据不一致;支持 x-max-lengthx-message-ttl 等。与经典队列/镜像队列对比:Quorum Queue 更可靠、适合高可用场景,但相比经典队列吞吐略低、内存占用较高。适合对可靠性要求高的场景。

Quorum Queue 用 Raft 复制保证高可用与一致,替代镜像队列,是可靠消息队列的推荐选择。

#
★★

10. RabbitMQ 的四种 Exchange 类型与路由键匹配规则(direct/topic/fanout/headers)

RabbitMQ 的四种 Exchange 类型与路由键匹配规则是什么?

  • direct/topic/fanout/headers
  • 路由匹配规则
  • 应用场景

RabbitMQ 四种 Exchange:direct——按 routing key 精确匹配(队列绑定 key 与消息 routing key 完全一致才路由);topic——按模式匹配(* 匹配一个词、# 匹配多个词),支持通配符路由;fanout——广播给所有绑定队列,忽略 routing key;headers——按消息 header 匹配(非 routing key),最灵活但少用。选型:精确路由用 direct,通配/分类路由用 topic,广播用 fanout,需按 header 过滤用 headers。四种类型决定消息路由到哪些队列。

Exchange + 绑定决定路由。direct 精确、topic 通配、fanout 广播、headers 按头,按路由需求选择。

#
★★

11. RabbitMQ 的死信队列(DLX/DLQ)在消息重试与失败处理场景的应用

RabbitMQ 的死信队列(DLX/DLQ)在消息重试与失败处理场景如何应用?

  • 死信队列(DLX/DLQ)
  • 死信触发条件
  • 重试与失败处理

死信队列(DLX/DLQ)接收无法被正常消费的消息。死信触发条件:消息被拒绝(basicNack/basicReject 且 requeue=false)、消息过期(TTL)、队列超过长度。应用:把失败消息路由到死信队列,还有专门的死信交换(DLX)把消息转投到死信队列,用于重试与失败处理——常见做法是"延迟重试队列 + 死信队列":消费失败将消息投到延时队列(TTL),到期重新进入原队列重试,重试次数用尽后进入死信队列,人工处理或告警。这样避免消息无限重试挤占队列,并保留失败消息供排查。

DLX/DLQ 是"失败消息的收容所",配合延时队列实现重试,配合人工/告警处理最终失败,是 RabbitMQ 可靠消费的重要组成部分。

#
★★

12. RabbitMQ 的消息确认(ack/nack/requeue)

RabbitMQ 的消息确认(ack/nack/requeue)如何工作?

  • basicAck 确认
  • basicNack 否定确认
  • requeue 重新入队

RabbitMQ 消费端确认:basicAck(确认成功,消息从队列删除);basicNack/basicReject(否定确认,可配合 requeue 参数)。requeue=true 表示消息重新入队(立即重试,可能死循环);requeue=false 表示不重新入队,消息进入死信队列(或丢弃)。配合 prefetch(预取数)与手动 ack(channel.basicAck)实现可靠消费。注意:ack 前消费者崩溃,消息会重新投递(at-least-once);nack+requeue=true 可能造成同一消息反复投递(需重试次数控制)。Spring AMQP 用 channel.basicAck/ChannelAwareMessageListener@RabbitListenerAckMode

ack 确认成功、nack 否定确认、requeue 决定是否重新入队。可靠消费需手动 ack + 合理 requeue(配合死信防死循环)。

#
★★

13. RocketMQ 5.x Pop 消费的消息粒度分配与消费者心跳/可见性超时(visibility timeout)的协同

RocketMQ 5.x Pop 消费的消息粒度分配与消费者心跳/可见性超时(visibility timeout)如何协同?

  • Pop 消费模型
  • 可见性超时(visibility timeout)
  • 消息粒度与心跳

RocketMQ 5.x 的 Pop 消费(CONSUMER_POP)改变传统拉取模型:消费者用 Pop 拉取消息,消息从队列"借出"并被赋予"可见性超时"(visibility timeout),在超时内可采用/确认;若超时未确认,消息重新可见(可被其他消费者消费),实现"先取后保"的消费语义。消息粒度分配:Pop 消费按消息粒度分配(而非按队列整个分配),支持更灵活的负载均衡;消费者通过心跳(heartbeat)向 broker 报告存活与消费进度,broker 根据可见性超时与心跳管理消息归属。协同:消费者心跳维持会话,可见性超时控制消息"借出-释放",超时未 ack 的消息重新投递,实现 at-least-once。

Pop 消费用"可见性超时"替代传统"取走即删除",超时未确认消息重新可见,配合心跳保证消费者存活与消息重投。

#
★★

14. RocketMQ 的 Topic/Queue/Message 模型

RocketMQ 的 Topic/Queue/Message 模型是怎样的?

  • Topic 逻辑分类
  • Queue 分区与并行
  • Message 结构

RocketMQ 以 Topic 组织消息分类,每个 Topic 划分为多个 Queue(分区),消息按写入顺序分布到 Queue,Queue 是并行与消费的最小单位(一个 Queue 响应一个消费者,保证 Queue 内有序)。Message 包含主题、标签(Tag)、key、body、属性等。生产者在写时按策略(key hash/轮询/指定)选择 Queue;消费者组内按 Queue 分配,一个 Queue 同一时刻被组内一个消费者消费(保证 Queue 内有序)。并行度 = Topic 的 Queue 数。与 Kafka 的 Partition 类似,Queue 是 RocketMQ 的并行与顺序单元。

Topic-Queue-Message 三层模型,Queue 是并行与顺序单元,支持顺序消费(MessageListenerOrderly)与高吞吐。

#
★★

15. RocketMQ 集群消费与广播消费的差异及顺序消费(MessageListenerOrderly)的锁语义

RocketMQ 集群消费与广播消费的差异,以及顺序消费(MessageListenerOrderly)的锁语义是什么?

  • 集群 vs 广播消费
  • 顺序消费的锁
  • 并行度与顺序

集群消费:消息被组内一个消费者消费(分摊),保证一条消息只被处理一次;广播消费:消息被组内每个消费者都消费(各自独立),适合"每个实例都需要"的场景。顺序消费(MessageListenerOrderly)用"队列锁"实现:消费者对某个 Queue 加锁,串行消费该 Queue 的消息,保证 Queue 内顺序;全局消费(MessageListenerConcurrently)并行处理,无全局顺序。顺序消费的锁语义:同一 Queue 的锁保证同一时刻只有一个线程消费该 Queue,维持分区内顺序;锁粒度是"队列",不同 Queue 可并行。顺序消费牺牲并行度换取顺序。

集群消费分摊、广播消费全量。顺序消费用队列锁串行消费 Queue 内消息保序,全局消费并行但无顺序。

#
★★

16. 消息丢失的三大场景(生产者、Broker、消费者)

消息丢失的三大场景(生产者、Broker、消费者)是什么?

  • 生产者丢失
  • Broker 丢失
  • 消费者丢失

消息丢失三大场景:一是生产者丢失——发送失败未重试、未用确认机制(如 Kafka acks=0、RabbitMQ 无 confirm),消息未达 broker;二是 Broker 丢失——消息到 broker 但未持久化(内存中),宕机丢失;副本不足(单副本)故障丢失;三是消费者丢失——消费后未确认/未提交偏移,崩溃后消息丢失(或 offset 提交过早导致消息未处理)。防护:生产者端用确认(acks=all/confirm)+ 重试;Broker 端持久化(刷盘、多副本、同步复制);消费者端处理成功后再 ack/提交,幂等。消息可靠 = 三端都保证。

消息可靠性要覆盖发送、存储、消费三端。三端各自用"确认 + 持久化 + 处理完再 ack"层层保证不丢。

#
★★

17. 消息可靠投递(生产端确认 + 消费端 ACK)

消息可靠投递如何实现(生产端确认 + 消费端 ACK)?

  • 生产端确认(confirm/acks)
  • 消费端 ACK
  • 可靠投递链路

消息可靠投递 = 生产端确认 + 消费端 ACK。生产端:Kafka 用 acks=all + 幂等/事务,RabbitMQ 用 Publisher Confirms,RocketMQ 用发送确认/重试,保证消息到达 broker 且持久化。消费端:消费成功后再确认(ack/提交 offset),失败不确认以触发重投,配合死信/重试处理。链路:生产确认(发送成功)-> broker 持久化(多副本)-> 消费确认(处理成功才 ack)。任一环节缺失都会导致丢失。配合重试与幂等,实现"不丢不重"(至少一次 + 幂等)。可靠投递是端到端保证,需三端配合。

可靠投递是"发送确认 + 存储持久 + 消费确认"的端到端链路,配合 at-least-once + 幂等实现不丢不重。

#

18. RabbitMQ Channel 的线程安全边界与连接复用(Channel 数限制)

RabbitMQ Channel 的线程安全边界与连接复用(Channel 数限制)是什么?

  • Channel 非线程安全
  • 连接复用与 Channel 数
  • 多线程使用

RabbitMQ 中一个 Connection 可复用多个 Channel(虚拟通道),但 Channel 不是线程安全的——同一 Channel 不能同时被多个线程并发使用(命令会交错导致协议错误)。用法:每个线程/每个任务用独立 Channel(从连接创建),或用连接池(如 CachingConnectionFactory)统一管理 Channel 的获取与归还。Channel 数限制:broker 有最大 Channel 数(channel_max,默认 2047),连接创建的 Channel 数受此限制,过多会报错。Spring AMQP 的 CachingConnectionFactory 缓存 Channel 复用,避免频繁创建。实践:一个连接 + 多个 Channel,每个 Channel 单线程使用(或池化归还)。

Connection 复用、Channel 非线程安全、Channel 数有上限。用连接池管理 Channel,保证每个 Channel 单线程访问。