宁语之溪
众里寻他千百度,蓦然回首,那人却在灯火阑珊处。 (宋·辛弃疾·青玉案)
自豪地使用 Typecho 建站搭配使用 🌻Sunny 主题博主 昨天 13:54 在线 · 当前 45 人活跃
宁语之溪
独在异乡为异客,每逢佳节倍思亲。 (唐·王维·九月九日忆山东兄弟)
自豪地使用 Typecho 建站搭配使用 🌻Sunny 主题博主 昨天 13:54 在线 · 当前 46 人活跃
文章

RocketMQ 异步消息链路中的降级思路

语之溪

·

Java

·

⚠️ 本文最后更新于2026年05月16日,已经过了100天没有更新,若内容或图片失效,请留言反馈

HAOVP 使用 RocketMQ 处理部分视频互动统计。点赞、收藏、播放、分享和弹幕等操作产生统计消息,消费者再更新视频主表和每日统计数据。

4 月接口验收时,RocketMQ 路由仍存在异常。内容服务增加了直接更新视频主表的降级路径,使弹幕等当前接口用例能够得到预期的核心计数结果。

这次处理让我更清楚:消息发送失败时做一次数据库回写,可以保住某个核心结果,但不等于整条异步消息链已经可靠。

正常消息链做了什么

一条简化的统计链路是:

业务操作完成
    -> 生产统计消息
    -> RocketMQ 接收和存储
    -> 消费者取得消息
    -> 更新视频主表计数
    -> 更新每日统计数据

使用消息队列的目的,是让主请求不必同步等待所有派生统计完成,同时把统计处理从业务操作中拆出来。

消费者处理的不只是一列总数。正常路径还会更新按日期保存的增量统计。因此,消息链承担了两部分工作:

核心展示计数
每日统计明细

降级时只更新主表,能够让页面上的核心计数发生变化,却不能自动补齐每日统计。

当前降级路径怎样工作

统计消息使用异步发送。发送成功时只记录结果;发送异常或异步回调失败时,生产者直接按动作类型更新视频主表。

可以抽象为:

messageQueue.asyncSend(message, new Callback() {
    @Override
    public void onSuccess() {
        log.info("message sent");
    }

    @Override
    public void onFailure(Throwable error) {
        updateMainRecord(message);
    }
});

这是脱敏结构示例,不包含实际主题、服务、字段和认证信息。

这条降级路径的优点很直接:消息组件不可用时,播放、点赞或弹幕等核心计数不必完全依赖消息消费才能变化。

它的边界也同样明确:

  • 只更新主表,不更新消费者负责的每日统计;
  • 失败只记录日志,没有持久化失败事件;
  • 应用退出后,内存中的异步回调无法继续;
  • 主请求成功与计数更新成功仍是两件事;
  • 消息状态不确定时直接降级可能带来重复更新。

发送失败有不同阶段

“发送失败”不能只看成一种结果。

调用发送方法前失败
连接 Broker 失败
Broker 未接受消息
Broker 已接受,但响应丢失
回调执行异常
消息发送成功,但消费者失败

前几种情况可能适合在生产者侧降级;最后一种已经不属于发送失败,生产者通常不知道消费者是否完成。

最难处理的是状态不确定:Broker 可能已经接收消息,但发送方没有得到明确成功响应。如果此时直接更新数据库,而消息随后又被消费者处理,同一个计数可能执行两次。

因此,可靠降级不能只靠捕获异常,还要回答:

  • 能否确认消息没有进入 Broker;
  • 降级更新和消息消费是否具备幂等性;
  • 是否保存了事件唯一编号;
  • 后续如何识别和修正重复结果。

为什么直接加一容易重复

统计消息通常携带动作和增量:

video = 某个资源
operation = like
increment = +1

如果同一事件被消费两次,主表会连续加两次。单纯的 count = count + 1 无法判断这次增量是否已经应用。

可以给事件增加唯一标识:

{
  "eventId": "example-event-id",
  "resourceId": "example-resource-id",
  "action": "like",
  "delta": 1
}

示例中的标识均为占位值。

消费者先检查事件是否处理过,再更新计数并记录处理结果。这里仍要考虑“更新业务数据”和“记录已消费”是否处于可靠事务中,否则两步之间失败也可能产生重复或遗漏。

主表降级不等于完整消费

