首页 / 科技 / 生产级实战:基于Redis Lua的分布式幂等框架设计与实现

生产级实战:基于Redis Lua的分布式幂等框架设计与实现

摸鱼不慌
摸鱼不慌

一、痛点分析:为什么你的接口不安全?

在支付、下单、账务核心链路中,我们常遇到:
  1. 用户重复提交:前端防抖失效,用户连续点击“支付”按钮。

  2. 超时重试:HTTP/RPC 客户端设置了超时重试机制(如 Feign Retry)。

  3. 消息队列重复消费:Kafka 的 Rebalance 或 RocketMQ 的 ACK 超时导致消息重新投递。

  4. 分布式事务回滚重试:Seata 或 TCC 模式的 Cancel 阶段重试。

核心诉求
  • 唯一性:同一个请求只执行一次。

  • 原子性:判断幂等Key是否存在、记录状态必须原子操作。

  • 高性能:不能因为幂等校验成为系统瓶颈。

  • 高可用:支持过期清理,防止Redis无限膨胀。

二、方案选型与设计

1. 幂等Key的设计

我们采用 业务唯一标识 + 令牌(Token) 的方式:
  • Token模式(推荐):前端在调用接口前先获取Token,提交时携带,Token用后即焚。

  • 业务Key模式order:create:{userId}:{productId}:{timestamp},适用于MQ消费。

2. 存储选型:Redis

利用 Redis 的 SETNX (Set if Not Exists) 特性。但单纯的 SETNX 无法解决 状态流转 问题(例如:请求正在处理中,还是已处理完成)。

3. 核心逻辑

我们将幂等记录分为三个状态:
  • 0: 处理中 (Processing)

  • 1: 处理成功 (Success)

  • 2: 处理失败 (Fail)


三、生产级代码实战

1. 定义幂等注解(AOP的核心)

通过注解实现无侵入式的幂等校验。
/**
 * 幂等注解
 * 标注在Controller方法上,用于自动进行幂等校验
 */@Target(ElementType.METHOD)@Retention(RetentionPolicy.RUNTIME)@Documentedpublic @interface Idempotent {    /**
     * 幂等Key的前缀
     */
    String keyPrefix();    /**
     * 获取幂等Key的SpEL表达式
     * 例如:#request.orderId 或 #token
     */
    String key();    /**
     * 过期时间(秒),默认1小时
     */
    int expire() default 3600;    /**
     * 错误信息
     */
    String message() default "请勿重复提交";    /**
     * 是否删除Key(true=请求完成后删除,false=保留至过期)
     * 一般查询类保留,写操作删除
     */
    boolean delKey() default false;
}

2. Lua脚本:保证原子性(核心)

这是整个方案的灵魂。我们使用 Lua 脚本来处理复杂的逻辑判断,避免并发下的竞态条件。
-- idempotent.lua-- KEYS[1]: 幂等Key-- ARGV[1]: 过期时间(秒)-- ARGV[2]: 当前时间戳(用于记录时间)-- 1. 尝试设置Key,如果不存在则设置值为 "0"(处理中),并设置过期时间local result = redis.call('SET', KEYS[1], '0', 'NX', 'EX', ARGV[1])if result then
    -- 设置成功,说明是第一次请求
    return 1end-- 2. 如果Key已存在,获取其值local value = redis.call('GET', KEYS[1])-- 3. 判断状态if value == '1' then
    -- 已经处理成功,返回重复提交标识
    return 0elseif value == '2' then
    -- 上次处理失败,允许重试(或者根据业务需求返回失败)
    -- 重置状态为处理中,并重置过期时间
    redis.call('SET', KEYS[1], '0', 'EX', ARGV[1])    return 1else
    -- value == '0',表示正在处理中,返回并发请求标识
    return -1end

3. Redis配置与Lua脚本加载

