幂等与分布式锁 

应用场景 

 

总结

针对共享资源,同时请求

1,一个人,重复操作,短时间内多次请求,只允许第一次请求成功,其他提示,请稍后重试

2,两个人,同时操作一个数据,比如订单号

代码示例

1,租赁重复提交

RLock rLock = redissonClient.getLock(Constant.LEASE_CONFIRM_DELIVERY_LOCK + request.getId());
try {
    boolean isLocked = rLock.tryLock();
    if (!isLocked) {
        throw new XXXServiceException(ErrConstant.OPERATION_FAILED, "请勿重复提交");
    }
} finally {
    if (rLock.isLocked() && rLock.isHeldByCurrentThread()) {
        rLock.unlock();
    }
}

2,订单中心,重复下发

   @Override
    @Transactional(rollbackFor = Throwable.class)
    @RedisDistributedLock(redisTemplateBeanName = "ylOrderRedisTemplate", handlerName = "DPS订单下发新增", lockName = "'" + DPS_TO_YL_ORDER + "'+#orderInfoDTO.outOrderNo")
    public String createOrderByEsb(OrderInfoDTO orderInfoDTO) {
        log.info("创建订单,开始执行,外部订单号: {}", orderInfoDTO.getOutOrderNo());
        String outOrderNo = orderInfoDTO.getOutOrderNo();
        TtOrder ttOrder = new TtOrder();
        ttOrder.setIsDel(false);
        ttOrder.setOutOrderNo(outOrderNo);
        // 先查询,没有再写入
        // 实际结果 还没有执行新增订单,就查到了
        // 后续都执行完成了,订单创建了,运单也异步创建了,这个再抛出一个异常,导致dps中间表状态更新失败,是2
        // 然后再重试,这时候订单是真有了,再重新下发,重试9次,依然失败
        log.info("创建订单,订单查询,开始执行,外部订单号: {}", orderInfoDTO.getOutOrderNo());
        List<TtOrder> ttOrders = selectAllByRecord(ttOrder);
        log.info("创建订单,订单查询,完成执行,外部订单号: {}", orderInfoDTO.getOutOrderNo());
        if (!CollectionUtils.isEmpty(ttOrders)) {
            List<TtOrder> orderList = ttOrders.stream().filter(x -> {
                String orderStatus = x.getOrderStatus();
                OrderStatus status = OrderStatus.getOrderStatus(orderStatus);
                return !OrderStatus.IS_CANCEL.equals(status) && !OrderStatus.FINISH.equals(status);
            }).collect(Collectors.toList());
            if (!CollectionUtils.isEmpty(orderList)) {
                log.info("创建订单,重复下发,下发订单信息: {},查询订单信息: {}", JSON.toJSONString(orderInfoDTO),JSON.toJSONString(ttOrders));
                throw new OrderServiceException(OrderServiceErrorCode.ORDER_IS_EXISTS);
            }
        }
@Target({ElementType.METHOD})
@Retention(RetentionPolicy.RUNTIME)
@Documented
public @interface RedisDistributedLock {
    String redisTemplateBeanName() default "redisTemplate";

    String handlerName() default "default distributed lock name";

    String lockName();

    int lockExpire() default 60;
}

切面逻辑,如果获取锁失败,则不执行业务逻辑 (这个理解有偏差,详细如下)

  @Around("redisDistributedLock()")
    public Object doAround(ProceedingJoinPoint point) throws Throwable {
        Signature signature = point.getSignature();
        MethodSignature methodSignature = (MethodSignature)signature;
        Method method = methodSignature.getMethod();
        RedisDistributedLock redisDistributedLock = (RedisDistributedLock)method.getAnnotation(RedisDistributedLock.class);
        String redisTemplateBeanName = redisDistributedLock.redisTemplateBeanName();
        RedisTemplate bean = (RedisTemplate)this.beanFactory.getBean(redisTemplateBeanName);
        String handlerName = redisDistributedLock.handlerName();
        String lockName = redisDistributedLock.lockName();
        if (lockName.contains("#")) {
            lockName = (String)this.parse(lockName, String.class, point);
        }

        int lockExpire = redisDistributedLock.lockExpire();
        Object proceed = null;
        NewRedisDistributedLock newRedisDistributedLock = new NewRedisDistributedLock(bean);
        RedisLock lock = newRedisDistributedLock.getLock(lockName, (long)lockExpire);

        try {
            if (lock != null) {
                log.debug("分布式锁处理器开始执行[{},{}]", handlerName, lockName);
                proceed = point.proceed();
            }
        } catch (Exception var18) {
            log.error("分布式锁处理器[{}]执行失败", handlerName, var18);
            throw new YlException("013", String.format("[%s]执行异常[%s]", handlerName, var18.getMessage()));
        } finally {
            log.debug("分布式锁处理器执行[{},{}]结束", handlerName, lockName);

            assert lock != null;

            lock.unlock();
        }

        return proceed;
    }

实际上,获取锁的逻辑有重试的代码

RedisLock lock = newRedisDistributedLock.getLock(lockName, (long)lockExpire);

过期时间60,重试480,每次休眠8s。

也就是说60s内,有重复下发,是获取锁失败的,之后进行重试。

实际上应该不需要进行重试,60s内重复下发不处理。如果60s外再次下发,分布式锁是可以获取锁的,兜底方案就是查询数据库订单,是否已经存在。

为什么用分布式锁呢?首先他解决的是同一个时刻,同一个变量的,并发下发,只需处理一个。

如果是查询数据库并发是控制不住的。

分布式锁需要重试的场景

如果上游有重试,比如网络原因,上游job下发9次/或者ESB下发重试,一般重试时间比较短,

这时候是需要重试,但是还是需要做幂等,重复创建的问题。

    public RedisLock getLock(String key, long expireSeconds) {
        return this.getLock(key, expireSeconds, 480, 8L);
    }
    public RedisLock getLock(final String key, final long expireSeconds, int maxRetryTimes, long retryIntervalTimeMillis) {
        final String value = UUID.randomUUID().toString();

        for(int i = 0; i < maxRetryTimes; ++i) {
            String status = (String)this.stringRedisTemplate.execute(new RedisCallback<String>() {
                public String doInRedis(RedisConnection connection) throws DataAccessException {
                    Jedis jedis = (Jedis)connection.getNativeConnection();
                    String status = jedis.set(key, value, "nx", "ex", expireSeconds);
                    return status;
                }
            });
            if ("OK".equals(status) && value.equals(this.stringRedisTemplate.opsForValue().get(key))) {
                logger.debug("线程[{}]获取了锁,Hash值[{}],尝试次数[{}]", Thread.currentThread().getName() + " :" + value, i + 1);
                return new NewRedisDistributedLock.RedisLockInner(this.stringRedisTemplate, key, value);
            }

            try {
                if (retryIntervalTimeMillis > 0L) {
                    Thread.sleep(retryIntervalTimeMillis);
                } else {
                    Thread.sleep(10L);
                }
            } catch (InterruptedException var11) {
                logger.error("线程[{}]中断,竞争锁失败,Hash值[{}]", Thread.currentThread(), value);
                break;
            }

            if (Thread.currentThread().isInterrupted()) {
                break;
            }
        }

        logger.debug("线程[{}]竞争锁失败,Hash值[{}],尝试次数[{}]", Thread.currentThread().getName() + " :" + value, maxRetryTimes);
        return null;
    }

3,运单中心

实际上只是锁定解锁

  @RedisDistributedLock(redisTemplateBeanName = "taskRedisTemplate", handlerName = "运单锁定/解锁分布式锁", lockName = "'" + BusiConstant.YL_ORDER_NO_LOCK + "'+#p0.orderNumber")
public void lockOrUnLockTask(DPTaskLockOrUnlockReqDTO dpTaskLockOrUnlockReqDTO) {
lockName = "'" + BusiConstant.YL_ORDER_NO_LOCK + "'+#p0.orderNumber"

相关key 

1,pc锁定解锁

2,更新运单状态

3,状态机触发

详细如下 

    @Override
    public Object fire(DPContext dpContext) {
        //开启分布式锁
        DPStateEnum initialState = dpContext.getInitialState();
        UntypedStateMachine stateMachine = stateMachineBuilder.newUntypedStateMachine(
                initialState, applicationContext);
        if (dpContext.isReentrantLock()) {//可重入锁,不需要重复加锁
            stateMachine.fire(dpContext.getEvent(), dpContext);
        } else {
            String lockName = BusiConstant.YL_ORDER_NO_LOCK + dpContext.getOrderNumber();
            distributedLockFire("地跑运单状态流转分布式锁", lockName, dpContext.getEvent(), dpContext, stateMachine, null);
        }
        return dpContext.getJghcTaskCreateForDP();
    }
    protected void distributedLockFire(String handlerName, String lockName, E event, C context, UntypedStateMachine stateMachine, RedisTemplate redis) {
        if (ToolUtil.isOneEmpty(handlerName, lockName, event, context, stateMachine)) {
            throw new YlTaskException(ErrConstant.INVALID_DATAFILED, "参数不全");
        }
        if (null == redis) {
            redis = redisTemplate;
        }
        NewRedisDistributedLock newRedisDistributedLock = new NewRedisDistributedLock(redis);
        RedisLock lock = newRedisDistributedLock.getLock(lockName, 60L);
        try {
            if (lock != null) {
                log.debug("分布式锁处理器开始执行[{}},{}]", handlerName, lockName);
                stateMachine.fire(event, context);
            }
        } finally {
            log.debug("分布式锁处理器执行[{},{}]结束", handlerName, lockName);
            assert lock != null;
            lock.unlock();
        }
    }

之前理解偏差,加了分布式锁,运单创建指派发运到达是排队执行的 

但实际是并发执行没有加锁,如下,可重复

if (dpContext.isReentrantLock()) {//可重入锁,不需要重复加锁
    stateMachine.fire(dpContext.getEvent(), dpContext);

不然pc锁定解锁,如果使用了分布式锁,其他并发操作,比如发运,到达,没有获取锁,是不会执行的,反而导致流程不正常。

4,订单中心,其他场景

 

思考

只是解决了并发情况下的,重复下单,如果过了几分钟又重复下单了呢,

方案是:代码层面根据订单号查询了订单信息,相当于又做了一次幂等校验。

其他方案呢:比如使用setNx做幂等,缓存订单号

会存在几个问题

1,缓存时间多久,才能保证后续不重复下发

2,业务层面如果是缓存订单,其实还需要判断订单状态,如果取消或者完成,还可以重复下单,只用setNx,阻塞了正常流程

所以,有必要查询数据库做幂等

更多推荐