
这部分内容摘自 JavaGuide 下面几篇文章的重点:
消息队列基础常见问题总结
Kafka 常见面试题总结
RocketMQ 常见面试题总结
RabbitMQ 常见面试题总结
Disruptor 常见面试题总结
回答消息队列问题时,建议把一条消息的生命周期串起来:生产者为什么发送消息,Broker 如何接收、存储和复制,消费者在什么时机确认,失败后如何重试,最终怎样发现和修复异常数据。只回答某个配置项,很容易在继续追问时断掉。
消息队列基础
什么是消息队列?
消息队列(Message Queue,MQ)是一类用于存储和转发消息的中间件。生产者把消息发送给 Broker,Broker 负责存储和投递,消费者从 Broker 获取消息并完成业务处理。
生产者和消费者不需要同时在线,也不需要相互知道对方的地址。生产者只关心消息能否可靠进入 MQ,消费者可以按照自己的处理能力异步消费。这种通信方式能够解耦服务,但也把原来的同步调用问题变成了消息可靠性、重复消费、顺序和积压问题。
队列通常具有先进先出的特点,但分布式 MQ 一般只能保证某个队列或分区内有序,不能直接理解为整个 Topic 全局有序。
⭐️消息队列有什么用?
消息队列常见作用如下:
异步处理:把不需要同步完成的动作放到后台执行,缩短主链路响应时间。
削峰填谷:高峰流量先写入 MQ,消费者按照自身处理能力消费。
系统解耦:生产者只负责发布事件,不需要直接依赖所有消费者。
顺序处理:让同一业务对象的事件按顺序进入同一队列或分区。
延时处理:实现订单超时关闭、通知提醒等延时任务。
最终一致性:通过事务消息、本地消息表和补偿任务协调跨服务状态。


以订单创建为例,订单入库和库存扣减通常属于核心流程;短信、积分、推荐和部分风控动作可以异步处理。拆分前要先确认哪些结果必须立即返回,哪些动作允许最终一致。
使用消息队列会带来哪些问题?
引入 MQ 后,系统会增加下面这些问题:
可用性依赖:MQ 故障可能阻塞生产或消费链路。
消息可靠性:生产、存储和消费阶段都可能丢失消息。
重复消费:确认丢失、消费超时和重平衡都可能让消息再次投递。
顺序变化:多分区、多消费者和重试会影响消息顺序。
消息积压:消费者处理能力低于生产速度时,延迟会持续扩大。
数据一致性:异步链路通常只能提供最终一致,需要对账和补偿。
排查难度:问题定位需要同时检查生产者、Broker、消费者、业务日志和 Trace。

