ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

企业级消息通知系统高可用架构设计与Spring Boot落地实践

企业级消息通知系统高可用架构设计与Spring Boot落地实践 1. 项目概述与核心需求拆解1.1 什么是企业级消息通知系统消息通知系统说白了就是解决“业务系统如何把一条消息可靠地送到用户手里”这个问题的中间层服务。你点外卖、收验证码、接到订单状态变更提醒背后都是一套通知服务在干活。很多团队初期图省事直接在业务代码里写死一个函数调短信接口或邮件服务等业务跑起来了才发现问题一大堆每天几百万条通知把第三方供应商的配额打爆了、其中一个通道挂了导致所有通知全部丢失、产品临时要新增一个钉钉机器人通知结果要改十几个业务方代码。我这篇文章要讲的就是如何从零搭建一个能支撑企业级规模的消息通知服务核心围绕三条主线高可用、可扩展、可观测。高可用解决通道故障和流量洪峰下的存活问题可扩展解决新通道接入和业务方接入的效率问题可观测解决出了问题能快速定位的问题。这套方案适用于电商平台、SaaS产品、企业内部协同工具的通用通知需求架构上以Spring Boot Redis RabbitMQ为底座不做任何商业组件依赖团队有Java后端基础就能落地。1.2 通知系统的业务边界与通用场景企业级通知服务一般要承担三类场景第一类是交易类通知比如订单支付成功、退款到账、物流状态变更这类通知对可靠性要求最高漏发一条可能引发客诉甚至合规风险第二类是营销类通知比如优惠券到期提醒、活动上线推送这类通知量大、时效要求相对宽松但必须考虑用户退订意愿第三类是运维类通知比如告警触发、工单流转这类通知通常走IM或Webhook要求分钟级触达。在设计系统时我习惯先把三类场景的SLA指标定下来交易类要求全年可用性不低于99.95%单条推送到达率不低于99.99%营销类允许批量发送但有频控约束运维类要求通道故障自动切换时间不超过30秒。很多团队一上来就写代码后面返工十有八九是因为没有事先定义这些指标。指标定清楚了架构选型才有依据。1.3 项目初始评估与技术选型依据在动手编码之前先评估一下核心依赖组件的选型。通知服务的本质是“消息转储 路由分发 渠道适配”所以关键组件是消息队列、缓存、数据库和任务调度框架。我选型的思路是这样的消息队列使用RabbitMQ而不是Kafka虽然Kafka吞吐更高但通知场景的核心诉求是“每条消息不丢、可重试”RabbitMQ的ack机制、死信队列、延迟队列生态更成熟中小团队维护成本更低。缓存统一用Redis既要支撑消息去重、计数限流也要承担频控令牌桶的存储Redis单实例性能足够生产环境用哨兵模式做主从高可用。数据库用MySQLPercona分支通知记录表虽然写入量大但业务简单分库分表不需要一上来就做先按天分表即可。任务调度用XXL-Job负责扫表补偿、定时批量发送这样的场景。这套组合的特点就是每个组件听上去都没有什么惊艳的地方但合在一起就是一套很皮实的架构也是我在多个项目里验证过的组合。2. 整体架构设计与关键决策2.1 分层架构与数据流向整个通知系统我把它拆成四层接入层、处理层、发送层、管理面。接入层负责接收业务方请求统一鉴权、参数校验、数据脱敏处理层负责消息的模板渲染、去重、频控、路由决策发送层负责与各个第三方通道短信、邮件、App Push等交互管理面是后台配置包括模板管理、通道配置、重试策略配置。数据流向是这样的业务侧调用HTTP接口或者写入MQ生产一条通知事件事件进入处理层后会先做幂等去重和频控校验通过后把通知任务持久化到MySQL同时投递到RabbitMQ的待发送队列。发送层的消费端拿到任务按路由规则选择通道执行发送执行结果回写数据库失败的任务进入重试队列。以上每一条链路都埋了TraceID方便全链路追踪。这里有一个很多人容易犯的错误把发送动作放在接入请求的同步链路里。用户请求支付成功回调同步去调短信接口等第三方响应结果第三方超时5秒用户的HTTP请求也卡了5秒整个支付链路被拖垮。通知服务必须做到接入异步化同步接口只保证“接单成功”不保证“发送成功”发送动作全发生在异步链路里。2.2 高可用设计的三重保险高可用不能靠口号得靠具体的冗余机制。第一重保险是存储层高可用MySQL一主一从半同步复制Redis哨兵模式三节点MQ做镜像队列。任何一个组件宕机系统都有自动切换的兜底能力不需要人工介入。第二重保险是无状态多副本所有应用节点都是无状态的可任意横向扩容前置负载均衡用Nginx或SLB某一个节点挂了流量自动摘除。第三重保险是通道层容灾短信、App Push这类渠道依赖第三方供应商我会给每个通道配置主备两个供应商路由策略里加入健康检查连续失败N次自动切换到备用通道。这三重保险是叠buff的关系层层兜底。我踩过最大的坑是只做了应用层高可用觉得节点多了就稳了结果短信供应商因为欠费直接停了服务系统里几千条通知积压。后来才加了通道健康检查和自动熔断这种情况才能自动切换。2.3 可扩展性设计的插件化思路可扩展性的核心是“新通道接入不改业务流程代码”。我的做法是定义一套统一的Channel接口所有通道短信、邮件、Push、Webhook、站内信都实现同一套接口各自封装各自的供应商对接逻辑。新接一个通道就是新增一个实现类并注册到Spring容器不需要动任何上层逻辑。public interface NotifyChannel { // 通道唯一编码如 sms、email、push_ios String getChannelCode(); // 发送消息 SendResult send(NotifyMessage message); // 健康检查 boolean healthCheck(); // 当前通道是否支持该消息类型 boolean support(NotifyType notifyType); }接口抽象出来之后还有一个很关键的点消息模型也要统一。业务方不用关心用什么通道发只需要提交一个“通知事件”里面带上业务类型、接收人、模板编号、模板参数具体走哪个通道由服务端路由决定。这样即使未来新增一个“企业微信”通道上游业务方一行代码都不用动。还有一个扩展点容易被忽略模板渲染。不同通道对内容的格式要求不同短信有70字限制邮件支持HTMLPush有标题加正文。模板需要按通道维度独立配置渲染引擎负责把业务参数填充到不同模板里这里我选择的是Java原生SPEL表达式模板不引入额外的模板框架保持轻量。3. 核心细节解析与实操要点3.1 消息数据模型与表结构设计先看数据库设计这是整个系统的地基。通知系统最核心的表是通知任务表我给它命名为notify_task。这个表有几个关键字段biz_type业务类型、biz_id业务幂等ID、template_id模板编号、channel_code路由后的通道编码、receiver接收人、params模板参数JSON、status状态、retry_count重试次数、next_retry_time下次重试时间、trace_id链路追踪ID。状态流转要设计得简洁清晰我用了这几个状态PENDING待处理、SENDING发送中、SUCCESS成功、FAILED失败、DEAD死亡即重试耗尽。每次发送动作只做状态之间的原子转换避免并发情况下状态错乱。建表时需要特别关注索引设计。实践中查询频率最高的是“按状态和下次重试时间扫补偿任务”以及“按业务ID查重”所以至少要建这两个组合索引idx_status_next_retry_time(status, next_retry_time)、idx_biz_type_biz_id(biz_type, biz_id)。顺序别反了反了会导致索引失效。模板和通道配置是低频数据单独建立notify_template、notify_channel_config、notify_channel_route三张表业务上支持动态修改。通道配置里要保存供应商的AK/SK、接口地址、限流阈值等信息这些字段需要加密存储项目里我用的是AES256加密后落库服务启动时解密加载到内存缓存。3.2 模板渲染与参数校验的细节处理模板渲染环节有不少细节讲究。每个模板可能被多个业务方复用参数命名必须收敛规范化。我规定参数名统一为驼峰格式模板里通过${userName}这样的占位符引用。渲染引擎实现时要注意两点一是参数缺失时不要直接抛异常而是用空字符串占位并记录一条WARN日志因为通知不是核心交易链路局部参数错误不能让整条任务失败二是渲染结果要按通道重新校验长度最典型的是短信超过70字要拆条拆条逻辑按字符合并而不是简单截断要处理换行和标点的边界。参数校验的层级要区分开。接入层校验的是参数类型和必填项处理层校验的是业务规则比如同一用户同一类型通知的频控发送层校验的是通道约束比如手机号格式、邮箱格式。这个分层不做的话后面会很痛苦我见过一个团队把所有校验堆在接入层结果加一个“邮件通道不支持发送手机号”的规则要改接入层代码发布上线才能生效完全放弃了动态性。模板数据还需要支持SuperVar多语言或动态配置的替换但这个视团队实际情况不需要过度设计。3.3 去重、频控与防刷的实现策略去重和频控是企业级通知和玩具项目之间最明显的分水岭。去重要做两层接入层的幂等去重和发送层的重复发送控制。幂等去重利用biz_type biz_id做唯一索引重复提交时直接返回原任务ID不产生新任务。发送层的去重更加精细比如短信通道同一个手机号5分钟内收到同一条模板消息的次数不能超过1次我用Redis的SETNX 过期时间来实现。频控则要区分用户维度和通道维度。用户维度限制的是“单个用户每天最多收到多少条营销通知”这个容易理解通道维度限制的是“短信供应商每秒最多接受多少条请求”这个经常被忽视。通道频控我用的Redis令牌桶每个通道一个Key每秒放入固定数量的令牌取不到令牌的消息进入等待队列而不是直接丢弃。这里多说一句防刷和频控如果不做被黑的成本极高。比如攻击者拿到开放接口后批量伪造手机号提交几万条短信费用几个小时内就能把预算打光。所以接入层必须有API级鉴权每个业务方一个AppId 密钥再加IP白名单和总调用量限制。3.4 通道路由与失败切换机制通道路由是通知系统的“大脑”它决定了一条消息应该走哪个通道。路由决策我推荐使用可配置的规则链先根据NotifyType确定通道候选集比如交易验证码优先短信、营销通知优先Push和邮件再根据用户偏好过滤用户在App里主动关闭了Push通知就要剔除Push通道然后按照通道健康状态过滤已熔断的通道剔除最后按优先级权重选择最终通道。路由规则要存储在数据库里而不是写在代码里。这样产品同学调整消息策略时不需要发版后台改配置即时生效。我见过一些团队把路由逻辑写在Java代码里改一次规则发一次版这种效率不可接受。失败切换的兜底逻辑是发送失败时如果存在备用通道且失败原因不是“消息内容本身有问题”比如模板渲染失败而是“通道不可用”超时、限流、供应商报错则自动切换到备用通道重试。切换的同时要发送一条内部告警提醒运维处理主通道故障。这个降级策略能够让交易类通知在短信通道故障时自动转走邮件或者站内信保住最低限度的触达能力。4. 实操过程与核心环节实现4.1 基础环境准备与组件部署我先梳理一下基础环境的部署以一台4核8G的Linux服务器为例生产建议至少三台同配置机器组成集群。操作系统使用Rocky Linux 9或Ubuntu 22.04均可Java环境用OpenJDK 17这部分不做过多赘述。MySQL、Redis、RabbitMQ三个中间件我倾向于用Docker Compose统一编排管理虽然是生产环境但只是这三个基础组件容器化的运维效率更高。需要特别注意的是Redis部分生产环境不要裸跑单实例至少要开启主从哨兵模式。我的Redis高可用部署方案是1个主节点、2个从节点、3个哨兵进程哨兵配置里quorum设为2表示至少2个哨兵同意才能判定主节点故障并触发切换。这个方案我跑过很长时间稳定性非常可靠。RabbitMQ开启镜像队列在Management UI或者通过Policy设置ha-mode: all保证任何一个MQ节点宕机消息不丢失。给通知服务单独建两个队列notify.task.queue用于待发送任务notify.dead.queue用于延迟重试和死信处理。4.2 应用工程结构与核心代码实现应用工程采用标准的Maven多模块结构我按职责拆成四个模块notify-api、notify-core、notify-channel、notify-admin。API模块给业务方提供SDK的Feign接口定义Core模块承载核心处理流程Channel模块放各个通道实现Admin模块是后台管理界面。核心处理流程的入口是NotifyService.submit()方法流程如下先做接口鉴权再校验参数然后根据biz_type biz_id到Redis做幂等检查通过后持久化任务并投递到MQ。关键代码如下实际项目里我还加了全局异常处理和日志埋点这里只展示核心逻辑。Service public class NotifyServiceImpl implements NotifyService { Autowired private NotifyTaskMapper taskMapper; Autowired private StringRedisTemplate redisTemplate; Autowired private RabbitTemplate rabbitTemplate; Transactional public Long submit(NotifyRequest request) { // 1. 鉴权校验 authService.checkAppAccess(request.getAppId()); // 2. 幂等去重 String idempotentKey request.getBizType() : request.getBizId(); Boolean firstSubmit redisTemplate.opsForValue() .setIfAbsent(idempotentKey, 1, Duration.ofHours(24)); if (Boolean.FALSE.equals(firstSubmit)) { throw new DuplicateSubmitException(重复提交的通知请求); } // 3. 持久化任务 NotifyTask task buildTask(request); taskMapper.insert(task); // 4. 发送到MQ待处理队列 rabbitTemplate.convertAndSend(notify.exchange, notify.task.queue, task); return task.getId(); } }4.3 发送端消费逻辑与重试机制发送端的消费逻辑要特别注意一定不要用默认的自动ack模式改用手动ack。手动ack的含义是MQ推送消息给消费者消费者处理完成后再告诉MQ“这条消息我处理好了”MQ才把消息从队列中删除如果没有ackMQ会把消息重新投递给其他消费者。这样能保证消息不会因为消费者宕机而丢失。消费处理器拿到任务后调用路由服务选择通道再调用通道发送。如果通道发送成功更新任务状态为SUCCESS并提交ack如果失败判断重试次数是否达到上限没达到则把任务投递到延迟队列延迟时间按照指数退避计算第一次30秒、第二次1分钟、第三次5分钟达到上限则标记为DEAD并走人工补偿流程。这块有个很容易出问题的点消费端是并发消费的同一个手机号可能有两条不同业务的通知同时进入发送流程如果通道实现不是线程安全的会产生数据竞争。比如某些短信供应商的SDK内部用一个非线程安全的连接池并发调用时会报错。我的建议是每个通道实例内部单独持有连接资源用ThreadPoolExecutor做隔离避免一个通道卡死影响其他通道。延迟队列在RabbitMQ里的实现原理是设置消息的expiration属性让消息在队列中待够指定时间才能被消费。这里有一个需要注意的坑同一个队列里的消息如果设置了不同的过期时间RabbitMQ只会看队首消息的过期时间所以如果队首消息过期时间是5分钟后面的消息过期时间即使是30秒也要等队首消息时间到了才会被扫描。解决办法是按照延迟级别拆成多个独立队列每个队列只处理固定的延迟时间。4.4 通道接入实站以阿里云短信为例我用短信通道作为例子完整走一遍通道接入流程。首先在notify-channel模块里新建AliyunSmsChannel类实现NotifyChannel接口。getChannelCode()返回sms_aliyunsend()方法里调用阿里云短信SDK发送发送之前要做手机号格式校验发送之后根据SDK返回结果构造统一的SendResult。Component public class AliyunSmsChannel extends AbstractSmsChannel { Value(${sms.aliyun.accessKeyId:}) private String accessKeyId; Value(${sms.aliyun.accessKeySecret:}) private String accessKeySecret; Value(${sms.aliyun.signName:}) private String signName; Override protected SendResult doSend(NotifyMessage message) { DefaultProfile profile DefaultProfile.getProfile(cn-hangzhou, accessKeyId, accessKeySecret); IAcsClient client new DefaultAcsClient(profile); CommonRequest request new CommonRequest(); request.setSysMethod(MethodType.POST); request.setSysDomain(dysmsapi.aliyuncs.com); request.setSysVersion(2017-05-25); request.setSysAction(SendSms); request.putQueryParameter(PhoneNumbers, message.getReceiver()); request.putQueryParameter(SignName, signName); request.putQueryParameter(TemplateCode, message.getTemplateCode()); request.putQueryParameter(TemplateParam, message.getTemplateParams()); try { CommonResponse response client.getCommonResponse(request); MapString, Object result JSON.parseObject(response.getData()); if (OK.equals(result.get(Code))) { return SendResult.success(); } return SendResult.fail(String.valueOf(result.get(Message))); } catch (Exception e) { return SendResult.fail(e.getMessage()); } } }通道接入的关键是把供应商的差异全部隔离在通道实现类内部。业务层只认SendResult的success和fail不关心通道内部调用了什么接口、用了什么认证方式。这样接新通道的工时基本可以控制在半天以内无非是照葫芦画瓢再实现一个类。4.5 消息积压与补偿机制的兜底方案即使架构设计得再完美在大促或者突发流量面前消费速度仍有可能赶不上生产速度造成消息积压。积压并不可怕可怕的是没有积压的感知和补偿机制。我做的方案是监控RabbitMQ队列中待消费消息的数量超过阈值就触发告警并自动扩容消费者实例手动扩容或K8s的HPA策略。同时维护一个补偿任务每隔5分钟扫描一次notify_task表中状态为PENDING且next_retry_time小于当前时间的任务重新投递到MQ队列。这个补偿任务能兜底处理一切因为系统重启、MQ宕机等原因遗漏的消息。补偿任务的设计上要注意不要和正常的消费任务冲突。扫描SQL必须带上next_retry_time字段条件并且一次最多取500条避免一把梭查全表导致数据库压力过大。处理完的任务要立即更新时间戳防止被下一次扫描重复捞起。5. 常见问题与排查技巧实录5.1 高频故障案例提炼我在多个项目的生产环境里总结了一些高频故障整理成速查表方便大家在项目上线后对照排查故障现象可能原因排查命令/工具解决方案消息大量积压、消费速率极低消费线程数过小或某一通道阻塞rabbitmqctl list_queues、查看线程池活跃度调大消费线程数通道线程池隔离通知延时明显增大模板渲染耗时过高或频控等待Tracing查看耗时Redis Key数量检查模板渲染加缓存频控改为异步等待短信丢失率高供应商限流被拒查看通道返回码检查限流阈值开启备用通道自动切换重复推送消费端业务处理超时MQ重投导致查看任务表biz_id是否有重复发送前增加Redis去重检查数据库压力大黄金时段频繁扫表补偿慢SQL日志、SHOW PROCESSLIST调整扫描周期分批处理通知内容乱码模板参数编码问题查看渲染日志确认字符串编码统一UTF-8编码并在接入层强校验5.2 排查链路追踪与日志规范排查问题最重要的工具是链路追踪和结构化日志。我的做法是在接入层生成一个trace_id用UUID或雪花算法然后在处理层、发送层每个关键模块的日志中都打印这个trace_id。这样一条消息从提交到成功全链路日志可以用一个ID串起来查问题效率提升非常明显。日志规范上我要求核心模块的WARN和ERROR日志统一用JSON格式输出带上时间戳、trace_id、biz_type、channel_code等字段方便接入ELK或者Loki做日志检索。一个时刻大量ERROR日志报警时用tail -f加grep trace_id的方式能快速定位到具体是哪一批消息在哪一步出了问题。5.3 高可用切换的真实场景复盘有一次生产环境中短信通道的底层供应商因为系统升级频繁超时成功率从99%掉到70%。我的通道健康检查逻辑检测到连续失败30次后自动熔断并切换到了备用短信供应商整个过程没有人工干预。事后复盘发现两个问题一个是健康检查的失败阈值设得偏保守导致的后果是熔断触发前已经有一小部分消息失败了另一个是切换过程中正在发送中的消息因为通道状态不一致丢失了几十条。后来我做了两个改进第一健康检查的窗口从“连续失败N次”改成“滑动窗口内失败率超过50%”这样能更快感知故障第二通道切换动作不直接在发送线程里做而是通过一个独立的管理线程完成切换前要等正在发送的消息返回结果或者超时保证发送中的消息不会被白白丢弃。5.4 性能压测与容量评估记录系统上线前一定要做压测不能拍脑袋觉得“够用”。我的压测方案用JMeter模拟业务方调用目标场景是“峰值每秒5000条通知提交持续10分钟”。测试结果基本符合预期单机4核8G部署的应用节点TPS稳定在3000左右CPU利用率约70%瓶颈出现在Redis的频控检查上单个Redis实例的GETSET操作在高并发下毛刺明显。优化方案是频控Key的拆分和本地缓存。每个通道的频控计数从“Redis每秒一个Key”改成“Redis每5秒一个Key 本地内存补一次”这样Redis压力直接降到原来的五分之一。另外压测时建议把日志级别调到WARN不然大量INFO日志的IO开销会占掉不少CPU。容量评估方面我给的参考公式是单机支撑能力TPS× 节点数 × 安全系数0.7 ≥ 预估峰值TPS。比如预估峰值5000单机3000那至少需要3台节点。安全系数0.7是为故障节点预留的冗余不能算得太满否则一台机器宕机整个集群就扛不住了。6. 项目落地与经验总结6.1 从0到1的落地路线图如果你所在团队想要落地这样一套通知服务我的建议是分四个阶段走。第一阶段搞定MVP实现HTTP接入、RabbitMQ异步处理、短信和邮件两个通道能够完成“提交通知任务 → 发送成功”的基础闭环。第二阶段完善可靠性加上Redis去重频控、失败重试、死信补偿确保消息不丢不重复。第三阶段提升可扩展性完善后台管理功能模板做到按通道独立配置路由规则做到动态可调整。第四阶段做高可用和可观测部署多个应用节点中间件全部高可用接入监控告警和链路追踪。每个阶段的验收标准也要提前想好。第一阶段是“能发出去”第二阶段是“发得可靠”第三阶段是“改得快”第四阶段是“挂不了”。6.2 关于n8n等自动化工具的融合思考在做完自研通知服务之后我其实还调研过n8n这类自动化工作流工具的企业级部署方案。它们擅长的是把不同系统的动作编排成自动化流程比如当CRM里新增一个线索自动发钉钉通知并建一条企微待办。这类工具和自研通知服务并不冲突甚至是一个很好的补充自研服务负责“中心化的通道能力和可靠性保障”n8n负责“轻量化的业务自动化编排”。如果你团队里有很多需要“事件触发 → 多步骤动作”的场景可以考虑在通知服务之上再套一层工作流引擎让业务人员也能自助编排通知行为。不过n8n的节点执行效率和一些高级特性可能与自研服务有差距大规模、高并发场景还是自研服务更可控。这个取舍就看团队业务的侧重点了。6.3 维护成本与团队能力要求最后聊一下维护成本和团队要求。这套系统落地后日常维护核心就是三块中间件的健康巡检、通道供应商的配额和价格监控、模板和路由规则的变更配置。代码层面的改动频率并不会太高主要集中在新通道接入和异常逻辑优化。对团队的能力要求是至少有一名成员熟悉RabbitMQ和Redis的原理和运维有两名成员能看懂核心业务代码并做代码Review。如果团队连这些基础能力都不具备我建议先不要上自研用云厂商的现成推送服务过渡会更稳妥。还有一个心得通知系统的代码量虽然不大但它的跨团队协作属性很强。上线前一定要和业务方、客服团队对齐“通知延迟超过XX分钟应该找谁”“短信发送失败用户投诉应该走什么流程”这些非技术问题。很多项目技术方案都很好最后因为协作流程没定清楚系统上线后被投诉淹没。我个人在实际项目里最大的体会是做通知系统细节决定成败。一个模板参数没配上、一个频控阈值设少了、一个通道没做健康检查生产上就是真实的资损和客诉。架构上可以借鉴成熟方案但真正的功力都藏在那些不起眼的重试策略、路由降级、日志规范里。希望这篇文章能让大家少走一些我之前走过的弯路把消息通知这个“小系统”做出真正可靠的企业级水准。
返回列表