在双 11 这种每秒数十万笔交易的极端流量下,分布式消息中间件 Kafka 在提供极致吞吐削峰能力的同时,对下游消费者提出了一个极其苛刻的技术前提:Kafka 默认只保证“至少一次(At-Least-Once)”投递。
在大促高压期间,网络偶发丢包、Broker 节点由于垃圾回收发生数毫秒的短暂假死、或者消费者实例因为内存紧张触发超时重平衡(Rebalance),都会导致大量消息的位点(Offset)无法及时成功提交。当新的消费者节点重新接管分区时,必然会将最后几百条甚至上千条已经消费过的消息重新拉取并重新投递一遍。
如果下游消费端是一个涉及资产变更的履约服务(例如:发放优惠券、账户充值、积分扣减、向物流商推送发货单),一旦重复消费了一条消息,就会导致用户收到双份红包、或者同一笔订单被扣了两次钱。这种直接造成严重资金损失与财务混乱的事故,在大厂属于不可触碰的 P0 级严重缺陷。
在百万 TPS 的洪峰下,如何设计一套既能抗住毫秒级高并发冲击、又能 100% 杜绝重复扣款的消费端防重(Idempotent Consumer)流水校验架构?
业界最严密的工业级实践,正是依托“Redis 极速前置过滤 + MySQL 唯一键终局兜底”的双层联合防重体系。
为什么单纯依赖 Redis 或单纯依赖 MySQL 都是不可取的
在防重设计上,很多团队容易走入两个极端,最终在生产大促中付出惨痛代价:
- 单纯依赖 Redis 分布式锁/缓存的“丢锁漏网”隐患:
有人认为只要在消费前执行redis.set(order_id, "1", "NX", "EX", 86400)就能搞定防重。然而:- Redis 是基于内存的 AP/弱一致性存储。如果在大促期间 Redis 发生了主从切换(Failover),未同步到从节点的防重 Key 会瞬间丢失,导致重复消息穿透;
- 即使 Key 写入成功,若后续数据库事务在提交时抛出了异常,Redis 里却已经标记了“已处理”,导致真实的业务订单反而被永远遗漏,造成严重的业务卡单。
- 单纯依赖 MySQL 唯一索引插入的“性能雪崩”:
有人主张每次消费都向一张 MySQL 防重表INSERT INTO t_idempotent (id, ...) VALUES (...),依靠数据库主键唯一约束来防重。
然而在大促峰值期,几十万 QPS 的消息并发同时向单张 MySQL 防重表执行插入与锁检查,数据库的行锁、B+ 树页分裂(Page Split)与死锁检测机制会被瞬间压垮,写入延迟从 1ms 暴增至数秒,数据库连接池瞬间打满,直接拖垮整个消费链路。
Redis 与 MySQL 联合双层防重架构拓扑
[ Kafka Consumer 拉取到消息 (包含唯一业务凭据: order_id) ] │ ▼ ┌─────────────────────────────────────────────────────────────┐ │ 【第一道关卡:Redis 极速布隆/分布式令牌前置校验 (Fast Filter)】 │ │ - 纯内存原子指令: SET order_id "PROCESSING" NX EX 120 │ └──────────────────────────┬──────────────────────────────────┘ │ ┌─────────────┴─────────────┐ │ (返回 false: Key 已存在) │ (返回 true: 抢到处理令牌) ▼ ▼ ┌──────────────────────────┐ ┌──────────────────────────────────┐ │ 判定为重复消息,直接安全 │ │ 进入核心业务本地数据库事务 │ │ 提交 Offset 并丢弃,零 DB │ │ ┌──────────────────────────────┐ │ │ 压力 (拦截 99.9% 重复) │ │ │ INSERT 业务订单表 │ │ └──────────────────────────┘ │ │ INSERT 终局防重流水表 (Unique)│ │ │ └──────────────────────────────┘ │ └────────────────┬─────────────────┘ │ ▼ ┌──────────────────────────────────┐ │ 事务成功提交 ──> 异步更新 Redis │ │ 状态为 COMPLETED 并延长 TTL │ └──────────────────────────────────┘该双层架构将性能与一致性完美融为一体:
- 第一层(内存极速过滤):利用 Redis 的
SET ... NX在 1 毫秒内拦截 99.9% 紧随而来的重复重试请求,将百万级的防重并发压力完全隔离在内存层,保护底层数据库不被击穿; - 第二层(数据库终局强一致保障):在业务落库的同一本地事务中,捆绑插入一条包含
business_id唯一索引的防重流水记录。依靠关系型数据库的 ACID 强一致性约束,守住绝对不发生重复资产操作的终极底线。
生产级双层联合防重核心代码实现
以下是我们在 Java 24 核心交易消息消费端落地的联合防重处理器模板:
package com.architect.kafka.idempotency; import org.springframework.data.redis.core.StringRedisTemplate; import org.springframework.jdbc.core.JdbcTemplate; import org.springframework.transaction.support.TransactionTemplate; import java.time.Duration; public class HighConcurrencyIdempotentConsumer { private final StringRedisTemplate redisTemplate; private final TransactionTemplate transactionTemplate; private final JdbcTemplate jdbcTemplate; public HighConcurrencyIdempotentConsumer( StringRedisTemplate redisTemplate, TransactionTemplate transactionTemplate, JdbcTemplate jdbcTemplate ) { this.redisTemplate = redisTemplate; this.transactionTemplate = transactionTemplate; this.jdbcTemplate = jdbcTemplate; } /** * 执行带双层防重保障的业务消费 */ public boolean consumeOrderMessage(String orderId, String messageBody) { String redisKey = "idempotent:order:" + orderId; // 1. 第一道关卡:Redis 内存原子占位 (设置 2 分钟防重处理中租期) Boolean acquired = redisTemplate.opsForValue().setIfAbsent(redisKey, "PROCESSING", Duration.ofMinutes(2)); if (Boolean.FALSE.equals(acquired)) { // 无法获取占位符,说明该订单正在处理中或已处理完毕,安全忽略重复消息 System.out.println("【Redis 防重拦截】检测到重复消息,直接丢弃: " + orderId); return true; } try { // 2. 第二道关卡:本地数据库事务 (业务写入 + 防重流水强一致绑定) boolean txSuccess = Boolean.TRUE.equals(transactionTemplate.execute(status -> { try { // A. 插入防重流水表 (依赖 order_id 唯一主键约束) String insertLockSql = "INSERT INTO t_idempotent_flow (order_id, create_time) VALUES (?, NOW())"; jdbcTemplate.update(insertLockSql, orderId); // B. 执行真实的业务落盘 (如更新订单状态、扣减库存) executeRealBusiness(orderId, messageBody); return true; } catch (org.springframework.dao.DuplicateKeyException dke) { // 命中了数据库唯一索引冲突,说明是极端情况下穿透了 Redis 的重复消息 status.setRollbackOnly(); System.err.println("【DB 防重终局兜底拦截】命中唯一键冲突: " + orderId); return true; // 视为幂等成功 } catch (Exception ex) { // 真实业务异常,事务回滚 status.setRollbackOnly(); throw new RuntimeException("业务处理失败,触发事务回滚", ex); } })); if (txSuccess) { // 3. 事务提交成功,将 Redis 状态更新为 COMPLETED 并赋予 24 小时兜底保护期 redisTemplate.opsForValue().set(redisKey, "COMPLETED", Duration.ofHours(24)); return true; } return false; } catch (Exception e) { // 发生非幂等性业务错误,必须立即删除 Redis 占位符,允许后续合法重试 redisTemplate.delete(redisKey); throw e; } } private void executeRealBusiness(String orderId, String payload) { // 执行真实履约落盘 } }消费端防重设计的最后三道安全红线
- 唯一防重凭证必须具备全局业务唯一性(Business Unique Key):严禁使用 Kafka 的
Topic + Partition + Offset作为防重凭据!因为一旦上游发生重发或重新推单,Offset 是全新的,系统无法识别其业务重复。防重键必须严格来源于业务实体本身的确定性凭据(如支付流水号pay_trade_no、订单号order_id)。 - 防重流水表必须采用按月分表与分区定期归档:随着大促进行,
t_idempotent_flow表的记录会迅速累积上亿条。该表必须按照月份或者日期进行物理分表(例如t_idempotent_flow_202610),并由自动化定时任务定期将 30 天前的历史防重流水迁移至冷归档库或直接物理截断,防止单表膨胀拖死数据库。 - 严格防范事务与 Redis 状态之间的乱序竞态:如果网络极其卡顿导致数据库事务耗时超过 2 分钟,Redis 占位符自动过期,新的重复消息可能会再次进入。因此,Redis 租期必须根据业务 P99 耗时留出 5 到 10 倍的安全冗余(通常设为 2 到 5 分钟),且绝对不能省略数据库层面的唯一键强校验。