调用方必须立即拿到下游处理结果时,同步 RPC 通常更直接;下游动作可以延后,并且业务接受最终一致时,消息队列更合适。
RPC 和消息队列有什么区别?
对比维度RPC消息队列通信方式调用方直接请求服务提供方生产者和消费者通过 Broker 通信返回结果通常需要在调用链中返回结果生产者通常不等待消费者完成时效性更适合需要立即结果的请求更适合允许延后处理的任务耦合关系调用方知道目标服务和接口生产者可以不知道具体消费者典型用途查询、校验、需要同步结果的业务异步处理、削峰、事件通知和最终一致性
对比维度RPC消息队列
对比维度
RPC
消息队列
通信方式调用方直接请求服务提供方生产者和消费者通过 Broker 通信
通信方式
调用方直接请求服务提供方
生产者和消费者通过 Broker 通信
返回结果通常需要在调用链中返回结果生产者通常不等待消费者完成
返回结果
通常需要在调用链中返回结果
生产者通常不等待消费者完成
时效性更适合需要立即结果的请求更适合允许延后处理的任务
时效性
更适合需要立即结果的请求
更适合允许延后处理的任务
耦合关系调用方知道目标服务和接口生产者可以不知道具体消费者
耦合关系
调用方知道目标服务和接口
生产者可以不知道具体消费者
典型用途查询、校验、需要同步结果的业务异步处理、削峰、事件通知和最终一致性
典型用途
查询、校验、需要同步结果的业务
异步处理、削峰、事件通知和最终一致性
两者可以同时出现在一条业务链路中。例如,创建订单时同步调用库存服务校验库存,订单提交后再通过 MQ 通知积分和消息服务。
消息可靠性与消费语义
⭐️消息队列常见的投递语义有哪些?
常见的消息投递语义有三种:
最多一次(At Most Once):消息最多投递一次,可能丢失,但不会重复。通常是先确认、后处理,或者发送失败后不重试。
至少一次(At Least Once):消息至少投递一次,不允许轻易丢失,但可能重复。生产者和消费者都可以重试,消费端必须实现幂等。这是业务系统中最常见的选择。
恰好一次(Exactly Once):从定义上看,每条消息的业务效果只发生一次。实际产品通常只在特定范围内提供这种能力,例如 Kafka 幂等生产者避免同一个生产者会话内因重试产生重复记录,Kafka 事务可以原子提交一批 Kafka 记录和消费位点。
不能只看到产品文档中的 Exactly Once 就认为整个业务链路天然只执行一次。消息处理只要涉及 MySQL、Redis、第三方支付等外部系统,就还需要唯一索引、状态机、幂等键、Outbox 或业务事务来保证最终效果。
⭐️如何保证消息不丢失?
需要分别处理生产者、Broker 和消费者三个阶段:
生产者阶段:开启发送确认;发送失败时按业务策略重试;可靠性要求高时,用本地消息表或 Outbox 记录待发送事件。
Broker 阶段:开启消息持久化和多副本复制,并根据业务要求设置确认副本数、刷盘策略和故障转移方案。
消费者阶段:业务处理成功后再提交 offset 或发送 ACK;处理失败时进入重试,超过阈值后进入死信队列或补偿流程。
消息可靠性不能只靠某一个配置。例如,Broker 已经持久化消息,但消费者先确认后处理,消费者宕机时仍然可能丢失业务结果;消费者改成处理后确认,又必须接受重复投递并实现幂等。
⭐️为什么消息队列会出现重复消费?如何处理?
常见原因包括:
消费者完成业务处理后,还没来得及确认就宕机。
ACK 或 offset 提交成功,但响应在网络中丢失。
消费超时、消费者重启或集群重平衡后,消息被重新分配。
生产者发送超时后重试,但前一次发送其实已经成功。
生产环境更常见的是 至少一次投递 + 消费端幂等。常见幂等手段如下:
使用订单号、支付流水号等业务唯一键去重。
通过数据库唯一索引拦截重复写入。
使用状态机限制状态只能沿合法方向变化。
建立消费记录表,保存消息 ID 和处理结果。
对允许短期去重的消息使用 Redis,并设置合理过期时间。
不要只依赖 MQ 的消息 ID。业务重试可能生成新的消息 ID,但业务动作仍然是同一次;支付、入账等场景通常需要业务唯一键。
⭐️如何保证消息顺序消费?
常见做法如下:
根据订单 ID、用户 ID 等业务键,把相关消息发送到同一个队列或分区。
同一个队列或分区由一个消费线程顺序处理,或者消费者内部按业务键串行化。
消费失败时不能简单跳过,需要重试、阻塞后续同一业务键消息,或者使用版本号和状态机拒绝乱序更新。

