ARTICLE DETAIL

资讯详情

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

消息中心高并发架构实战:分层设计、幂等去重与重试补偿

消息中心高并发架构实战:分层设计、幂等去重与重试补偿 消息中心这套东西乍一看没什么技术含量不就是发个通知吗但真把它当成一个独立系统来做你会发现坑比想象中深得多。我从最早在业务代码里直接调短信 SDK到后来负责一个日均发送量几千万的消息平台中间踩过的坑足够写一本小册子。这篇笔记就把消息中心从需求拆解、架构分层、数据建模到高并发落地这条链路完整讲一遍重点放在那些文档里不会写、但线上一定会遇到的问题上。适合正在做统一消息平台的同学也适合业务开发想搞清楚我的消息为什么发丢了的排查思路。全文不谈虚的所有结论都对应到具体的表结构、参数和代码片段上。1. 先把消息中心的能力边界说清楚1.1 业务代码里直连短信 SDK 到底会烂在哪几乎每个团队都经历过这个阶段注册要发验证码直接在注册逻辑里 import 一个短信 SDK下单成功要发通知在订单服务里再调一次。刚开始只有两三个发送点感觉挺清爽。等到业务铺开问题就集中爆发了。第一个问题是渠道散落。短信、邮件、站内信、App 推送各自有各自的 SDK 和配置密钥散在十几个服务的配置文件里运维轮换一次密钥要改十几个仓库。第二个问题是没有统一的重试和补偿。短信通道偶尔抖动代码里 catch 一下就吞掉了用户没收到验证码客服那边完全查不到记录。第三个问题是业务方无法感知发送结果。运营想知道昨天那批营销消息触达率多少开发只能去翻日志因为根本没有落库。更隐蔽的是重复发送。业务逻辑重试、MQ 消息重投、定时任务重复触发三条路径叠在一起用户可能收到三遍同样的短信。这种问题在测试环境几乎测不出来上线之后才被投诉。这些问题的根源其实是一个发送这个动作被当成了业务逻辑的一部分而不是一个独立的基础能力。只要它还寄生在业务代码里就没法统一治理。所以消息中心存在的第一个理由就是把发送链路的复杂度从业务里彻底剥离出去。1.2 一个消息中心该收哪些活、不该收哪些活边界划不清楚消息中心最后会变成一个什么都能塞的垃圾桶。我的经验是收进来的必须满足与具体业务无关的通用发送能力这一条。该收的活包括多渠道统一发送短信、邮件、站内信、推送、微信模板消息等、模板管理与变量渲染、发送记录落库与查询、失败重试与降级、频控与限流、发送结果回调、灰度与 A/B 投放。不该收的活或者说要想清楚再收的业务侧的内容策略。比如下单满 100 元才发这条券这是业务规则应该在业务服务里判断完只把该给谁发什么模板告诉消息中心。如果让消息中心去查订单表它就和业务耦合死了改一个营销规则要动消息中心的代码。还有一个容易过界的点是用户触达偏好。用户有没有关闭营销推送、是否在免打扰时段这类偏好属于用户中心或者营销平台消息中心应该通过接口查询而不是自己维护一份。我见过一个平台把偏好数据冗余在消息中心结果用户关闭推送之后还是会收到因为用户中心的数据同步延迟了半小时。这种数据一致性问题能避就避。提示划边界的时候问自己一句——这个规则换一个业务方还成立吗成立就放消息中心不成立就留在业务侧。2. 整体架构怎么分层每一层的取舍在哪2.1 接入层、调度层、渠道层、存储层的职责划分分层的第一原则是每一层只做一件事且可以被独立替换。我最终落地的结构大概是四层从上到下说。接入层负责两件事协议转换和参数校验。业务方通过 HTTP 或者 RPC 提交发送请求接入层把不同协议统一成内部的发送指令做基本校验必填字段、模板是否存在、接收方格式是否合法然后立刻返回一个 msgId。注意这里是异步返回不等实际发送结果。同步等发送结果会把接入层的 RT 拖到几百毫秒甚至秒级通道一抖动接入层线程池直接打满形成雪崩。业务方要结果通过查询接口或者回调拿。调度层是核心。它从消息队列里消费发送指令做频控判断、路由决策这条消息走哪个通道、拆分成渠道任务然后投递给渠道层。调度层是唯一有状态判断逻辑的地方也是最需要做幂等的地方。渠道层是通道适配器。每个通道一个实现对外暴露统一的 send 接口内部处理各自 SDK 的差异。渠道层只负责把这条消息发出去并返回结果不做业务判断。这样一来新增一个通道只需要加一个实现类不动上层。存储层分两块消息元数据谁在什么时候给谁发了什么落在 MySQL发送轨迹和回执存在 MySQL 或者 ES 里供查询。不要把所有东西都塞进 MySQL回执数据量是元数据的几倍用 ES 存轨迹查询会更舒服。这四层之间用什么通信接入层到调度层用消息队列削峰调度层到渠道层既可以直接 RPC 调用因为要拿同步结果也可以再走一层队列。我倾向于调度到渠道直接 RPC因为渠道层的并发是可控的而重试逻辑在调度层统一做会更清晰。2.2 消息队列与存储引擎的选型对比队列这块社区里讨论最多的就是 Kafka 和 RocketMQ。我实际两个都用过结论是消息中心这个场景 RocketMQ 更省心原因有几点。第一是延迟消息。重试需要延迟投递比如第一次失败 1 分钟后重试第二次 5 分钟第三次 30 分钟。RocketMQ 原生支持任意精度的延迟消息新版本支持自定义时间戳Kafka 要靠外部调度或者时间轮自己实现多一层组件就多一个故障点。第二是事务消息。有些场景需要落库和发消息原子化RocketMQ 的事务消息能直接解决Kafka 得靠本地消息表方案绕一圈。第三是消息轨迹RocketMQ 自带排查问题的时候能看到消息到底卡在哪个环节。当然 Kafka 的吞吐更高如果你的场景是纯日志型、海量、允许少量丢失Kafka 更合适。但消息中心的短信验证码是绝对不能丢的投递语义要求更高这时候 RocketMQ 的优势就体现出来了。存储选型上MySQL 是主力。消息元数据用 MySQL索引建在(receiver_id, create_time)和(biz_id, create_time)上运营查询和用户查询都能覆盖。表量大的话按create_time做冷热分离三个月以前的数据归档到历史库。组件适用场景我的取舍理由RocketMQ发送指令、延迟重试原生延迟消息 事务消息重试链路简单Kafka回执日志、埋点吞吐高允许少量丢失用于分析MySQL消息元数据、任务状态强一致支持事务和精确查询Redis频控计数、幂等去重高性能计数配合 Lua 保证原子性Elasticsearch发送轨迹、全文检索海量回执查询支持多条件组合注意不要为了技术先进上 Kafka 做核心发送链路除非你的团队有能力自己维护延迟调度和时间轮。组件越多故障排查的路径越长。3. 数据模型与消息状态机3.1 三张核心表的设计思路表设计是消息中心的地基地基建歪了后面改起来非常痛苦。我的方案是三张表消息主表、渠道任务表、发送轨迹表。消息主表记录一条业务消息的整体信息一个业务请求对应一行。核心字段包括msg_id全局唯一雪花算法生成、biz_id业务方标识用于频控和统计、template_id、receiver接收方手机号或用户 ID、params模板变量JSON 存储、channel_mask需要发送的渠道位图、status整体状态、create_time、expire_time过期不再发送。这里channel_mask用位图存储是个小技巧。假设 bit0 是短信、bit1 是邮件、bit2 是站内信、bit3 是推送那短信站内信就是二进制 0101十进制 5。这样一条消息多渠道只占一个字段查询时用位运算就能筛出所有包含短信的消息比存逗号分隔的字符串高效得多。渠道任务表是把一条消息拆成多个渠道任务之后的记录。一条消息发了三个渠道这里就有三行。核心字段包括task_id、msg_id、channel、status、retry_count已重试次数、next_retry_time下次重试时间、send_time。这张表是重试调度的核心扫描status FAILED AND next_retry_time now()就能捞出所有该重试的任务。发送轨迹表记录每一次发送尝试的结果包括失败原因。同一个 task 重试三次这里就有三行。这张表写入量最大只追加不更新适合放 ES 或者按月分表。三张表的关系是主表 1 条任务表 N 条N 等于渠道数轨迹表 M 条M 大于等于 N因为有重试。这个结构看起来很啰嗦但好处是职责清晰主表回答这条消息整体成功了吗任务表回答某个渠道发成功了吗轨迹表回答某次尝试为什么失败。排查问题的时候层层下钻非常顺畅。3.2 消息状态流转与补偿机制状态机是防止状态错乱的关键。主表的状态我定义了六个INIT已接收、DISPATCHED已拆分、SENDING发送中、SUCCESS全部渠道成功、PARTIAL部分渠道成功、FAILED全部失败。状态流转有严格的规则只允许从INIT到DISPATCHED从DISPATCHED到SENDINGSENDING到终态。任何状态都不允许从终态回到中间态这个约束一定要在代码里强制校验否则重试逻辑可能把一条已经成功的消息重新拉回到发送中造成重复发送。任务表的状态更简单INIT、SUCCESS、FAILED、DEAD超过最大重试次数进死信。FAILED是可重试的会带上next_retry_timeDEAD是不可重试的需要人工介入或者告警。补偿机制有两个层面。同步补偿是在调度层消费消息的时候如果发现这条消息的任务表记录已经存在直接跳过避免消息队列重投造成重复拆分。异步补偿是定时任务扫描长时间停留在SENDING状态的任务比如超过 10 分钟主动去查渠道层的回执如果查不到就判定为失败并进入重试。这个兜底非常必要因为渠道层的回调可能丢网络也可能断不能完全依赖同步返回。状态的最终一致性靠一张对账表保证。每天凌晨跑一次对账用主表的msg_id去渠道侧拉取发送结果比对状态是否一致不一致的修正。这个动作看起来很土但它是数据准确性的最后一道防线。4. 高并发下真正难啃的几个点4.1 幂等与去重别让用户收到三遍验证码重复发送是用户投诉最多的问题必须从三个层面同时防。入口层幂等。业务方提交发送请求时带一个request_id业务方自己生成全局唯一消息中心用 Redis 做SETNXkey 是msg:req:{request_id}过期时间设 10 分钟。如果 key 已存在说明是重复提交直接返回上次的 msgId不新建消息。这一步挡住了业务方的重试和网络重传。调度层幂等。消息队列消费的时候用msg_id channel作为幂等 key同样用 Redis 做SETNX。为什么要在这里再做一次因为 MQ 本身是 at-least-once 语义同一条消息可能被消费两次如果不做幂等任务表会插入两行用户收到两条。发送层幂等。渠道层调 SDK 的时候把task_id作为厂商侧的幂等 key 传过去如果厂商支持的话。大部分短信厂商都支持这个参数能有效避免网络超时但其实已经发出去了这种最恶心的重复。三层幂等叠加重复率能压到极低。但要注意Redis 的可靠性。如果 Redis 挂了幂等就失效了。所以关键业务场景我建议在 MySQL 里也建一张幂等表Redis 作为第一道快筛MySQL 作为兜底。写入前先查 MySQL虽然慢一点但保险。实操心得幂等 key 的过期时间不能太短。如果设成 1 分钟业务方 2 分钟后的重试就挡不住了。我的经验值是至少 10 分钟验证码这种场景可以设成和验证码有效期一致。4.2 频控、限流与降级的三层防线频控是面向用户的防止骚扰。比如同一用户 1 分钟内最多 1 条短信同一用户 1 天最多 5 条营销短信。频控用 Redis 的计数器实现key 是freq:{biz_id}:{receiver}:{channel}用INCREXPIRE组合。这里必须用 Lua 脚本保证原子性否则并发下先 INCR 后 EXPIRE 之间如果进程挂了key 就永远不会过期用户被永久限制。频控的粒度要仔细设计。验证码这种强需求频控要宽松1 分钟 1 条但允许用户主动请求重发营销消息要严格1 天 3 条且要和用户偏好叠加。限流是面向系统的防止通道被打爆。每个通道的 QPS 上限是固定的比如短信通道 5000 QPS超过就要排队。用令牌桶算法桶容量是通道的突发承受能力令牌生成速率是稳态 QPS。超过令牌的请求进入等待队列队列满了就降级。降级是最后一道防线。当某个通道不可用时策略要分场景验证码类消息可以切换到备用通道比如从通道 A 切到通道 B营销类消息直接丢弃并记录不要重试避免通道恢复后瞬间涌入大量积压消息把通道再打挂。这个区分非常重要很多平台在通道恢复后因为积压的营销消息一起重试导致二次故障。防线作用对象实现方式触顶后的动作频控单个用户Redis 计数器 Lua直接拒绝记录原因限流单个通道令牌桶 等待队列排队队列满则降级降级系统整体通道切换 / 消息丢弃验证码切通道营销丢弃4.3 重试策略与死信处理重试次数和间隔不能拍脑袋定要按通道特性和消息类型分别配置。验证码短信对时效要求极高重试间隔要短1 秒、5 秒、30 秒最多三次超过就告知用户发送失败请稍后重试。营销消息对时效要求低重试间隔可以拉长1 分钟、5 分钟、30 分钟这样能避开通道的短时抖动。站内信基本不会失败重试一次就够了。重试的退避算法建议用指数退避加随机抖动。为什么要加随机抖动因为瞬时的通道故障会让大量消息在同一时刻失败如果重试间隔完全一样它们会在同一时刻再次集中重试形成重试风暴。加上随机抖动比如 ±20%能把重试请求打散。超过最大重试次数的任务进死信。死信不是终点需要有死信告警和人工处理后台。我在实际项目里给死信设了分级告警验证码类消息的死信立刻告警给值班同学营销类消息汇总成日报。死信任务在后台支持手动重发重发前要先确认通道状态正常。重试调度用什么拉定时任务每分钟扫一次会有一个问题任务量大的时候扫描压力大。我的做法是用分片扫描按task_id % 分片数把扫描任务分散到多个 worker每个 worker 只扫自己负责的分片。同时给(status, next_retry_time)建联合索引扫描效率能提升一个数量级。5. 从建表到链路打通的实际操作5.1 表结构落地先看消息主表的建表语句字段做了精简只保留核心部分。CREATE TABLE t_message ( msg_id BIGINT UNSIGNED NOT NULL COMMENT 消息ID雪花算法, biz_id VARCHAR(64) NOT NULL COMMENT 业务方标识, template_id VARCHAR(64) NOT NULL COMMENT 模板ID, receiver VARCHAR(128) NOT NULL COMMENT 接收方, params JSON DEFAULT NULL COMMENT 模板变量, channel_mask INT UNSIGNED NOT NULL DEFAULT 0 COMMENT 渠道位图, status TINYINT NOT NULL DEFAULT 0 COMMENT 0INIT 1DISPATCHED 2SENDING 3SUCCESS 4PARTIAL 5FAILED, expire_time DATETIME DEFAULT NULL COMMENT 过期时间, create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, PRIMARY KEY (msg_id), KEY idx_receiver (receiver, create_time), KEY idx_biz (biz_id, create_time) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4 COMMENT消息主表;params用 JSON 类型而不是 TEXT好处是 MySQL 5.7 以后可以直接用params-$.orderNo这样的语法查询运营排查问题时能直接按订单号定位消息。channel_mask用 INT 存位图一个 int 支持 32 个渠道足够用了。再看渠道任务表这张表是重试的核心。CREATE TABLE t_channel_task ( task_id BIGINT UNSIGNED NOT NULL COMMENT 任务ID, msg_id BIGINT UNSIGNED NOT NULL COMMENT 消息ID, channel VARCHAR(16) NOT NULL COMMENT 渠道SMS/EMAIL/INBOX/PUSH, status TINYINT NOT NULL DEFAULT 0 COMMENT 0INIT 1SUCCESS 2FAILED 3DEAD, retry_count TINYINT NOT NULL DEFAULT 0 COMMENT 已重试次数, next_retry_time DATETIME DEFAULT NULL COMMENT 下次重试时间, send_time DATETIME DEFAULT NULL COMMENT 实际发送时间, PRIMARY KEY (task_id), KEY idx_msg (msg_id), KEY idx_retry (status, next_retry_time) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4 COMMENT渠道任务表;idx_retry这个联合索引是给重试调度用的扫描条件status 2 AND next_retry_time NOW()能走索引。注意索引的顺序不能颠倒status在前面因为它的区分度更低但过滤性更强先过滤掉大量成功任务再按时间筛。分表策略上主表按msg_id取模分 16 张任务表按task_id取模分 32 张。为什么任务表分得多因为任务表行数是主表的好几倍。分片数选择上是 2 的幂方便后期翻倍扩容。5.2 核心发送链路代码发送链路的骨架不复杂关键在细节处理。下面是一个接入层的核心处理逻辑用 Java 写。public SendResult send(SendRequest request) { // 1. 入口幂等 String idemKey msg:req: request.getRequestId(); Boolean first redis.opsForValue().setIfAbsent(idemKey, 1, 10, TimeUnit.MINUTES); if (Boolean.FALSE.equals(first)) { Long existMsgId queryMsgIdByRequestId(request.getRequestId()); return SendResult.duplicate(existMsgId); } // 2. 参数校验 if (!templateService.exists(request.getTemplateId())) { throw new BizException(模板不存在); } if (!receiverValidator.valid(request.getReceiver())) { throw new BizException(接收方格式错误); } // 3. 频控检查 FreqResult freq freqService.check(request.getBizId(), request.getReceiver(), request.getChannelMask()); if (!freq.isPass()) { return SendResult.rejected(freq.getReason()); } // 4. 落库主表 long msgId idGenerator.nextId(); MessageDO message buildMessage(msgId, request); messageMapper.insert(message); // 5. 投递到 MQ让调度层异步处理 mqProducer.send(TOPIC_MSG_DISPATCH, msgId, message); return SendResult.accepted(msgId); }注意这里的顺序先落库再投递 MQ。如果先投 MQ 后落库MQ 消费速度比落库快的话调度层可能查不到主表记录。另外第 5 步用msgId作为 MQ 的消息 key能保证同一 message 的消息投到同一个队列保持顺序。调度层的任务拆分逻辑是这样的RocketMQMessageListener(topic TOPIC_MSG_DISPATCH, consumerGroup cg_dispatch) public class DispatchConsumer implements RocketMQListenerMessage { Override public void onMessage(Message message) { long msgId message.getMsgId(); // 幂等检查是否已拆分 if (taskMapper.existsByMsgId(msgId)) { log.warn(消息已拆分跳过 msgId{}, msgId); return; } ListInteger channels ChannelMask.parse(message.getChannelMask()); for (Integer channel : channels) { ChannelTask task new ChannelTask(); task.setTaskId(idGenerator.nextId()); task.setMsgId(msgId); task.setChannel(ChannelEnum.of(channel)); task.setStatus(TaskStatus.INIT); taskMapper.insert(task); // 立刻发起首次发送 channelSender.sendAsync(task); } messageMapper.updateStatus(msgId, MsgStatus.SENDING); } }这里existsByMsgId是幂等的关键但要注意并发场景如果两条相同的 MQ 消息同时消费都查到不存在就会插入两遍。解决办法是给t_channel_task的(msg_id, channel)加唯一索引插入冲突时捕获异常跳过。这个细节很容易被忽略但线上是会出问题的。5.3 容量估算与参数计算容量估算决定了你的资源投入不能拍脑袋。我按一个中等规模的平台来算一遍。假设日均发送 2000 万条消息平均每条消息发 1.5 个渠道那么任务表每天新增 3000 万行主表 2000 万行。峰值集中在早 9 点到晚 9 点峰值系数取 3那么峰值 QPS 大约是2000 万 / 12 小时 / 3600 秒 × 3 ≈ 1400 QPS。这个量级不算大单组 MySQL 加合理分表能扛住。存储估算主表每行平均 500 字节一天 2000 万行就是 10GB一年 3.6TB。任务表每行 200 字节一天 3000 万行就是 6GB一年 2.2TB。所以冷热分离是必须的三个月热数据保留在 MySQL冷数据归档到对象存储或者点击流系统。Redis 的容量估算主要看频控 key。假设活跃用户 500 万每个用户每个业务方每天一个频控 key10 个业务方就是 5000 万个 key每个 key 大概 100 字节总共 5GB。这个量级一台 16G 的 Redis 集群能装下但要考虑主从复制的内存放大。通道的并发估算短信通道假设配置 5000 QPS但真实峰值可能需要 8000这时候限流的等待队列就派上用场了。队列长度设置为峰值QPS × 平均发送耗时比如 8000 × 0.1 秒 800。队列超过这个长度就应该触发降级而不是无限堆积。实操心得容量估算一定要留 50% 的余量。我见过因为大促流量翻三倍任务表分片没提前扩容导致单个分片写入热点整组数据库响应变慢的案例。分片扩容提前做比事后补救容易得多。6. 线上问题排查实录6.1 几个典型故障的复盘故障一验证码大面积延迟。某天上午突然接到大量投诉验证码要等一两分钟才到。排查发现通道侧的响应时间从 100ms 涨到了 2 秒而调度层的线程池配置是固定的 200 个线程每个线程等 2 秒整体 QPS 从 2000 掉到了 100。问题定位很快但根因是线程池没有按超时时间动态调整。解决方案是给渠道调用加上超时熔断响应时间超过 500ms 直接切备用通道同时把线程池改成根据队列长度动态扩容。故障二用户收到重复短信。排查发现是 MQ 重投导致的。某个消费者的处理逻辑抛了异常MQ 触发重试但任务表已经插入了一行第二次消费时幂等检查用的是 Redis 缓存恰好在那一刻缓存过期了于是又插了一行。修复方式是给任务表加唯一索引并且把 Redis 幂等 key 的过期时间从 5 分钟延长到 30 分钟同时在 Redis 之外加了本地缓存做二级保护。故障三营销消息积压后集中爆发。某个营销活动一批 500 万条的消息因为通道限流被积压在队列里。第二天早上通道恢复正常积压的消息瞬间全部重试把通道又打挂了。这个问题的教训是重试必须分批和限速。修复方案是给重试队列加上速率限制每秒最多重试 1000 条同时给积压消息加上超过 6 小时不再重试的过期判断。这三个故障有一个共同点问题都不出在功能实现上而出在边界条件的处理上。功能测试的时候一切正常只有在真实的流量压力和异常场景下才暴露出来。6.2 常见问题速查表与避坑心得现象可能原因排查方向解决方式消息发不出去通道密钥失效 / 签名错误看渠道层返回码更新密钥检查签名用户收到重复消息MQ 重投 / 重试逻辑有 bug查轨迹表同一 task 的记录数加唯一索引 幂等消息延迟严重线程池打满 / 通道响应慢看线程池队列长度和通道 RT加超时熔断切备用通道频控误判Redis key 未设置过期查 Redis 里 key 的 TTLLua 脚本保证原子性状态不对状态机缺少校验看状态流转日志加状态流转白名单查询慢索引缺失 / 数据量大explain 看执行计划补索引冷热分离最后分享几个我踩过坑之后总结的经验。第一任何依赖外部的东西都要有超时不管是 HTTP 调用、Redis 操作还是数据库查询没有超时的调用等于把系统的命门交给了别人。第二日志要打全但不要打太多关键是打三个点入口收到请求、中间关键决策、出口发送结果这样出问题能快速定位。第三上线前一定要做混沌测试手动把 Redis 停掉、把通道响应改慢、把 MQ 消费停掉看看系统会不会雪崩。这些测试在测试环境跑一遍能提前发现 80% 的隐患。消息中心这个系统做出来不难做好很难。难点从来不是写多少个接口、建多少张表而是把那些可能会出问题的地方提前想到、提前防住。等你真正处理过几次线上故障回头看这些设计会发现每一条约束背后都是血泪教训。
返回列表