消息队列:知识点
以下笔记整理自个人学习,覆盖堆积、可靠性、重复消费、顺序消息等主题,便于检索与复习。
消息队列如果发生消息堆积怎么办?(字节面试:如何解决MQ消息积压问题?MQ(Message Queue)消息积压问题指的是在消息队列中累积了大量未处理的消 - 掘金)
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 | |
生产者确认机制
我个人认为,生产者确认机制对性能影响较大,无特殊需要不要开启
一般情况下,只要生产者与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 | |
这里publisher-confirm-type有三种模式可选:
- none:关闭confirm机制
- simple:同步阻塞等待MQ的回执
- correlated:MQ异步回调返回回执
添加之后需要去配置一个ReturnCallback,每个RabbitTemplate只能配置一个ReturnCallback,因此我们可以在配置类中统一设置。我们在publisher模块定义一个配置类(如果返回 nack则会调用这个方法:
1 | |
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 | |
重启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)
- 消息存储:
- RabbitMQ:内存为主,磁盘为辅。消息被消费后通常会立刻删除。它不适合存储大量积压的消息,积压会导致性能急剧下降。
- Kafka:基于磁盘的日志文件 (Log)。消息是追加写入(Append Only)的,并且会持久化保留(比如保留 7 天),无论是否被消费。这使得 Kafka 可以随时“回放”历史消息。
- 消费模型:
- RabbitMQ:推模式 (Push) 为主。Broker 主动把消息推送给消费者(需要做 QoS 限流)。
- Kafka:拉模式 (Pull)。消费者自己去 Broker 拉取消息,消费速度由消费者自己控制(配合 offset 记录消费进度)。
- 路由机制:
- RabbitMQ:非常强大。有 Exchange (交换机) 概念(Direct, Topic, Fanout),可以在 Broker 端实现极其复杂的路由规则。
- Kafka:非常简单。只有 Topic 和 Partition。生产者直接发给 Topic,没有复杂的路由逻辑。
三、 性能与吞吐量 (Performance)
- Kafka:吞吐量之王。
- 单机可以轻松达到 10万+ 甚至百万级 QPS。
- 秘诀:零拷贝 (Zero Copy)、磁盘顺序写 (Sequential Write)、批量发送 (Batching)。
- RabbitMQ:吞吐量一般。
- 单机通常在 万级 QPS。
- 虽然吞吐量不如 Kafka,但它的延迟极低(微秒级),适合对实时响应要求极高的在线业务。
四、 选型建议(面试总结)
“如果面试官问怎么选,我会这么说:
- 选 RabbitMQ 的场景:
- 需要处理复杂的路由逻辑(如根据 Header 分发)。
- 对延迟极度敏感(如金融交易)。
- 数据量不大,且需要极高的数据可靠性(不支持消息丢失)。
- 选 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,主要基于两个考虑:
- 业务需要可靠性高,比如订单、支付链路需要保证消息不丢失、支持ACK机制, 以及支持延时队列功能,在消费失败能成功人工介入;
- 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个数,就是全局中位数。
消息消费失败了怎么办?如果回过头来做? 怎么做到不丢消息?(可以用到项目中)
要做的:
消息不丢失:消费失败不能导致消息彻底消失。
可重复消费(重试):失败后可以重新消费。
幂等性:同一条消息消费多次,结果一致。
可追溯性:知道哪条消息失败了,能查、能修复。
手动 ACK(不要自动确认):如果你的消费逻辑失败了,不要确认(ack)消息,让消息队列重投。
死信队列:消费失败消息,超过最大重试次数,就会被送入 死信队列(DLQ)。死信队列可用于:后续分析 / 人工处理/补偿系统定时拉出再处理
1
消息消费失败 → 保存失败记录到 DLX → 定时任务补偿 → 再次投递消息 or 人工处理消息落库(失败消息存数据库):如果你的消息非常关键,比如支付扣款、订单发货,可以在消费失败时把消息保存到数据库。后台启动一个补偿程序(定时任务)去扫数据库,重新消费这些失败消息。
1
消息消费失败 → 保存失败记录到 MySQL → 定时任务补偿 → 再次投递消息 or 人工处理保持幂等性(防止重复消费出问题)
补偿机制:定时任务重试失败记录:写一个定时任务,定时从失败消息表中读取消息,重新投递或重新处理.
消息超过 N 次失败,标记为“毒消息”:毒消息(Poison Message) 是指:即使重试 10 次,还是失败的消息。如果定时任务扫描N次还是失败,标记为“待人工干预”,存入数据库,触发 告警(钉钉 / 微信 / 邮件)。
如果设计个web server你会怎么设计,tcp你如何处理呢,多路复用你如何来做,不是原理,TCP优化你能想到什么,除了多路复用
##生产者消费者模型你怎么实现,用什么数据结构,如果这个消息包非常大,你如何处理
集群如何保证一致性
讲一下如果让你设计一个jvm,如何管理内存的申请和释放,不要那么复杂的结构,申请,释放过程是怎样的,用的什么数据结构,复杂度是多少,有没有更简单的结构,不是OS内存是进程里面如何设计,如果一个大对象如何分配内存。
假如核心线程还没满,有几个空闲的核心线程,又来了一个新任务,此时会怎么做?(答错了 应该是新建核心线程)
如何实现一个无锁化的并发数据结构?
了解分布式事务吗?聊聊 2PC、3PC、Seata、TCC?(浅谈分布式事务及解决方案 | 京东物流技术团队1 背景 在讲述分布式事务的概念之前,我们先来回顾下事务相关的一些概念。 - 掘金)(终于有人把“TCC分布式事务”实现原理讲明白了! - JaJian - 博客园_)

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

那现在有一个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 次,对业务系统的结果都是一样的。