企业微信的回调服务上线第一周通常都很顺 —— 收到消息、处理、回复,一次 HTTP 请求里做完,代码短得让人怀疑为什么还要引入消息队列。问题都发生在之后:接了个几百人的活动群、下游数据库慢了一下、模型接口抖了一秒,整条链路一起塌。队列不是为吞吐加的,它买的是「这三类抖动不至于变成事故」。下面按 wecomapi 的回调接入方式,讲同步处理的三种崩法、队列该按什么切分,以及重放与死信这两件最容易做错的事。
同步处理的三种崩法
在回调请求的生命周期里把业务处理完,出事的方式有三种,按代价从便宜到贵排下来。它们都不是高并发才会遇到的问题 —— 前两种在日均几万条消息的量级就足以复现。
第一种是重试雪崩。回调方普遍有超时上限,超时或收到非 2xx 就会重投(具体阈值与重试次数以文档为准)。你的处理逻辑一慢,第一次投递超时,平台重投,重投的请求和原请求一起占着连接和线程,系统更慢,更多请求超时。重试在这里是放大器,不是补偿机制 —— 等你察觉的时候,在途请求量已经和真实消息量没关系了。
第二种是故障域连坐。回调是一个共享入口,所有账号、所有会话的事件都从这里进。同步处理意味着一个热点群的消息能把工作线程占满,其他几十个账号的消息全排在后面等 —— 一个客户的活动搞崩了所有客户的服务。这类故障最难向业务解释,因为出问题的模块和受影响的模块看起来毫无关系。
第三种最贵:静默丢失。你按最佳实践先返回了 2xx 再处理,然后进程重启、下游超时、某行代码抛了个没接住的异常,这条消息就没了 —— 平台那边显示投递成功,你这边没有任何记录。它不报警、不体现在任何指标上,通常是几天后客户问「我那条消息你们怎么没回」才被发现。前两种崩是响的,这一种是哑的。
很多团队照着「先快速 ACK 再异步处理」改造完,其实只解决了第一种,第三种反而被放大了 —— ACK 变快之后,丢失的窗口更隐蔽。这句建议少说了一个前提,见下一节。
ACK 之前唯一该做的事是持久化
「先返回 2xx 再处理」省略的前提是:ACK 之前必须先把原始报文可靠地落下来。落到哪都行 —— 消息中间件、数据库一张表、追加写的日志文件 —— 但必须是持久的,而且要落成功了才 ACK。落失败就返回非 2xx,让平台按它的重试策略重投,那正是重试机制该发挥作用的地方。
把顺序反过来(先 ACK 再往内存队列里塞)看起来快了几毫秒,代价是丢掉了平台的重试保障:入队那一步一旦失败,这条消息在世界上就不存在了。几毫秒不值这个价。
这条边界定下来之后,回调处理函数就只剩三件事:验签、持久化、返回。任何涉及业务判断、外部调用、模型推理的代码都不该出现在这个函数里,包括那些「反正很快」的查询 —— 它们平时确实很快。
// 示意逻辑,验签算法与事件字段名以线上文档为准
app.post("/wecom/callback", async (req, res) => {
if (!verifySignature(req)) return res.sendStatus(401);
try {
// 落盘成功才 ACK:原始报文 + 接收时间 + 去重键 + 分区键
await inbox.append({
raw: req.body,
dedupeKey: eventIdOf(req.body), // 事件唯一标识,字段名见 wecomapi 文档
partition: sessionIdOf(req.body), // 按会话分区,保证会话内有序
receivedAt: Date.now(),
});
} catch (e) {
return res.sendStatus(500); // 落盘失败就别 ACK,让平台重投
}
res.sendStatus(200);
});队列怎么切分:横着按类型,竖着按会话
一个大 topic 全塞进去,等于把队列当成了一次更慢的同步调用 —— 故障域没变,只是延后了。企业微信消息队列的切分要按两个维度同时做:横着按事件类型分消费组,竖着按会话分区。
按事件类型分消费组,是因为不同事件的处理耗时和失败率差着数量级。文本消息触发 AI 回复要几秒且依赖外部服务,wecomapi 的群成员变更事件只是更新一行数据。放在同一个消费组里,慢的会把快的堵死,而快的那一类往往是维持状态一致性的关键路径。
按会话分区,是因为有序性的真实需求只到会话这一层。全局有序没有业务价值、还堵死横向扩展;账号级有序看着更保险,代价是账号里一个热点群就能拖慢这个账号的全部会话。分区键取会话标识,同一会话内串行、不同会话之间并行,这是并发度和有序性之间最划算的那个点。
- AI 生成、文件下载这类秒级任务单独走慢通道,不要和毫秒级的状态更新共用消费组
- 分区数一开始就设得比预期消费者数大一截 —— 事后扩分区会打乱既有会话的顺序
- 每个消费组独立配置并发度、重试次数与超时;共用一套参数等于按最慢的那类来配
量小的时候不需要引入专门的消息中间件。一张带状态列的表加一个轮询 worker,同样能拿到解耦、重试、可重放这三样东西。队列在这里是一种处理形态,不是某个具体组件 —— 别为了架构图好看先上一套自己运维不动的东西。
幂等要做在产生副作用的那一层
重复投递一定会发生,这点大家都知道,常规做法是在消费入口用事件唯一标识去重。这一层是必要的,但只做这一层挡不住真正会被客户感知到的重复。
因为消费失败重试也会产生重复,而这次事件标识是同一个、消费入口的去重记录可能已经写下了。危险的位置在链路末端:消息已经发出去了,紧接着写业务库失败,任务被判失败进入重试 —— 重试时入口去重放行,客户就收到了第二条一模一样的回复。
所以幂等键要一路带到产生外部副作用的那一步上。调 wecomapi 的发送消息接口之前,先用「会话 + 触发事件 + 动作类型」拼出一个业务幂等键,在自己的发送记录表上按这个键占位;同一个键第二次进来直接返回上次的记录,不再真的发出去。这层去重只能在你这边做,别假设发送接口会替你识别重复请求。
- 消费入口去重:挡平台的重复投递,成本低,先做
- 发送侧幂等:挡重试与并发导致的重复外发,这一层才是客户体感的最后一道闸
- 去重记录带过期时间,永久保留既没必要,也会让这张表变成新的瓶颈
重放和死信怎么设计
重放能力的前提在第二节已经具备了:你存了原始报文。剩下要设计的是重放时要不要重新产生外发动作,而这个问题的答案几乎总是「要能选」。
真实发生过的事故长这样:修完一个解析 bug,把过去两小时的消息重放一遍补状态,结果自动回复对着几百个客户又发了一轮昨天的话术。重放接口必须支持一个只重建内部状态、屏蔽全部外发调用的模式,并且把这个模式设成默认值 —— 要真的发出去,必须显式打开。
死信队列同样不能当垃圾桶。消息进死信时必须带上失败分类,因为处置方式完全不同:可重试的(下游超时、被限流)应该退避后自动回来;不可重试的(报文不合预期、权限不足)重试多少次都是浪费,还会挤占真正该重试的那些消息的处理能力;毒消息(能稳定让消费者崩溃的那一条)必须立刻隔离,否则它会卡死整个分区,后面排队的消息一条都出不去。
- 重试用指数退避加次数上限,不可重试的错误直接进死信,别拿重试去刷频控
- 死信条目要能一键回放到原消费组,也要能标记为「已人工处理」后关闭
- 把死信数量做成告警指标 —— 它比错误日志更能反映系统的真实健康度
本文讲的是队列形态与工程取舍。回调的验签算法、事件结构与精确字段以 wecomapi 线上接口文档为准,示意代码不要直接照抄上生产。
常见问题
- 消息量不大,一天几千条,也要上队列吗?
- 要,但不必上专门的消息中间件。这里队列买的是故障隔离、失败重试和可重放,不是吞吐;一张带状态列的表加一个轮询 worker 就能拿到这三样。真正不该省的是顺序:原始报文落盘成功之后再 ACK。量小的时候丢一条消息的概率不低,只是没人发现。
- 应该先返回 2xx 再入队,还是先入队再返回?
- 快速 ACK 是对的,但「快速」指的是不在 wecomapi 的回调里跑业务,不是跳过落盘。ACK 之前该做的只有两件:验签,以及把原始报文可靠落下来;落失败就返回非 2xx,让平台按重试策略重投。反过来做的话,入队失败时消息会静默消失 —— 平台侧显示投递成功,你这边没有任何记录,也不会报警。
- 队列的分区键该取账号还是取会话?
- 默认取会话。有序性的真实需求只到会话这一层,账号级有序会让一个热点群拖慢该账号下的所有会话,并发度也上不去。只有当业务确实依赖账号内的全局顺序时才提升到账号级,并为此接受吞吐下降 —— 这种需求比大多数人以为的少得多。
准备好动手了?
精确字段、鉴权与端点以线上文档为准;可在控制台创建密钥后联调。
