ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

多实例同时扫描任务如何避免重复发送?条件抢占、租约与僵尸任务回收

多实例同时扫描任务如何避免重复发送?条件抢占、租约与僵尸任务回收 一、业务背景一条积分消息为什么必须靠扫表发出去电商下单成功之后有一串不那么核心但必须做的收尾动作给用户加积分、发短信、发 App 推送。这些动作如果跟创建订单写在同一个同步流程里会有两个后果接口响应被拖长。短信网关抖动一下用户就得在提交订单按钮上多等两秒故障域被放大。短信服务挂了订单服务跟着报错——用户明明下单成功了却因为他收不到短信而看到下单失败。所以引入 RocketMQ 做异步解耦订单落库成功之后把积分、短信、推送这些事丢给消息队列由消费者慢慢处理。但异步立刻带来一个更棘手的问题订单落库成功和消息发出去是两次独立的写操作中间没有任何事务能同时覆盖它们。先发 MQ 再落库 → MQ 发出去了结果本地事务回滚了 → 消费者给一个根本不存在的订单加积分先落库再发 MQ → 订单成功了但发 MQ 时网络抖动/进程被 kill → 消息静默丢失积分永远不到账。这就是经典的本地消息表Transactional Outbox要解决的问题把业务数据和待发消息塞进同一个本地事务落库。事务内只往t_event_record插一条statusWAIT的记录不发 MQ再由一个独立的发布器异步扫表补发。这样一来订单成功但事件丢失这个窗口就被彻底堵死了——订单和事件要么一起成功要么一起回滚。代价也很明确发布器必然是一个定时扫描任务。而order-service为了保证吞吐线上一定是多实例部署的。于是问题来了order-service 实例 A ──┐ order-service 实例 B ──┼── 每5秒都在扫同一张 t_event_record 表 order-service 实例 C ──┘同一条 WAIT 事件会被 A、B、C 同时捞到。如果不做任何处理用户下一单积分会被加三次。下面要讲的就是怎么让多实例并发扫描变成只有一个人能发。二、整体数据流查询只是候选归属权由 UPDATE 影响行数裁决先看全景。核心思想只有一句话SELECT出来的只是候选不是归属谁能改成功谁才算抢到。RocketMQEventRetrySchedulert_event_recordOrderServiceRocketMQEventRetrySchedulert_event_recordOrderService★ 事务内不发 MQ提交即代表事件一定不会丢抢到租约这条事件归我发超过 maxRetry(3) → 标记 FAIL死信 告警alt[发送成功][发送失败]已被别的实例抢走安静跳过不是异常alt[影响行数 1][影响行数 0]loop[每个实例每 5s 跑一轮]④ reclaimStuck()回收发送中僵尸UPDATE SET statusWAIT, claim_ownerNULLWHERE statusSENDING AND claimed_at NOW() - 10min★ 判定基准是【领取时间 claimed_at】不是 next_retry_time同事务 INSERT (statusWAIT, next_retry_timeNOW())① 捞候选无锁SELECT WHERE statusWAIT AND next_retry_timeNOW() LIMIT 20② 条件抢占 UPDATE SET statusSENDING, claim_owner?, claimed_atNOW() WHERE id? AND statusWAIT③ 发送 MQUPDATE statusSUCCESS带 claim_owner 条件retry_count1、next_retry_time 退避后移退回 WAIT第 4 步是这套设计里最容易被漏掉、也最容易在线上炸的一环抢占成功之后、MQ 真正发出去之前进程被 kill 了怎么办那条记录会永远卡在SENDING——扫描器只捞WAIT它再也回不来了。从业务视角看这条消息静默丢失了而且没有任何报错。所以claim_owner claimed_at 超时回收这套机制本质就是一个用数据库实现的租约lease抢占 领取租约记下谁领的、什么时候领的租约到期10 分钟还没变成 SUCCESS 领租约的人八成已经死了回收 把租约强行收回状态退回WAIT重新参与补发。三、代码实现3.1 调度器主体/** * Outbox 事件独立发布器 * * p设计要点 * 1. 多实例安全 —— 靠 UPDATE 条件抢占 影响行数闸门不用分布式锁 * 2. 崩溃自愈 —— SENDING 超时视为僵尸由 reclaimStuck 回收重发 * 3. 多实例可水平扩展 —— 抢占粒度是行实例越多吞吐越高不会互相打架 */Slf4jComponentRequiredArgsConstructorpublicclassEventRetryScheduler{/** 单轮最多捞多少条避免一次拉太多导致长时间持有连接 / 内存膨胀 */privatestaticfinalintBATCH_SIZE20;/** 最大重试次数超过进死信停止自动重试等人工介入 */privatestaticfinalintMAX_RETRY3;/** 租约时长SENDING 超过这个时间未完成视为发送中僵尸 */privatestaticfinalintLEASE_MINUTES10;privatestaticfinalStringEVENT_TOPICORDER_EVENT_TOPIC;privatefinalEventRecordMappereventRecordMapper;privatefinalRocketMQTemplaterocketMQTemplate;/** 实例唯一标识形如 order-8083-a1b2c3落库到 claim_owner 便于排查是谁发飞的 */privatefinalInstanceIdProviderinstanceIdProvider;/** * 主循环每 5 秒扫一次待发事件。 * * p用 fixedDelay 而不是 fixedRate —— fixedRate 是每 5 秒启动一轮 * 上一轮没跑完就会并发出第二轮任务堆积fixedDelay 是上一轮跑完再等 5 秒。 */Scheduled(fixedDelay5_000)publicvoidpublishPendingEvents(){StringownerinstanceIdProvider.get();ListEventRecordcandidateseventRecordMapper.selectWaitForPublish(BATCH_SIZE);if(CollectionUtils.isEmpty(candidates)){return;}for(EventRecordevent:candidates){try{claimAndPublish(event,owner);}catch(Exceptione){// ★ 单条异常绝不能中断整批这一轮其它事件的补发必须继续log.error([Outbox] 事件处理异常 eventId{},event.getId(),e);}}}privatevoidclaimAndPublish(EventRecordevent,Stringowner){// ① 条件抢占影响行数是归属权的唯一裁决者intclaimedeventRecordMapper.claim(event.getId(),owner);if(claimed0){// 已被同集群其它实例抢走 —— 这是设计内的正常现象用 debug 级别别打 errorlog.debug([Outbox] 事件已被其它实例抢占跳过 eventId{},event.getId());return;}// ② 抢到租约才有资格发 MQtry{rocketMQTemplate.syncSend(EVENT_TOPIC,buildMessage(event));// ③ 成功标记 SUCCESSSQL 里带 claim_owner 条件见 3.2introwseventRecordMapper.markSuccess(event.getId(),owner);log.info([Outbox] 事件发送成功 eventId{} bizNo{} rows{},event.getId(),event.getBizNo(),rows);}catch(Exceptione){// ④ 失败退避重试 / 超限进死信handleFailure(event,owner,e);}}privatevoidhandleFailure(EventRecordevent,Stringowner,Exceptioncause){intnextRetryevent.getRetryCount()1;if(nextRetryMAX_RETRY){eventRecordMapper.markFail(event.getId(),owner,nextRetry);// 死信必须打 error 并接告警这是需要人介入的信号不是重试一下就好log.error([Outbox] 事件重试超限进入死信 eventId{} retry{},event.getId(),nextRetry,cause);return;}// 退避递增retryCount² × 5s → 5s / 20s / 45s// 下游抖动时如果还按固定 5s 死磕等于持续给已经喘不过气的服务加压intbackoffSecondsnextRetry*nextRetry*5;eventRecordMapper.releaseForRetry(event.getId(),owner,nextRetry,backoffSeconds);log.warn([Outbox] 事件发送失败{}s 后重试 eventId{} retry{},backoffSeconds,event.getId(),nextRetry,cause);}/** * 回收发送中僵尸。 * * p场景实例抢占成功statusSENDING后、MQ 还没发出去时进程崩溃 / 被 kill。 * 没有这一步这条记录会永远卡在 SENDING —— 扫描器只捞 WAIT它再也回不来 * 消息静默丢失且没有任何报错。这是 Outbox 方案里最隐蔽的坑。 */Scheduled(fixedDelay60_000)publicvoidreclaimStuck(){introwseventRecordMapper.reclaimStuck(LEASE_MINUTES);if(rows0){// 有回收就说明发生过实例异常退出值得留一条 warn 供排查log.warn([Outbox] 回收发送中僵尸事件 {} 条,rows);}}privateMessageStringbuildMessage(EventRecordevent){returnMessageBuilder.withPayload(event.getPayload()).setHeader(eventId,event.getId()).setHeader(eventType,event.getEventType()).setHeader(bizNo,event.getBizNo()).build();}}3.2 Mapper每一句 SQL 都在守一条边界MapperpublicinterfaceEventRecordMapper{/** * 捞候选。注意这里【不加锁、不占用归属】。 * 多个实例捞到同一条是必然的也是允许的 —— 归属交给下面的 claim 决定。 */Select(SELECT * FROM t_event_record WHERE status WAIT AND next_retry_time NOW() ORDER BY next_retry_time ASC LIMIT #{limit})ListEventRecordselectWaitForPublish(Param(limit)intlimit);/** * 条件抢占 —— 整个方案的核心。 * * pWHERE statusWAIT 是闸门并发的 UPDATE 会被 InnoDB 行锁串行化 * 第一个事务把 status 改成 SENDING 并提交后后面的 UPDATE 匹配不到行 * 影响行数只能是 0。所以谁能改成功天然唯一不需要任何分布式锁。 */Update(UPDATE t_event_record SET status SENDING, claim_owner #{owner}, claimed_at NOW() WHERE id #{id} AND status WAIT)intclaim(Param(id)Longid,Param(owner)Stringowner);/** * 标记成功。 * * p额外带 claim_owner 条件防止原实例的超时回执误伤 —— * 若租约已被 reclaimStuck 回收、事件又被别的实例抢走发送 * 原实例迟到的成功回执影响行数会是 0不会把别人的记录改坏。 */Update(UPDATE t_event_record SET status SUCCESS, update_time NOW() WHERE id #{id} AND status SENDING AND claim_owner #{owner})intmarkSuccess(Param(id)Longid,Param(owner)Stringowner);/** * 发送失败释放租约退回 WAIT并推后下次重试时间。 * 同样必须带 claim_owner —— 只能释放【自己的】租约 * 否则 A 实例会把 B 实例正在发送的事件抢回队列人为制造重复发送。 */Update(UPDATE t_event_record SET status WAIT, claim_owner NULL, claimed_at NULL, retry_count #{retryCount}, next_retry_time DATE_ADD(NOW(), INTERVAL #{backoffSeconds} SECOND) WHERE id #{id} AND status SENDING AND claim_owner #{owner})intreleaseForRetry(Param(id)Longid,Param(owner)Stringowner,Param(retryCount)intretryCount,Param(backoffSeconds)intbackoffSeconds);/** 超过重试上限 → 死信停止自动补发 */Update(UPDATE t_event_record SET status FAIL, retry_count #{retryCount}, update_time NOW() WHERE id #{id} AND status SENDING AND claim_owner #{owner})intmarkFail(Param(id)Longid,Param(owner)Stringowner,Param(retryCount)intretryCount);/** * 回收僵尸。★ 判定基准是 claimed_at领取时间不是 next_retry_time。 * * p为什么不能用 next_retry_time一条事件可能是排在队列里等重试WAIT * 也可能是已经被领走在途SENDING。如果用 next_retry_time 判定 * 那些刚被排队领走、还在正常发送中的事件会被误判回收直接造成重复补发。 * claimed_at 记录的是谁、什么时候把租约领走的只有它才能准确表达在途多久了。 */Update(UPDATE t_event_record SET status WAIT, claim_owner NULL, claimed_at NULL, update_time NOW() WHERE status SENDING AND claimed_at DATE_SUB(NOW(), INTERVAL #{minutes} MINUTE))intreclaimStuck(Param(minutes)intminutes);}3.3 索引两条扫描路径各配一个复合索引调度器每 5 秒扫一次、回收器每 60 秒扫一次这两个查询是全项目执行次数最高的 SQL必须有索引兜住-- 补发扫描WHERE statusWAIT AND next_retry_time NOW()-- 等值列 status 在前范围列 next_retry_time 在后 —— 复合索引列顺序的基本功CREATEINDEXidx_status_next_retryONt_event_record(status,next_retry_time);-- 僵尸回收WHERE statusSENDING AND claimed_at ?CREATEINDEXidx_status_claimed_atONt_event_record(status,claimed_at);四、异常与并发边界这几个坑不填线上必翻车坑 1SELECT到同一条是必然的不能把它当冲突处理三个实例同时SELECT捞到同一条 WAIT 记录是设计内的正常现象。如果在这里用异常或者告警去处理重复捞到日志会被刷爆。正确姿势SELECT 只负责缩小范围归属权完全交给后面那句条件 UPDATE 的影响行数。抢不到影响行数 0就走debug日志安静跳过。坑 2为什么不用SELECT ... FOR UPDATE也不用 Redisson 分布式锁SELECT ... FOR UPDATE锁要一直持有到事务提交而发 MQ这个网络 IO 绝不能放在事务里——MQ 卡 3 秒行锁就持有 3 秒扫描线程全被拖死Redisson 分布式锁锁的是整个扫描任务粒度太粗。它只解决同一时刻只有一个实例在扫解决不了扫到一半进程崩了怎么办而且锁释放和事务提交之间有时间差锁已释放、事务还没提交并发读可能读到旧状态。更关键的是加了全局锁之后多实例扩容就完全失去意义——扫表变成了单点。条件抢占把并发粒度降到行级实例越多能并行处理的事件越多是真正能水平扩展的。坑 3标记成功 / 释放租约都必须带claim_owner不带claim_owner条件的 UPDATE 是跨实例误伤的源头markSuccess不带 owner → 租约已被回收、事件已被别的实例重发原实例迟到的成功回执会把记录改成 SUCCESS掩盖了真实发送状态releaseForRetry不带 owner → A 实例把 B 实例正在发送的事件强行抢回队列凭空制造一次重复发送。坑 4租约回收的本质是承认至少一次不是恰好一次即使所有边界都填好了仍然存在这个窗口实例 A 抢占成功 → 发 MQ网络阻塞11 分钟才返回 ↓10分钟时租约到期 → reclaimStuck 回收 → 实例 B 抢到 → 又发了一次 ↓ 实例 A 终于返回成功 → markSuccess 影响行数0owner 已变结果MQ 里出现了两条消息。这不是 bug而是分布式系统里不引入分布式事务就不可能做到恰好一次的必然结果。正确的应对是两头一起做发送端MQ 发送必须配置远小于租约时长10 分钟的超时把阻塞 11 分钟这种极端情况压到几乎不可能消费端下游必须幂等。这个项目的消费者用msg_id前置检查 uk_msg_id唯一索引做三层防线业务表还有uk_order_type唯一索引兜底——重复投递不会重复加积分。配套的消费端设计见同系列文章《MQ 重复投递不可怕前置检查、数据库事务与唯一索引三层幂等》。坑 5死信是给人看的必须告警超过MAX_RETRY3 次标记成FAIL之后自动补发就彻底停了。此时如果没有告警这批消息就真的静默烂在表里了。所以死信分支必须打error级别日志 接监控告警 留人工补偿入口。坑 6fixedDelay而不是fixedRatefixedRate是每 5 秒启动一轮上一轮没跑完会并发出下一轮任务不断堆积fixedDelay是上一轮跑完再等 5 秒天然自带背压。五、简历 / 面试亮点这段代码可以直接写进简历因为它一次性覆盖了三个面试高频考点最终一致性方案、无锁并发控制、崩溃自愈。可量化的表述基于 Transactional Outbox 实现订单事件最终一致性事务内只落库不发消息消除业务成功但消息丢失窗口 发布器采用UPDATE 条件抢占 影响行数闸门替代分布式锁支持 order-service 多实例水平扩展抢占冲突以 debug 日志静默跳过不产生误告警 设计claim_owner claimed_at 租约机制10 分钟超时自动回收因进程崩溃卡在 SENDING 的发送中僵尸事件杜绝消息静默丢失 失败按retryCount² × 5s退避递增重试3 次超限落FAIL死信并告警 两条高频扫描 SQL 分别用idx_status_next_retry、idx_status_claimed_at复合索引覆盖EXPLAIN 从typeALL全表扫描rows≈9873优化为rangerows≈3171扫描行数下降 67% 测试覆盖 70 个用例其中包含调度器超时退回 WAIT 重试、且重试不双还的场景验证。规范 Commit Message 建议feat(outbox): 发布器改用条件抢占替代分布式锁支持多实例水平扩展 fix(outbox): 补 claimed_at 超时回收修复进程崩溃导致事件永久卡 SENDING fix(outbox): markSuccess/releaseForRetry 补 claim_owner 条件防止跨实例误伤 docs(readme): 补充租约机制与至少一次语义说明六、面试场景题暴击场景题 1租约设为 10 分钟但 MQ 发送阻塞了 11 分钟才返回成功会发生什么怎么设计才能既不丢、又尽量不重先答发生什么租约在第 10 分钟被reclaimStuck回收状态退回WAIT实例 B 抢到并成功发送第 11 分钟实例 A 返回成功执行markSuccess因为带了claim_owner条件此时 owner 已不是 A影响行数 0——这条迟到的回执被安全丢弃。最终结果是 MQ 里有一条消息B 发的A 的那次发送在 MQ 侧其实也成功了的话就是两条 →重复投递由消费端幂等兜住。再答怎么设计按优先级MQ 发送超时 租约时长。这是第一道也是最有效的一道发送超时设 3 秒租约 10 分钟阻塞 11 分钟这种事根本不可能发生。租约时长应该约等于单次发送的最坏耗时 × 巨大安全系数而不是拍脑袋定。租约续期renew。长任务场景可以让执行方定期UPDATE claimed_at NOW() WHERE id? AND claim_owner?续租只有续租失败才认为租约真丢了。适合处理逻辑本身耗时不确定的场景本项目发送是秒级 IO用不上。接受至少一次把幂等下沉到消费端。这是分布式系统的正解——只要不引入分布式事务“恰好一次在工程上做不到能做的是至少一次 消费端幂等”。所以消费端必须有msg_id前置检查 唯一索引。加可观测性。回收次数、死信条数、SENDING状态的 P99 驻留时长上监控看板。回收次数突然飙升 有实例在频繁异常退出这是比消息重复更严重的信号。场景题 2为什么不直接给扫描方法加一把 Redisson 分布式锁只让一个实例扫表面看很合理加锁之后同一时刻只有一个实例在扫天然不会重复发送代码还更简单。但有三个致命问题锁解决不了崩溃。实例 A 抢到锁、抢占成功、发 MQ 前被 kill锁因为过期自动释放记录却永远卡在SENDING。锁只能保证同一时刻一个人扫保证不了扫到一半死了怎么办——那得靠租约回收。扩容失去意义。扫表变成单点任务order-service扩到 10 个实例补发吞吐还是 1 个实例的量。而条件抢占的粒度是行实例越多一轮能并行处理的事件越多。锁的粒度与生命周期都不匹配。分布式锁是任务级、秒级持有的业务需要的是行级、可能长达分钟级的归属权——用锁就得为每条事件单独加锁并设置超长 TTL那还不如直接用数据库的状态字段当租约少一个中间件依赖且归属状态和业务数据在同一个事务里天然一致。一句话收口「能用数据库一次 UPDATE 解决的事不要引入分布式锁。锁是有状态的中间件依赖会带来锁超时、锁续期、锁与事务提交不同步这些新问题而UPDATE ... WHERE statusWAIT影响行数闸门是无状态的——实例随时可以重启、扩容、被杀正确的那个实例总能抢到。这就是无锁化的思路把并发控制下沉到数据库的行锁让架构更简单、更能水平扩展。」推荐标签RocketMQ、分布式事务、Spring Boot、微服务、Transactional Outbox
返回列表