顺序范围需要先说清楚。大多数系统只要求同一订单、同一账户或同一设备的消息有序,没有必要让整个 Topic 全局有序。顺序范围越大,并行能力越低。
消息积压怎么办?
先判断积压发生在哪个环节:
生产流量是否突然增加。
消费者是否出现慢 SQL、外部接口超时、锁竞争或频繁 GC。
消费者实例是否异常退出。
Broker 的磁盘、网络或分区是否出现瓶颈。
分区数、队列数和消费者并发是否匹配。
确认原因后,可以临时扩容消费者、提高批量消费能力、修复慢逻辑,或者把积压消息转存后使用临时消费者处理。增加消费者之前要检查分区数或队列数;如果只有少量分区,新增的消费者无法获得分区,也不会提升消费速度。
积压处理完成后,还要补充监控:消息堆积量、最老消息年龄、生产消费速率差、消费失败率和死信数量,比只监控“队列里有多少条消息”更有用。
消费失败后应该如何重试?
重试应区分异常类型:
网络抖动、临时超时等瞬态错误,可以使用指数退避和随机抖动重试。
参数错误、业务状态不允许、权限失败等永久错误,继续重试通常没有意义。
下游已经持续异常时,应限制重试速率,避免重试流量进一步压垮下游。
超过最大重试次数的消息可以进入死信队列,随后告警、人工处理或由补偿任务再次投递。消费端仍然要保证幂等,因为“业务处理成功但 ACK 失败”时,重试消息会再次到达。
什么是死信队列?什么时候需要死信队列?
死信队列(Dead Letter Queue,DLQ)用于接收无法继续正常消费的消息。常见来源包括:
消息达到最大重试次数。
消费者明确拒绝消息,并且不再重新入队。
消息超过 TTL。
队列达到长度限制,消息被丢弃或转移。
死信队列的作用是隔离异常消息,避免一条永久失败的消息无限重试、持续占用消费资源。进入死信队列后,还需要告警、查询原始消息和失败原因、人工修复或受控重放。
死信队列不是失败消息的终点。没有监控和处理流程的 DLQ,只是把问题从主队列挪到了另一个队列。
Kafka
Kafka 的核心概念有哪些?
Producer:生产并发送消息。
Broker:Kafka 服务节点。
Topic:消息的逻辑分类。
Partition:Topic 的物理分片,也是并行处理和局部顺序的基本单位。
Replica:分区副本,用于故障恢复。
Consumer Group:一组协作消费消息的消费者;同一分区在一个消费组内同一时刻只分配给一个消费者。
Offset:消费者在分区中的消费位置。
Topic 增加分区可以提高并行能力,但会增加副本、文件、选主和重平衡成本。分区也不是越多越好,应根据吞吐、消费者并发和运维能力规划。

⭐️Kafka 的多副本机制和 ISR 是什么?
Kafka 的每个 Partition 可以配置多个副本,其中一个是 Leader,其余是 Follower。生产者和消费者主要与 Leader 交互,Follower 从 Leader 复制日志。Leader 故障后,Kafka 会从符合条件的 Follower 中选出新 Leader。
ISR(In-Sync Replicas)保存与 Leader 保持同步的副本集合。Follower 长时间没有跟上 Leader 时会被移出 ISR,恢复同步后可以重新加入。高可靠写入通常要同时考虑:
replication.factor:一个分区有多少个副本。
replication.factor
acks=all:生产者等待 Leader 按当前 ISR 状态完成确认。
acks=all
min.insync.replicas:ISR 至少保留多少个副本时才接受高可靠写入。
min.insync.replicas
unclean.leader.election.enable:是否允许落后副本成为 Leader。开启后可用性可能提高,但存在丢失已确认消息的风险。
unclean.leader.election.enable
复制进度中还经常出现两个概念:LEO(Log End Offset)表示副本下一条待写记录的位置;HW(High Watermark)表示已经被同步确认、普通消费者可以安全读取到的位置。Follower 复制推进后,HW 才会继续向前。
多副本提高了容灾能力,但也增加磁盘、网络复制和故障恢复成本。副本数、ISR 和确认策略必须一起回答,只说 acks=all 不能完整说明可靠性。
acks=all
Kafka 现在还依赖 ZooKeeper 吗?
Kafka 早期使用 ZooKeeper 保存集群元数据并辅助 Controller 选举。Kafka 2.8 引入 KRaft,使用 Kafka 自己实现的 Raft 元数据仲裁管理控制面;Kafka 3.3 开始将 KRaft 标记为适合新集群生产使用。
Kafka 4.0 已经移除 ZooKeeper 模式,只支持 KRaft。 新集群不再需要部署 ZooKeeper,Broker 和 Controller 可以分角色部署,也可以在开发环境中合并角色。老集群升级时要按官方迁移流程转换元数据,不能直接删除 ZooKeeper 后重启。
⭐️Kafka 如何保证消息有序?
Kafka 只能保证 单个 Partition 内的日志顺序,不保证一个 Topic 的所有 Partition 全局有序。常见做法如下:
生产者使用订单 ID、用户 ID 等业务键选择 Partition,确保同一业务对象的消息进入同一 Partition。
同一 Consumer Group 内,一个 Partition 同一时刻只交给一个消费者实例。
消费者处理时避免把同一 Partition 的消息无约束地交给多个线程并发执行。
发生失败重试时,要避免后发批次先于失败批次写入。Kafka 生产者启用幂等后会结合序列号处理重试顺序;如果关闭幂等并允许多个未确认请求并发,重试可能造成乱序。
如果业务处理需要并发,可以按业务键分发到不同的本地串行队列,但 ACK、进程故障和内存积压都要自行处理。数据库侧最好再加版本号或状态机,防止重放和异常重试把状态改回旧值。
⭐️Kafka 的 Exactly Once 是如何实现的?
⭐️Kafka 的 Exactly Once 是如何实现的?
Kafka 的 Exactly Once 主要依赖 幂等生产者和事务:
幂等生产者使用 Producer ID、Epoch 和序列号识别重试批次,避免网络超时重试在 Kafka 日志中写入重复记录。Kafka 3.0 起,在没有冲突配置时默认启用 enable.idempotence。
enable.idempotence
事务生产者配置 transactional.id,可以把写入多个 Partition 的消息以及消费位点作为一个事务提交或回滚。
transactional.id
下游消费者配置 isolation.level=read_committed,只读取已经提交的事务消息,跳过被中止的记录。
isolation.level=read_committed
这套机制适合“从 Kafka 读取、处理、再写回 Kafka”的链路。事务范围不自动包含外部 MySQL 或第三方接口:如果消费逻辑还要写数据库,仍然需要数据库事务、幂等约束、Outbox 或 CDC 等方案。
Kafka 为什么吞吐量高?
Kafka 的高吞吐来自多个设计共同作用:
使用追加写和顺序 I/O,减少随机磁盘访问。
利用操作系统 Page Cache 缓存热点数据。
通过零拷贝能力减少内核态和用户态之间的数据复制。
使用分区实现并行读写和水平扩展。
Producer 批量发送并压缩消息,降低网络往返和协议开销。
Consumer 按批拉取消息,提高处理效率。
这些优化有适用条件。批次越大,吞吐通常越高,但消息等待成批的时间也可能增加;分区越多,并行度更高,集群元数据和重平衡成本也会上升。
⭐️Kafka 如何降低消息丢失风险?
生产者、Broker 和消费者都要配置:
Producer 使用 acks=all,配置合理的重试,并启用幂等生产者。
acks=all
Topic 设置多个副本。
Broker 配置合适的 min.insync.replicas,避免 ISR 过少时仍然接受高可靠写入。
min.insync.replicas
禁止把不在 ISR 中的副本直接选为 Leader,除非业务明确接受数据丢失风险。
Consumer 在业务处理成功后再提交 offset。

