系列目录
- 从单线程到线程池:云盘转 Java 后的第一堂并发课
- 线程池不是 new 出来就完事:参数、队列与快慢接口隔离
- 数据库连接池与 Spring 声明式事务:把 Node.js 的坑填上
- ConcurrentHashMap 与锁:文件元数据的并发读写
- Future 与 CountDownLatch:一个接口聚合一堆下游
- 用 Kafka 做异步化与削峰:从热点上报到全局事务(本篇)
第 3 篇梳理过 Spring 声明式事务的范围:它保证单个数据库内的 ACID,写库与多端通知这类跨资源操作需要单独处理。当时我已有 Node.js、C# 等后端经验,转入 Java 项目前先预判了同步调用、资源竞争和跨资源一致性可能带来的问题,提前学习了消息队列和最终一致性的处理方式。云盘随后用 Kafka 异步发送通知,并通过重试、幂等和对账处理最终一致性。下面记录两个实际链路。
热点上报占用主接口资源
去年做的云盘点击热点图,最终链路是「采集 → Kafka → Java 消费 → MySQL 分钟级入库」。那篇文章记录的是最终形态,最初的实现是同步写库。
热点上报接口收到坐标后,直接 insert 到 MySQL,再返回成功。UED 使用的数据没有实时要求,但接口实现与普通业务接口相同,同步写库是当时最直接的选择。
流量增加后,限制开始显现。热点 SDK 注册在 document 的 click 事件上,用户在云盘页面每点击一次就会产生一次上报,一个活跃用户浏览一段时间可以产生几十条。白天工作时段全量用户叠加,上报接口的流量接近主接口量级(具体数字按惯例模糊)。MySQL 热点表的写入 RT 升高,连接被大量短事务占用。第 3 篇讨论过的连接池竞争再次出现:热点接口与普通接口共用连接池,热点写入变慢时,其他接口在借连接时排队。
为了给 UED 提供辅助数据而占用主站接口资源,成本过高。当时先做了降级:上报接口增加开关,高峰期收到请求后直接返回并丢弃数据。这样可以保护主链路,代价是高峰期缺少热点数据。
这个需求的约束明确:
- 热点数据允许少量丢失。它是统计数据,分钟级聚合后少几条不会改变图形;
- 数据允许在数分钟后入库。UED 看的是趋势,没有实时要求;
- 上报处理不能影响主站接口。
消息队列适合允许延迟、允许少量丢失且需要隔离主链路的场景。
上传后的多端通知
第二个链路在核心路径上,影响范围更大。
云盘是多端产品:web、PC、Mac、Android、iPhone 都在线。用户在 web 端上传文件后,其他端需要尽快感知文件列表和配额变化。第一版采用同步调用:上传接口写入元数据后,在同一个请求中依次调用各端通知通道。
起初只有两个端,增加两次调用的成本可以接受。通知通道增加后,接口逐渐接近第 5 篇描述的情况,RT 受多个下游调用的累计耗时影响。通知属于不可靠的外部调用,任一通道抖动或超时都会拉长上传接口 RT。元数据已经写入成功而某个通知超时时,客户端会收到「上传失败」并重传,服务端还要处理重复上传。
元数据写入成功后,上传结果已经确定。通知属于后续动作,可以稍后送达。PC 端列表晚一秒刷新,在当时的产品体验中可以接受。
两个链路的处理方式相同:把同步链路中允许延迟的步骤拆出,异步执行。云盘当时已经有 Kafka。转 Java 时的架构设计写过「通过 Kafka 做异步消息通知、保证系统最终一致」,公司也提供了 mafka(内部中间件。这里仅描述使用方式:申请 topic、取得生产者和消费者客户端封装)。我先把 Kafka 用在热点上报,再将多端通知迁移到这条链路。
消息队列解决的问题与条件
这两个链路说明了 MQ 可以解决什么,也说明了使用条件。
- 削峰。 高峰流量到来时,消息暂存在 Kafka,消费者按自身处理能力消费。数据库不必直接承受全部点击峰值。热点上报改造后,MySQL 的写入速率由消费者处理能力和积压量决定,不再随用户点击同时变化。
- 异步。 接口完成关键路径后返回,后续动作通过消息交给消费者。上传接口改造后,RT 仍包含写入元数据和投递消息的成本,但不再等待各通知通道完成。
- 解耦。 生产方发布业务事件,消费方按需订阅。新增通知渠道或统计需求时,可以增加消费组,不需要修改上传接口。前提是新消费方理解事件格式,并能处理历史消息和重复消息。
设计前还需要确认以下限制:
- 延迟。 消费时间由消费者能力、积压和重试决定。热点数据数分钟后入库、多端通知数秒后到达可以接受,有些操作仍需保留在关键路径。
- 顺序。 Kafka 仅保证单个 partition 内的追加顺序。同一文件的事件需要使用稳定分区键,例如文件 ID,才能进入同一个 partition;同一 consumer group 中,一个 partition 同一时刻只分配给一个消费者。不同 partition 的消息没有全局顺序,重试和消费失败也会带来等待或重复。按文件处理顺序时,消费端还要依据该文件的单调版本号判断事件新旧,时间戳不足以完成这个判断。
- 投递失败和重复处理。 Kafka 0.8/0.9 没有生产者幂等投递和跨 Kafka、数据库的原子提交。生产端在未确定 broker 是否接收消息时重试,可能写入重复记录。消费者在业务处理后、提交 offset 前崩溃,恢复后也会再次处理该消息。至少一次处理语义要求业务处理成功后再提交 offset,并由消费端实现幂等。自动提交 offset 可能在业务处理完成前推进位点,使当前消费组无法重新读取该消息。
MQ 的适用范围取决于业务能否接受延迟、重复和分区边界外的无序。支付性质的配额扣减校验仍应放在同步事务中。拆出的操作需要明确允许的延迟、顺序要求和补偿方式。
最终一致性的处理步骤
链路拆开后,多个资源不会在同一时刻一致。这里的目标是最终一致:业务数据提交后,相关事件经过投递、消费、重试和补偿,使下游状态与业务状态对齐。实现由四部分组成:
- 本地事件表。 本地事务无法覆盖 Kafka。业务数据提交后再投递消息时,进程可能在两者之间退出并遗漏通知。我们在同一个数据库事务中写入业务数据和一条「待发送」事件记录,后台任务扫描事件表并投递至 Kafka。收到生产者确认后,将记录更新为「已发送」。投递成功后状态更新失败时,扫描任务会再次投递。因此 outbox 保证业务数据和待投递事件的原子写入,重复投递仍需由后续环节处理。
- 消息重试。 消费端遇到下游通道超时或数据库故障时,不提交该消息的 offset,使它在后续轮次重新处理。重试需要限制次数并设置间隔。达到上限后,消息及失败原因需要保存到可处理的位置,再由补偿任务或人工处理;完成接管后才可以推进对应 offset。
- 消费幂等。 生产端和消费端都可能产生重复。消费逻辑需要保证同一消息处理一次与处理 N 次得到一致的业务结果。具体做法在改造部分说明。
- 定时对账。 前三层无法证明所有下游都已对齐。定时任务检查长期停留在「待发送」或发送中的事件,比较生产侧和消费侧可比较的数据,并对差异告警或补偿。对账频率和阈值决定发现问题的时间,不能提供实时一致性。
我们后来半开玩笑地称这套组合为「全局事务」。它不提供跨资源的原子提交,而是将全局操作拆成可重试的本地步骤,通过幂等、补偿和对账追平结果。它比 XA 弱,适合当时可以接受最终一致的云盘场景。
生产与消费规范
两个链路落地后,我们将这些约束整理为团队规范。
生产侧:
- 事件表与业务数据在同一事务中写入。需要可靠通知的业务事件通过后台扫描任务投递,避免只依赖业务代码中的一次直接发送;
- 消息体携带唯一 ID(事件表主键)、事件发生时间和版本号,供消费端去重和判断新旧。版本号需要按有顺序要求的业务对象递增;
- topic 按事件类型划分,不按消费方划分。消费方的增减由消费侧处理。
消费侧的重点是幂等。关键代码如下(脱敏简化):
@Transactional // 去重插入与业务处理在同一个本地事务里提交
public void onMessage(String msg) {
FileEvent event = parse(msg);
// 幂等第一道闸:消息唯一键 + 消费者名做去重表主键,重复插入影响行数为 0
int rows = consumedDao.insertIgnore(event.getMsgId(), CONSUMER_NAME);
if (rows == 0) {
return; // 已消费过,直接确认,不重复执行
}
doConsume(event); // 真正的业务处理:多端通知 / 热点入库 / …
// 业务失败则抛出异常,整个事务回滚。回滚后去重记录不存在,重投自然再进来;
// 不确认 offset,等待重试
}这段代码的事务只能同时覆盖去重表和 doConsume 所在数据库的写入。doConsume 调用外部通知通道时,事务回滚无法撤销已经发送的通知,通知接口自身仍要能按消息 ID 去重,或者接受重复通知。调用方应在本地事务提交成功后提交 offset;代码中的「确认」和 offset 提交不在该方法内,需要由实际消费者客户端按此顺序实现。
配套约束如下:
- 重试三次为限,间隔递增;三次仍失败的消息保存到死信表并触发告警,由对账任务或人工处理。消息转入死信处理时需要保存可恢复记录,避免提交 offset 后丢失失败上下文;
- 去重表和业务处理尽量落在同一个库。去重记录插入和业务写入在一个本地事务中提交,可以避免「去重记录写了、业务没写」的状态;
- 消费逻辑不假设跨 partition 顺序。需要按对象排序的场景使用消息中的递增版本号裁决,晚到的旧版本不能覆盖新版本;
- 对账任务每天运行:事件表中「待发送」超过阈值未清时告警,生产和消费两侧可比较的计数差异超过阈值时告警。
改造后,热点高峰的消息会在 Kafka 积压,由消费者按能力处理,MySQL 不再在请求路径中承受每次点击的写入。上传接口不再等待通知通道完成。后来新增通知渠道和统计订阅时,主链路没有修改。这些结果依赖消费者容量、Kafka 积压监控和失败补偿持续可用。
当时对 Kafka 可靠性的了解范围
还需要说明当时的了解范围。方案将 Kafka 作为可信基础设施使用:业务侧假设已确认的消息由集群保存,但我当时没有真正搞清副本机制、ISR 和 broker 故障后的 leader 选举。例如,ISR 收缩时哪些配置组合会扩大数据丢失窗口、min.insync.replicas 应如何设置,当时答不上来。
公司 mafka 由中间件团队维护,集群副本配置和容量规划不由业务团队负责。业务侧仍要保证投递失败可以重试、消费可以幂等,并通过对账发现差异。这些措施处理的是业务事件从本地事务到下游处理之间的失败窗口,Kafka 集群可靠性仍取决于集群配置。
下一篇记录一次 Full GC 排查:线上服务没有发版、流量没涨,却开始周期性卡顿。
参考资料
- Kafka 0.8/0.9 官方文档(生产者与消费者投递语义部分)
- 云盘服务端从nodejs 专项 java 相关复盘(架构 2 的设计,原文措辞为「通过 kafka 进行全局事务维护,保证整个系统是最终一致的」)
- 云盘点击热点图

