消息队列:知识点

以下笔记整理自个人学习,覆盖堆积、可靠性、重复消费、顺序消息等主题,便于检索与复习。

消息队列如果发生消息堆积怎么办?(字节面试:如何解决MQ消息积压问题?MQ(Message Queue)消息积压问题指的是在消息队列中累积了大量未处理的消 - 掘金)

1.消息堆积问题的危害

  • 系统性能下降:当消息堆积量过大时,消息队列的处理能力可能受到影响,甚至导致系统的整体性能下降。
  • 消息延迟:随着队列中消息的增多,新消息的消费速度变慢,导致消息的处理延迟增加,用户体验变差。
  • 内存溢出和磁盘空间耗尽:如果消息堆积持续增长,队列可能会占用过多的内存或磁盘空间,最终导致内存溢出或磁盘空间耗尽。
  • 消息丢失:某些消息队列系统可能在消息堆积到一定程度后丢弃旧消息,从而造成数据丢失。

2.消息堆积的类型

  • 突发性消息堆积:生产者在短时间内产生大量的消息,导致消费者一时无法处理。通常是由于系统的流量激增或特殊事件(如促销活动)导致的。
  • 持续性消息堆积:生产者产生消息的速度持续超过消费者的消费速度,导致消息逐渐积压。此类问题通常是系统设计或配置不合理的结果。
  • 周期性消息堆积:消息堆积呈现周期性变化,可能与系统的负载峰谷相关,例如白天访问量大、夜间访问量低。

3.消息生产速度分析

消息堆积的一个常见原因是生产者发送消息的速度过快,超过了消费者的处理能力。在这种情况下,消费者无法及时消费所有的消息,导致消息不断积压。

排查思路:

  1. 检查生产者的消息发送速率。
  2. 确定生产者发送的消息是否异常,例如是否有某些业务逻辑导致短时间内发送大量重复或无效消息。
  3. 分析生产者的流量是否超出了系统的设计范围,是否有突然增加的流量高峰。

解决方案:

  • 通过限流手段控制生产者的发送速率,防止短时间内发送过多消息。
  • 使用流量整形(Traffic Shaping)手段,将消息发送的速度平滑化,避免突发流量。
  • 如果是系统流量激增导致的堆积,考虑通过弹性扩展增加消费者实例,缓解堆积。

消息去重与有效性

某些情况下,生产者可能会因为代码问题或业务逻辑错误,产生大量重复或无效的消息。这些无效的消息占用了消息队列的资源,影响了消费者的正常工作。

排查思路:

  • 检查消息的唯一性标识(如 messageId),确定是否有大量重复消息。
    分析消息内容的有效性,确保没有生产空消息或无效消息。
    解决方案:

  • 在生产者端增加去重机制,确保同一条消息不会被重复发送。
    对消息内容进行校验,避免发送无效消息。

消息生产失败与重试机制

如果消息在发送过程中失败,生产者通常会自动进行重试。然而,如果重试机制设计不合理,可能会导致生产者发送大量重复消息,从而导致消息堆积。

排查思路:

  • 分析生产者的重试机制,检查是否存在频繁重试的现象。
  • 确定消息发送失败的原因,是否是由于网络问题、队列满等情况导致的。

解决方案:

  • 优化生产者的重试机制,设置合理的重试间隔和重试次数,避免频繁重试。
  • 使用幂等机制,确保消息发送的重复操作不会影响业务逻辑。

4.消费者处理能力不足

消费者的处理能力不足是导致消息堆积的主要原因之一。如果消费者处理消息的速度低于生产者的发送速度,消息就会不断堆积。

解决方案:

  • 扩展消费者实例,增加消费线程或部署更多的消费者实例,提升消息的处理速度。
  • 优化消费者的代码逻辑,减少不必要的计算或 I/O 操作,提高消费效率。
  • 使用多线程或异步处理的方式,提升消费者的并发处理能力。

消费者重试机制问题

在某些情况下,消费者处理消息失败时会进行重试。如果重试机制设计不合理,可能会导致消费者反复处理失败的消息,进一步加剧消息堆积问题。

