尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

一张无门槛券被发两次:Kafka 重复消费下幂等设计的三种防线

一张无门槛券被发两次:Kafka 重复消费下幂等设计的三种防线 title: 一张无门槛券被发两次Kafka 重复消费下幂等设计的三种防线date: 2026-08-31tags: [幂等设计, Kafka, 分布式事务, 后端, 高可用]大促当天客服群炸了去年双十一零点刚过 9 分钟客服群里开始有人 我「用户投诉收到两张 5 元无门槛券但他只下了一次单。」一开始我以为是前端重复点了查了下下单接口有防重复提交。再往后投诉从 3 条涨到 40 多条我们才意识到事情没那么简单。拉出发券服务的日志同一笔订单号20261111-8842在 12 秒内出现了两条发券记录而且两条都成功了。更离谱的是这两条记录对应的 Kafka 消息 offset 不一样——一个在 partition-3 的 offset 10233一个在 partition-3 的 offset 10234是两条「不同」的消息但业务含义完全相同。这就是典型的「消息重复投递」。我们的发券链路是下单成功 → 发一条OrderPaidEvent到 Kafka → 发券消费者订阅该 topic 发券。Kafka 的投递语义是at-least-once至少一次重复是设计内的常态不是异常。重复从哪来重平衡 手动提交 offset 的时序坑很多人以为「我用了手动提交 offset 就不会重复」这是误解。看一段我们当时消费者的伪代码骨架KafkaListener(topics order-paid, groupId coupon-sender) public void onMessage(ConsumerRecordString, String record, Acknowledgment ack) { OrderPaidEvent event JSON.parseObject(record.value(), OrderPaidEvent.class); couponService.send(event.getOrderId(), event.getUserId()); // ① 先发券 ack.acknowledge(); // ② 再提交 offset }逐行解释- 第 2 行把 Kafka 的 JSON 字符串反序列化成事件对象record.value()是消息体OrderPaidEvent是订单支付事件携带orderId和userId。- 第 3 行调用发券逻辑按orderId给用户发券。注意这里先执行业务、再提交位移。- 第 4 行ack.acknowledge()是 Spring-Kafka 的手动提交把当前 offset 提交给 broker。问题就在这两行的顺序。大促时消息堆积Kafka 触发了 consumer group 的重平衡rebalance某个消费者实例被踢掉、位移还没提交分区被分配给另一个实例新实例从「上一次已提交的 offset」重新拉取——于是这条消息被重新消费一次。我们当时max.poll.interval.ms设的是默认 5 分钟而发券逻辑里同步调用了三个外部服务用户中心、营销中心、通知中心一次耗时 800ms~2s高峰期偶发超过 5 分钟重平衡就频繁触发。换句话说只要「业务执行」和「offset 提交」不是原子的at-least-once 就必然可能重复。想靠提交方式消除重复是做不到的幂等必须在业务层兜底。防线一去重表最稳的兜底我们最终的主防线是一张去重表用order_id做唯一键靠数据库的唯一约束天然挡掉重复。Transactional public void sendSafe(OrderPaidEvent event) { int inserted dedupMapper.insertIfAbsent(event.getOrderId(), COUPON); if (inserted 0) { // 唯一键冲突说明已经处理过 log.warn(duplicate order {}, event.getOrderId()); return; // 直接返回不再发券 } try { couponService.send(event.getUserId(), event.getOrderId()); } catch (Exception e) { dedupMapper.delete(event.getOrderId()); // 业务失败要回滚去重标记否则会「假成功」 throw e; } }逐行解释- 第 3 行insertIfAbsent往去重表插一行order_id上有唯一索引。第一次插入返回 1重复插入触发DuplicateKeyException或被捕获后返回 0。- 第 4 行如果inserted 0代表这条订单已经发过券直接 return整个发券逻辑短路。- 第 7 行真正发券第 10 行是关键细节——业务失败时必须删掉刚插入的去重记录否则这条订单会永远被「假去重」挡住再也发不出券。我们第一版就漏了这行导致重平衡失败后重试的订单全部静默丢失。insertIfAbsent对应的 SQL 是INSERT IGNORE INTO dedup(order_id, biz_type) VALUES(?, ?)靠UNIQUE KEY uk_order (order_id, biz_type)保证。注意我加了biz_type维度同一笔订单可能在「发券」和「发货」两个业务里都要幂等共用一个唯一键会互相误伤。防线二状态机避免「已发券」和「已退款」打架去重表解决的是「同一条消息别处理两次」但解决不了「处理到一半、状态乱套」的问题。比如发券流程里有「锁定库存」「扣减预算」「写用户券包」三步如果第一步成功、第三步超时重试时第一步又跑了一遍库存就被多锁一次。这时候需要状态机把订单的发放过程建模成有限状态public enum CouponState { INIT, LOCKED, BUDGET_DEDUCTED, GRANTED, FAILED; private static final SetCouponState TERMINAL Set.of(GRANTED, FAILED); public boolean canTransitTo(CouponState next) { return switch (this) { case INIT - next LOCKED || next FAILED; case LOCKED - next BUDGET_DEDUCTED || next FAILED; case BUDGET_DEDUCTED - next GRANTED || next FAILED; case GRANTED, FAILED - false; // 终态不可再变更 }; } }逐行解释- 第 2 行定义发放过程的状态枚举初始化 → 库存锁定 → 预算扣减 → 已发放 → 失败。- 第 4 行把GRANTED和FAILED标记为终态进到这两个状态后不再接受任何转移。- 第 6 行canTransitTo是状态机的核心校验只有「当前态 → 合法下一态」才放行。比如已经GRANTED的订单再来一次发放请求GRANTED.canTransitTo(LOCKED)返回false直接拒绝。每次推进状态都要用UPDATE coupon_state SET state? WHERE order_id? AND state?加「乐观锁」条件保证并发下只有一条能成功推进。状态机和去重表是互补的去重表挡「重复触发」状态机管「内部步骤乱序」。防线三前端 Token别把它当主防线很多教程把「前端下单拿 Token、提交时校验 Token」当作幂等主方案。我们早期也这么做过但它有两个硬伤public boolean consumeToken(String orderToken) { // SET key value NX EX 300只有第一次设置成功才返回 true return redisTemplate.opsForValue() .setIfAbsent(idemp:token: orderToken, 1, Duration.ofMinutes(5)); }这行 RedisSET NX能挡住「用户手抖连点两次」的场景但它防不住服务端的消息重放、定时任务重试、跨机器并发。Token 是客户端维度的补充不是服务端幂等的替代。把它当唯一防线一旦 Token 没传比如网关重试、内部 RPC 调用防线就空了。三种防线怎么选防线防的重复来源优点缺点适合场景去重表消息重投、任务重试简单、可靠、数据库保证多一次 DB 写几乎所有写操作状态机步骤乱序、并发推进状态清晰可审计建模成本高多步骤长事务Token用户重复点击零侵入业务防不住服务端表单/接口入口复盘那次事故的数字受影响用户 312 人重复发券 327 张部分用户有两张以上券面额 5 元直接资损 1635 元但更严重的是「用户以为平台乱扣乱发」的信任成本。从发现到止血我们花了 47 分钟其中 30 分钟在争论「到底是前端还是 Kafka 的问题」——如果当时第一道防线就在这 30 分钟根本不会存在。我的取舍我不建议把「前端 Token」当作幂等的主防线它顶多算入口处的礼貌性拦截。也更不建议用「先查后插」代替去重表的唯一键——SELECT再INSERT在高并发下根本挡不住重复我们压测时 200 并发就能复现重复插入。我的判断是去重表 数据库唯一键是 90% 写接口的最低标配状态机留给真正有多步骤、需要审计的长流程。如果你的接口一天写不了几万次去重表就够了别为了「架构感」上全套。另外去重表记得加expire_at字段定期清理否则它会长成第二张主表。我们设的是 30 天自动过期和退款窗口对齐。我们最终落地的生产配置三道防线要真正生效参数得落到具体版本和具体值上不能只停留在「思想上」。我们这套发券链路的运行时是 Spring-Kafka 2.8.0 kafka-clients 3.4.0 MySQL 8.0.32关键配置如下// 消费者侧先把重平衡概率压下去再去谈业务幂等 KafkaListener(topics order-paid, groupId coupon-sender, properties { max.poll.interval.ms300000, // 拉取间隔放宽到 5 分钟避免消费慢被踢 max.poll.records20, // 单次少拉点单批处理时间可控 enable.auto.commitfalse // 关掉自动提交位移完全由业务 ack 控制 }) public void onMessage(ConsumerRecordString, String record, Acknowledgment ack) { sendSafe(JSON.parseObject(record.value(), OrderPaidEvent.class)); ack.acknowledge(); }逐行解释- 第 4 行max.poll.interval.ms300000是这次事故后我们改的最关键一项——之前用默认 5 分钟还偶发超时其实根因是单批拉了 500 条、处理太久改成max.poll.records20第 5 行后单批最多 20 条、处理时间稳定在 40 秒内重平衡基本不再触发。- 第 6 行enable.auto.commitfalse关掉自动提交位移只在我们业务跑完、去重表插入成功后才由ack.acknowledge()提交从源头消除了「业务没跑完位移先提交了」的漏洞。去重表的清理我们用了定时任务每天凌晨扫expire_at now()的行30 天窗口外的标记一次性删掉避免它长成第二张主表拖慢写入。思考题你现在负责的写接口里哪些只靠「前端防重复点击」在扛如果 Kafka 重平衡或 RPC 超时重试发生它们扛得住吗试着给对方一个订单号重复发两次请求看数据库里会不会多一条脏数据。本文为 Round 5 重写稿与 R3 幂等设计篇使用不同的事故场景与代码未复用旧文。
返回列表