acks=all 并不表示所有副本都确认,而是由 Leader 等待当前 ISR 中满足要求的副本确认。副本数、ISR、min.insync.replicas 和集群故障情况要结合起来分析。
acks=all
min.insync.replicas
Kafka 重平衡有什么影响?如何减少影响?
消费组成员、订阅 Topic 或分区数发生变化时,Kafka 需要重新分配分区,这个过程叫 Rebalance。重平衡期间,部分分区会暂时停止消费;如果发生得过于频繁,会增加消费延迟,还可能带来重复处理。
常见优化方向包括:
避免消费者频繁重启或长时间卡住消费线程。
合理设置 session.timeout.ms、heartbeat.interval.ms 和 max.poll.interval.ms。
session.timeout.ms
heartbeat.interval.ms
max.poll.interval.ms
控制单批消息的处理时长,耗时任务可以与拉取线程解耦,但要自行管理 offset 和背压。
使用静态成员资格减少短暂重启引发的重平衡。
在客户端和 Broker 版本支持时评估协作式重平衡,减少一次性撤销全部分区的影响。
参数不能单独照抄。max.poll.interval.ms 设置得很大虽然能减少误判,也会让真正失效的消费者更晚被替换。
max.poll.interval.ms
RocketMQ
RocketMQ 的核心组件有哪些?
RocketMQ 的核心组件包括:
NameServer:保存 Topic 路由信息。Broker 定期注册和发送心跳,Producer、Consumer 从 NameServer 获取路由。
Broker:接收、存储和投递消息,维护 CommitLog、ConsumeQueue、消费进度等数据。
Producer:从 NameServer 获取路由并选择 MessageQueue 发送消息。
Consumer:订阅 Topic,按照集群消费或广播消费模式处理消息。
Proxy:RocketMQ 5.x 增加的代理层,可以为 gRPC 等访问方式提供统一入口;传统客户端也可以直接访问 Broker。
NameServer 不存储业务消息,也不参与每次消息转发。多个 NameServer 节点彼此独立,Broker 向所有 NameServer 注册;客户端可以连接其中任意可用节点获取路由。
RocketMQ 的事务消息解决什么问题?
RocketMQ 事务消息用于协调 本地事务和消息发送,使两者尽量保持一致。典型流程如下:
生产者先发送半消息,消费者暂时不可见。
Broker 确认半消息后,生产者执行本地事务。
本地事务成功则提交消息,失败则回滚消息。
Broker 长时间没有得到明确结果时,会回查生产者的本地事务状态。
事务消息适合“本地事务成功后必须可靠触发下游动作”的场景。它保证的是本地事务与消息发送之间的协调,消费者仍然可能重复收到消息,因此消费端幂等、失败重试和业务补偿依然需要保留。
RocketMQ 如何实现延时消息?
RocketMQ 支持延时消息。生产者指定延时时间后,消息不会立即投递给消费者,适合订单超时关闭、预约提醒等场景。
延时消息不能替代完整的业务状态机。订单取消消费者处理消息时,需要再次检查订单状态,只允许“待支付”订单转为“已取消”;订单已经支付时应直接忽略延时消息。这样才能处理消息延迟、重复投递和支付取消并发。
⭐️RocketMQ 如何保证顺序消息?
RocketMQ 通常保证同一业务分组内有序,而不是让整个 Topic 全局有序:
生产者按照订单 ID、用户 ID 等业务键选择同一个 MessageQueue。RocketMQ 5.x 的 FIFO 消息使用 MessageGroup 表达这一分组。
同一业务分组的消息需要按顺序发送;多个线程并发发送时,Broker 无法推断它们原本的业务顺序。
消费者按照队列存储顺序串行处理,并在当前消息处理完成后再处理后续消息。
某条顺序消息失败时,后续同组消息通常要等待它重试完成,否则会破坏顺序。
顺序保证会降低并行能力。业务键不能过度集中,例如所有订单都使用同一个分组,否则这个分组对应的队列会成为热点。消费端仍然建议使用状态机或版本号防止重放和异常路径导致的乱序更新。
RocketMQ 为什么读写性能较高?
RocketMQ 的消息主体统一顺序写入 CommitLog,避免为每个 Topic 随机写不同文件。消费者不会扫描整个 CommitLog,而是先读取对应队列的 ConsumeQueue 索引,再根据物理偏移量定位消息;按 Key 或时间查询则可以使用 IndexFile。
这套结构把“大文件顺序写”和“按队列读取”拆开了:
CommitLog 负责保存完整消息和元数据。
ConsumeQueue 保存 CommitLog 偏移量、消息大小和 Tag Hash 等固定长度索引。
IndexFile 提供按消息 Key 和时间范围查询的能力。
RocketMQ 还会利用 Page Cache、内存映射和批量处理降低磁盘与网络开销。不过,顺序写不等于数据已经安全落盘,最终可靠性还取决于刷盘和副本复制策略。
RocketMQ 同步刷盘、异步刷盘有什么区别?
同步刷盘:Broker 把消息写入磁盘后再向生产者返回成功,消息持久性更强,但写入延迟和磁盘压力更大。
异步刷盘:消息写入 Page Cache 后即可返回,由后台线程批量刷盘,吞吐和延迟更好,但机器掉电时可能丢失尚未落盘的数据。
副本复制也有同步和异步之分。同步复制需要从节点确认后再返回,能够降低主节点故障造成的数据丢失风险;异步复制由主节点先返回,再把消息复制给从节点,性能更好,但故障切换时风险更高。
刷盘和复制是两个维度。同步刷盘只说明本机磁盘已经写入,不代表副本已经同步;同步复制也不能替代消费端 ACK、幂等和补偿。
RabbitMQ
RabbitMQ 的 Exchange、Queue 和 Routing Key 分别是什么?
RabbitMQ 的 Exchange、Queue 和 Routing Key 分别是什么?
RabbitMQ 基于 AMQP 路由模型:
Producer:把消息发送到 Exchange。
Exchange:根据类型、Routing Key 和 Binding 把消息路由到一个或多个 Queue。
Queue:存储等待消费的消息。
Routing Key:生产者发送消息时携带的路由键。
Binding Key:Exchange 和 Queue 建立绑定时使用的匹配规则。
Consumer:从 Queue 获取并处理消息。

