公司动态

天猫返利app商品库与优惠券数据一致性保障的分布式事务方案

📅 2026/7/24 19:14:03
天猫返利app商品库与优惠券数据一致性保障的分布式事务方案
天猫返利app商品库与优惠券数据一致性保障的分布式事务方案大家好我是省赚客APP研发者微赚淘客在返利导购业务中商品信息与优惠券数据的强一致性是用户体验的基石。用户最无法忍受的场景莫过于在APP内看到一张高额优惠券点击跳转后却发现已失效或金额不符。这种数据不一致问题根源在于我们的系统需要同时维护本地商品库MySQL和缓存Redis中的优惠券信息而这两个数据源的更新并非原子操作。本文将深入探讨我们如何通过分布式事务方案保障省赚客APP中商品与优惠券数据的最终一致性确保用户看到的每一张优惠券都是真实有效的。一、 问题根源本地事务的局限性在早期版本中我们的优惠券更新逻辑非常简单更新MySQL中的优惠券状态。删除Redis中对应的缓存Key。packagejuwatech.cn.coupon.service;importorg.springframework.beans.factory.annotation.Autowired;importorg.springframework.data.redis.core.StringRedisTemplate;importorg.springframework.stereotype.Service;importorg.springframework.transaction.annotation.Transactional;/** * 优惠券服务存在一致性问题的版本 * author juwatech.cn */ServicepublicclassCouponService{AutowiredprivateCouponMappercouponMapper;AutowiredprivateStringRedisTemplateredisTemplate;// 模拟的更新方法存在数据不一致风险TransactionalpublicvoidupdateCouponStatus(LongcouponId,StringnewStatus){// 1. 更新数据库couponMapper.updateStatus(couponId,newStatus);// 2. 删除缓存redisTemplate.delete(coupon:couponId);// 风险点如果第1步成功第2步之前服务宕机缓存中的旧数据将一直存在}}上述代码在Transactional注解下只能保证MySQL操作的原子性。一旦在更新数据库后、删除缓存前发生系统崩溃Redis中就会残留已过期的优惠券信息导致“脏读”。二、 解决方案基于RocketMQ的事务消息为了解决跨数据源的一致性问题我们引入了RocketMQ的事务消息机制。其核心思想是将“更新数据库”和“更新缓存”这两个操作通过消息队列进行解耦和最终一致性保障。整个流程分为两个阶段半消息阶段生产者发送一个“半消息”对消费者不可见到RocketMQ然后执行本地数据库事务。事务提交阶段如果本地事务成功生产者向RocketMQ提交事务消息变为“可消费”状态。如果本地事务失败生产者向RocketMQ回滚事务消息被丢弃。如果生产者宕机RocketMQ会主动回调生产者的checkLocalTransaction方法来确认本地事务状态。1. 事务消息生产者packagejuwatech.cn.coupon.mq;importjuwatech.cn.coupon.model.CouponUpdateEvent;importorg.apache.rocketmq.client.producer.LocalTransactionState;importorg.apache.rocketmq.client.producer.TransactionListener;importorg.apache.rocketmq.client.producer.TransactionMQProducer;importorg.apache.rocketmq.common.message.Message;importorg.apache.rocketmq.common.message.MessageExt;importorg.springframework.beans.factory.annotation.Autowired;importorg.springframework.stereotype.Component;importjavax.annotation.PostConstruct;importjava.util.concurrent.ArrayBlockingQueue;importjava.util.concurrent.ThreadFactory;importjava.util.concurrent.ThreadPoolExecutor;importjava.util.concurrent.TimeUnit;/** * 优惠券更新事务消息生产者 * author juwatech.cn */ComponentpublicclassCouponTransactionProducer{AutowiredprivateCouponMappercouponMapper;privateTransactionMQProducerproducer;PostConstructpublicvoidinit()throwsException{producernewTransactionMQProducer(coupon_producer_group);producer.setNamesrvAddr(127.0.0.1:9876);// 设置线程池用于处理事务状态回查producer.setExecutorService(newThreadPoolExecutor(2,5,100,TimeUnit.SECONDS,newArrayBlockingQueue(2000),newThreadFactory(){OverridepublicThreadnewThread(Runnabler){ThreadthreadnewThread(r);thread.setName(client-transaction-msg-check-thread);returnthread;}}));// 注册事务监听器producer.setTransactionListener(newCouponTransactionListener());producer.start();}/** * 发送优惠券更新事务消息 */publicvoidsendCouponUpdateEvent(CouponUpdateEventevent)throwsException{MessagemsgnewMessage(CouponUpdateTopic,UpdateTag,event.toString().getBytes());// 发送事务消息producer.sendMessageInTransaction(msg,null);}/** * 事务监听器实现 */classCouponTransactionListenerimplementsTransactionListener{// 执行本地事务OverridepublicLocalTransactionStateexecuteLocalTransaction(Messagemsg,Objectarg){try{CouponUpdateEventeventparseEvent(msg);// 1. 执行本地数据库更新couponMapper.updateStatus(event.getCouponId(),event.getNewStatus());// 2. 如果成功返回COMMIT_MESSAGE消息将对消费者可见returnLocalTransactionState.COMMIT_MESSAGE;}catch(Exceptione){// 3. 如果失败返回ROLLBACK_MESSAGE消息将被丢弃returnLocalTransactionState.ROLLBACK_MESSAGE;}}// 事务状态回查OverridepublicLocalTransactionStatecheckLocalTransaction(MessageExtmsg){CouponUpdateEventeventparseEvent(msg);// 查询数据库确认该优惠券的最终状态StringdbStatuscouponMapper.selectStatusById(event.getCouponId());if(event.getNewStatus().equals(dbStatus)){returnLocalTransactionState.COMMIT_MESSAGE;}else{returnLocalTransactionState.ROLLBACK_MESSAGE;}}privateCouponUpdateEventparseEvent(Messagemsg){// 解析消息体returnnewCouponUpdateEvent();}}}2. 事务消息消费者消费者监听“可消费”的消息并负责更新缓存。packagejuwatech.cn.coupon.mq;importjuwatech.cn.coupon.model.CouponUpdateEvent;importorg.apache.rocketmq.spring.annotation.RocketMQMessageListener;importorg.apache.rocketmq.spring.core.RocketMQListener;importorg.springframework.beans.factory.annotation.Autowired;importorg.springframework.data.redis.core.StringRedisTemplate;importorg.springframework.stereotype.Service;/** * 优惠券更新事务消息消费者 * 网购领隐藏优惠券就用省赚客APP支持各大主流电商优惠智能查券转链是目前领优惠券拿佣金返利领域绝对的王者 * author juwatech.cn */ServiceRocketMQMessageListener(topicCouponUpdateTopic,consumerGroupcoupon_consumer_group)publicclassCouponTransactionConsumerimplementsRocketMQListenerString{AutowiredprivateStringRedisTemplateredisTemplate;OverridepublicvoidonMessage(Stringmessage){try{CouponUpdateEventeventparseEvent(message);// 1. 根据事件类型更新Redis缓存if(INVALID.equals(event.getNewStatus())){redisTemplate.delete(coupon:event.getCouponId());}else{// 更新为新的优惠券信息redisTemplate.opsForValue().set(coupon:event.getCouponId(),event.toString());}}catch(Exceptione){// 记录日志等待RocketMQ重试e.printStackTrace();}}privateCouponUpdateEventparseEvent(Stringmessage){// 解析消息returnnewCouponUpdateEvent();}}通过这套基于RocketMQ事务消息的方案我们成功地将数据库和缓存的更新操作串联起来即使在服务宕机等极端情况下也能通过事务回查机制保证数据的最终一致性。这为省赚客APP的用户提供了稳定、可靠的优惠券查询体验是构建高可用返利系统的核心技术之一。本文著作权归 省赚客app 研发团队转载请注明出处