排查思路:

  • 分析消费者的重试机制,确定是否存在频繁重试的现象。
  • 检查消息处理失败的原因,是否由于网络问题、数据不一致等导致。
    解决方案:
  • 设置合理的重试间隔和重试次数,避免频繁重试。
    对于某些无法立即处理的消息,可以将其放入死信队列,避免影响其他正常消息的处理。

5.如何防止消息堆积问题

除了针对性的解决方案外,合理的系统设计和配置可以有效预防消息堆积问题的发生。以下是一些常见的防止消息堆积的策略:

5.1 使用限流机制

通过限流机制,可以控制生产者的消息发送速度,避免在高并发场景下生产者发送过多的消息,导致消息堆积。

5.2 动态扩容

在流量激增的情况下,可以通过动态扩容的方式增加消费者实例,提升系统的处理能力。

5.3 使用死信队列

对于无法正常处理的消息,可以将其放入死信队列(Dead Letter Queue,DLQ),避免影响其他消息的正常消费。

5.4 消息过期与优先级队列

对于一些时效性较强的消息,可以设置过期时间,确保消息在一定时间后自动失效,避免长期积压。此外,优先级队列可以确保高优先级的消息被优先处理。

5.5 合理设置消费者重试策略

对于消费者的重试机制,需要设置合理的重试次数和间隔,避免频繁重试导致系统性能下降。