@Configuration@Slf4jpublic class RedisIdempotentConfig {    @Bean
    public DefaultRedisScript<Long> idempotentScript() {
        DefaultRedisScript<Long> redisScript = new DefaultRedisScript<>();
        redisScript.setLocation(new ClassPathResource("lua/idempotent.lua"));
        redisScript.setResultType(Long.class);        return redisScript;
    }    @Bean
    public RedisTemplate<String, Object> redisTemplate(RedisConnectionFactory factory) {
        RedisTemplate<String, Object> template = new RedisTemplate<>();
        template.setConnectionFactory(factory);        // 使用String序列化Key
        template.setKeySerializer(new StringRedisSerializer());        // 使用Jackson序列化Value
        Jackson2JsonRedisSerializer<Object> serializer = new Jackson2JsonRedisSerializer<>(Object.class);        ObjectMapper mapper = new ObjectMapper();
        mapper.activateDefaultTyping(LaissezFaireSubTypeValidator.instance, ObjectMapper.DefaultTyping.NON_FINAL);
        serializer.setObjectMapper(mapper);
        template.setValueSerializer(serializer);
        template.setHashKeySerializer(new StringRedisSerializer());
        template.setHashValueSerializer(serializer);
        template.afterPropertiesSet();        return template;
    }
}

4. AOP切面:拦截请求

@Aspect@Component@Slf4jpublic class IdempotentAspect {    @Autowired
    private RedisTemplate<String, Object> redisTemplate;    @Autowired
    private DefaultRedisScript<Long> idempotentScript;    @Pointcut("@annotation(com.example.idempotent.Idempotent)")
    public void pointCut() {}    @Around("pointCut()")
    public Object around(ProceedingJoinPoint joinPoint) throws Throwable {        MethodSignature signature = (MethodSignature) joinPoint.getSignature();        Method method = signature.getMethod();        Idempotent idempotent = method.getAnnotation(Idempotent.class);        // 1. 解析SpEL表达式获取幂等Key
        String key = parseKey(idempotent.key(), method, joinPoint.getArgs());        String redisKey = idempotent.keyPrefix() + ":" + key;        // 2. 执行Lua脚本
        Long result = redisTemplate.execute(
                idempotentScript,
                Collections.singletonList(redisKey),
                String.valueOf(idempotent.expire())
        );        // 3. 处理返回结果
        if (result == null || result == -1) {            // -1: 正在处理中(并发请求)
            log.warn("Request is processing, key: {}", redisKey);            throw new BusinessException(idempotent.message());
        } else if (result == 0) {            // 0: 已处理成功(重复提交)
            log.warn("Duplicate request detected, key: {}", redisKey);            // 这里可以根据业务返回缓存的结果,或者抛异常
            // 例如:return getCachedResult(redisKey);
            throw new BusinessException("重复提交,该请求已处理成功");
        }        // 4. 执行业务逻辑
        Object proceed;        try {
            proceed = joinPoint.proceed();            // 5. 业务成功,更新状态为 1
            redisTemplate.opsForValue().set(redisKey, "1", idempotent.expire(), TimeUnit.SECONDS);            return proceed;
        } catch (BusinessException e) {            // 6. 业务失败,更新状态为 2(允许重试)
            redisTemplate.opsForValue().set(redisKey, "2", idempotent.expire(), TimeUnit.SECONDS);            throw e;
        } catch (Exception e) {            // 系统异常,删除Key,允许重试(视业务而定,也可以标记为失败)
            log.error("System error, removing idempotent key: {}", redisKey, e);
            redisTemplate.delete(redisKey);            throw e;
        } finally {            // 7. 如果配置了删除Key(非查询类操作),在事务提交后删除
            if (idempotent.delKey()) {                // 注意:这里需要确保事务提交后再删除,可以使用TransactionSynchronizationManager
                // 简化版:直接删除(在高并发下可能有极短的窗口期问题,生产环境建议用事务同步)
                // redisTemplate.delete(redisKey);
            }
        }
    }    /**
     * 解析SpEL表达式
     */
    private String parseKey(String keyExpression, Method method, Object[] args) {        LocalVariableTableParameterNameDiscoverer discoverer = new LocalVariableTableParameterNameDiscoverer();
        String[] paramNames = discoverer.getParameterNames(method);        if (paramNames == null || paramNames.length == 0) {            return keyExpression;
        }        ExpressionParser parser = new SpelExpressionParser();        StandardEvaluationContext context = new StandardEvaluationContext();        for (int i = 0; i < paramNames.length; i++) {
            context.setVariable(paramNames[i], args[i]);
        }        Expression expression = parser.parseExpression(keyExpression);        return expression.getValue(context, String.class);
    }
}

5. 业务接口应用

