ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

高可用企业级通知系统架构设计与实战:消息驱动、削峰填谷全解析

高可用企业级通知系统架构设计与实战:消息驱动、削峰填谷全解析 有时候后台系统里最不起眼的模块反而是线上事故最容易爆发的点。消息通知服务就是一个典型——你平时觉得它简单不就是把一条消息塞给用户吗可一旦业务量上来短信、邮件、站内信、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万条请求。如果不经过MQ50万条发送请求同时打给下游短信网关通道直接被打爆发送大量失败。接入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) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;这里几个关键设计点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) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;主表状态为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-elsepublic 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注入的ListNotificationChannel根据消息的渠道类型动态选择对应的实现类。新增一个渠道比如钉钉机器人只需要新增一个实现类注册为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) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;请求进来时根据sceneCode channel找到对应模板用参数替换占位符。用Spring自带的SimplePlaceholder或者正则替换都可以public String render(String template, MapString, Object params) { for (Map.EntryString, 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会被自动调度到其他健康节点服务不中断。同时配置HPAHorizontalPodAutoscaler当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的关键配置brokerClusterNameDefaultCluster brokerNamebroker-a brokerId0 brokerRoleSYNC_MASTER flushDiskTypeSYNC_FLUSHSYNC_MASTER要求消息同步到从节点才确认SYNC_FLUSH要求消息落盘才确认这两个都开启消息可靠性最高但吞吐量会打折。对通知系统来说可靠性优先于吞吐宁可发送慢一点不能丢消息。Kafka对应的配置是replication.factor3、min.insync.replicas2、acksall这三个参数组合含义是每个分区有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?useSSLfalse username: notify_rw driver-class-name: com.mysql.cj.jdbc.Driver replica: jdbc-url: jdbc:mysql://172.16.1.11:3306/notification?useSSLfalse 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的CircuitBreakerCircuitBreakerConfig 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 OptionalString 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%以上允许少量重复但绝不丢失。建议你做完通知系统后也做一次这种混沌实验这比任何代码评审都更能暴露问题。
返回列表