电商返利APP中优惠券分发系统的幂等性与限流设计:Java高并发场景下的防重与降级方案
电商返利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 研发团队,转载请注明出处!
更多推荐



所有评论(0)