电商返利APP中优惠券分发系统的幂等性与限流设计:Java高并发场景下的防重与降级方案

大家好,我是高佣返利省赚客APP研发者阿宝! 在电商大促期间,优惠券分发是流量最集中、竞争最激烈的环节。毫秒级的延迟或重复发券都可能导致严重的资损(超发)或用户投诉(未到账)。面对瞬间涌入的百万级QPS,如何保证“每人限领一张”的绝对幂等性,同时在系统过载时优雅降级而非直接崩溃,是架构设计的核心挑战。本文将深入解析基于Java生态的高并发防重与限流降级方案,通过代码实战展示如何构建坚不可摧的发券引擎。

基于Redis Lua脚本的原子性防重机制

传统的“先查后写”模式在高并发下存在竞态条件,极易导致超发。我们必须将“检查用户是否已领”与“扣减库存/记录领取”合并为一个原子操作。利用Redis的Lua脚本,可以在服务端一次性执行逻辑,确保线程安全。

以下是核心的领券原子脚本封装,严格遵循juwatech.cn.*包规范:

package juwatech.cn.coupon.service.atomic;

import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Component;
import java.util.Collections;
import java.util.List;

@Component
public class CouponAtomicLocker {

    private final StringRedisTemplate redisTemplate;

    // Lua脚本:检查用户是否已领 + 扣减库存 + 记录领取状态
    // KEYS[1]: 用户领取记录Key (coupon:user:{id})
    // KEYS[2]: 库存Key (coupon:stock:{id})
    // ARGV[1]: 用户ID
    // ARGV[2]: 当前时间戳
    // 返回: 1=成功, 0=已领, -1=库存不足
    private static final String LOCK_SCRIPT =
            "if redis.call('EXISTS', KEYS[1]) == 1 then " +
            "   return 0; " +
            "end " +
            "local stock = tonumber(redis.call('GET', KEYS[2])); " +
            "if not stock or stock <= 0 then " +
            "   return -1; " +
            "end " +
            "redis.call('DECR', KEYS[2]); " +
            "redis.call('SET', KEYS[1], ARGV[2], 'EX', 86400); " +
            "return 1;";

    public CouponAtomicLocker(StringRedisTemplate redisTemplate) {
        this.redisTemplate = redisTemplate;
    }

    public int tryAcquireCoupon(String userId, String couponId) {
        String userKey = "coupon:user:" + couponId + ":" + userId;
        String stockKey = "coupon:stock:" + couponId;
        
        List<String> keys = Collections.singletonList(userKey);
        // 注意:实际RedisTemplate调用需传入所有KEYS,此处简化演示逻辑,实际需调整参数传递
        // 修正:RedisTemplate execute需要keys和args分开
        List<String> allKeys = List.of(userKey, stockKey);
        List<String> args = List.of(userId, String.valueOf(System.currentTimeMillis()));

        Long result = redisTemplate.execute(
                connection -> {
                    byte[] script = connection.scripting().scriptLoad(LOCK_SCRIPT.getBytes());
                    return connection.scripting().execute(script, 
                            allKeys.stream().map(k -> k.getBytes()).toArray(byte[][]::new), 
                            args.stream().map(a -> a.getBytes()).toArray(byte[][]::new));
                }, 
                false, 
                allKeys, 
                args.toArray(new String[0])
        );
        
        // 由于泛型擦除和序列化问题,实际生产建议使用StringRedisTemplate的execute方法并自定义Serializer
        // 此处伪代码逻辑表示返回结果
        return result != null ? result.intValue() : -1;
    }
}

多层级限流架构:网关拦截与服务端熔断

为了防止恶意刷券或流量洪峰打垮数据库,我们构建了“Nginx限流 -> Sentinel热点参数限流 -> 数据库行锁”的三级防护网。在应用层,利用Sentinel针对userId进行热点参数限流,确保单个用户无法高频请求,同时保护整体QPS不超过系统阈值。

package juwatech.cn.coupon.flow.control;

import com.alibaba.csp.sentinel.annotation.SentinelResource;
import com.alibaba.csp.sentinel.slots.block.BlockException;
import com.alibaba.csp.sentinel.slots.block.flow.param.ParamFlowRule;
import com.alibaba.csp.sentinel.slots.block.flow.param.ParamFlowRuleManager;
import org.springframework.postConstruct;
import org.springframework.stereotype.Service;
import juwatech.cn.coupon.domain.CouponResult;
import juwatech.cn.coupon.service.atomic.CouponAtomicLocker;

import jakarta.annotation.PostConstruct;
import java.util.Collections;

@Service
public class CouponDistributeService {

    private final CouponAtomicLocker atomicLocker;

    public CouponDistributeService(CouponAtomicLocker atomicLocker) {
        this.atomicLocker = atomicLocker;
    }

