有时候后台系统里最不起眼的模块,反而是线上事故最容易爆发的点。消息通知服务就是一个典型——你平时觉得它简单,不就是把一条消息塞给用户吗?可一旦业务量上来,短信、邮件、站内信、App推送,各渠道商接口不稳定,上游业务方一窝蜂发通知,偶发连不上数据库、渠道商响应超时、消息积压把内存打爆,各种问题一夜之间全来了。
我刚接手公司通知系统重构时,线上已经在用一套很“朴素”的实现:订单状态一变,业务代码里直接调sendMail()、sendSms(),发送失败就打印一条错误日志。初期没问题,后来一天几百万条通知,数据库连接被占满,渠道商接口频繁限流,系统互相拖垮,数据库一慢整个订单服务全挂。后来我花了两个多月把它重构成了高可用、可扩展的企业级通知服务,今天这篇就把整个设计思路、核心代码、踩过的坑全部整理出来,给准备做或正在做类似系统的同学一个可以直接参考的落地路径。
这篇文章适合两类人看:一类是后端开发,想了解企业级通知系统的架构长什么样;另一类是技术负责人或架构师,正在做系统拆分的选型评估。内容会涉及核心技术选型的底层逻辑、消息模型设计、投递引擎实现、高可用部署方案,以及我在真实生产环境里碰到的问题和排查方法。
1. 整体设计与需求拆解
1.1 企业级通知系统的核心需求是什么
在动手写代码之前,必须先把需求讲清楚。很多人做通知系统第一步就错了——上来就挑框架、写发送逻辑,结果做到后面发现扩展性差、扛不住流量、事故频发,再回头重构成本极高。
企业级通知系统,核心要解决四件事:
可靠性。这是最硬核的指标。消息只要进了我们的服务,就尽量不能丢。比如用户下单成功就得收到通知,漏发一条可能引发投诉甚至资损。我们需要做到状态可追踪、失败可重试、重试有上限。
发送能力可扩展。通知渠道天然是多样化的,短信、邮件、站内信、App Push、Webhook,今天只有邮件,明天可能就要接钉钉、企业微信,后天可能又要接第三方国际短信。系统设计要允许“加新渠道像加插件一样简单”,而不是改一堆老代码。
业务接入简单。通知这件事对上游业务方来说只是“捎带手的事”,不能要求人家懂我们的内部细节。最好提供一个统一的API,业务方传一个场景码和一批用户ID就行。不同场景的模板、渠道、限流规则、重试策略,全部在通知服务里配置化处理。
抗压能力。双十一大促、秒杀活动、营销推送,一个活动通知可能瞬间产生几百万条发送任务。系统要能削峰填谷,不能上游一抖我们就崩,我们一崩又把数据库拖崩。
这四条不是并列关系,可靠性是底线,扩展性决定系统能走多远,接入简单决定业务方愿不愿意用,抗压能力决定系统上线后运维稳不维稳。整套架构的设计都是围绕这四个目标展开的。
1.2 从单机定时任务到消息驱动架构
我们最初的通知服务是什么形态呢?一张notification_record表,一个Spring定时任务每30秒扫一次状态为PENDING的记录,然后逐条发送,发送成功后UPDATE状态。听起来还行,实际上三个问题:
一是扫描式处理是有上限的。每30秒扫一次,数据库里积压的记录越多,单次扫描越慢,CPU和数据库连接都消耗在无关数据的查询上。一旦单日消息量上百万,这张表的数据量到了千万级,定时任务基本就废了。
二是发送和业务强耦合。系统里直接依赖了邮件服务器、短信网关。一旦某个渠道商接口变更,或者需要增加一种新渠道,只能改代码、发版、重启,线上流程又长又慢。
三是可靠性没法保证。进程重启、机器宕机、数据库连接池满,都会导致正在处理的消息状态丢失,遗漏发送,而且没人知道漏了哪些。
企业级方案普遍选择消息驱动加事件驱动的架构,把“通知请求”和“通知发送”解耦。业务方调用我们的API,我们立刻落库,同时把任务ID扔进消息队列。发送器从队列拉取任务,按渠道分发,异步推进状态。核心链路变成:
业务方HTTP调用 -> 消息落库 -> 写入MQ -> 发送Worker消费 -> 调用渠道商 -> 更新发送状态很多人不理解,为什么多引入一个MQ,链路长了,延迟高了,不是更麻烦吗?关键在于,MQ不是用来传“消息内容”的,它解决的是流量削峰和故障隔离。
典型的场景是营销通知。某天上午10点业务方开启全员推送,10分钟内涌入50万条请求。如果不经过MQ,50万条发送请求同时打给下游短信网关,通道直接被打爆,发送大量失败。接入MQ之后,Worker按固定的消费速率去处理,比如每秒200条,多余的消息在MQ中排队等待。下游渠道商看到的流量是平稳可控的,不会出现过载。即便某个渠道商故障,消息也只是在队列里堆积,不丢消息,等渠道恢复后继续消费。
1.3 技术选型背后的思考
选型这事没有银弹,完全取决于团队技术栈和部署环境。我这次重构基于Spring Boot 3.x,这是当前Java后端最主流的选择,生态成熟、资料丰富,团队招人也好招。
消息队列我选了RocketMQ,主要是看重它的事务消息能力——通知服务有“落库 + 发MQ”的原子性需求,RocketMQ的事务消息能保证这两步要么都成功、要么都失败,不会出现库里有记录但MQ里没有任务的情况。如果团队已经重度使用Kafka,用Kafka也能做,但“本地消息表 + 定时对账”这类方案需要自己实现,复杂度会高一些。
存储层用的MySQL,部署上做了主从同步和读写分离。后文会详细讲为什么需要配合Redis做二级缓存,以及短期窗口优化策略。
Redis在这里承担三个角色:分布式锁(保证同一个任务的并发幂等)、短时去重(防止同一用户短时间重复收到相同通知)、发送频率控制(限制同一用户或同一场景的发送速率)。
至于高可用,Redis必须上哨兵或Cluster模式,不能单机裸奔,这个后面会作为重点展开。
2. 消息模型与存储设计
2.1 统一消息模型:不绑定任何具体渠道
很多通知系统越做越难扩展,根源在于一开始就把消息模型和发送渠道绑死了。表里有email字段、phone字段、push_token字段,新增一个渠道就往表里加一列,最后表膨胀到几十列,代码里全是if-else判断渠道类型。
我做设计时强制统一消息模型,核心只有一张宽表,加一个内容JSON字段:
CREATE TABLE `notification_message` ( `id` bigint NOT NULL AUTO_INCREMENT, `biz_id` varchar(64) NOT NULL COMMENT '业务方幂等ID', `scene_code` varchar(64) NOT NULL COMMENT '场景码,如 ORDER_PAID', `channels` varchar(128) NOT NULL COMMENT '需要发送的渠道,逗号分隔: SMS,EMAIL,IN_APP', `title` varchar(255) DEFAULT NULL, `content` text COMMENT '统一渲染后的内容', `template_id` varchar(64) DEFAULT NULL, `template_params` text COMMENT '模板参数JSON', `status` tinyint NOT NULL DEFAULT '0' COMMENT '0待发送 1发送中 2部分成功 3成功 4失败 5已取消', `max_attempts` tinyint NOT NULL DEFAULT '3', `attempted` tinyint NOT NULL DEFAULT '0', `next_retry_time` datetime DEFAULT NULL, `trace_id` varchar(64) DEFAULT NULL, `created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, `updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, PRIMARY KEY (`id`), KEY `idx_status_retry` (`status`, `next_retry_time`), KEY `idx_biz` (`biz_id`), KEY `idx_created` (`created_at`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;这里几个关键设计点:
biz_id是业务方传入的幂等ID,比如订单号加通知类型的组合。数据库对这个字段建唯一索引,重复请求直接丢弃或返回已有结果。
scene_code是解耦的核心,业务方不用关心发什么渠道、用哪个模板。这些映射规则在配置中心维护,ORDER_PAID该走短信还是邮件,由通知服务决定。
template_params不存渲染后文本,存原始参数,这样后续要换模板、要多语言,都能基于参数重新渲染,不需要业务方改代码。
status从0到5的流转是状态机设计的核心,后面展开。
这张表解决的是“消息从哪里来、到哪里去、现在什么状态”的问题。但它本身存储的还是“一条通知请求”,而不是“每个渠道的发送状态”。同一个消息需要发短信和邮件,这两个发送任务是并行的,状态不同步更新。因此还需要一张渠道明细表。
CREATE TABLE `notification_channel_record` ( `id` bigint NOT NULL AUTO_INCREMENT, `message_id` bigint NOT NULL, `channel` varchar(32) NOT NULL, `status` tinyint NOT NULL DEFAULT '0' COMMENT '0待发送 1发送中 2成功 3失败', `channel_message_id` varchar(128) DEFAULT NULL COMMENT '渠道商返回的消息ID', `error_code` varchar(64) DEFAULT NULL, `error_msg` varchar(512) DEFAULT NULL, `attempted` tinyint NOT NULL DEFAULT '0', `send_at` datetime DEFAULT NULL, PRIMARY KEY (`id`), KEY `idx_message` (`message_id`), KEY `idx_channel_status` (`channel`, `status`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;主表状态为2部分成功时,明细表里可能短信成功、邮件失败,下次重试只重试失败的邮件。这个设计后面讨论重试策略时还会用到。
2.2 Redis在通知链路里的三个不可替代的作用
MySQL是最终数据底座,但直接在业务高峰期把大量状态更新打在MySQL上,数据库很快成为瓶颈。Redis在通知链路中承担的职责,很多系统都没用好。
第一个作用是分布式锁。通知系统是多实例部署的,同一个MQ消息可能被两个Worker同时消费。如果不加锁,同一条消息会被发送两次。用Redis的SET NX EX实现一个简单的分布式锁:
public boolean tryLock(String key, String requestId, long expireSeconds) { String result = redisTemplate.opsForValue() .setIfAbsent(key, requestId, Duration.ofSeconds(expireSeconds)); return Boolean.TRUE.equals(result); }消费消息时以messageId为锁的key,拿到锁才处理;处理完删除锁。锁的过期时间根据处理耗时设定,比如短信发送最长10秒,锁过期设15秒,避免处理时间过长锁自动释放导致并发问题。
第二个作用是短时去重。业务方接口写得不严谨,同一个事件可能重复回调。比如订单支付成功,支付回调处理程序重试了三次,通知请求也发了三次。去重方案依赖数据库的唯一索引,但每次请求都去数据库查一次,浪费连接资源。更好的做法是Redis缓存最近N分钟已处理的bizId,命中缓存直接返回上一次的处理结果:
public boolean isDuplicate(String bizId) { return Boolean.TRUE.equals(redisTemplate.hasKey("NOTIFY:DUP:" + bizId)); } public void markProcessed(String bizId) { redisTemplate.opsForValue().set("NOTIFY:DUP:" + bizId, "1", Duration.ofMinutes(30)); }实际处理时,bizId的去重是双保险:先查Redis,没有再查数据库唯一索引,两条路都过了才真正创建消息记录。
第三个作用是频率控制。有些用户对通知很敏感,10分钟内连续下三单,每单发一条短信,用户很容易烦躁,甚至投诉渠道商造成通道被封。系统需要支持“频控”能力,按用户维度限制发送量:
public boolean rateLimit(String userId, int maxCount, int seconds) { String key = "NOTIFY:RATE:" + userId; Long count = redisTemplate.opsForValue().increment(key); if (count != null && count == 1) { redisTemplate.expire(key, Duration.ofSeconds(seconds)); } return count != null && count <= maxCount; }这里INCR加EXPIRE两步不是原子的,极端情况可能出现并发问题,但对通知系统这种非强一致场景,误放行一两条比误拦截一两条更可接受。如果要求更强保障,可以用Lua脚本原子执行。
Redis在高可用方案上的选择要单独说。企业环境我不会建议只部署单机Redis。单点Redis一旦宕机,分布式锁失效、去重失效、频控失效,高并发下垃圾消息直接击穿数据库。推荐至少上Redis Sentinel(哨兵)模式,一主两从三哨兵,自动故障切换。数据量特别大、需要横向扩展的场景就上Redis Cluster。
我当时线上用的是一主两从加三哨兵,主节点挂了哨兵自动把从节点提升为主,整个切换过程对业务方基本无感。真正要小心的反而是客户端配置,必须配置正确的哨兵地址和masterName,别把哨兵地址当成Redis地址直连,那是新手最容易踩的坑。
2.3 数据库防击穿:短期窗口优化策略
前面说Redis承担了去重和频控,但还有一个隐患:业务尖峰时刻,大量请求同时打到数据库,MySQL的连接池不够用,连接等待导致接口RT飙升,接着触发上游服务超时重试,重试又带来更多流量,最终拖垮整个服务。为了避免这种情况,有几个非常实用的策略:
写请求异步化,接口快速返回。创建通知记录这步,本身是必须落库的。但并不意味着业务方必须同步等数据库写入完成。我在入口层做了两段式处理:第一段,收到请求后先校验参数、做去重检查,然后直接返回成功;第二段,通过MQ异步通知Worker真正落库和发送。这里有一个取舍——如果下游数据库宕机,接口已经返回业务方成功了,但消息并没有真正落库,存在丢消息风险。所以我落库的方式是全同步,接口必须等记录写入成功才返回。但通过“批量插入”来减少数据库交互次数——业务方可能同时给1000个用户发通知,循环单条插入会有1000次网络往返,改用INSERT ... VALUES (), (), ()一条SQL批量插入,性能提升非常明显。
对查询场景做强隔离。通知系统的读场景主要来自后台管理页面,运营想查“这个订单有没有给用户发短信”。这种查询不能和写入抢连接池。数据源拆成两个:写库专用连接池配置稍小,读库连接池配置稍大。MyBatis或MyBatis-Plus的@DS注解可以方便地实现读写分离。
短期窗口合并写。这是我从交易系统借鉴过来的经验。通知记录表的高频写入集中在几个固定时段。我把某段时间内针对同一场景的写操作合并为一个事务批量提交,虽然对整体并发提升有限,但显著减少了行锁竞争和日志落盘次数,高峰期数据库负载能降低近三成。
3. 投递引擎与渠道网关的架构实现
3.1 两级队列加动态分发:让每个渠道互不拖累
前面讲过,MQ的核心价值是削峰和故障隔离。但这里要明确:一个MQ Topic给所有渠道共用的方案,耦合非常严重。短信通道慢、邮件通道快,如果共用一个Topic,慢渠道的消息会阻塞快渠道的处理。
我采用的是一级Topic加二级工作队列的分发架构:
业务消息进入 NOTIFY_JOB_TOPIC | v 分发Worker(按渠道类型) | +--> 短信任务队列(独立消费组) +--> 邮件任务队列(独立消费组) +--> 站内信任务队列(独立消费组)分发Worker消费NOTIFY_JOB_TOPIC里的消息,根据消息携带的渠道列表,分别投递到对应的渠道队列。每个渠道使用独立的消费组,消费组之间互不干扰。
这样设计带来的直接好处是:某个渠道商故障导致短信消费组积压,邮件消费组完全不受影响,站内信照常发送。同时,每个消费组可以配置不同的消费线程数——短信渠道商限流严,消费线程调小;邮件服务性能高,消费线程调大。能实现这一点,依赖的正是RocketMQ的Tag过滤能力,或者Kafka的独立Topic方案。
Worker处理消息时,真正的发送逻辑全部通过策略模式实现,不写一长串if-else:
public interface NotificationChannel { String channelType(); SendResult send(NotificationMessage message); } @Component public class SmsChannel implements NotificationChannel { @Override public String channelType() { return "SMS"; } @Override public SendResult send(NotificationMessage message) { // 调用短信服务商API } } @Component public class EmailChannel implements NotificationChannel { @Override public String channelType() { return "EMAIL"; } @Override public SendResult send(NotificationMessage message) { // 调用邮件服务商API } }在Worker里通过Spring注入的List<NotificationChannel>,根据消息的渠道类型动态选择对应的实现类。新增一个渠道(比如钉钉机器人),只需要新增一个实现类,注册为Spring Bean,不改任何老代码。
3.2 发送状态机:从待发送到终态的完整流转
通知记录的状态不能乱跳,必须有明确的状态机约束。我们的设计是:
待发送(0) -> 发送中(1) -> 成功(3) | -> 部分成功(2) -> 重试 -> 发送中(1) -> 成功(3) / 失败(4) | -> 失败(4) -> 重试 -> 发送中(1) -> 成功(3) / 失败(4) |-> 已取消(5)为什么要区分部分成功?因为一条消息会发多个渠道,短信成功、邮件失败,整体状态不能算成功也不能算失败。我做了一个状态聚合逻辑:每个渠道发送完成后上报结果,主表汇总所有渠道结果,只要有一个渠道成功,状态就是部分成功;所有渠道都成功,才是成功;所有渠道都失败,则是失败。
这个消息状态机在重试逻辑里扮演了核心角色,也是高可用架构里可靠性目标的具体载体。关于重试,很多初学同学会把重试写成“发送失败后立刻重新发送”。这是不对的。如果渠道商都返回服务不可用,立刻重试大概率还是失败,反而给对方增加压力。必须用指数退避加抖动:
public long calcRetryDelay(int attempt, long baseDelayMs) { long expDelay = baseDelayMs * (long) Math.pow(2, attempt - 1); long jitter = ThreadLocalRandom.current().nextLong(0, 1000); return expDelay + jitter; }第一次重试等待30秒,第二次1分钟,第三次2分钟,最多重试3次。每次重试之间加一个随机抖动,避免大量重试消息同时触发,形成重试风暴。
3.3 平平无奇的限流实现,如何拦住99%的渠道商封禁
渠道商都不是无限吞吐的。很多系统上线后突然被短信服务商暂停服务,原因就是没有做发送限流,短时间请求量太大,触发了通道风控。
限流必须做在通知服务内部,而不是依赖渠道商的风控。我们按两个维度限流:
全局通道速率限制。每个渠道商都有一个建议的QPS上限,短信服务商给的是每秒200条。我们使用Redis配合令牌桶或滑动窗口算法实现全局限流。
令牌桶的实现非常简单,但需要注意多实例并发时,需要把计数器放在Redis里,而不是本地内存:
public boolean acquireToken(String channel, int qps) { String key = "NOTIFY:LIMIT:" + channel; Long current = redisTemplate.opsForValue().increment(key); if (current != null && current == 1) { redisTemplate.expire(key, Duration.ofSeconds(1)); } return current != null && current <= qps; }这里的语义是“窗口内计数不超过QPS”,实现的是固定窗口算法,业务量特别大时,边界处会有突刺(比如正好跨秒边界,前一秒最后100条、后一秒最开始的100条连在一起,形成瞬时200条)。要求更平滑的场景需要上Lua脚本实现滑动窗口或令牌桶。
业务方维度限流。某个业务方的接口突然异常,疯狂调用通知服务,不能让它拖垮整个平台。为每个业务方配置调用配额,超过配额直接返回“频率超限”错误码。类似之前按用户限流的实现,只是key从用户ID换成业务方AppId。
public boolean acquireAppQuota(String appId, int quota) { String key = "NOTIFY:APP_QUOTA:" + appId; // 秒级计数,允许单个业务方每秒最多quota条 Long current = redisTemplate.opsForValue().increment(key); if (current != null && current == 1) { redisTemplate.expire(key, Duration.ofSeconds(1)); } return current != null && current <= quota; }3.4 模板渲染与多场景配置化
通知内容的上游是业务方,业务方不该关心“短信要拼接成什么文本”。我们把内容渲染放到通知服务内部,业务方只传模板参数。
模板存在DB表中:
CREATE TABLE `notification_template` ( `id` bigint NOT NULL AUTO_INCREMENT, `scene_code` varchar(64) NOT NULL, `channel` varchar(32) NOT NULL, `content_template` varchar(1000) NOT NULL COMMENT '内容模板,如:您的订单${orderId}已支付成功', `status` tinyint NOT NULL DEFAULT '1', PRIMARY KEY (`id`), UNIQUE KEY `uk_scene_channel` (`scene_code`, `channel`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;请求进来时,根据sceneCode + channel找到对应模板,用参数替换占位符。用Spring自带的SimplePlaceholder或者正则替换都可以:
public String render(String template, Map<String, Object> params) { for (Map.Entry<String, Object> entry : params.entrySet()) { template = template.replace("${" + entry.getKey() + "}", entry.getValue().toString()); } return template; }模板配置支持版本管理,可以看历史变更记录。改文案不用发代码,运营自己改配置就能生效,这对企业级系统的敏捷性非常重要。
4. 高可用与系统部署方案
4.1 基于K8s的多实例部署与故障转移
之前部署方式是单机运行jar包,一台机器挂了服务就不可用,这在高可用设计里完全不合格。现代企业级系统,我强烈建议容器化部署到Kubernetes。
Kubernetes带来的高可用能力有几个层次:
多副本自动调度。通知服务是无状态应用(所有状态都在MySQL和Redis里),可以水平扩容。Deployment配置replicas: 3,三副本调度到不同节点上,一台机器宕机,Pod会被自动调度到其他健康节点,服务不中断。同时配置HPA(HorizontalPodAutoscaler),当Pod的CPU使用率超过70%或者内存使用率超80%,自动扩容副本数量。
apiVersion: autoscaling/v2 kind: HorizontalPodAutoscaler metadata: name: notification-hpa spec: scaleTargetRef: apiVersion: apps/v1 kind: Deployment name: notification-service minReplicas: 3 maxReplicas: 10 metrics: - type: Resource resource: name: cpu target: type: Utilization averageUtilization: 70健康检查与优雅停机。配置readinessProbe和livenessProbe,探针路径是/actuator/health。如果应用失联,K8s自动重启或摘除流量。要做到优雅停机,JVM收到SIGTERM信号后要停止接收新请求,处理完正在进行的发送任务再退出。Spring Boot 2.3+原生支持优雅停机:
server: shutdown: graceful spring: lifecycle: timeout-per-shutdown-phase: 30s三master高可用控制面。K8s集群本身也要高可用。生产环境推荐三台master节点的架构,etcd三节点集群,API Server通过负载均衡对外提供服务。kube-vip或HAProxy做API Server的VIP,一台master宕机不影响集群管理面。如果对K8s运维不熟,用kubekey这类工具可以比较轻松地搭建一套三master高可用集群,它内部已经集成了etcd集群和负载均衡的配置逻辑。
4.2 消息队列的高可用配置
RocketMQ或Kafka,部署时都必须开启副本机制。RocketMQ用主从同步模式(SYNC_MASTER),生产消息写入主节点后要同步到从节点才返回成功。这样主节点宕机,从节点还能继续消费,最大限度降低消息丢失概率。
RocketMQ的关键配置:
brokerClusterName=DefaultCluster brokerName=broker-a brokerId=0 brokerRole=SYNC_MASTER flushDiskType=SYNC_FLUSHSYNC_MASTER要求消息同步到从节点才确认,SYNC_FLUSH要求消息落盘才确认,这两个都开启,消息可靠性最高,但吞吐量会打折。对通知系统来说,可靠性优先于吞吐,宁可发送慢一点,不能丢消息。
Kafka对应的配置是replication.factor=3、min.insync.replicas=2、acks=all,这三个参数组合含义是:每个分区有3个副本,至少2个副本同步成功才算写入成功。
另外,消费端必须开启手动ACK,不能消费完消息自动确认,必须在渠道发送成功后才ACK,否则消息还没发送成功就确认了,进程崩溃就丢消息了。
RocketMQ的消费配置:
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> { for (MessageExt msg : msgs) { try { // 1. 从消息体解析出messageId // 2. 根据messageId查询数据库中消息状态 // 3. 如果状态已经是成功,直接跳过 // 4. 否则执行真正的发送逻辑 // 5. 发送成功后更新数据库状态 } catch (Exception e) { // 返回 RECONSUME_LATER,让MQ稍后重投递 return ConsumeConcurrentlyStatus.RECONSUME_LATER; } } return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; });这里有一个极其重要的幂等性问题:消息队列的投递是“至少一次”的,消费端可能对一个消息重复收到。如果收到重复消息,必须先去数据库检查状态,已经是终态的直接跳过。这条逻辑必须在消息处理和状态更新中贯穿始终。
4.3 数据库高可用与读写分离
MySQL数据库的高可用,生产环境我推荐一主多从加MHA或Orchestrator自动故障切换。应用层面使用读写分离,写操作走主库,读操作走从库。
如果对自动化切换不熟悉,退一步做“主从热备 + VIP漂移”也可以。主库宕机后,DBA手动把VIP漂移到从库,提升从库为主库。MTTR 5分钟内,绝大多数场景够用。但如果是7x24高可用要求严格的服务,建议直接上MHA,自动检测主库故障,自动选主、自动切换,切换时间秒级。
Spring Boot层面配置多数据源:
spring: datasource: primary: jdbc-url: jdbc:mysql://172.16.1.10:3306/notification?useSSL=false username: notify_rw driver-class-name: com.mysql.cj.jdbc.Driver replica: jdbc-url: jdbc:mysql://172.16.1.11:3306/notification?useSSL=false username: notify_ro driver-class-name: com.mysql.cj.jdbc.Driver应用中通过@DS("replica")注解把查询路由到从库,@DS("primary")把写入路由到主库。注意从库的数据同步存在秒级延迟,刚写入的通知记录,立即查详情可能查不到。所以我们规定:状态查询走主库,列表查询走从库。状态一致性比性能更重要。
4.4 服务高可用场景下的后端编码准则
仅仅依赖基础组件的高可用还不够,应用中大量低质量的代码会把高可用架构变得形同虚设。我在重构通知系统的过程中,总结了几条在多实例部署场景下的硬性编码准则:
不要使用本地内存缓存做全局判断。之前有个同事把一个用户的频控计数放在HashMap里。单实例当然没问题,多实例部署后同一个用户被不同实例处理,每个实例各自计数,限流形同虚设。凡是需要全局一致的计数、锁、去重,必须放Redis。
不要在进程内维护有状态信息。比如保存一份“正在处理的消息ID列表”在内存里。如果实例被K8s杀死重启,这个列表就丢了,相关消息会被重新消费。该查数据库就查数据库。
设置全局的Outbox模式。有些消息是先更新业务库,再发通知。如果两步之间应用崩溃,通知就漏了。我们的做法是:业务方调用通知API之前,先在业务库里记录一条outbox记录,通知服务返回成功后再更新outbox状态。同时有个定时任务扫描超过2分钟未完成的outbox记录,重新调用通知API。这是一种兜底补偿机制,防止极端情况下的消息丢失。
线程池必须显式命名和设置拒绝策略。通知发送线程池,队列不能是无界的。无界队列会耗尽内存。设置有界队列,队列满后执行CallerRunsPolicy,让提交线程自己执行任务,天然的背压机制。
ThreadPoolExecutor sendPool = new ThreadPoolExecutor( 10, 20, 60, TimeUnit.SECONDS, new ArrayBlockingQueue<>(1000), new ThreadFactoryBuilder().setNameFormat("notify-send-%d").build(), new ThreadPoolExecutor.CallerRunsPolicy() );4.5 降级与熔断:渠道商故障不拖垮主链路
渠道商是不可控的,邮件服务商偶尔延迟、短信服务商偶尔拒绝。高可用架构里必须有降级和熔断能力。
我们的做法是在渠道发送层集成Resilience4j的CircuitBreaker:
CircuitBreakerConfig config = CircuitBreakerConfig.custom() .failureRateThreshold(50) .waitDurationInOpenState(Duration.ofSeconds(30)) .permittedNumberOfCallsInHalfOpenState(10) .slidingWindowSize(20) .build();含义是:最近20次调用中,如果超过50%失败,熔断器打开,后续所有请求直接返回失败,不再调用渠道商接口;30秒后进入半开状态,允许10次探测调用,如果成功率达到要求,熔断器关闭恢复。
熔断打开期间的降级策略:
- 短信失败-> 重新排队,进入重试队列
- 邮件失败-> 降低优先级,延迟10分钟再发
- 站内信失败-> 直接写库,用户下次登录时主动拉取
这套降级策略写在渠道策略实现中,可以配置化调整。
5. 常见问题与排查技巧实录
5.1 消息积压了,怎么快速定位瓶颈
通知系统最常遇到的线上问题就是消息积压。MQ里的消息堆积量持续增长,用户收不到通知。排查思路按三步走:
第一步,看消费组状态。RocketMQ提供的mqadmin consumerStatus命令可以查看消费组每个消费者的实时消费速率和堆积量:
sh mqadmin consumerStatus -g notify-sms-group -n 172.16.1.20:9876如果消费速率低于生产速率,说明消费端存在瓶颈。
第二步,看日志中渠道商响应耗时。渠道商接口响应慢是消费速率低的最主要原因。我们每发送一条都记录耗时,通过Prometheus监控P99耗时。如果P99超过5秒,基本可以判定是渠道商侧的性能问题,而不是消费代码的问题。
第三步,看数据库连接池使用率。每次发送结束后需要更新数据库状态,数据库连接池被打满也会导致消费线程阻塞。定期监控HikariPool的活跃连接数。
关键点:不要把消息积压的原因一股脑归结于“代码有bug”。我遇到过多次,原因就是下游渠道商的某个API在晚高峰响应时间从200ms涨到10秒,消费速率被外部拖慢,积压在所难免。这时候要做的是降级、快速限流、联系渠道商,而不是瞎调消费线程数。
5.2 消息丢失的排查思路
通知系统最怕的是“消息丢了用户没收到”,排查手段要有章法。我们的核心思路是:每条消息都有全链路追踪ID,从入口到渠道商响应全部打日志。
消息的traceId在入口生成,贯穿消息落库、MQ消息体、消费日志、渠道商请求头。日志采集到ELK后,排查时直接按traceId搜全链路日志:
入口接收 -> 消息落库 -> 投递MQ -> 消费拉取 -> 调用渠道商 -> 收到渠道商响应 -> 更新数据库状态哪一步缺失,问题就出在哪一步。常见场景:
- 入口有日志,但消息没落库 -> 数据库写入失败,看是否有主键冲突或连接池打满
- 落库成功但MQ没有消息 -> 事务提交前MQ发送失败,用事务消息修复
- MQ有消息但消费日志没打印 -> 消费组疑似rebalance或者消息被其他实例消费,查看该消费组的消费者列表
- 消费有日志但渠道商没收到 -> 渠道商API调用前报错了,看异常堆栈
而且所有消费端的异常日志不能只打error级别,必须包含messageId和traceId,否则后续排查要大海捞针。
5.3 渠道商接口被限流怎么办
业务大促期间,短信渠道商返回限流错误码是家常便饭。此时要做两件事:
内部降速。当前消费线程数是20,发现连续10条返回“限流”后,通过动态配置中心把消费线程数降为5,给渠道商喘息时间。这种运行时动态调整,Spring Cloud Config配合@RefreshScope就能实现。
逃生通道。短信被限流不等于通知就不发了。我们做了渠道自动降级:短信发送遇到限流时,把消息重新投递到站内信队列,同时标记一条“降级记录”,用户在App端会收到站内信,而不是完全丢失通知。
代码上,这在策略模式里就能优雅实现:
public class SmsChannel implements NotificationChannel { @Override public SendResult send(NotificationMessage message) { try { // 调用短信服务商 } catch (RateLimitException e) { // 返回降级标记 return SendResult.degraded("RATE_LIMIT"); } } }Worker统一处理degraded结果,将消息改投站内信。
5.4 发给同一用户的消息太多,如何聚合收敛
营销场景最容易出现这个问题:用户一晚上连续下单5次,就收到5条短信,体验很差。企业级通知系统一定要有“消息聚合”能力。
我们在投递前加了一个聚合逻辑:以用户+场景+时间段(比如10分钟)为维度,如果该窗口内已有相同场景的消息,新消息不立即发送,而是更新已有消息的content字段,附加“您有3条新通知,请点击查看”。
这个策略在Redis里实现比较简单,用Hash记录每个用户的聚合信息:
public Optional<String> aggregateIfNeeded(String userId, String sceneCode, String content) { String key = "NOTIFY:AGG:" + userId + ":" + sceneCode; // 如果10分钟窗口内已有聚合,更新内容,然后返回不发送 // 否则创建新的聚合窗口,正常发送 }通过这样一层聚合,营销消息的发送量可以直接减少30%到50%,也降低了对渠道商的压力。
6. 写在最后的实践经验
整套系统重构上线后,服务稳定性有了质的提升。原来高峰期每周至少一两次因为渠道商慢导致上游接口超时的投诉,重构后半年多了,最严重的一次是RocketMQ某个broker磁盘满了,因为有主从同步和自动故障恢复,业务基本无感。
有几点体会非常深:
一是通知系统看似简单,实际是分布式系统的“小型博物馆”。幂等、削峰、限流、熔断、降级、状态机、分布式锁、消息可靠性,这些分布式系统里的核心议题,在通知系统里全都能碰到。认真做完这套系统,对整个后端知识体系的提升非常大。
二是不要迷信某个中间件,要根据场景做取舍。我见过有人为了追求高可靠,把所有通知都走RocketMQ事务消息,导致链路极长,一条站内信消息延迟好几秒。实际上站内信这种场景,即使丢失几条,用户下次登录拉取也能看到,完全不需要事务消息。不同渠道不同等级的消息,可靠性要求是不一样的。
三是可观测性的建设不能滞后。通知系统的排查复杂度和业务系统不是一个量级,链路长、依赖多。Prometheus监控指标和ELK日志系统,最开始就要同步建设好。没有监控就上线,出了问题定位问题所花费的时间,可能比重构整个系统还长。
最后分享一个小技巧:为了验证消息不丢,我在重构完成后专门写了一个模拟故障的测试用例,随机杀掉消费者进程,看消息消费情况。连续跑了48小时,最终确认消息不丢失的可靠性达到99.9%以上(允许少量重复,但绝不丢失)。建议你做完通知系统后也做一次这种混沌实验,这比任何代码评审都更能暴露问题。