ARTICLE DETAIL

资讯详情

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

消息中心架构设计:分层模型、投递链路与稳定性保障

消息中心架构设计:分层模型、投递链路与稳定性保障 消息中心这四个字听起来像是个后台小模块但凡在一线做过几年业务系统的人都知道它几乎是整个公司里最容易被低估、也最容易在半夜把人叫起来的那套系统。我在一次内部架构评审上被要求用四十分钟讲清楚一版消息中心的架构设计台下坐着的既有做交易的老手也有做增长和投放的同学。讲完之后最大的感受是消息中心的难点从来不在“把消息发出去”而在于“该发的发出去、不该发的别发、发出去的能追踪、出问题能兜住”。这几件事串起来才是完整的消息中心架构设计。这篇笔记就是那次评审前后整理的思考过程包括模型怎么定、链路怎么切、存储怎么选、容量怎么算、线上出了事怎么查。内容偏实战适合正在从零搭消息中心、或者接手了一坨历史遗留发送代码准备重构的同学。不管你是刚入行的后端还是做过几年业务但没系统梳理过这块的人应该都能从里面抠出点能直接用的东西。我不会讲太多空泛的“高可用高并发”口号更多是把每个决策背后的取舍摊开来说。1. 从业务盘点到整体分层先想清楚要解决什么问题1.1 把消息类型盘一遍才知道复杂度在哪任何架构设计的第一步都不是画图而是把需求摊开。消息中心的复杂度完全取决于你承接的消息类型。我通常按四个维度先做分类谁触发、发给谁、有多急、有多重要。按触发方分有交易类支付成功、发货提醒、退款到账、互动类点赞、评论、关注、系统类登录提醒、安全通知、运营类活动推送、优惠券到期。这四类的特征差异极大。交易类要求必达且时序严格用户收不到会直接找客服互动类量大但可容忍丢失晚几分钟甚至丢几条问题都不大系统类对时效要求高但量小运营类则是典型的“量大、可批量、可延迟、但对成本极其敏感”。按收件方分有 C 端用户、商家、内部员工。C 端用户是海量商家是中等规模但逻辑复杂一个商家下面还有子账号、角色权限内部员工量最小但通常要求最严比如故障告警漏了一条可能就是一次线上事故。按紧急程度分有实时验证码、支付结果、准实时订单状态变更、延迟可接受日报、周报。验证码这类对延迟要求在秒级但它的量是脉冲式的晚高峰一个短信通道被打满是很常见的事。把这张表列清楚之后你会发现“一套代码发所有消息”是不现实的。真正合理的做法是底层共用一套投递能力和存储上层按业务域做策略隔离。这个判断会直接决定后面的分层设计。1.2 自建还是接入外部通道一个真实的决策过程渠道层面短信、邮件、App Push 通常都要接第三方通道这一点没什么好纠结的。真正需要拍板的是“编排层要不要自建”。市面上有一些成熟的推送服务平台接入快几天就能跑通。我们当时的评估结论是能做但不能只做这个。原因有三个。第一是数据归属。消息发送记录里含有大量用户行为信息哪些用户收到过什么、什么时候点开过这些数据沉淀在别人那里后续做频控优化、做用户触达分析都会受制于人。第二是策略定制。用户免打扰时段、每用户每日推送上限、同一活动多波次去重这些规则每家公司的业务逻辑都不一样靠平台配置改不动。第三是成本核算。通道费用按条计费谁在什么时候、因为什么原因多发了十万条如果没有自己的编排层做归因这笔账永远算不清。所以最终的形态是通道层用外部服务编排层、策略层、存储层全部自建。这个划分的好处是边界清晰——通道挂了我们换通道策略要调整不依赖任何外部排期。1.3 整体分层五层结构各管什么画架构图的时候我习惯用五层来描述从上到下依次是接入层、编排层、投递层、存储层、治理层。每一层的职责必须单一跨层的直接调用要严格控制。接入层负责统一入口。业务方不应该直接去调通道接口而是通过统一 SDK 或 HTTP 接口提交一个“发送请求”。这一层要做的事包括参数校验、调用方身份识别、限流、返回 msgId。注意这里是“返回 msgId”而不是“返回发送结果”这是异步化的关键——接入层只负责收单不负责送达。编排层是整个系统的大脑。它接收请求后要完成模板渲染、渠道选择、策略校验频控、免打扰、黑白名单、消息落库、投递任务生成。这一层的逻辑最复杂也是最容易出 bug 的地方因为它是有状态的。投递层负责把消息真正送到通道。每种通道对应一个适配器适配器屏蔽掉各家通道的参数差异对外提供统一的send(msg)接口。这一层要处理并发控制、超时、重试、熔断。存储层承载消息本体、收件箱、未读数、投递状态四类数据。后面会单独展开因为选型和分片策略是这套系统里最容易埋雷的地方。治理层不参与主链路但决定了系统的下限。它包括频控中心、审计日志、监控告警、运营后台、灰度开关。很多团队一开始不重视这层等到运营同学要临时停掉某类推送、或者要查“这条消息到底发没发出去”的时候才发现没地方点。分层定下来之后有一条纪律必须守住上层可以调下层下层绝不能反向调上层。我们见过太多把频控逻辑写到通道适配器里的代码结果是换一个通道就得重写一遍策略维护成本极高。2. 核心模型设计消息、模板、任务、收件箱怎么落2.1 四层数据模型与建表思路模型设计是消息中心的地基。我推荐四层模型模板template→ 任务task→ 消息message→ 投递记录delivery。很多人会漏掉“任务”这一层结果就是运营要发一条全员推送只能循环调接口既没法统计也没法撤回。模板层存的是内容和渲染规则。一条模板包含模板编码、渠道类型、标题模板、正文模板、变量定义、审核状态。变量用占位符表示比如尊敬的${name}您的订单${orderNo}已发货。这里有个经验变量必须显式声明类型和是否必填不能靠运行时兜底。我在生产环境见过因为一个变量为空导致整条消息渲染成尊敬的null的事故用户截图发到社交平台上非常难堪。任务层描述的是一次批量发送行为。它记录任务名、模板编码、目标人群、发送时间、限流配置、创建人。任务的价值在于可追溯——运营问“昨天那波推送发了多少、点了多少”直接查任务维度的报表就够了。消息层是真正落地的每一条消息。这里给一个参考建表语句字段看着多但每一个都是踩坑之后加的CREATE TABLE msg_message_00 ( id bigint unsigned NOT NULL AUTO_INCREMENT, msg_id varchar(64) NOT NULL COMMENT 全局唯一ID雪花算法生成, biz_id varchar(128) NOT NULL COMMENT 调用方幂等键, task_id bigint unsigned NOT NULL DEFAULT 0, template_codevarchar(64) NOT NULL, receiver_id bigint unsigned NOT NULL COMMENT 接收者用户ID, channel tinyint NOT NULL COMMENT 1站内信 2push 3短信 4邮件, status tinyint NOT NULL DEFAULT 0, priority tinyint NOT NULL DEFAULT 5, content_ref varchar(256) NOT NULL COMMENT 正文存储引用, plan_time datetime(3) NOT NULL COMMENT 计划发送时间, create_time datetime(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3), PRIMARY KEY (id), UNIQUE KEY uk_biz (biz_id,receiver_id,channel), KEY idx_receiver_time (receiver_id,create_time) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;两个关键点uk_biz这个唯一索引是防重复的第一道防线业务方传相同的 biz_id 进来第二次插入直接失败从数据库层面保证不重发content_ref不直接存正文而是存一个引用正文放在对象存储或者单独的正文表里。这个拆分的理由是正文长度不可控有的模板渲染出来几百字如果跟在主表里行会变得很宽索引效率会明显下降。投递记录层记录每一次投递尝试。一条消息可能有多次投递重试每次都要留痕包括通道返回码、通道侧消息 ID、耗时、失败原因。这一层是排查问题的核心依据务必设计成独立表不要试图塞进消息主表。2.2 状态机与幂等把“不重不漏”落到代码上消息的生命周期必须用状态机管理不能散落成一堆 if-else。我把状态定义成下面这些状态值状态名含义是否终态0INIT已创建待处理否10QUEUED已进入投递队列否20SENDING投递中否30SENT通道已接收否40DELIVERED确认送达部分通道可回执是50READ用户已读是60FAILED最终失败是70CANCELLED被撤回是状态流转必须单向不允许从终态回到非终态。实现上所有状态更新都用带条件的 UPDATE// 只有当前状态是 QUEUED 或 SENDING 时才允许推进到 SENT int rows mapper.updateStatus(msgId, FROM_STATUS_LIST, SENT, channelMsgId); if (rows 0) { // 说明状态已被其他线程改过本次投递结果丢弃避免状态回退 log.warn(status conflict, msgId{}, msgId); }这个“带前置状态的乐观更新”是我认为最值得抄的一招。它在不引入分布式锁的前提下解决了多消费者并发更新同一条消息状态的问题。关于幂等要分两个层面看。发送幂等靠uk_biz唯一索引重复请求直接被数据库挡住。投递幂等要靠通道侧的幂等键很多短信通道支持传一个自定义的 outId重复的 outId 通道会拒绝这能有效防止重试导致的重复下发。如果通道不支持那就只能靠本地投递记录的“是否已成功”来做前置判断牺牲一点性能换正确性。注意幂等键的粒度建议是bizId receiverId channel。有些团队只用 bizId结果一个 bizId 对应多个接收人时全部被去重掉这是个很低级但很常见的坑。2.3 收件箱聚合与未读数最容易出一致性问题的两块站内信、互动消息这类需要用户能看到历史记录的必须要有收件箱模型。收件箱的本质是“以用户维度重新组织的消息视图”。同一份消息数据业务维度存一份谁发的、什么模板用户维度存一份这个用户收到过什么。收件箱表按receiver_id分片按时间倒序排列。查询只需要一个索引(receiver_id, create_time desc)翻页用游标而不是 offset因为深翻页在分片表上性能会急剧恶化。未读数的处理是另一个老大难。我试过三种方案最后选的是 Redis 计数加定期校准。第一种是实时 count 数据库用户每次打开 App 都 count 一次。量小的时候没问题用户量一上来count 语句会成为数据库最大的压力源直接放弃。第二种是纯 Redis 计数所有变更都走 Redis。性能和实时性都很好但 Redis 会丢数据主从切换、内存淘汰一旦丢了未读数就对不上了。第三种是 Redis 计数加离线校准。读路径走 Redis写路径走 Redis 加异步落库同时每隔十分钟用离线任务扫一遍增量消息和 Redis 里的值做校正不一致就覆盖。这套方案兼顾了性能、实时性和最终一致性我认为是目前最平衡的选择。计数用的 Lua 脚本长这样核心是保证判断和更新在一个原子操作里完成-- KEYS[1]: 用户未读数 Hash 的 key如 unread:10086 -- ARGV[1]: 分类字段如 chat -- ARGV[2]: 增量1 或 -1 -- ARGV[3]: 当前时间戳 local cur tonumber(redis.call(HGET, KEYS[1], ARGV[1]) or 0) local next cur tonumber(ARGV[2]) if next 0 then next 0 end redis.call(HSET, KEYS[1], ARGV[1], next) redis.call(EXPIRE, KEYS[1], 604800) return next用 Hash 而不是多个独立 key 有个好处一个用户所有分类的未读数聚在一个 key 里可以整体设置过期时间避免冷用户产生大量短生命周期的小 key。这里next 0的保护是必须的因为已读回执和消息删除可能乱序到达没有这个保护就会出现负数未读数前端会显示成-3非常难看。3. 高并发写入与投递链路把峰值扛住的关键动作3.1 接入层限流三道闸门的设计消息中心被上游打挂是常态因为上游业务方通常在写代码时不会考虑你的容量。所以限流必须在接入层做而且要分三个维度。全局闸门保护系统本身比如整个接入层限制 10 万 QPS超过直接拒绝并返回特定错误码。这个值要基于压测结果留 30% 余量来定不能拍脑袋。调用方闸门保护公平性。给每个业务方分配一个配额交易域 3 万、增长域 2 万、其他域共享 1 万。没有这个隔离一个做活动的同学写了个循环调用能把所有人的短信都挤掉。配额存在配置中心支持动态调整出问题的时候运维可以一键降某方的量。用户闸门保护用户体验。同一个用户一秒内收到五条推送体验是灾难性的。这里除了 QPS 限制更重要的是频控策略同一模板对同一用户一天最多一次、营销类消息一天最多三条、免打扰时段比如 22:00 到次日 8:00不投递营销消息。令牌桶的实现建议用 Redis Lua而不是本地 Guava RateLimiter。原因是接入层是多实例部署本地限流只能做到单机维度集群总量会失控。-- KEYS[1]: 限流 key -- ARGV[1]: 桶容量 ARGV[2]: 每秒补充速率 ARGV[3]: 当前毫秒时间戳 local bucket redis.call(HMGET, KEYS[1], token, ts) local token tonumber(bucket[1]) or tonumber(ARGV[1]) local ts tonumber(bucket[2]) or tonumber(ARGV[3]) local delta math.max(0, tonumber(ARGV[3]) - ts) / 1000 * tonumber(ARGV[2]) token math.min(tonumber(ARGV[1]), token delta) if token 1 then redis.call(HMSET, KEYS[1], token, token, ts, ARGV[3]) return 0 end redis.call(HMSET, KEYS[1], token, token - 1, ts, ARGV[3]) redis.call(PEXPIRE, KEYS[1], 60000) return 13.2 削峰与落库为什么接入层要“先返回再处理”接入层处理完校验和限流之后不该直接把消息写进数据库而是应该投递到消息队列然后立刻返回 msgId。这个设计的核心理由是数据库的写入能力是有上限的而消息队列可以缓冲。大促期间瞬时写入量可能是平时的十倍如果同步落库数据库连接池瞬间打满整个服务跟着挂。队列的选型上我倾向用支持分区顺序的组件。分区键选receiver_id这样同一个用户的消息进同一个分区天然保证了单用户的消息顺序。顺序这件事对交易类消息很重要用户先收到“订单已支付”再收到“订单已发货”逻辑上说得通反过来就会让人困惑。落库环节用批量写。单条 insert 在这个量级下完全不可取正确姿势是消费者攒批比如积累 500 条或等待 200 毫秒就 flush 一次。这里有个细节攒批的大小要可配置因为写入 RT 和吞吐是矛盾的大促期间调大批量、平时调小批量这个动态调整能力在实战中很有用。3.3 渠道适配与并发控制投递层对外的接口必须统一。每个渠道实现一个适配器屏蔽掉参数差异public interface ChannelAdapter { /** 渠道类型 */ ChannelType type(); /** 单次投递 */ SendResult send(SendRequest req); /** 批量投递默认降级为循环单发 */ default ListSendResult batchSend(ListSendRequest reqs) { return reqs.stream().map(this::send).collect(Collectors.toList()); } }并发控制有两层考虑。自身出口限流每个通道通常有配额限制比如短信通道每秒最多 5000 条那投递线程池的核心线程数就不能超过这个数的合理倍数否则大量请求会被通道侧拒绝白白浪费重试次数。通道隔离不同通道用不同的线程池短信通道慢不能影响邮件通道这是最基本的故障隔离。通道的优先级路由也要设计好。一条通知类消息优先走免费的 App Push如果用户长时间没打开 App比如 30 分钟未送达降级到短信。这个降级逻辑能显著降低短信成本我们实测下来能省掉接近四成的短信量。判断“是否送达”依赖通道回执如果是 Push 就要用厂商的回执接口或者用“App 下次启动时上报”这种近似方案。3.4 重试、死信与补偿的三层兜底投递失败的处理必须是分层的不能只有“重试三次”这一招。第一层是即时重试针对网络抖动这类瞬时错误。采用指数退避第一次 1 秒、第二次 5 秒、第三次 30 秒。但要注意重试必须异步化不能在投递线程里 sleep那会把线程池堵死。正确做法是发一条延迟消息到队列到期再消费。第二层是死信队列。重试次数用尽后进死信死信里的消息会被定时任务捞起来分析。这里有个关键判断如果失败原因归类为“不可重试”比如手机号格式错误、用户已注销、模板审核未通过就不要进重试直接标记 FAILED 并记录原因否则会浪费大量资源在注定失败的任务上。第三层是人工补偿。死信消息会进运营后台运营同学可以按批次筛选、修改参数后重新提交。这一层看起来原始但很多事故的最后一道防线就是它。我给的建议是后台的补偿功能一定要支持按 msgId 单条重发和按任务批量重发两种模式并且操作必须留审计日志。实操心得重试次数不是越多越好。我们曾经把营销消息的重试次数设成 8 次结果是一条通道故障持续了二十分钟期间积压的失败消息排队重试把通道恢复后的第一波容量全吃掉了正常的验证码反而发不出去。后来改成关键消息 5 次、营销消息 1 次问题就消失了。4. 存储选型与容量估算算不清楚就会在上线三个月后爆掉4.1 四类数据各自的存储归属消息中心的数据形态差异很大用一套存储硬扛是不明智的。我整理了一张对照表可以直接参考数据类型特征推荐存储理由消息主体写多读少按用户和时间查询分库分表 MySQL事务可靠索引灵活运维成熟收件箱按用户维度倒序翻页分库分表 MySQL 缓存查询模式固定可用用户ID分片未读数高频读写强实时Redis Hash单点性能高数据结构贴合投递记录写多主要用于排查和统计MySQL 归档到列式存储热数据量小冷数据需聚合分析全文检索可选按内容搜索历史消息搜索引擎倒排索引MySQL 做不了报表统计多维聚合T1数仓不影响线上有人问能不能全部用文档型数据库解决。可以做但要注意事务边界。消息状态更新和未读数更新如果不在一个事务里就会出现消息状态是 SENT 但未读数没加的情况。我的做法是未读数走异步最终一致用定时校准兜底不强求强一致因为它对用户体验的影响是可控的而消息状态和投递记录必须强一致因为它们直接影响计费和排查。4.2 容量估算的完整计算过程这部分是我认为最值得认真做的功课。用一组假设数据来演示整个推导过程。假设日活用户 5000 万人均每天收到 8 条消息其中推送 5 条、站内信 3 条。那么日增消息量是 4 亿条。峰值 QPS 估算推送有明显的时间聚集性假设 20% 的量集中在晚高峰的一个小时里那这一小时的量是 8000 万条即 8000 万 / 3600 ≈ 2.2 万 QPS。再考虑余量按 3 倍预留写入侧需要支撑 6.6 万 QPS。这个数字直接决定了分片数量和批量写入的批次大小。存储容量估算单条消息的行数据含索引开销大约 300 字节4 亿 × 300B 120 GB/天。保留 90 天热数据就是 10.8 TB。这还没算投递记录投递记录按平均 1.2 次投递计算量和消息主体同量级所以整体热数据大约 20 TB 出头。分片数计算单表控制在一行 1000 万以内比较舒服90 天的数据总量是 4 亿 × 90 360 亿条需要 360 张表。取 2 的幂次方好做取模所以定为 512 张表。再按 16 张表一个库就是 32 个库。最终形态是32 库 × 16 表 512 片分片键用receiver_id。这个分片数不是拍出来的是按数据量和单表容量反推的。很多团队上线时只分了 8 张表三个月后就得做痛苦的数据迁移前期多分一点基本没有额外成本。缓存容量估算未读数用 Redis Hash每个活跃用户一个 key假设每个 key 加上内部编码开销大约 120 字节5000 万活跃用户就是 6 GB。考虑到 Redis 集群通常需要预留一倍内存做 rehash 和过期清理规划 16 GB 内存的集群就够了成本完全可控。4.3 冷热分离与归档策略90 天以前的数据不是没用只是访问频率极低。归档的目标是让热库保持在一个稳定的量级同时保证冷数据还能查。我的做法是按月分表 归档。每个月切换一次写入表历史表在 T30 之后开始归档到列式存储。归档任务用离线工具跑按create_time分批扫描每批 5 万条避免长事务。归档完成后热库里只保留消息的索引信息和归档位置正文按需回源。这里有个必须注意的点归档不能影响线上查询。做法是给查询接口加一个时间范围判断超过 90 天的查询自动路由到归档表。用户侧的界面也要配合只展示最近三个月的消息更早的引导到“查询历史记录”这个独立入口避免默认查询就打到大表上。5. 稳定性保障降级、隔离和看得见的监控5.1 多渠道隔离与分级降级开关稳定性设计的核心思想是“局部故障不影响全局”。消息中心最容易出问题的就是通道所以隔离的第一刀就切在通道上。每个通道独立线程池、独立超时配置、独立熔断器。短信通道超时设 800 毫秒、Push 设 1500 毫秒、邮件设 5 秒这些值来自各通道的历史 P99 数据再加一点余量。熔断用滑动窗口统计错误率错误率超过 50% 且请求数超过阈值就打开熔断半开状态试探恢复。降级开关要分级设计。我把开关分成三级单通道开关、单业务域开关、全局开关。出故障时的操作顺序是从细到粗先关单通道观察影响面不行再关业务域。开关必须放在配置中心支持秒级推送生效而且要有权限控制——曾经有同学误关过全局推送开关导致所有用户十分钟收不到任何消息这个教训之后我们给全局开关加了二次确认和审批。5.2 监控指标只盯这几个就够了监控指标不在于多在于能不能快速定位问题。我建议按“量、率、时延、错误”四个维度建看板。量维度看三个数接入 QPS、各通道发送量、队列积压量。队列积压是最灵敏的指标一旦开始持续增长说明消费能力跟不上通常五到十分钟后就会影响到用户。率维度看成功率和重试率。成功率要分通道看短信成功率掉到 95% 以下就要警惕Push 的成功率因为受厂商影响本身波动就大所以要看趋势而不是绝对值。时延维度看端到端时延即从接入层收单到用户收到消息的时间。这个指标要按消息优先级分层统计关键消息 P99 应该在 3 秒以内营销类允许到分钟级。错误维度要按错误码聚合。通道返回的错误码要建立字典表把原始错误码翻译成人能看懂的原因比如SMS_1002翻译成“手机号在黑名单”。这一步做过和没做过排查效率差距是数量级的。5.3 压测与故障演练的实操方法压测不能只压接入层。我见过很多团队压测只打接入接口接口 QPS 轻松过十万结果一上生产就挂因为真正的瓶颈在下游的数据库和通道。完整的压测要分三段打第一段打接入层验证限流和收单能力第二段打投递层只统计入队不真实发送验证消费和落库能力第三段做小流量真实发送用测试号或者内部员工号验证通道的真实吞吐。故障演练建议每季度做一次场景包括主通道故障切换、Redis 集群主从切换、数据库主库切换、消息队列分区不可用。演练的关键是要有明确的验收标准比如“主通道故障后 30 秒内自动切换到备用通道切换期间失败率不超过 5%”。没有标准的演练就是走过场。注意故障演练一定要在业务低峰期做而且提前通知客服和运营。我们有次演练没提前打招呼客服接到了一堆用户反馈“收不到验证码”的工单白白浪费了两个人力。6. 常见问题排查实录把踩过的坑整理成速查表6.1 消息重复三种成因和对策重复是最高频的问题。成因通常有三类。第一类是上游重复提交。业务方在重试逻辑里没有做幂等控制网络超时后重试了一次结果两条都到达。对策是在接入层强校验 bizId配合数据库唯一索引兜底。这类问题占比最高接口文档里一定要把幂等要求写清楚并且给业务方提供 SDK 里的自动幂等封装。第二类是消费端重复消费。消息队列在 at-least-once 语义下消费成功但 ack 丢失会导致重投。对策是消费逻辑必须幂等处理前先查投递记录表里是否已有该 msgId 的成功记录。这里要注意查表判断本身有并发风险所以要么加分布式锁要么依赖数据库唯一索引。第三类是通道侧重复。这个最麻烦因为控制权不在我们手上。对策是尽量使用通道提供的幂等键如果没有就在通道适配器里加一层本地缓存记录最近五分钟内发过的receiver content hash重复的直接拦掉。这个方案只能减少重复不能根除。6.2 消息延迟从队列积压到通道限速延迟排查有个固定的思考顺序从下游往上游倒推效率最高。先看通道侧。查通道的响应时间曲线如果 P99 明显升高说明是通道问题这时候应该考虑切换备用通道或者临时降级非关键消息。曾经有一次短信通道在晚上八点集体变慢排查发现是通道侧在做扩容这个信息在通道的公告里所以养成看通道状态公告的习惯能省下大量排查时间。再看消费端。观察消费者线程池的活跃度和队列长度。如果是消费能力不足要看是不是有慢查询或者慢接口拖累。常见的坑是在消费链路里同步调了一个外部接口而这个接口偶尔会慢到几秒把整个消费线程堵住。对策是把这类调用改成异步或者加严格超时。最后看写入侧。批量写入的批次大小如果设置得太大单次 flush 耗时变长会造成周期性的延迟毛刺。这时候应该观察落库耗时的分布如果出现明显的尖峰就把批次调小。6.3 未读数不一致校准任务的实现要点未读数对不上的表现是“小红点和实际消息数不符”。这个问题的根本原因是计数和消息落库不在一个事务里。校准任务的实现思路是以消息表为准按用户聚合出真实的未读数和 Redis 里的值比对。全量校准代价太大所以要做增量——只扫描最近 24 小时内有过消息变动的用户。具体做法是在消息落库成功后往一个“待校准用户集合”里塞 user_id可以用 Redis Set校准任务定期取出这个集合里的用户逐个比对。比对的时候要注意消息的“已读”状态可能在校准过程中发生变化所以要给校准加一个时间快照以快照时刻的状态为准。另外校准任务要有幂等和限速避免一下子把 Redis 打满。我们线上跑了半年未读数不一致的投诉从每天几十条降到了基本为零。6.4 排查速查表把上面这些整理成一张表出问题的时候可以直接对号入座现象可能原因快速排查手段处置动作用户收到两条相同消息上游重复提交 / 重复消费按 bizId 查消息表检查上游幂等补唯一索引消息长时间未送达队列积压 / 通道限速看队列长度和通道 RT扩容消费者或切通道未读数显示负数已读回执乱序抽查 Redis 中的值触发校准任务确认 Lua 保护逻辑部分用户完全收不到频控误伤 / 黑名单查该用户的频控计数和黑名单调整频控规则解除误加黑名单短信成功率骤降通道故障 / 内容被拦截看通道错误码分布切换通道检查模板内容投递记录缺失异步写失败 / 队列丢消息对比消息表和投递记录数补投递记录检查队列配置端到端时延升高慢查询 / 外部依赖阻塞看各段耗时埋点加超时改异步消息顺序错乱分片键选错检查是否按 receiverId 分区改分区键补顺序补偿逻辑7. 上线之后踩过的坑以及这套架构还能往哪走7.1 我踩过的最疼的三次第一次是模板变量没做严格校验。运营配置了一个模板变量名拼错了一个字母渲染引擎找不到变量就原样输出占位符那一波推送全部带着${userName}发出去了。事后复盘问题不在于渲染引擎而在于模板发布前缺少一次试渲染。后来我们在模板审核流程里强制加了一步必须选一个真实用户做试渲染渲染结果人工确认后才能发布。这一步多花两分钟但挡住了后面所有类似的错误。第二次是批处理的批次大小没做动态调整。大促期间流量涨了八倍批量写入的批次还是 500 条结果单次 flush 耗时从 30 毫秒涨到 400 毫秒数据库连接被占满整个写入链路雪崩。后来做了两件事批次大小按当前队列积压量动态调整积压多的时候按 200 条快写、积压少的时候按 1000 条慢写以及给数据库连接池加了获取超时超时后快速失败而不是无限等待。第三次是全局降级开关的权限没收紧。这个前面提过了影响是十分钟内所有推送都没发出去。修复方案除了加二次确认还加了一个“影响面预估”功能——点开关之前会弹出提示告诉你这个开关影响多少条消息、多少用户看到“影响 5000 万用户”这个数字手就会停一下。7.2 这套架构后续可以怎么演进如果业务量再涨一个数量级我认为最先需要动的不是存储而是编排层。因为编排层是有状态的也是逻辑最复杂的它横向扩容的成本远高于存储。可行的方向是把编排层拆成“策略计算”和“消息组装”两个独立的无状态服务策略计算结果做缓存让高频的重复计算能被复用。另一个方向是把频控和人群圈选下沉到独立的服务。现在很多团队的频控逻辑散落在各个业务代码里规则不统一出了问题也不知道该找谁。把频控做成一个独立的决策服务输入是用户和场景输出是“允许/拒绝/延迟”这样规则能集中管理也方便做 A/B 实验。还有一个容易被忽略的方向是成本归因。消息中心的通道费用通常不小但很少有团队能精确回答“上个月的钱花在哪个业务域、哪个活动上”。把费用按调用方维度做归因输出的报表可以直接驱动业务方优化自己的发送策略这比技术团队自己去限制别人有效得多。至于更远的事情比如智能发送时机预测、基于用户行为的个性化频控阈值我觉得都是可以探索的。但有一点我比较坚持在基础能力没打牢之前不要过早引入智能化的东西。消息中心的价值首先体现在“稳定、可查、可控”这三点上锦上添花的功能永远排在它们后面。我个人在实际推进这套东西的时候最大的体会是消息中心的每一次事故回头看几乎都能归到一个“当时觉得没必要做”的细节上——一个唯一索引、一个超时配置、一次试渲染。做架构设计最难的不是想出多花哨的方案而是把那些朴素但必要的环节一个都不落下。
返回列表