diff --git a/urbanops-module-garden/src/main/java/com/zteits/urbanops/module/garden/dal/redis/IntegralLockCoreRedisDAO.java b/urbanops-module-garden/src/main/java/com/zteits/urbanops/module/garden/dal/redis/IntegralLockCoreRedisDAO.java new file mode 100644 index 0000000..db50caa --- /dev/null +++ b/urbanops-module-garden/src/main/java/com/zteits/urbanops/module/garden/dal/redis/IntegralLockCoreRedisDAO.java @@ -0,0 +1,75 @@ +package com.zteits.urbanops.module.garden.dal.redis; + +import jakarta.annotation.Resource; +import lombok.extern.slf4j.Slf4j; +import org.redisson.api.RLock; +import org.redisson.api.RedissonClient; +import org.springframework.stereotype.Component; +import org.springframework.util.Assert; + +import java.util.concurrent.TimeUnit; + +/** + * 积分操作分布式锁通用工具类 + */ +@Slf4j +@Component +public class IntegralLockCoreRedisDAO { + + // 积分操作锁前缀 + private static final String INTEGRAL_LOCK_PREFIX = "integral:operate:%s:%s"; + // 锁持有超时时间(30秒) + public static final long INTEGRAL_LOCK_TIMEOUT_MILLIS = 30 * 1000L; + + @Resource + private RedissonClient redissonClient; + + /** + * 加锁执行积分操作逻辑 + * @param unitId 单位ID + * @param statDate 统计日期 + * @param timeoutMillis 锁持有超时时间 + * @param runnable 要执行的核心业务逻辑 + */ + public void lock(String unitId, String statDate, Long timeoutMillis, Runnable runnable) { + // 空值校验 + Assert.hasText(unitId, "单位ID不能为空"); + Assert.hasText(statDate, "统计日期不能为空"); + Assert.notNull(timeoutMillis, "锁超时时间不能为空"); + Assert.notNull(runnable, "执行逻辑不能为空"); + + // 构建细粒度锁Key + String lockKey = String.format(INTEGRAL_LOCK_PREFIX, unitId, statDate); + RLock lock = redissonClient.getLock(lockKey); + boolean lockAcquired = false; + + try { + // 非阻塞获取锁:3秒等待超时,timeoutMillis自动过期 + lockAcquired = lock.tryLock(3, timeoutMillis, TimeUnit.MILLISECONDS); + if (!lockAcquired) { + String errorMsg = String.format("获取积分操作锁失败,单位ID:%s,统计日期:%s", unitId, statDate); + log.error(errorMsg); + throw new RuntimeException(errorMsg); + } + log.info("成功获取积分操作锁,锁Key:{}", lockKey); + + // 执行核心业务逻辑 + runnable.run(); + + } catch (InterruptedException e) { + log.error("获取积分操作锁被中断,单位ID:{},统计日期:{}", unitId, statDate, e); + Thread.currentThread().interrupt(); + throw new RuntimeException("积分操作锁获取被中断,请重试", e); + } finally { + // 安全释放锁 + if (lockAcquired && lock.isHeldByCurrentThread()) { + try { + lock.unlock(); + log.info("成功释放积分操作锁,锁Key:{}", lockKey); + } catch (Exception e) { + log.error("释放积分操作锁异常,锁Key:{}", lockKey, e); + } + } + } + } +} diff --git a/urbanops-module-garden/src/main/java/com/zteits/urbanops/module/garden/service/unitintegraltotal/UnitIntegralServiceImpl.java b/urbanops-module-garden/src/main/java/com/zteits/urbanops/module/garden/service/unitintegraltotal/UnitIntegralServiceImpl.java index 7c03d27..e50cc5e 100644 --- a/urbanops-module-garden/src/main/java/com/zteits/urbanops/module/garden/service/unitintegraltotal/UnitIntegralServiceImpl.java +++ b/urbanops-module-garden/src/main/java/com/zteits/urbanops/module/garden/service/unitintegraltotal/UnitIntegralServiceImpl.java @@ -1,5 +1,6 @@ package com.zteits.urbanops.module.garden.service.unitintegraltotal; +import cn.hutool.extra.spring.SpringUtil; import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; import com.baomidou.mybatisplus.core.toolkit.Wrappers; import com.zteits.urbanops.framework.common.util.object.BeanUtils; @@ -12,10 +13,12 @@ import com.zteits.urbanops.module.garden.dal.dataobject.unitintegraladd.UnitInte import com.zteits.urbanops.module.garden.dal.dataobject.unitintegraltotal.UnitIntegralTotalDO; import com.zteits.urbanops.module.garden.dal.mysql.unitintegraladd.UnitIntegralAddMapper; import com.zteits.urbanops.module.garden.dal.mysql.unitintegraltotal.UnitIntegralTotalMapper; +import com.zteits.urbanops.module.garden.dal.redis.IntegralLockCoreRedisDAO; import com.zteits.urbanops.module.garden.service.unitintegraladd.UnitIntegralAddService; import com.zteits.urbanops.module.garden.service.unitintegralsub.UnitIntegralSubService; import jakarta.annotation.Resource; import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.ObjectUtils; import org.apache.commons.lang3.StringUtils; import org.redisson.api.RLock; import org.redisson.api.RedissonClient; @@ -55,6 +58,9 @@ public class UnitIntegralServiceImpl implements UnitIntegralService{ private UnitIntegralAddService unitIntegralAddService; @Resource private RedissonClient redissonClient; + @Resource + private IntegralLockCoreRedisDAO integralLockCoreRedisDAO; + // 锁前缀(区分加/减积分锁) private static final String LOCK_PREFIX_INTEGRAL = "integral:"; @@ -125,85 +131,78 @@ public class UnitIntegralServiceImpl implements UnitIntegralService{ * @param reqVO 积分增加请求(注意:totalIntegral 为「本次增加积分」) */ @Override - @Transactional(rollbackFor = Exception.class) public void addUnitIntegral(UnitIntegralSaveReqVO reqVO){ // 1. 参数校验 + if (reqVO == null) { + throw exception0(BAD_REQUEST.getCode(), "积分增加请求参数不能为空"); + } if (reqVO.getTotalIntegral() <= 0) { - throw exception0(BAD_REQUEST.getCode(),"本次增加积分必须为正数"); + throw exception0(BAD_REQUEST.getCode(), "本次增加积分必须为正数"); + } + if (ObjectUtils.isEmpty(reqVO.getUnitId()) || ObjectUtils.isEmpty(reqVO.getStatDate())) { + throw exception0(BAD_REQUEST.getCode(), "单位ID和统计日期不能为空"); } - // 2. 构建细粒度锁Key:单位ID + 统计日期 + 操作类型 - String lockKey = buildLockKey(reqVO.getUnitId(), reqVO.getStatDate()); - RLock lock = redissonClient.getLock(lockKey); - - // 3. 获取锁(非阻塞,避免线程等待;3秒获取超时,10秒自动过期防死锁) - boolean lockAcquired = false; - try { - lockAcquired = lock.tryLock(3, 10, TimeUnit.SECONDS); - if (!lockAcquired) { - String errorMsg = String.format("获取积分增加锁失败,单位ID:%s,统计日期:%s(请稍后重试)", reqVO.getUnitId(), reqVO.getStatDate()); - log.error(errorMsg); - throw exception0(INTERNAL_SERVER_ERROR.getCode(), errorMsg); - } - - // 2. 构建查询条件 - LambdaQueryWrapper queryWrapper = Wrappers.lambdaQuery(UnitIntegralTotalDO.class) - .eq(UnitIntegralTotalDO::getUnitType, reqVO.getUnitType()) - .eq(UnitIntegralTotalDO::getStatDate, reqVO.getStatDate()) - .eq(UnitIntegralTotalDO::getUnitId, reqVO.getUnitId()) - .eq(UnitIntegralTotalDO::getDeleted, 0); - - // 3. 可加行锁(FOR UPDATE)防止多线程同时扣减,保证数据一致性 - UnitIntegralTotalDO unitIntegralTotalDO = unitIntegralTotalMapper.selectOne(queryWrapper); - - // 4. 初始化变量(语义清晰,避免冗余) - int beforeIntegral = 0; // 本次增加前的累计增加积分 - int afterIntegral = 0; // 本次增加后的累计增加积分 - int currentAdd = reqVO.getTotalIntegral(); // 本次增加积分(重命名,语义更准) + // 2. 调用通用锁工具类,加锁执行核心逻辑(参考executeNotify的Runnable风格) + integralLockCoreRedisDAO.lock( + reqVO.getUnitId(), + reqVO.getStatDate(), + IntegralLockCoreRedisDAO.INTEGRAL_LOCK_TIMEOUT_MILLIS, + () -> { + // 3. 锁内二次校验(关键!避免分布式锁并发问题,参考executeNotify的dbTask校验) + // 场景:两个线程同时通过前置校验,第一个执行完后,第二个拿到锁仍会执行,需校验数据库状态 + LambdaQueryWrapper checkWrapper = Wrappers.lambdaQuery(UnitIntegralTotalDO.class) + .eq(UnitIntegralTotalDO::getUnitId, reqVO.getUnitId()) + .eq(UnitIntegralTotalDO::getStatDate, reqVO.getStatDate()) + .eq(UnitIntegralTotalDO::getUnitType, reqVO.getUnitType()) + .eq(UnitIntegralTotalDO::getDeleted, 0); + UnitIntegralTotalDO dbTotalDO = unitIntegralTotalMapper.selectOne(checkWrapper); - // 5. 新增/更新总积分表(核心逻辑) - if (unitIntegralTotalDO != null) { - // 5.1 已有记录:更新积分 - beforeIntegral = unitIntegralTotalDO.getTotalAdd(); - int newTotalAdd = beforeIntegral + currentAdd; - // 更新累计增加积分 + 当前总积分(累计增加 - 累计扣减) - unitIntegralTotalDO.setTotalAdd(newTotalAdd); - unitIntegralTotalDO.setTotalIntegral(newTotalAdd - unitIntegralTotalDO.getTotalSub()); - unitIntegralTotalDO.setUpdateTime(LocalDateTime.now()); - unitIntegralTotalDO.setUpdater(String.valueOf(getLoginUserId())); - unitIntegralTotalMapper.updateById(unitIntegralTotalDO); - afterIntegral = newTotalAdd; + // 4. 调用带事务的核心执行方法(自注入避免事务失效) + getSelf().executeAddUnitIntegral0(reqVO, dbTotalDO); + } + ); + } + /** + * 积分增加核心执行方法(带事务,内部调用) + */ + @Transactional(rollbackFor = Exception.class) + public void executeAddUnitIntegral0(UnitIntegralSaveReqVO reqVO, UnitIntegralTotalDO unitIntegralTotalDO) { + int beforeIntegral = 0; // 本次增加前的累计增加积分 + int afterIntegral = 0; // 本次增加后的累计增加积分 + int currentAdd = reqVO.getTotalIntegral(); // 本次增加积分(重命名,语义更准) - log.info("更新单位积分成功,单位ID:{},统计日期:{},本次增加:{},累计增加:{}", - reqVO.getUnitId(), reqVO.getStatDate(), currentAdd, newTotalAdd); - } else { - // 5.2 无记录:插入新记录 - UnitIntegralTotalDO saveDO = new UnitIntegralTotalDO(); - BeanUtils.copyProperties(reqVO, saveDO); - saveDO.setTotalAdd(currentAdd); // 累计增加 = 本次增加 - saveDO.setTotalSub(0); // 初始扣减为0 - saveDO.setTotalIntegral(currentAdd); // 初始总积分 = 本次增加 - unitIntegralTotalMapper.insert(saveDO); - afterIntegral = currentAdd; + // 5. 新增/更新总积分表(核心逻辑) + if (unitIntegralTotalDO != null) { + // 5.1 已有记录:更新积分 + beforeIntegral = unitIntegralTotalDO.getTotalAdd(); + int newTotalAdd = beforeIntegral + currentAdd; + // 更新累计增加积分 + 当前总积分(累计增加 - 累计扣减) + unitIntegralTotalDO.setTotalAdd(newTotalAdd); + unitIntegralTotalDO.setTotalIntegral(newTotalAdd - unitIntegralTotalDO.getTotalSub()); + unitIntegralTotalDO.setUpdateTime(LocalDateTime.now()); + unitIntegralTotalDO.setUpdater(String.valueOf(getLoginUserId())); + unitIntegralTotalMapper.updateById(unitIntegralTotalDO); + afterIntegral = newTotalAdd; - log.info("新增单位积分成功,单位ID:{},统计日期:{},本次增加:{}", - reqVO.getUnitId(), reqVO.getStatDate(), currentAdd); - } - // 5. 新增积分明细 - createAddDetail(reqVO, beforeIntegral, afterIntegral, currentAdd); - log.info("积分增加成功,单位ID:{},本次增加:{}", reqVO.getUnitId(), currentAdd); + log.info("更新单位积分成功,单位ID:{},统计日期:{},本次增加:{},累计增加:{}", + reqVO.getUnitId(), reqVO.getStatDate(), currentAdd, newTotalAdd); + } else { + // 5.2 无记录:插入新记录 + UnitIntegralTotalDO saveDO = new UnitIntegralTotalDO(); + BeanUtils.copyProperties(reqVO, saveDO); + saveDO.setTotalAdd(currentAdd); // 累计增加 = 本次增加 + saveDO.setTotalSub(0); // 初始扣减为0 + saveDO.setTotalIntegral(currentAdd); // 初始总积分 = 本次增加 + unitIntegralTotalMapper.insert(saveDO); + afterIntegral = currentAdd; - } catch (InterruptedException e) { - log.error("获取锁被中断,单位ID:{}", reqVO.getUnitId(), e); - Thread.currentThread().interrupt(); // 恢复中断状态 - throw exception0(INTERNAL_SERVER_ERROR.getCode(), "积分增加操作被中断,请重试"); - } finally { - // 6. 释放锁(仅当前线程持有锁时释放,避免误删) - if (lockAcquired && lock.isHeldByCurrentThread()) { - lock.unlock(); - log.debug("释放积分增加锁,锁Key:{}", lockKey); - } + log.info("新增单位积分成功,单位ID:{},统计日期:{},本次增加:{}", + reqVO.getUnitId(), reqVO.getStatDate(), currentAdd); } + // 5. 新增积分明细 + createAddDetail(reqVO, beforeIntegral, afterIntegral, currentAdd); + log.info("积分增加成功,单位ID:{},本次增加:{}", reqVO.getUnitId(), currentAdd); } /** * 封装增加明细创建逻辑 @@ -222,81 +221,78 @@ public class UnitIntegralServiceImpl implements UnitIntegralService{ * @param reqVO 积分扣减请求(totalIntegral 为「本次扣减积分」) */ @Override - @Transactional(rollbackFor = Exception.class) public void subUnitIntegral(UnitIntegralSaveReqVO reqVO){ // 1. 前置参数校验 + if (reqVO == null) { + throw exception0(BAD_REQUEST.getCode(), "积分扣减请求参数不能为空"); + } if (reqVO.getTotalIntegral() <= 0) { - throw exception0(BAD_REQUEST.getCode(),"本次增加积分必须为正数"); + throw exception0(BAD_REQUEST.getCode(), "本次扣减积分必须为正数"); // 修复文案错误 + } + if (ObjectUtils.isEmpty(reqVO.getUnitId()) || ObjectUtils.isEmpty(reqVO.getStatDate())) { + throw exception0(BAD_REQUEST.getCode(), "单位ID和统计日期不能为空"); } - // 2. 构建细粒度锁Key:单位ID + 统计日期 - String lockKey = buildLockKey(reqVO.getUnitId(), reqVO.getStatDate()); - RLock lock = redissonClient.getLock(lockKey); + // 2. 调用通用锁工具类,加锁执行核心逻辑(参考executeNotify的Runnable风格) + integralLockCoreRedisDAO.lock( + reqVO.getUnitId(), + reqVO.getStatDate(), + IntegralLockCoreRedisDAO.INTEGRAL_LOCK_TIMEOUT_MILLIS, + () -> { + // 3. 锁内二次校验(关键!避免分布式锁并发问题,参考executeNotify的dbTask校验) + LambdaQueryWrapper queryWrapper = Wrappers.lambdaQuery(UnitIntegralTotalDO.class) + .eq(UnitIntegralTotalDO::getUnitType, reqVO.getUnitType()) + .eq(UnitIntegralTotalDO::getStatDate, reqVO.getStatDate()) + .eq(UnitIntegralTotalDO::getUnitId, reqVO.getUnitId()) + .eq(UnitIntegralTotalDO::getDeleted, 0); // 排除已删除的积分记录 + UnitIntegralTotalDO unitIntegralTotalDO = unitIntegralTotalMapper.selectOne(queryWrapper); + // 4. 校验积分记录是否存在(锁内校验,避免并发删除场景) + if (unitIntegralTotalDO == null) { + String errorMsg = String.format("单位积分记录不存在,单位ID:%s,统计日期:%s", + reqVO.getUnitId(), reqVO.getStatDate()); + log.error(errorMsg); + throw exception0(BAD_REQUEST.getCode(), errorMsg); + } - // 3. 获取锁(非阻塞,避免线程等待;3秒获取超时,10秒自动过期防死锁) - boolean lockAcquired = false; - try { - lockAcquired = lock.tryLock(3, 10, TimeUnit.SECONDS); - if (!lockAcquired) { - String errorMsg = String.format("获取积分增加锁失败,单位ID:%s,统计日期:%s(请稍后重试)", reqVO.getUnitId(), reqVO.getStatDate()); - log.error(errorMsg); - throw exception0(INTERNAL_SERVER_ERROR.getCode(), errorMsg); - } - // 2. 构建查询条件 - LambdaQueryWrapper queryWrapper = Wrappers.lambdaQuery(UnitIntegralTotalDO.class) - .eq(UnitIntegralTotalDO::getUnitType, reqVO.getUnitType()) - .eq(UnitIntegralTotalDO::getStatDate, reqVO.getStatDate()) - .eq(UnitIntegralTotalDO::getUnitId, reqVO.getUnitId()) - .eq(UnitIntegralTotalDO::getDeleted, 0); // 排除已删除的积分记录 + // 5. 调用带事务的核心扣减方法 + getSelf().executeSubUnitIntegral0(reqVO, unitIntegralTotalDO); + } + ); + } + /** + * 扣减积分核心执行方法(带事务,内部调用) + */ + @Transactional(rollbackFor = Exception.class) + public void executeSubUnitIntegral0(UnitIntegralSaveReqVO reqVO, UnitIntegralTotalDO unitIntegralTotalDO) { + // 1. 初始化核心变量 + int currentSub = reqVO.getTotalIntegral(); // 本次扣减积分 + int beforeSub = unitIntegralTotalDO.getTotalSub(); // 扣减前累计扣减积分 + int afterSub = beforeSub + currentSub; // 扣减后累计扣减积分 + int currentTotalIntegral = unitIntegralTotalDO.getTotalAdd() - afterSub; // 扣减后总积分 - // 3. 可加行锁(FOR UPDATE)防止多线程同时扣减,保证数据一致性 - //UnitIntegralTotalDO unitIntegralTotalDO = unitIntegralTotalMapper.selectOne(queryWrapper.last("FOR UPDATE")); - UnitIntegralTotalDO unitIntegralTotalDO = unitIntegralTotalMapper.selectOne(queryWrapper); - // 4. 校验积分记录是否存在 - if (unitIntegralTotalDO == null) { - String errorMsg = String.format("单位积分记录不存在,单位ID:%s,统计日期:%s", - reqVO.getUnitId(), reqVO.getStatDate()); - log.error(errorMsg); - throw exception0(BAD_REQUEST.getCode(),errorMsg); - } + // 2. 业务规则校验:扣减后总积分不能为负数 + if (currentTotalIntegral < 0) { + String errorMsg = String.format("单位积分扣减失败,扣减后总积分为负!单位ID:%s,当前累计增加:%s,本次扣减:%s,累计扣减:%s", + reqVO.getUnitId(), unitIntegralTotalDO.getTotalAdd(), currentSub, beforeSub); + log.error(errorMsg); + throw exception0(BAD_REQUEST.getCode(),errorMsg); + } - // 5. 初始化核心变量 - int currentSub = reqVO.getTotalIntegral(); // 本次扣减积分 - int beforeSub = unitIntegralTotalDO.getTotalSub(); // 扣减前累计扣减积分 - int afterSub = beforeSub + currentSub; // 扣减后累计扣减积分 - int currentTotalIntegral = unitIntegralTotalDO.getTotalAdd() - afterSub; // 扣减后总积分 + // 3. 更新积分总表(核心扣减逻辑) + unitIntegralTotalDO.setTotalSub(afterSub); // 更新累计扣减积分 + unitIntegralTotalDO.setTotalIntegral(currentTotalIntegral); // 更新当前总积分 + unitIntegralTotalDO.setUpdateTime(LocalDateTime.now()); + unitIntegralTotalDO.setUpdater(String.valueOf(getLoginUserId())); + unitIntegralTotalMapper.updateById(unitIntegralTotalDO); - // 6. 业务规则校验:扣减后总积分不能为负数 - if (currentTotalIntegral < 0) { - String errorMsg = String.format("单位积分扣减失败,扣减后总积分为负!单位ID:%s,当前累计增加:%s,本次扣减:%s,累计扣减:%s", - reqVO.getUnitId(), unitIntegralTotalDO.getTotalAdd(), currentSub, beforeSub); - log.error(errorMsg); - throw exception0(BAD_REQUEST.getCode(),errorMsg); - } + log.info("单位积分扣减成功,单位ID:{},统计日期:{},本次扣减:{},扣减前累计扣减:{},扣减后累计扣减:{},扣减后总积分:{}", + reqVO.getUnitId(), reqVO.getStatDate(), currentSub, beforeSub, afterSub, currentTotalIntegral); - // 7. 更新积分总表(核心扣减逻辑) - unitIntegralTotalDO.setTotalSub(afterSub); // 更新累计扣减积分 - unitIntegralTotalDO.setTotalIntegral(currentTotalIntegral); // 更新当前总积分 - unitIntegralTotalDO.setUpdateTime(LocalDateTime.now()); - unitIntegralTotalDO.setUpdater(String.valueOf(getLoginUserId())); - unitIntegralTotalMapper.updateById(unitIntegralTotalDO); - - log.info("单位积分扣减成功,单位ID:{},统计日期:{},本次扣减:{},扣减前累计扣减:{},扣减后累计扣减:{},扣减后总积分:{}", - reqVO.getUnitId(), reqVO.getStatDate(), currentSub, beforeSub, afterSub, currentTotalIntegral); - - // 8. 新增扣减明细 - // 新增减积分明细 - createSubDetail(reqVO, beforeSub, afterSub, currentSub); - log.info("单位{}减积分成功,本次减{},累计减{},剩余{}", - reqVO.getUnitId(), currentSub, afterSub, currentSub); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - throw exception0(INTERNAL_SERVER_ERROR.getCode(), "积分操作被中断,请重试"); - } finally { - if (lockAcquired && lock.isHeldByCurrentThread()) { - lock.unlock(); - } - } + // 8. 新增扣减明细 + // 新增减积分明细 + createSubDetail(reqVO, beforeSub, afterSub, currentSub); + log.info("单位{}减积分成功,本次减{},累计减{},剩余{}", + reqVO.getUnitId(), currentSub, afterSub, currentSub); } /** @@ -313,4 +309,14 @@ public class UnitIntegralServiceImpl implements UnitIntegralService{ log.info("新增积分扣减明细成功,单位ID:{},明细ID:{}", updateReqVO.getUnitId(), detailVO.getId()); } + + /** + * 获得自身的代理对象,解决 AOP 生效问题 + * + * @return 自己 + */ + private UnitIntegralServiceImpl getSelf() { + return SpringUtil.getBean(getClass()); + } + }