常见 Exchange 类型如下:
类型路由方式directRouting Key 与 Binding Key 精确匹配fanout忽略 Routing Key,发送到所有绑定队列topic使用 * 和 # 对分段 Routing Key 做模式匹配headers根据消息 Header 匹配
类型路由方式
类型
路由方式
directRouting Key 与 Binding Key 精确匹配
direct
Routing Key 与 Binding Key 精确匹配
fanout忽略 Routing Key,发送到所有绑定队列
fanout
忽略 Routing Key,发送到所有绑定队列
topic使用 * 和 # 对分段 Routing Key 做模式匹配
topic
使用 * 和 # 对分段 Routing Key 做模式匹配
*
#
headers根据消息 Header 匹配
headers
根据消息 Header 匹配

RabbitMQ 中的死信是怎么产生的?
RabbitMQ 中的消息出现下面几种情况时,可以被发送到 Dead Letter Exchange(DLX),再由 DLX 路由到死信队列:
消费者执行 basic.reject 或 basic.nack,并设置 requeue=false。
basic.reject
basic.nack
requeue=false
消息 TTL 到期。
队列超过长度限制,消息按溢出策略被丢弃。
Quorum Queue 中消息重复投递次数超过 delivery-limit。
delivery-limit
DLX 本身仍然是普通 Exchange,需要提前声明并绑定目标队列。生产环境更推荐用 Policy 配置 DLX 和 Routing Key,后续调整时不必删除并重新声明队列。
RabbitMQ 如何实现延迟消息?
常见方案有两种:
TTL + DLX:消息先进入设置了 TTL 的队列,过期后通过死信交换器路由到真正的消费队列。该方案不需要额外插件,但按消息设置不同 TTL 时可能出现队头阻塞:后面的短延时消息要等前面的长延时消息离开队头后才能被处理。
延迟消息插件:使用 rabbitmq-delayed-message-exchange 提供的 x-delayed-message Exchange,根据消息头设置延迟时间。使用前要确认插件版本兼容性、延迟规模和故障恢复要求。
rabbitmq-delayed-message-exchange
x-delayed-message
订单超时关闭等场景无论使用哪种方案,消费者都要重新检查订单状态。消息可能延迟、重复或与支付请求并发到达,不能收到延时消息就直接取消订单。
⭐️RabbitMQ 如何保证消息可靠性?
可以沿着消息链路回答:
生产者到 Broker:使用 Publisher Confirms 确认 Broker 是否接收消息;通过 mandatory 和 Return Listener 发现消息到达 Exchange 后没有路由到队列的情况。
mandatory
Broker 存储:声明持久化 Exchange 和 Queue,消息使用持久化投递;需要副本时,根据业务选择 Quorum Queue 或 Streams。
Broker 到消费者:消费者使用手动 ACK,业务处理成功后再确认;失败消息按异常类型决定重试、拒绝或进入死信队列。
业务处理:消费端实现幂等,并通过监控和补偿任务处理异常消息。