@RestController@RequestMapping("/order")public class OrderController {    @Autowired
    private OrderService orderService;    /**
     * 创建订单接口
     * 使用幂等注解
     * @param token 前端获取的令牌
     * @param request 下单请求
     */
    @PostMapping("/create")
    @Idempotent(
        keyPrefix = "idempotent:order:create",
        key = "#token", // 使用请求参数中的token作为幂等Key
        expire = 300,   // 5分钟有效期
        message = "订单正在处理中,请勿重复提交"
    )
    public Response<OrderVO> createOrder(@RequestParam("token") String token,                                         @RequestBody CreateOrderRequest request) {        // 业务逻辑
        OrderVO order = orderService.create(request);        return Response.success(order);
    }    /**
     * 支付接口
     * 使用订单号作为幂等Key
     */
    @PostMapping("/pay")
    @Idempotent(
        keyPrefix = "idempotent:order:pay",
        key = "#request.orderNo", // 使用订单号
        expire = 600
    )
    public Response<String> pay(@RequestBody PayRequest request) {
        orderService.pay(request);        return Response.success("支付成功");
    }
}

6. 全局异常处理器

@RestControllerAdvice@Slf4jpublic class GlobalExceptionHandler {    @ExceptionHandler(BusinessException.class)
    public Response<Void> handleBusinessException(BusinessException e) {
        log.warn("Business exception: {}", e.getMessage());        return Response.fail(e.getCode(), e.getMessage());
    }
}

四、源码级深度剖析

Lua脚本为何能保证原子性?

Redis 是单线程执行命令的。当 Lua 脚本被调用时,Redis 会将其作为一个 整体 执行,期间不会被其他命令打断。这就解决了以下并发问题:
  1. 线程A判断Key不存在。

  2. 线程B同时判断Key不存在。

  3. 线程A设置Key,线程B设置Key(导致重复执行)。

在Lua脚本中,SET NX 和后续的 GET 操作是连续的,中间不会插入其他Redis指令,因此保证了判断和设置的原子性。

状态机设计

我们的Lua脚本实现了一个简单的 状态机
  • 初始状态:Key不存在。

  • 迁移1SET NX 成功 -> 0 (Processing)。

  • 迁移2:业务成功 -> 1 (Success)。

  • 迁移3:业务失败 -> 2 (Fail)。

  • 迁移4:Fail状态下再次请求 -> 重置为 0 (允许重试)。

这种设计比单纯的 SETNX 更强大,因为它区分了“处理中”和“处理完成”,有效防止了 “悬挂请求”(即第一个请求很慢,第二个请求进来时第一个还没写完结果)导致的数据不一致。

五、生产环境避坑指南

  1. Redis Key 爆炸:务必设置合理的 expire 时间。对于Token模式,建议在业务完成后主动删除Key(设置 delKey=true),减少Redis内存占用。

  2. SpEL 表达式性能:虽然 SpEL 很方便,但在超高并发下(QPS > 10万),解析表达式会有微小开销。如果对性能极度敏感,可以改为在注解中直接指定参数名,通过反射获取,或者使用 ThreadLocal 传递幂等Key。

  3. Redis 集群模式:Lua 脚本要求所有 Key 必须落在同一个 Slot 上。我们的 Key 设计是 prefix:id,只要 prefix 相同,就会路由到同一个 Slot,因此该方案天然支持 Redis Cluster。

  4. 事务一致性:如果业务方法包含数据库事务,Lua脚本的执行是在事务之外的。这意味着:Redis中标记为成功,但数据库事务回滚了。解决方案:使用 TransactionSynchronizationManager 在事务提交后再更新Redis状态,或者采用 最终一致性 方案(定时任务核对Redis与DB状态)。

  5. Token生成策略:如果是前端获取Token,Token必须是 一次性 的。生成Token时也要利用 Redis 的原子操作。

六、总结

本文设计了一套基于 Redis + Lua + AOP 的通用幂等框架。
核心优势
  • 通用性:通过注解和SpEL表达式,适用于任何接口。

  • 安全性:Lua脚本保证了原子性,状态机设计防止了并发问题。

  • 高性能:Redis内存操作,微秒级响应。

  • 可维护:AOP实现无侵入,业务代码零耦合。

生产建议
  • 对于核心链路(如支付),建议配合 分布式锁 使用(先拿锁,再校验幂等)。

  • 定期监控 Redis 中幂等Key的数量,防止内存泄漏。

  • 在网关层(如Spring Cloud Gateway)也可以集成类似的幂等逻辑,做第一道拦截。