MQ怎么保证消息的可靠性的(RabbitMQ进阶–保证消息的可靠性_rabbitmq 消息可靠性-CSDN博客

1.保证MQ消息可靠性的三种方式

1.1发送者可靠性

生产者重试机制

1
2
3
4
5
6
7
8
9
10
11
12
13
通过在配置文件中添加相关配置打开重试机制
spring:
rabbitmq:
connection-timeout: 1s # 设置MQ的连接超时时间
template:
retry:
enabled: true # 开启超时重试机制
initial-interval: 1000ms # 失败后的初始等待时间
multiplier: 1 # 失败后下次的等待时长倍数,下次等待时长 = initial-interval * multiplier
max-attempts: 3 # 最大重试次数
注意:当网络不稳定的时候,利用重试机制可以有效提高消息发送的成功率。不过SpringAMQP提供的重试机制是阻塞式的重试,也就是说多次重试等待的过程中,当前线程是被阻塞的。

如果对于业务性能有要求,建议禁用重试机制。如果一定要使用,请合理配置等待时长和重试次数,当然也可以考虑使用异步线程来执行发送消息的代码。

生产者确认机制

我个人认为,生产者确认机制对性能影响较大,无特殊需要不要开启

一般情况下,只要生产者与MQ之间的网路连接顺畅,基本不会出现发送消息丢失的情况,因此大多数情况下我们无需考虑这种问题。
不过,在少数情况下,也会出现消息发送到MQ之后丢失的现象,比如:

  • MQ内部处理消息的进程发生了异常
  • 生产者发送消息到达MQ后未找到Exchange
  • 生产者发送消息到达MQ的Exchange后,未找到合适的Queue,因此无法路由

针对上述情况,RabbitMQ提供了生产者消息确认机制,包括Publisher Confirm和Publisher Return两种。在开启确认机制的情况下,当生产者发送消息给MQ后,MQ会根据消息处理的情况返回不同的回执。

  • 当消息投递到MQ,但是路由失败时,通过Publisher Return返回异常信息,同时返回ack的确认信息,代表投递成功
  • 临时消息投递到了MQ,并且入队成功,返回ACK,告知投递成功
  • 持久消息投递到了MQ,并且入队完成持久化,返回ACK ,告知投递成功
  • 其它情况都会返回NACK,告知投递失败

其中ack和nack属于Publisher Confirm机制,ack是投递成功;nack是投递失败。而return则属于Publisher Return机制。

默认两种机制都是关闭状态,需要通过配置文件来开启。
在生产者中添加配置

1
2
3
4
5
spring:
rabbitmq:
publisher-confirm-type: correlated # 开启publisher confirm机制,并设置confirm类型
publisher-returns: true # 开启publisher return机制

这里publisher-confirm-type有三种模式可选:

  • none:关闭confirm机制
  • simple:同步阻塞等待MQ的回执
  • correlated:MQ异步回调返回回执

添加之后需要去配置一个ReturnCallback,每个RabbitTemplate只能配置一个ReturnCallback,因此我们可以在配置类中统一设置。我们在publisher模块定义一个配置类(如果返回 nack则会调用这个方法:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
@Slf4j
@Configuration
@RequiredArgsConstructor
public class MqConfig {

private final RabbitTemplate rabbitTemplate;

@PostConstruct
public void init(){
rabbitTemplate.setReturnsCallback(returnedMessage -> {
log.error("监听到消息return callback");
log.debug("exchange:{}", returnedMessage.getExchange());
log.debug("routingKey:{}", returnedMessage.getRoutingKey());
log.debug("message:{}", returnedMessage.getMessage());
log.debug("replyCode:{}", returnedMessage.getReplyCode());
log.debug("replyText:{}", returnedMessage.getReplyText());
});
}
}

MQ可靠性

交换器的持久化,队列的持久化,消息的持久化

消费者可靠性

由于消息回执的处理代码比较统一,因此SpringAMQP帮我们实现了消息确认。并允许我们通过配置文件设置ACK处理方式,有三种模式:

  • none:不处理。即消息投递给消费者后立刻ack,消息会立刻从MQ删除。非常不安全,不建议使用
  • manual:手动模式。需要自己在业务代码中调用api,发送ack或reject,存在业务入侵,但更灵活
  • auto:自动模式。SpringAMQP利用AOP对我们的消息处理逻辑做了环绕增强,当业务正常执行时则自动返回ack. 当业务出现异常时,根据异常判断返回不同结果:
    • 如果是业务异常,会自动返回nack;
    • 如果是消息处理或校验异常,自动返回reject;

消费者重连机制

当消费者出现异常后,消息会不断requeue(重入队)到队列,再重新发送给消费者。如果消费者再次执行依然出错,消息会再次requeue到队列,再次投递,直到消息处理成功为止。

极端情况就是消费者一直无法执行成功,那么消息requeue就会无限循环,导致mq的消息处理飙升,带来不必要的压力:当然,上述极端情况发生的概率还是非常低的,不过不怕一万就怕万一。为了应对上述情况Spring又提供了消费者失败重试机制:在消费者出现异常时利用本地重试,而不是无限制的requeue到mq队列。

修改consumer服务的application.yml文件,添加内容:

1
2
3
4
5
6
7
8
9
10
11
spring:
rabbitmq:
listener:
simple:
retry:
enabled: true # 开启消费者失败重试
initial-interval: 1000ms # 初识的失败等待时长为1
multiplier: 1 # 失败的等待时长倍数,下次等待时长 = multiplier * last-interval
max-attempts: 3 # 最大重试次数
stateless: true # true无状态;false有状态。如果业务中包含事务,这里改为false

重启consumer服务,重复之前的测试。可以发现:

  • 消费者在失败后消息没有重新回到MQ无限重新投递,而是在本地重试了3次
  • 本地重试3次以后,抛出了AmqpRejectAndDontRequeueException异常。查看RabbitMQ控制台,发现消息被删除了,说明最后SpringAMQP返回的是reject

结论:

  • 开启本地重试时,消息处理过程中抛出异常,不会requeue到队列,而是在消费者本地重试
  • 重试达到最大次数后,Spring会返回reject,消息会被丢弃

失败处理策略

在之前的测试中,本地测试达到最大重试次数后,消息会被丢弃。这在某些对于消息可靠性要求较高的业务场景下,显然不太合适了。

因此Spring允许我们自定义重试次数耗尽后的消息处理策略,这个策略是由MessageRecovery接口来定义的,它有3个不同实现:

  • RejectAndDontRequeueRecoverer:重试耗尽后,直接reject,丢弃消息。默认就是这种方式
  • ImmediateRequeueMessageRecoverer:重试耗尽后,返回nack,消息重新入队
  • RepublishMessageRecoverer:重试耗尽后,将失败消息投递到指定的交换机

消息堆积问题(消息队列(MQ)消息堆积问题排查与解决思路_mq消息堆积问题排查-CSDN博客)(字节面试:如何解决MQ消息积压问题?MQ(Message Queue)消息积压问题指的是在消息队列中累积了大量未处理的消 - 掘金

消息队列怎么保证有序,业务层面怎么保证订单到达的顺序(消息队列(五)如何保证消息的顺序性?_Java_奈何花开_InfoQ写作社区)

四种队列的对比(Kafka、RabbitMQ、RocketMQ与ActiveMQ:消息队列比较与选型指南-CSDN博客)

一、 核心设计定位 (The “Why”)

  • RabbitMQ
    • 定位传统的消息中间件
    • 核心Erlang 语言编写,实现了 AMQP(高级消息队列协议)。
    • 强项复杂路由、低延迟、高可靠性。它像一个精密的“邮局”,非常擅长处理复杂的业务逻辑流转(如订单状态机、金融支付)。
  • Kafka
    • 定位分布式的流处理平台
    • 核心Scala/Java 编写。
    • 强项超高吞吐、大数据存储、日志处理。它更像一个巨大的“日志管道”,甚至可以被视为一个分布式存储系统。

二、 架构模式与存储 (Architecture)

  1. 消息存储
    • RabbitMQ:内存为主,磁盘为辅。消息被消费后通常会立刻删除。它不适合存储大量积压的消息,积压会导致性能急剧下降。
    • Kafka基于磁盘的日志文件 (Log)。消息是追加写入(Append Only)的,并且会持久化保留(比如保留 7 天),无论是否被消费。这使得 Kafka 可以随时“回放”历史消息。
  2. 消费模型
    • RabbitMQ推模式 (Push) 为主。Broker 主动把消息推送给消费者(需要做 QoS 限流)。
    • Kafka拉模式 (Pull)。消费者自己去 Broker 拉取消息,消费速度由消费者自己控制(配合 offset 记录消费进度)。
  3. 路由机制
    • RabbitMQ:非常强大。有 Exchange (交换机) 概念(Direct, Topic, Fanout),可以在 Broker 端实现极其复杂的路由规则。
    • Kafka:非常简单。只有 TopicPartition。生产者直接发给 Topic,没有复杂的路由逻辑。

三、 性能与吞吐量 (Performance)

  • Kafka吞吐量之王
    • 单机可以轻松达到 10万+ 甚至百万级 QPS。
    • 秘诀:零拷贝 (Zero Copy)、磁盘顺序写 (Sequential Write)、批量发送 (Batching)。
  • RabbitMQ:吞吐量一般。
    • 单机通常在 万级 QPS。
    • 虽然吞吐量不如 Kafka,但它的延迟极低(微秒级),适合对实时响应要求极高的在线业务。

四、 选型建议(面试总结)

“如果面试官问怎么选,我会这么说:

  1. 选 RabbitMQ 的场景
    • 需要处理复杂的路由逻辑(如根据 Header 分发)。
    • 延迟极度敏感(如金融交易)。
    • 数据量不大,且需要极高的数据可靠性(不支持消息丢失)。
  2. 选 Kafka 的场景
    • 大数据日志采集(ELK)。
    • 实时流计算(Flink/Spark Streaming)。
    • 超高并发的削峰填谷。
    • 需要回溯/重放历史消息的场景。”

1.kafka消息队列(【消息队列】——Kafka入门一篇就够了!_kafka消息队列-CSDN博客)(一文看懂-Kafka消息队列 - Flamings - 博客园)(https://blog.csdn.net/weixin_45366499/article/details/106943229)

  • RabbitMq:基于AMQP协议,可靠性很强,支持ACK, 延迟消息,消息路由,适合对数据一致性和可靠性要求高的场景,比如订单系统、通知系统。我实际用它做过电商项目的下单异步处理。
  • Kafka:吞吐量非常高,是一个分布式的日志系统,适合大数据、日志采集、流式计算场景。它支持分区、顺序性,偏向高性能批量处理,但不支持延迟消息。
  • RocketMQ: 阿里开源,天然支持分布式部署,原生支持事务消息、延迟消息,在国内金融、电商领域用得很多,适合需要高吞吐+事务保障的场景。

如果面试官继续追问:“你们为什么选 RabbitMQ?

我们当时选用的是 RabbitMQ,主要基于两个考虑:

  1. 业务需要可靠性高,比如订单、支付链路需要保证消息不丢失、支持ACK机制, 以及支持延时队列功能,在消费失败能成功人工介入;
  2. RabbitMq管理控制台比较直观。

实际上,我觉得选 MQ 不是“谁好用谁上”,而是要看:

  • 业务对可靠性 / 性能 / 延迟 / 顺序性的要求;
  • 团队熟悉程度和社区活跃度;
  • 是否需要事务 / 延迟 / 死信队列等功能。

你消息队列用哪种比较多?它在实现高性能方面做了哪些优化或者设计

  • RabbitMQ 的核心是基于 Erlang 构建的 Actor 模型,天生支持高并发和异步消息传递
  • 利用 Exchange(交换机)+ Binding + RoutingKey 的灵活组合,支持多种路由策略(Direct, Topic, Fanout, Headers),实现精准高效路由
  • RabbitMQ 会优先将消息缓存在内存中(page cache),等到达到阈值才写磁盘,提升性能。磁盘顺序写优化(尤其对于持久化队列),保证即使重启也不会丢数据。
  • 支持 手动/自动 ACK 确认机制,可以结合批量确认来减少网络交互、提高吞吐。

如果要提升mq的消费吞吐能力,要怎么做呢?

了解RabbitMQ的实现细节吗?RabbitMQ怎么实现持久化的?

设计短链系统,给了场景,需要压缩的链接百万级,提示不建议用算法,考虑好开始说(字节三面:如何设计一个高性能短链系统?-腾讯云开发者社区-腾讯云

改一下:当用自增id时,用户访问短链时,系统解码短链字符串,还原出ID。

文件中有100亿个无序整数,内存100M,找中位数

【1】第一次扫描(分桶统计)

  • 定义一个范围划分(比如按照数值大小区间分桶)。
  • 比如整个数值可能落在[-2^31, 2^31-1]。
  • 切成很多区间,比如1亿个区间,每个区间负责一小段值域。
  • 扫描一遍数据,每个数丢到对应的区间里,只统计各个桶的数量(用一个小数组存每个桶里的数量,不存具体数字)。

➡️ 目的是为了快速定位中位数在哪个区间

【2】确定中位数所在桶

  • 扫描完后,加上每个桶的计数,累加到第50亿个数,看中位数落在哪个桶。
  • 比如中位数落在第K个桶。

【3】第二次扫描(精确处理中位数桶)

  • 重新扫一遍文件,只把落到第K个桶的数读出来(这部分数据量远小于全部)。
  • 这部分数据量可能几百万、几千万,内存能放得下了(或者继续分批放)。

【4】排序+找中位数

  • 把第K个桶的数据拿到内存里,进行快速排序(或者其他简单的排序算法)。
  • 排好后取第k个数,就是全局中位数。

消息消费失败了怎么办?如果回过头来做? 怎么做到不丢消息?(可以用到项目中)

要做的:

消息不丢失:消费失败不能导致消息彻底消失。

可重复消费(重试):失败后可以重新消费。

幂等性:同一条消息消费多次,结果一致。

可追溯性:知道哪条消息失败了,能查、能修复。

  1. 手动 ACK(不要自动确认):如果你的消费逻辑失败了,不要确认(ack)消息,让消息队列重投。

  2. 死信队列:消费失败消息,超过最大重试次数,就会被送入 死信队列(DLQ)。死信队列可用于:后续分析 / 人工处理/补偿系统定时拉出再处理

    1
    消息消费失败 → 保存失败记录到 DLX → 定时任务补偿 → 再次投递消息 or 人工处理
  3. 消息落库(失败消息存数据库):如果你的消息非常关键,比如支付扣款、订单发货,可以在消费失败时把消息保存到数据库。后台启动一个补偿程序(定时任务)去扫数据库,重新消费这些失败消息。

    1
    消息消费失败 → 保存失败记录到 MySQL → 定时任务补偿 → 再次投递消息 or 人工处理
  4. 保持幂等性(防止重复消费出问题)

  5. 补偿机制:定时任务重试失败记录:写一个定时任务,定时从失败消息表中读取消息,重新投递或重新处理.

  6. 消息超过 N 次失败,标记为“毒消息”:毒消息(Poison Message) 是指:即使重试 10 次,还是失败的消息。如果定时任务扫描N次还是失败,标记为“待人工干预”,存入数据库,触发 告警(钉钉 / 微信 / 邮件)。

如果设计个web server你会怎么设计,tcp你如何处理呢,多路复用你如何来做,不是原理,TCP优化你能想到什么,除了多路复用

##生产者消费者模型你怎么实现,用什么数据结构,如果这个消息包非常大,你如何处理

集群如何保证一致性

讲一下如果让你设计一个jvm,如何管理内存的申请和释放,不要那么复杂的结构,申请,释放过程是怎样的,用的什么数据结构,复杂度是多少,有没有更简单的结构,不是OS内存是进程里面如何设计,如果一个大对象如何分配内存。

假如核心线程还没满,有几个空闲的核心线程,又来了一个新任务,此时会怎么做?(答错了 应该是新建核心线程)

如何实现一个无锁化的并发数据结构?

了解分布式事务吗?聊聊 2PC、3PC、Seata、TCC?(浅谈分布式事务及解决方案 | 京东物流技术团队1 背景 在讲述分布式事务的概念之前,我们先来回顾下事务相关的一些概念。 - 掘金)(终于有人把“TCC分布式事务”实现原理讲明白了! - JaJian - 博客园_)

image-20250908204336758

RPC和Http通信的区别,各自的优缺点(RPC 和 HTTP 理解,看完这一篇就够了 - 星火燎原智勇 - 博客园)

image-20250912171437086

那现在有一个10g的大文件,里面存着不同的ip地址,你只有1g的内存,如何找出里面出现最多ip地址的top10?

ip地址的,可以先分批遍历一遍文件,按ip的哈希值取模分成100个文件,然后再维护每个文件的前10

TCP 连接复用,chrome 打开新的标签页会使用 TCP 连接复用吗

具体来说,Chrome 在处理新标签页的连接时,会根据不同的 HTTP 协议版本采取不同的策略:

  • HTTP/1.1 环境下: Chrome 会为每个域名维护一个连接池,通常允许同时建立多个(例如 6-8 个)并行的 TCP 连接。当您打开一个新的标签页访问同一个域名时,Chrome 会从这个连接池中选取一个空闲的连接来发送请求,而不是重新创建一个新的连接。这大大提高了资源加载的并行度,缩短了页面整体加载时间。
  • HTTP/2 环境下: HTTP/2 协议引入了“多路复用”(Multiplexing)这一关键特性,从根本上改变了连接管理的方式。在 HTTP/2 下,Chrome 只需要为每个域名建立一个单一的 TCP 连接,就可以在该连接上同时并发地处理多个请求和响应,而不会相互阻塞。当您打开一个新的标签页访问支持 HTTP/2 的网站时,新的请求会通过这个已经建立的连接以新的“流”(Stream)的形式发送,从而实现了更高效的连接利用和更低的延迟。

值得注意的是, 这种连接复用机制是针对同一域名(源)的。如果您在新标签页中访问的是一个完全不同的网站,那么 Chrome 还是需要为其建立一个新的 TCP 连接。此外,像 WebSocket 这样的特殊协议会独占一个连接,无法被其他 HTTP 请求复用。

总而言之,Chrome 浏览器通过先进的连接管理策略,包括在 HTTP/1.1下的连接池技术和在 HTTP/2 下的多路复用技术,有效地实现了 TCP 连接的复用。这使得在打开新的标签页访问相同网站时,能够获得更快的加载速度和更优的网络性能。

RabbitMq的实现原理

RabbitMQ 的底层基于 Erlang 语言和 AMQP 协议 实现。

它的核心架构非常灵活,引入了 Exchange(交换机) 这一层。生产者不直接对接队列,而是将消息发给交换机,由交换机根据路由规则分发到具体的队列。这实现了生产和消费的完全解耦。

在通信层面,它使用了 Channel(信道) 复用 TCP 连接,解决了频繁建连的开销。

在存储层面,它采用内存索引 + 磁盘日志的方式。对于持久化消息,会顺序写入磁盘;对于非持久化消息,优先存内存,内存不足时再换出到磁盘。这种设计保证了它在低延迟和高可靠性之间的平衡。

如何考虑一个消息队列的实现?有多少要注意的地方?

如果要设计一个消息队列(MQ),这是典型的系统设计题。你不能只说“收发消息”,必须从高可用、高性能、高可靠三个维度,全面考虑一个中间件该有的素养。

至少有 7 个核心注意点 需要考虑:

1. 通信协议 (Protocol)

  • 问题:客户端和服务端怎么说话?
  • 设计
    • TCP:必须要用长连接,避免频繁握手。
    • 自定义协议:像 Kafka/RocketMQ 那样,设计一套紧凑的二进制协议(Header+Body),减少网络开销,支持多语言 SDK。

2. 消息存储 (Storage) —— 最核心

  • 问题:消息存在哪?内存还是磁盘?
  • 设计
    • 内存:快,但易丢(如 Redis Pub/Sub)。
    • 磁盘:可靠,但慢。优化手段是 顺序写 (Sequential Write),像 Kafka 的 CommitLog 那样追加写入,速度接近内存。
    • 索引:为了快速读,需要维护一个稀疏索引(如 ConsumeQueue),通过 Offset 快速定位文件位置。

如果要设计一个消息队列(MQ),这是典型的系统设计题。你不能只说“收发消息”,必须从高可用、高性能、高可靠三个维度,全面考虑一个中间件该有的素养。

至少有 7 个核心注意点 需要考虑:

1. 通信协议 (Protocol)

  • 问题:客户端和服务端怎么说话?
  • 设计
    • TCP:必须要用长连接,避免频繁握手。
    • 自定义协议:像 Kafka/RocketMQ 那样,设计一套紧凑的二进制协议(Header+Body),减少网络开销,支持多语言 SDK。

2. 消息存储 (Storage) —— 最核心

  • 问题:消息存在哪?内存还是磁盘?
  • 设计
    • 内存:快,但易丢(如 Redis Pub/Sub)。
    • 磁盘:可靠,但慢。优化手段是 顺序写 (Sequential Write),像 Kafka 的 CommitLog 那样追加写入,速度接近内存。
    • 索引:为了快速读,需要维护一个稀疏索引(如 ConsumeQueue),通过 Offset 快速定位文件位置。

3. 高可用架构 (High Availability)

  • 问题:机器挂了怎么办?
  • 设计
    • 主从复制 (Replication):数据必须有多副本。
    • 选主机制:Master 挂了,Slave 怎么上位?是用 Zookeeper/Raft 自动选主(Kafka),还是依赖 NameServer(RocketMQ)?

4. 消费模型 (Consumption Model)

  • 问题:怎么给消费者?
    • Push (推):实时性高,但容易压垮消费者。
    • Pull (拉):消费者自己掌握节奏,适合批量拉取,但可能有延迟。
    • 长轮询 (Long Polling):结合两者优点。消费者拉取,如果没消息,服务端挂起请求,等有消息了立马返回。

5. 可靠性保障 (Reliability)

  • 问题:怎么保证不丢消息?
    • 生产端:Ack 确认机制。
    • 服务端:刷盘策略(同步刷盘 vs 异步刷盘)。
    • 消费端:手动 Ack。

6. 高级特性 (Advanced Features)

  • 延时队列:通过时间轮算法实现。
  • 事务消息:支持分布式事务(两阶段提交 + 回查)。
  • 死信队列:重试 N 次失败后隔离消息。

7. 零拷贝 (Zero-Copy) —— 性能杀手锏

  • 设计:在读取消息发给消费者时,利用 sendfile 或 mmap 技术,直接把数据从磁盘文件拷贝到网卡,绕过 CPU 和用户态内存,极大提升吞吐量。

  • 问题:机器挂了怎么办?

  • 设计

    • 主从复制 (Replication):数据必须有多副本。
    • 选主机制:Master 挂了,Slave 怎么上位?是用 Zookeeper/Raft 自动选主(Kafka),还是依赖 NameServer(RocketMQ)?

4. 消费模型 (Consumption Model)

  • 问题:怎么给消费者?
    • Push (推):实时性高,但容易压垮消费者。
    • Pull (拉):消费者自己掌握节奏,适合批量拉取,但可能有延迟。
    • 长轮询 (Long Polling):结合两者优点。消费者拉取,如果没消息,服务端挂起请求,等有消息了立马返回。

5. 可靠性保障 (Reliability)

  • 问题:怎么保证不丢消息?
    • 生产端:Ack 确认机制。
    • 服务端:刷盘策略(同步刷盘 vs 异步刷盘)。
    • 消费端:手动 Ack。

6. 高级特性 (Advanced Features)

  • 延时队列:通过时间轮算法实现。
  • 事务消息:支持分布式事务(两阶段提交 + 回查)。
  • 死信队列:重试 N 次失败后隔离消息。

7. 零拷贝 (Zero-Copy) —— 性能杀手锏

  • 设计:在读取消息发给消费者时,利用 sendfile 或 mmap 技术,直接把数据从磁盘文件拷贝到网卡,绕过 CPU 和用户态内存,极大提升吞吐量。

kafka怎么保证消息不丢失?同一消费组的三个消费者消费三个分区,如果有一个消费端挂了,其他消费者来消费,怎么保证不重复消费?这个时候消息id在哪

场景描述:一个 Topic 有 3 个 Partition(P0, P1, P2),消费组里有 3 个 Consumer(C0, C1, C2),正好一对一消费。突然,C2 挂了。

1. 发生了什么?(Rebalance 机制)

  • 当 C2 挂掉(长时间未发送心跳)时,Kafka 集群的 Group Coordinator(组协调器) 会敏锐地察觉到。
  • 它会触发一次 Rebalance(重平衡)
  • Rebalance 的结果是:原本由 C2 消费的 P2 分区,会被重新分配给存活的 C0 或 C1(比如分配给了 C0)。

2. 为什么会产生“重复消费”?

  • 假设 C2 挂掉之前,已经从 P2 拉取了一批消息(Offset 100 到 150),并且业务逻辑已经处理到了 Offset 120。
  • 但是!C2 还没来得及向 Broker 提交 Offset 120 就宕机了
  • Broker 那边记录的 P2 分区的最新消费位移依然是上一批提交的 Offset 100
  • 当 Rebalance 结束后,C0 接管了 P2 分区。C0 问 Broker:“我该从哪里开始消费 P2?” Broker 会告诉它:“从 Offset 100 开始”。
  • 于是,C0 又把 100 到 120 的消息拉下来重新处理了一遍。这就是经典的、因 Consumer 宕机导致的重复消费场景。

3. 怎么保证不重复消费?(业务层面的幂等性设计)

敲黑板:Kafka 本身无法 100% 避免这种场景下的重复投递(它提供的是 At Least Once 语义)。解决重复消费的终极防线,必须落在业务代码的“幂等性(Idempotence)”设计上。

所谓幂等性,就是同一条消息,无论你处理 1 次还是 100 次,对业务系统的结果都是一样的

敲黑板:Kafka 本身无法 100% 避免这种场景下的重复投递(它提供的是 At Least Once 语义)。解决重复消费的终极防线,必须落在业务代码的“幂等性(Idempotence)”设计上。

所谓幂等性,就是同一条消息,无论你处理 1 次还是 100 次,对业务系统的结果都是一样的


消息队列:知识点
https://kyy-logs.github.io/2026/05/13/消息队列/消息队列-知识点/
作者
Yangyang Kong
发布于
2026年5月13日
许可协议