RabbitMQ 全套复盘 + Nacos+ES+MyBatis-Plus 梳理
RabbitMQ 消息队列(异步通信、可靠性全套解决方案)
在短信项目里:上万条批量营销短信不会直接同步调用第三方通道,而是先丢进 MQ,消费者慢慢消费,防止瞬间流量打垮服务。
RabbitMQ 核心组成组件
1.生产者 Producer
业务服务,负责发送消息。在项目里就是营销活动模块、定时短信任务模块,批量生成短信任务后作为生产者发送消息。
2.Broker 服务节点
RabbitMQ 服务本体,内部包含 Exchange(交换机)、Queue(队列)。
3.Exchange 交换机
消息的路由中转站,生产者不会直接把消息发到队列,先发给交换机,交换机根据路由规则把消息分发到对应队列。 交换机 4 种类型:
- Direct:精准匹配路由键,一对一分发,短信业务普通下发队列使用;
- Topic:模糊匹配路由键,用来区分不同短信渠道、不同营销活动;
- Fanout:广播模式,一条消息同步发给所有绑定队列,用于平台通知、告警推送;
- Headers:基于消息头匹配,项目极少使用。
4.Binding 绑定关系
交换机和队列之间的绑定,绑定的时候指定 routingKey 路由键,决定消息分发规则。
5.Queue 消息队列
存储消息的容器,消费者从队列拉取消息。队列支持持久化、设置最大长度、消息过期时间。
6.消费者 Consumer
监听队列、处理消息的服务,项目中专门的短信推送消费服务,读取消息后调用第三方短信通道完成下发。
7.VirtualHost 虚拟主机
MQ 内的环境隔离机制,不同业务线、开发 / 测试环境分开,互不干扰,项目开发环境单独一套 vhost。
三大核心设计作用(异步、解耦、削峰)
RabbitMQ 是消息中间件,核心作用:异步解耦、流量削峰、最终一致性。
流量削峰
营销活动、节日大促会一次性生成几十万条短信任务,如果同步循环调用第三方通道,瞬间大量请求打满服务线程池、第三方接口超限封禁。
MQ 可以把瞬时海量消息缓冲在队列中,消费者匀速慢慢消费,抹平流量高峰。
系统解耦
短信下发逻辑和活动创建逻辑完全分离。活动模块只负责生产消息,不用关心短信下发是否成功、通道是否故障;后续更换短信渠道、新增风控校验,只改消费端,不用改动活动创建代码。
异步通信,提升接口响应速度
运营后台批量创建上万条短信任务,同步执行会接口超时;丢入 MQ 后接口直接返回成功,消息后台异步处理,前端不用长时间等待。
消息可靠性底层机制(防止消息丢失核心全套原理)
线上三大核心故障:消息丢失、重复消费、消息堆积,配套死信队列做兜底处理。
1.消息丢失
消息丢失分三个链路,每个链路都有独立保障机制,缺一不可:
链路 1:生产者 → Broker 之间丢失消息
场景:生产者发送消息过程中服务宕机、网络中断,消息没到达 Broker 就消失。
底层机制:生产者确认机制 Publisher Confirms
- 开启 confirm 模式,每条消息发送成功后,Broker 会返回 ack 确认;
- 如果返回 nack 或者长时间无响应,生产者本地重试发送,或者记录本地日志定时补发。
配套操作:发送消息时设置持久化标识 deliveryMode=2。
链路 2:Broker 内部丢失消息
场景:MQ 服务器断电、重启,内存中未落地磁盘的消息全部清空。
底层机制:交换机持久化 + 队列持久化 + 消息持久化
- 交换机持久化:重启后交换机不消失;
- 队列持久化:重启后队列不消失;
- 消息持久化:消息写入磁盘,而非仅存内存。 三者同时配置,才能保证 Broker 重启消息不丢失。
链路 3:Broker → 消费者 丢失消息
场景:消息推送给消费者,消费者还没处理完程序宕机,消息直接被删除。
底层机制:消费者手动 ACK 确认机制
- 自动 ACK(默认):消息一推送给消费者,Broker 立刻删除消息,风险极高,线上禁用;
- 手动 ACK:消费者完整处理完业务逻辑(短信下发成功、日志入库),手动发送 ack 指令,Broker 才删除消息;
- 处理失败:发送 nack 指令,消息重新放回队列重试,多次失败转入死信队列。
2.重复消费问题完整原理
产生根源
消费者业务处理成功(短信已经下发完成),但是网络波动,ACK 指令没能传递到 Broker。 Broker 收不到 ack,认为这条消息没有处理完成,一段时间后重新投递,导致同一条短信下发两次,造成用户收到重复短信。
底层解决方案:消费幂等
每条短信任务生成全局唯一 messageId,存入 Redis 做幂等标记。
消费逻辑流程:
- 拿到消息先提取 messageId;
- 查询 Redis,判断该 ID 是否已消费;
- 已存在直接 ACK 丢弃,不执行下发;
- 不存在,执行短信下发,下发成功后写入 Redis 缓存,再发送 ack。
3.消息堆积问题底层原理
产生原因
生产速度 > 消费速度:
- 批量营销活动瞬间生成几十万消息,生产者速度极快;
- 消费者处理慢:调用第三方短信通道有网络延迟、数据库写入慢;
- 消费者数量过少,单节点处理能力有限;
- 消费逻辑阻塞、频繁重试,拖慢消费效率。
堆积带来的危害
- 队列消息过多占用服务器内存,MQ 服务卡顿、崩溃;
- 新消息无法写入,生产者阻塞报错;
- 消息长期积压,过期失效,业务数据丢失。
分层解决方案
- 横向扩容:增加消费者实例,分摊消息处理压力;
- 队列拆分:按渠道、活动拆分多个独立队列,避免单一队列消息过载;
- 限制队列最大长度,超长消息直接转入死信;
- 优化消费内部逻辑:减少同步 IO、批量操作数据库、异步附属逻辑。
4.死信队列 DLX 完整机制
什么消息会进入死信队列
- 消息重试达到最大次数,依旧消费失败;
- 消息过期;
- 队列长度超限,新消息被丢弃转入死信。
项目落地多级重试 + 死信架构
- 业务队列消费失败,先转入短期重试队列,等待 1 分钟重试;
- 重试 3 次仍失败,转入长期重试队列,间隔 5 分钟重试;
- 两轮重试全部失败,投递至死信队列;
- 后台定时任务监听死信队列,统一记录异常短信,运营人工排查黑名单、通道额度、手机号错误等问题。
价值:失败消息不阻塞正常业务队列,统一归档排查,不丢失异常数据。
核心问题解决方案对比表格
| 故障场景 | 底层产生原因 | 全套落地解决方案(短信平台) |
|---|---|---|
| 消息丢失 | 生产者无确认、队列未持久化、自动 ACK | 1. 开启 Publisher-Confirm 生产者确认;2. 交换机、队列、消息三层持久化;3. 消费者手动 ACK |
| 重复消费 | 业务处理成功,ACK 网络丢失,Broker 重发消息 | 每条消息生成唯一 messageId,消费前查询 Redis 校验,实现幂等,拦截重复下发 |
| 消息堆积 | 批量活动瞬时大量消息,消费处理速度跟不上生产 | 1. 扩容消费者节点;2. 按渠道拆分多队列;3. 优化消费内部 IO 逻辑;4. 设置队列最大长度超限转死信 |
| 消费持续失败 | 黑名单手机号、通道欠费、号码格式错误 | 多级延迟重试队列,重试耗尽转入死信队列,后台统一统计异常短信 |
短信 MQ 完整业务流程流程图
流程文字简化口述版
- 运营创建营销活动 / 定时短信任务,生产者生成唯一 messageId 封装消息;
- 开启生产者确认、消息持久化,发送到 Topic 交换机,按渠道路由分到对应业务队列;
- 消费者手动 ACK 拉取消息,先查 Redis 校验 messageId 实现幂等,重复消息直接丢弃;
- 校验通过后执行黑名单、频次风控,校验失败直接进入重试队列;
- 风控通过调用第三方通道下发,下发成功记录幂等标识、写入日志,手动 ACK 删除消息;
- 下发失败则多次延迟重试,重试耗尽转入死信队列;
- 定时任务读取死信消息归档至 ES,运营统一查看异常短信。
Nacos 注册中心 + 配置中心
Nacos 整体定位
Nacos 是阿里开源的微服务组件,同时提供两大核心能力:服务注册发现、动态配置管理,完全替代 Eureka + Spring Cloud Config,项目微服务体系核心底座。
在短信平台作用:
- 所有后端微服务(营销活动服务、短信推送服务、风控服务、AI 文案服务)统一注册到 Nacos;
- 第三方短信渠道密钥、营销限流阈值、QLExpress 规则开关、AI 模型密钥全部统一托管在配置中心。
服务注册与发现(AP 架构)
核心概念
- 服务提供者:短信推送服务、风控服务等业务服务,启动时向 Nacos 上报自身 IP、端口、服务名;
- 服务消费者:营销后台服务,需要调用短信推送接口时,从 Nacos 拉取所有可用服务实例;
- 健康检测:Nacos 定时发送心跳,长时间无心跳的服务实例自动剔除,不会转发请求到故障节点。
AP 架构特性(服务注册选用 AP)
CAP 理论中,AP 代表高可用 + 分区容错,牺牲强一致性。
为什么注册中心选 AP:
营销高峰期,哪怕短暂数据不一致,也不能让服务注册功能瘫痪;多节点集群,一台 Nacos 宕机,其余节点仍能正常提供注册、查询服务,保证业务不中断。
项目落地场景
短信平台集群部署多台推送服务,流量负载均衡分发;某一台推送服务宕机,Nacos 自动剔除实例,Feign 远程调用不会路由到故障节点,避免大量短信下发失败。
动态配置中心(CP 架构)
核心功能
集中管理项目所有配置,不用分散在每个服务 yml 文件;支持配置分环境隔离(开发 / 测试 / 生产)、配置分组区分业务模块。
关键注解@RefreshScope:配置修改后,无需重启服务,Spring Bean 自动刷新配置,实时生效。
CP 架构特性(配置中心选用 CP)
CP 代表强一致性 + 分区容错,牺牲部分可用性。
配置(渠道密钥、额度、风控规则)属于核心敏感数据,必须保证所有服务读取到的配置完全一致,不允许出现 A 服务读取旧配置、B 服务读取新配置的情况,因此采用 CP 模式保证数据统一。
项目落地配置内容
- 第三方短信渠道账号、API 密钥、请求地址;
- 短信下发限流阈值、单日用户最大发送次数;
- QLExpress 规则引擎开关、黑白名单拦截阈值;
- SpringAI 大模型接口地址、Token 消耗上限、AI 降级开关;
- Seata、Sentinel 中间件参数。
注册中心 vs 配置中心 核心区别
| 维度 | Nacos 注册中心(AP) | Nacos 配置中心(CP) |
|---|---|---|
| 核心职责 | 管理微服务实例,实现远程调用负载均衡 | 统一管理项目所有业务、中间件配置 |
| CAP 选型 | AP,优先保证高可用 | CP,优先保证数据强一致 |
| 更新方式 | 服务心跳自动上报实例状态 | 后台手动修改配置,推送变更事件 |
| 故障影响 | 单节点宕机,集群仍可正常查询服务 | 集群半数节点不可用,暂时无法修改配置 |
| 业务价值 | 微服务远程调用不路由故障机器 | 修改渠道密钥、风控规则不用重启服务 |
@RefreshScope 底层简单原理
- Nacos 配置变更后,推送事件给服务;
- Spring 监听器捕获配置变更,刷新对应作用域下的 Bean;
- 带有
@RefreshScope注解的类,会重新从配置中心读取最新参数; - 无注解 Bean 不会自动刷新,必须重启服务才能加载新配置。
Nacos 服务注册与发现流程
服务注册流程:
短信推送、风控等微服务启动后注册到 Nacos,持续上报心跳;营销服务调远程接口时,从 Nacos 拉取健康实例做负载均衡;服务宕机心跳停止,Nacos 自动剔除,不会转发请求。
Nacos 配置动态刷新流程
配置刷新流程:
后台修改渠道密钥、限流规则等配置,Nacos 推送变更消息给所有服务;程序通过 @RefreshScope 刷新 Bean,直接读取新配置,不用重启服务。
Elasticsearc
业务背景:早期千万级短信下发日志存在 MySQL,多条件检索、按手机号 / 活动 / 时间范围查询很慢,后面迁移 ES 做日志检索,MySQL 只存核心业务主数据
ES 是什么,核心定位
Elasticsearch 是分布式全文检索引擎,底层基于 Lucene 封装。
核心能力:全文检索、多维条件筛选、海量数据近实时查询。
⚠️重点区分:
ES不适合高频事务写入,不替代 MySQL;适合海量日志、报表、检索类查询。
在短信平台职责:存储海量短信下发日志,支撑运营后台按手机号、活动 ID、下发状态、时间范围多维度追溯短信记录。
核心底层:倒排索引
正向索引(MySQL 的存储思路)
文档 ID → 内容
例:
1 号日志:手机号 138xxxx,状态成功,2026-08-15 下发
2 号日志:手机号 139xxxx,状态失败,2026-08-15 下发
如果查「所有失败短信」,MySQL 要全表扫描 / 走索引,数据量大之后很慢。
倒排索引(ES 核心)
词条 → 文档 ID 列表
把字段内容拆成词条,建立词条和文档的映射
示例:
词条【下发失败】→ [2 号日志]
词条【138xxxx】→ [1 号日志]
词条【139xxxx】→ [2 号日志]
优势:检索时直接根据词条找到对应的文档,不用遍历全表,海量数据下多条件查询速度远快 MySQL。
补充:短信日志不会做分词,手机号、状态这类字段我们设置为keyword精确匹配,不分词。
ES 基础核心概念
| ES 概念 | 对标 MySQL 概念 | 短信项目举例 |
|---|---|---|
| Index(索引) | Database 数据库 | sms_log_index 短信日志索引 |
| Type(旧版,7.x 后废弃) | Table 表 | 7.x 不再使用,了解即可 |
| Document 文档 | Row 一行数据 | 单条短信下发记录(一条日志) |
| Field 字段 | Column 列 | 手机号、活动 ID、下发状态、通道、创建时间 |
| Mapping 映射 | Table 表结构(字段类型) | 定义手机号 keyword、时间 date、消息内容 text |
| Shard 分片 | 分库分表 | 大索引拆分多个分片,分布式存储,水平扩容 |
| Replica 副本 | 数据备份 | 分片副本,节点宕机不丢数据、保证查询可用 |
版本提醒:生产一般 ES7+,已经移除 Type
写入与检索底层简单流程(短信日志场景)
(1)写入流程(短信下发成功后,同步 / 异步写入 ES 日志)
短信消费端下发完成 → 组装短信日志文档 → 写入 ES
- 文档先写入内存缓冲区(Index Buffer),同时写 translog 事务日志(宕机恢复用)
- 定时刷新 refresh:缓冲区生成段文件 segment,进入文件缓存,近实时可查(默认 1s)
- 段文件不断合并(merge),最终刷入磁盘
关键点:ES 不是写入立刻磁盘持久化,translog 保障宕机不丢;refresh 默认 1s,所以叫近实时检索,不是强实时。
(2)检索流程(运营后台查询短信记录)
运营传入条件:手机号 + 时间范围 + 下发状态
- 协调节点接收查询请求,路由到对应分片
- 在分片内基于倒排索引快速匹配符合条件的文档
- 各个分片结果汇总、排序、分页,返回给业务服务
为什么千万级短信日志要从 MySQL 迁移 ES?
- 短信日志属于海量、只查询、很少修改的数据,写入量大,不会更新;MySQL 大表多条件联合查询,B + 树索引效率急剧下降;
- 业务经常不定条件检索:手机号、活动、状态、通道、时间自由组合,MySQL 很难提前建好所有联合索引,索引过多会严重拖慢写入;
- ES 倒排索引天然适合多维检索,灵活组合条件,查询性能稳定;
- 冷热分离:MySQL 只保留近期核心业务数据,历史海量短信日志归档 ES,减轻 MySQL 压力。
限制点:ES 不适合强事务、高频更新场景,所以核心业务数据(商户、活动任务)依旧放 MySQL。
流程图
短信日志写入 ES 流程
运营后台检索短信日志流程
MyBatis-Plus
MyBatis-Plus 是什么
MyBatis-Plus(简称 MP)是 MyBatis 的增强工具,只增强、不修改,完全兼容原生 MyBatis。
核心目标:减少基础 CRUD 重复代码,不用手写简单的 Insert、Update、Delete、基础 SQL。
在短信平台里:商户信息、营销活动任务、黑名单、渠道配置这类 MySQL 业务表,全部使用 MP 开发。
核心常用能力
1.通用 CRUD 封装
BaseMapper内置 selectById、insert、updateById、deleteById,单表基础操作不用写 XML。
2.Lambda 查询构造器
LambdaQueryWrapper、LambdaUpdateWrapper,用 Java 实体类的方法引用,避免硬编码字段名
好处:编译期就能校验字段,字段改名直接编译报错,防止手写字符串字段名出错(比如 "phone" 手敲成 "tel")
3.分页插件
内置分页能力,配置分页插件后,直接selectPage,自动拼接分页 SQL,不用手写 limit。
4.主键策略、自动填充
比如创建时间、更新时间,@TableField(fill = FieldFill.INSERT_UPDATE)自动填充,不用每次 set。
LambdaQueryWrapper 通俗举例(短信项目场景)
需求:查询某个活动下、手机号在黑名单之外、状态为启用的营销任务
- 原生 MyBatis:写 XML,手写字段字符串,容易写错字段
- MP Lambda 写法:直接引用实体方法
SmsActivity::getActivityId、SmsActivity::getStatus编译阶段校验,字段改名直接报错,线上不会出现因为字段名写错导致查不出数据。
MP 适用边界
✅ 适合:单表简单 CRUD、单表分页、简单条件筛选(活动、商户、黑名单这类单表业务)
❌ 不适合:复杂多表联查、复杂统计、自定义复杂 SQL
规范:复杂 join、复杂统计我们仍然写原生 MyBatis XML,MP 只用来简化单表操作。