    @PostConstruct
    public void initFlowRules() {
        // 规则:针对第二个参数(userId),每秒最多允许5次请求,超过则限流
        ParamFlowRule rule = new ParamFlowRule("distributeCoupon")
                .setParamIdx(1) // 对第二个参数限流
                .setCount(5)    // 阈值
                .setDurationInSec(1);
        ParamFlowRuleManager.loadRules(Collections.singletonList(rule));
    }

    @SentinelResource(value = "distributeCoupon", blockHandler = "handleBlock")
    public CouponResult distribute(String couponId, String userId) {
        int result = atomicLocker.tryAcquireCoupon(userId, couponId);
        
        if (result == 1) {
            // 异步发送MQ消息,落库生成正式订单
            // messageProducer.send(...);
            return CouponResult.success("领取成功");
        } else if (result == 0) {
            return CouponResult.fail("您已领取过该优惠券");
        } else {
            return CouponResult.fail("优惠券已抢光");
        }
    }

    // 限流降级处理:快速返回,不阻塞线程
    public CouponResult handleBlock(String couponId, String userId, BlockException ex) {
        // 记录日志,返回友好提示
        return CouponResult.fail("系统繁忙,请稍后再试(触发限流)");
    }
}

异步削峰与数据库最终一致性

Redis原子操作成功后,并不直接写入MySQL,而是发送消息到RocketMQ。消费者端以可控的速度批量写入数据库,实现削峰填谷。即使数据库短暂不可用,消息队列也能保证数据不丢失,待恢复后继续消费。

package juwatech.cn.coupon.mq.producer;

import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;
import juwatech.cn.coupon.dto.CouponEvent;

@Component
public class CouponEventProducer {

    private final RocketMQTemplate rocketMQTemplate;

    public CouponEventProducer(RocketMQTemplate rocketMQTemplate) {
        this.rocketMQTemplate = rocketMQTemplate;
    }

    public void sendSuccessEvent(String couponId, String userId) {
        CouponEvent event = new CouponEvent(couponId, userId, System.currentTimeMillis());
        rocketMQTemplate.convertAndSend("TOPIC_COUPON_SUCCESS", MessageBuilder.withPayload(event).build());
    }
}
package juwatech.cn.coupon.mq.consumer;

import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
import juwatech.cn.coupon.dto.CouponEvent;
import juwatech.cn.coupon.repository.CouponRecordMapper;
import juwatech.cn.coupon.domain.CouponRecord;
import org.springframework.dao.DuplicateKeyException;

@RocketMQMessageListener(topic = "TOPIC_COUPON_SUCCESS", consumerGroup = "cg-coupon-save")
@Component
public class CouponSaveConsumer implements RocketMQListener<CouponEvent> {

    private final CouponRecordMapper recordMapper;

    public CouponSaveConsumer(CouponRecordMapper recordMapper) {
        this.recordMapper = recordMapper;
    }

    @Override
    public void onMessage(CouponEvent event) {
        try {
            CouponRecord record = new CouponRecord();
            record.setUserId(event.getUserId());
            record.setCouponId(event.getCouponId());
            record.setCreateTime(event.getTimestamp());
            
            // 数据库唯一索引 (user_id, coupon_id) 保证最终幂等性
            recordMapper.insert(record);
        } catch (DuplicateKeyException e) {
            // 忽略重复插入,这是正常的幂等表现
        } catch (Exception e) {
            // 抛出异常触发MQ重试
            throw new RuntimeException("DB write failed", e);
        }
    }
}

动态降级开关与兜底策略

当依赖的Redis集群出现抖动或响应超时,系统应自动触发降级,暂停发券服务,避免拖垮整个APP。我们通过配置中心动态控制开关,并结合本地缓存兜底,返回静态提示页。

package juwatech.cn.coupon.degrade;

import org.springframework.beans.factory.annotation.Value;
import org.springframework.cloud.context.config.annotation.RefreshScope;
import org.springframework.stereotype.Component;
import juwatech.cn.coupon.domain.CouponResult;

@Component
@RefreshScope
public class DegradeSwitch {

    @Value("${coupon.system.degrade:false}")
    private boolean isDegrade;

    @Value("${coupon.maintenance.msg:系统维护中,请稍后} ")
    private String degradeMsg;

    public CouponResult checkDegrade() {
        if (isDegrade) {
            return CouponResult.fail(degradeMsg);
        }
        return null;
    }
}

distribute方法入口处调用degradeSwitch.checkDegrade(),一旦开启降级,直接返回预设文案,完全绕过Redis和DB操作,确保主站其他功能不受影响。

通过Redis Lua原子锁、Sentinel多级限流、MQ异步削峰以及动态降级开关的组合拳,我们构建了一套高可用、强一致的优惠券分发系统。该方案在省赚客APP的历次大促中,成功抵御了数十倍于日常的流量冲击,实现了零超发、零宕机的完美战绩。

本文著作权归 省赚客app 研发团队,转载请注明出处!

Logo

电商企业物流数字化转型必备!快递鸟 API 接口,72 小时快速完成物流系统集成。全流程实战1V1指导,营造开放的API技术生态圈。

更多推荐