Publisher Confirm 只说明 Broker 处理了这次发布,不一定说明消息成功进入业务队列。路由失败还要通过 mandatory 返回或备用交换器处理。
mandatory
RabbitMQ 如何保证高可用?
RabbitMQ 集群中的节点可以共享拓扑元数据和客户端连接能力,但普通集群不等于队列消息已经复制。消息是否有副本取决于队列类型:
Classic Queue:RabbitMQ 4.x 中是非复制队列,适合不要求副本保护、临时或高频创建删除的队列。
Quorum Queue:基于 Raft 复制日志,需要多数副本可用,适合订单、支付、库存等强调数据安全的长期队列。
Stream:适合保留、回放和大吞吐场景,消费模型与普通队列不同。
部署 Quorum Queue 时通常使用奇数副本,并把副本分散到不同故障域。它能处理少数节点故障,但多数副本不可用时会停止提供服务;副本越多,磁盘和网络成本也越高。
客户端还要配置多个节点地址、自动恢复和合理的连接重试。只部署多个 RabbitMQ 节点,却让所有连接只指向一个地址,仍然会留下接入单点。
RabbitMQ 生产环境要监控哪些指标?
至少要关注:
消息状态:Ready、Unacked、消息进入和确认速率、最老消息年龄。
消费者状态:消费者数量、消费速率、ACK/NACK、重投递和 Prefetch 是否合理。
节点资源:内存水位、磁盘剩余空间、磁盘 I/O、文件句柄和 Erlang 进程数。
连接资源:Connection、Channel 数量以及频繁创建销毁速率。
可靠性异常:Publisher Nack、无法路由消息、死信数量和 Quorum Queue 副本状态。
只监控队列长度不够。队列长度不变,可能是生产和消费同时停了;Unacked 持续增长通常说明消费者拿到消息后处理过慢或没有及时确认。告警阈值应结合业务吞吐和允许延迟设置,不能照搬固定数字。
消息队列选型
⭐️Kafka、RocketMQ 和 RabbitMQ 如何选择?
⭐️Kafka、RocketMQ 和 RabbitMQ 如何选择?
可以先比较业务语义,再看团队运维能力:
场景常见选择重点关注日志、埋点、流式处理、大吞吐Kafka分区扩展、消息保留与回放、流处理生态订单、交易、业务事件RocketMQ事务消息、延时消息、顺序消息和 Java 生态路由规则复杂、业务队列RabbitMQExchange 路由、确认机制、队列类型和消费模型
场景常见选择重点关注
场景
常见选择
重点关注
日志、埋点、流式处理、大吞吐Kafka分区扩展、消息保留与回放、流处理生态
日志、埋点、流式处理、大吞吐
Kafka
分区扩展、消息保留与回放、流处理生态
订单、交易、业务事件RocketMQ事务消息、延时消息、顺序消息和 Java 生态
订单、交易、业务事件
RocketMQ
事务消息、延时消息、顺序消息和 Java 生态
路由规则复杂、业务队列RabbitMQExchange 路由、确认机制、队列类型和消费模型
路由规则复杂、业务队列
RabbitMQ
Exchange 路由、确认机制、队列类型和消费模型
选型时还要确认:
是否需要消息回放和较长时间保留。
是否需要事务消息、延时消息或灵活路由。
峰值吞吐、单条消息大小和可接受延迟。
顺序范围、消费语义和容灾目标。
团队是否具备对应产品的部署、升级、监控和故障处理经验。
技术特性只是选型的一部分。公司已经稳定运行某种 MQ,并且现有功能满足需求时,继续使用通常比为了某个局部特性引入新集群更稳妥。
Disruptor 和分布式消息队列有什么区别?
Disruptor 是 JVM 进程内的高性能事件处理框架,Kafka、RocketMQ 和 RabbitMQ 是跨进程、跨机器的分布式消息中间件。
对比维度Disruptor分布式消息队列使用范围单个 JVM 内部跨服务、跨进程、跨机器数据存储内存中的 RingBuffer通常支持持久化和多副本主要目标低延迟、减少锁竞争和对象分配可靠投递、削峰、解耦和消息回放故障恢复需要业务自行处理进程故障由 Broker、副本和消费位点提供支持
对比维度Disruptor分布式消息队列
对比维度
Disruptor
分布式消息队列
使用范围单个 JVM 内部跨服务、跨进程、跨机器
使用范围
单个 JVM 内部
跨服务、跨进程、跨机器
数据存储内存中的 RingBuffer通常支持持久化和多副本
数据存储
内存中的 RingBuffer
通常支持持久化和多副本
主要目标低延迟、减少锁竞争和对象分配可靠投递、削峰、解耦和消息回放
主要目标
低延迟、减少锁竞争和对象分配
可靠投递、削峰、解耦和消息回放
故障恢复需要业务自行处理进程故障由 Broker、副本和消费位点提供支持
故障恢复
需要业务自行处理进程故障
由 Broker、副本和消费位点提供支持

Disruptor 通过 RingBuffer、预分配、CAS 和缓存行填充减少锁竞争、GC 和伪共享。它不能替代分布式 MQ,更适合异步日志、撮合和单进程内的高频事件流水线。
写在最后
感谢你能看到这里,也希望这篇文章对你有点用。
JavaGuide 坚持更新 6 年多,近 6000 次提交、600+ 位贡献者一起打磨。如果这些内容对你有帮助,非常欢迎点个免费的 Star 支持下(完全自愿,觉得有收获再点就好):GitHub | Gitee。
如果你想要付费支持/面试辅导(比如简历优化、一对一提问、高频考点突击资料等)的话,欢迎了解我的知识星球。已经坚持维护六年,内容持续更新,虽白菜价(0.4元/天)但质量很高,主打一个良心!
