经典系统设计案例04:即时通讯怎么处理消息顺序
围绕即时通讯在经典系统设计案例场景下的生产落地,拆解线上痛点、底层矛盾、方案演进、Java 代码、隐性坑点、架构取舍和故障复盘。
1. 业务背景与线上痛点
即时通讯在秒杀、支付、订单、IM、短链、Feed、推荐、网关这类高频系统设计场景里不是一个孤立技术点。真实生产环境里,它通常挂在核心入口、交易状态、异步恢复和观测告警之间。经典系统设计不能背模板,真实系统要同时处理并发压力、状态流转、数据一致性、失败补偿和可观测性。
这类链路最常见的线上隐患有三个。第一,业务峰值超过预估后,线程池、连接池、缓存分片或数据库主库先被打满,用户看到的不是少量失败,而是整条核心链路超时。第二,服务间调用存在网络抖动和局部失败,单机里能靠事务或内存状态解决的问题,到了集群里会变成重复请求、部分成功和数据延迟。第三,缺少明确的降级、补偿和告警边界,值班人员只能靠临时 SQL 和人工判断止血,恢复时间不可控。
本篇只解决即时通讯在经典系统设计案例中的架构落地问题:围绕消息顺序、已读回执、离线推送、消息漫游,把入口保护、状态推进、异常恢复和上线观测做成闭环。它不解决所有业务建模问题,也不替代领域规则设计;如果业务状态本身定义混乱,再强的架构也只能把混乱传播得更快。
2. 底层核心原理深挖
系统设计的本质是围绕业务不变量构建状态机、资源保护、异步恢复和可验证的数据闭环。单机方案里,很多判断可以放在 JVM 内存、单库事务或本地队列中完成;但分布式集群下,节点之间没有共享内存,网络调用没有确定结果,超时不等于失败,重试也不等于成功。
从并发模型看,请求不是均匀进入系统,而是按渠道、活动、租户和热点资源形成尖峰。只要某个热点资源把一个资源池耗尽,后续请求会因为排队等待继续放大延迟。从数据库事务看,本地 ACID 只能覆盖单库单连接,跨服务状态推进必须额外处理幂等、消息可靠投递、补偿和对账。从中间件内核看,Redis、MQ、搜索集群、配置中心都可能出现主从延迟、分区迁移、消费者再均衡、连接闪断和节点异构性能差异。
因此,即时通讯的关键不是选一个组件,而是定义清楚几个不变量:什么状态只能推进一次,什么请求可以丢弃,什么结果允许延迟,什么异常必须进入补偿,什么指标一旦异常就要停止放量。没有这些不变量,架构图再完整也只是组件堆砌。
3. 方案迭代演进 + 横向对比
第一代简陋原生方案通常是直接在业务服务里同步处理即时通讯,依赖数据库事务和接口超时兜底。它开发快,但遇到峰值流量、下游抖动或重复请求时,问题会直接打到核心库和核心线程池。
第二代优化过渡方案会引入本地缓存、简单限流、状态字段和异步任务,把一部分压力从主链路移开。这一代能缓解峰值,但策略分散在多个服务里,缺少统一观测和补偿闭环。
第三代生产稳定方案会把策略中心、幂等表、Outbox 事件、补偿任务、灰度开关和链路指标打通。核心链路只做必要状态推进,非核心动作进入异步恢复路径。
第四代高可用终极方案会继续补齐跨机房流量调度、自动故障摘除、容量水位控制、多级降级和演练机制。它能支撑大促和多活,但成本、复杂度、团队要求都会明显上升。
| 方案阶段 | 性能吞吐 | 开发成本 | 运维难度 | 适配流量量级 | 故障风险 | 业务适配范围 |
|---|---|---|---|---|---|---|
| 简陋原生方案 | 依赖单服务和单库能力,峰值下容易排队 | 低 | 低 | 日常低峰、内部系统 | 超时和重复提交会直接污染状态 | 低频、非核心链路 |
| 优化过渡方案 | 能承接中等峰值,但策略分散 | 中 | 中 | 部门级业务、普通活动 | 补偿和观测不足,故障定位慢 | 有一定峰值但一致性要求适中 |
| 生产稳定方案 | 主链路短,异步恢复能力强 | 中高 | 中高 | 中大型分布式业务 | 依赖策略、补偿和监控质量 | 核心交易、订单、库存、履约 |
| 高可用终极方案 | 能承接大促和跨机房流量 | 高 | 高 | 千万级并发、异地多活 | 复杂度高,误配置风险更大 | 大促、交易核心、强稳定性系统 |
4. 必备可视化图形
整体架构部署图如下,重点看入口层、核心服务、数据层、异步补偿和观测平台之间的边界。即时通讯不要只放在一个服务方法里,它必须有策略入口、状态落点和恢复链路。
核心请求流转图如下。正常路径要尽快完成关键状态推进,异常路径要能判断是拒绝、重试、补偿、降级还是人工介入。
这两张图的重点不是画得复杂,而是把真实线上链路画完整:入口怎么进,策略在哪里生效,状态写到哪里,消息从哪里发,失败后谁接手,监控从哪里判断恢复。
5. 生产级可落地代码示例
下面代码只保留核心链路,不包含 Controller、DTO 和启动类。重点是幂等、策略判断、事务边界、Outbox 事件、重试耗尽后的补偿,以及中间件异常时的人工兜底入口。
@Service
@RequiredArgsConstructor
public class JavaSystemDesignCase04InstantMessagingSystemService {
private final ScenarioPolicyRepository policyRepository;
private final IdempotencyRepository idempotencyRepository;
private final CompensationTaskRepository compensationTaskRepository;
private final TransactionTemplate transactionTemplate;
private final RetryTemplate boundedRetryTemplate;
private final RabbitTemplate rabbitTemplate;
private final MeterRegistry meterRegistry;
public ScenarioResult handle(ScenarioCommand command) {
String bizKey = command.bizKey();
String scene = "java-system-design-case-04-instant-messaging-system";
ScenarioPolicy policy = policyRepository.getEnabledPolicy(scene);
IdempotencyRecord existed = idempotencyRepository.findSuccess(bizKey);
if (existed != null) {
meterRegistry.counter("scenario.idempotent.hit", "scene", scene).increment();
return ScenarioResult.fromCached(existed.resultPayload());
}
if (!policy.allow(command.tenantId(), command.priority())) {
meterRegistry.counter("scenario.rejected", "scene", scene).increment();
return ScenarioResult.degraded("当前链路水位过高,已进入降级保护");
}
try {
return boundedRetryTemplate.execute(ctx -> transactionTemplate.execute(status -> {
boolean locked = idempotencyRepository.tryStart(bizKey, Duration.ofMinutes(15));
if (!locked) {
return ScenarioResult.processing("请求正在处理,避免重复推进状态");
}
ScenarioResult result = executeCoreBusiness(command, policy);
idempotencyRepository.markSuccess(bizKey, result.toJson());
rabbitTemplate.convertAndSend("scenario.exchange", "system.event", ScenarioEvent.of(scene, bizKey, result));
meterRegistry.counter("scenario.success", "scene", scene).increment();
return result;
}), ctx -> {
compensationTaskRepository.save(CompensationTask.retryLater(scene, bizKey, command.toJson(), ctx.getLastThrowable()));
meterRegistry.counter("scenario.retry.exhausted", "scene", scene).increment();
return ScenarioResult.accepted("核心状态未确认,已进入补偿队列");
});
} catch (DataAccessResourceFailureException | AmqpException ex) {
compensationTaskRepository.save(CompensationTask.manualReview(scene, bizKey, command.toJson(), ex));
meterRegistry.counter("scenario.fallback", "scene", scene).increment();
throw new ScenarioUnavailableException("即时通讯链路进入保护,已记录补偿任务", ex);
}
}
private ScenarioResult executeCoreBusiness(ScenarioCommand command, ScenarioPolicy policy) {
if (policy.isReadOnlyFallback()) {
return ScenarioResult.degraded("只读兜底已开启,跳过非核心写入");
}
// 这里放当前场景的核心状态推进,例如库存预扣、路由放行、配置发布、数据同步或故障切换。
return ScenarioResult.success(command.bizKey());
}
}
这段代码有几个风险点。第一,幂等锁的过期时间必须大于核心处理和补偿扫描周期,否则会出现同一请求被二次推进。第二,RetryTemplate 必须限制次数和退避时间,不能让重试风暴压垮下游。第三,事务内不要做不可控的远程调用,事件投递最好使用 Outbox 或本地消息表兜底。第四,补偿任务必须带原始请求、业务键、失败原因和重试次数,否则无法支持人工核验。
6. 生产隐性坑点 + 兜底补偿方案
第一个隐性坑是冷热点切换。平时低频的租户、商品、渠道或接口,在大促入口曝光后会突然变成热点。兜底方案是提前做热点预热、入口限流按资源维度拆分,并让策略中心支持单资源快速降级。
第二个隐性坑是重试风暴。调用方、网关、RPC 框架、MQ 消费者如果都配置重试,实际放大倍数可能远超压测模型。兜底方案是统一重试预算,核心链路用短超时和有限退避,重试耗尽后进入补偿队列,不在同步链路硬顶。
第三个隐性坑是异步数据偏差。主链路返回成功后,消息投递、搜索索引、缓存刷新、报表同步可能延迟,用户会看到短暂不一致。兜底方案是定义读写一致性等级:核心状态读主库或读状态表,非核心查询允许最终一致,并通过校验任务修复延迟数据。
第四个隐性坑是跨版本兼容。灰度期间新旧服务同时消费消息或读写同一张表,字段语义不一致会引入隐蔽错误。兜底方案是先兼容读、再双写、最后切读,并保留回滚窗口内的数据修复脚本。
7. 架构权衡取舍
中小单体业务推荐使用数据库唯一约束、事务状态机、少量本地缓存和明确的失败返回。坚决舍弃一开始就引入复杂多活、重型消息编排和过度平台化治理。短板是峰值能力有限,但胜在简单、可维护、故障面小。
中大型分布式业务推荐生产稳定方案:策略中心、幂等表、Outbox、本地消息表、补偿任务、统一指标和灰度开关要配齐。坚决舍弃“所有远程调用都同步等待”和“失败只靠人工 SQL 修”的做法。短板是研发和运维成本上升,需要团队具备稳定性治理能力。
大促超高并发 / 异地多活业务推荐高可用终极方案:入口分级限流、单元化路由、跨机房容灾、异步削峰、自动故障摘除、容量水位控制和常态化演练。坚决舍弃单中心强依赖、跨机房同步强一致和无演练的纸面预案。短板是复杂度高,策略误配本身也会成为事故源。
这套架构能解决的是即时通讯在高峰、失败、重复和延迟场景下的可恢复问题;永远无法解决的是错误的业务规则、缺失的组织协作、没有负责人维护的监控,以及没有演练过的应急预案。
8. 真实线上故障完整复盘
故障现象:高峰期用户重复提交请求,业务缺少幂等和状态推进约束,导致部分记录进入中间态,客服和运营无法快速修复。用户侧表现为请求长时间等待、部分结果重复刷新后才可见,客服开始收到状态不一致反馈。
监控指标异常表现:business_success_rate 指标快速越过告警阈值,核心接口 P99 延迟抬升,错误率和超时数同时上涨,MQ 堆积或补偿任务数量开始增加。链路追踪显示慢点集中在策略判断、状态写入或下游依赖处。
快速应急止血手段:第一时间关闭非核心入口,降低灰度比例,开启只读或降级开关;对异常渠道做临时限流;暂停会放大压力的自动重试;把已进入中间态的数据写入补偿队列,避免人工直接改核心表。
深层根因定位:复盘发现压测只覆盖了单接口成功路径,没有覆盖重复请求、下游抖动、消息延迟和热点资源集中访问。策略配置没有按租户或资源拆分,导致局部热点把共享资源池打满。
短期代码修复:补上业务幂等键、状态推进保护、有限重试、异常分类和补偿任务落库;对核心依赖增加短超时和熔断保护;把关键失败原因打到结构化日志里,支持按业务键检索。
长期架构根治:将即时通讯纳入统一策略中心和稳定性看板,补齐整体架构部署图里的异步恢复链路;压测模型加入混合流量、热点资源、失败注入和补偿积压;发布前必须验证回滚开关和降级预案。
后续常态化预防措施:每次大促前做容量复核和故障演练;每周检查补偿任务滞留;每次发布后观察业务成功率、P99、错误预算和中间态数量;事故复盘必须形成可追踪的架构改进项。
9. 中高阶架构面试追问
- 即时通讯在单机事务里能成立的方案,为什么到分布式集群里会失效?请从网络超时、重复请求和状态推进三个角度回答。
- 如果消息顺序、已读回执、离线推送、消息漫游同时遇到峰值流量和下游抖动,你会先限流、降级、排队还是扩容?判断依据是什么?
- 这套方案里的幂等键应该选业务唯一键、请求流水号还是状态版本号?不同选择会带来什么一致性风险?
- 当 MQ 已经堆积、数据库水位也很高时,补偿任务应该继续跑、暂停还是降速?如何避免补偿反向压垮主链路?
- 如果要把即时通讯扩展到异地多活,哪些状态不能跨机房强同步?哪些数据可以最终一致?哪些链路必须保留人工兜底?