正常消费者同时维护主表和每日统计,降级路径只改主表后,会出现:

视频主表计数:已经增加
每日统计记录:没有对应增量

页面可能看起来正常,但运营统计或趋势查询出现缺口。

可以选择的补偿方向包括:

  • 将失败事件写入可靠的本地表,稍后重放;
  • 定期用互动明细重算主表和每日统计;
  • 对主表与统计表做差异检查;
  • 恢复消息组件后补发未完成事件;
  • 让统计查询明确区分实时总数和每日聚合结果。

当前代码只支持直接主表更新,不能据此声称统计数据已经完整一致。

降级是否应该影响主请求

不同业务需要不同选择。

对于可修复的展示计数,可以接受:

核心明细成功
统计消息失败
    -> 主请求仍成功
    -> 记录失败事件并补偿

对于不能丢失的关键状态,如果异步链失败后又没有可靠落盘,就不应简单返回成功。

判断依据包括:

  • 主业务事实是否已经保存;
  • 派生数据是否能从明细重建;
  • 用户是否依赖计数立即准确;
  • 失败事件是否可以可靠重试;
  • 重复执行是否安全;
  • 补偿需要多长时间。

“接口可用”不能以静默丢失不可重建数据为代价。

生产者、Broker 和消费者分别监控什么

消息链要可排查,至少需要三段指标。

生产者

  • 发送请求数、成功数和失败数;
  • 同步与异步发送耗时;
  • 失败原因和超时;
  • 触发降级的动作类型;
  • 本地失败事件数量。

Broker

  • 节点是否可达;
  • 主题和路由是否存在;
  • 消息写入和积压;
  • 存储、磁盘和复制状态。

消费者

  • 消费成功、失败和重试;
  • 消费延迟与积压;
  • 重复事件数量;
  • 死信消息;
  • 主表和每日统计更新结果。

只看接口返回 200,无法判断消息到底走了正常消费还是数据库降级。

怎样验证降级路径

我会分别验证正常、发送失败、消费失败和恢复四类场景。

正常发送与消费

  • 消息发送成功;
  • 消费者只处理一次;
  • 主表计数正确;
  • 每日统计同步更新;
  • 日志可以用事件标识串联。

发送失败

  • 能明确触发生产者降级;
  • 主表只更新一次;
  • 失败事件可被发现;
  • 每日统计缺口有补偿方式;
  • 主请求的返回符合业务约定。

消费失败

  • Broker 中消息仍可重试;
  • 重试不会重复计数;
  • 多次失败后进入可观察状态;
  • 失败原因修复后能够继续消费。

组件恢复

  • 恢复后积压消息是否继续处理;
  • 已经走过降级的事件是否再次消费;
  • 主表和统计表是否出现重复;
  • 对账能否发现并修复差异。

没有这些场景,就不能用“有异常回调”证明降级链路已经完整。

重试不能没有上限

发送或消费失败后立即无限重试,可能进一步压垮正在恢复的组件。

重试需要考虑:

最大次数
退避时间
随机抖动
错误是否可重试
事件幂等性
最终失败去向

参数要根据真实故障和容量确定。当前还没有经过验证的重试次数和间隔,因此我暂时不写固定数字。

当前能够确认的结果

4 月验收能够确认:RocketMQ 路由当时仍有异常,内容服务增加主表降级更新后,弹幕发送等指定用例中的核心计数能够变化。

当前不能由此确认:

  • 主题路由已经完全恢复;
  • 所有消息都会可靠送达;
  • 消费失败会按预期重试;
  • 重复消息不会重复计数;
  • 每日统计没有缺口;
  • 组件恢复后不会再次处理已降级事件;
  • 消息链达到了零丢失。

这次降级的价值,是在消息路由异常时保住当前核心计数;它暴露出的下一步也很清楚:事件需要可识别、失败需要可持久化、消费需要幂等、主表与统计表需要对账。

异步消息链的降级不只是“发送失败就改数据库”,而是要在业务可用、重复风险和数据完整性之间划出明确边界。只有正常、失败、重试和恢复都能验证,才能说这条链路具备可靠性。

现在已有 9 次阅读,0 条评论,0 人点赞
评论:共0条
发表
搜索 消息 足迹 排行
你还不曾留言过..
你还不曾留下足迹..
博主 不再显示
博主
未知作品 歌曲封面
立即安装