know-engine

引入分布式锁解决文档处理的并发问题

在我们的课程中,有文档处理的过程,比如文档的转换、切分、嵌入和存储。因为我们为了解耦和提升吞吐量,我们做了基于事件驱动的方式做流程的串联。 因为有了事件驱动,而事件驱动有可能失败,所以我们引入了XXL JOB定时任务做兜底,那么在极…

TL;DR

在我们的课程中,有文档处理的过程,比如文档的转换、切分、嵌入和存储。因为我们为了解耦和提升吞吐量,我们做了基于事件驱动的方式做流程的串联。 因为有了事件驱动,而事件驱动有可能失败,所以我们引入了XXL JOB定时任务做兜底,那么在极…

在我们的课程中,有文档处理的过程,比如文档的转换、切分、嵌入和存储。因为我们为了解耦和提升吞吐量,我们做了基于事件驱动的方式做流程的串联。 因为有了事件驱动,而事件驱动有可能失败,所以我们引入了XXL-JOB定时任务做兜底,那么在极端情况下,就可能会出现定时任务和listener并发的情况。(这个方案,如果用MQ的话也一样,因为MQ也可能会重投) 所以,我们需要有一种方案,来解决这个并发的问题。那么,分布式锁是比较合适的。

Redis部署

## Redis 官方镜像支持 ARM64,直接运行即可(不加 --platform)
docker run -d --name redis \
  -p 6379:6379 \
  -v redis-data:/data \
  --restart always \
  redis redis-server --appendonly yes --requirepass "knowengine"


## 我的Mac 是 Apple Silicon 芯片(M1/M2/M3),架构是 ARM64,但 Docker 默认拉取的是 AMD64(x86) 架构的 Redis 镜像。所以要用下面这个命令:

docker run -d --name redis \
  -p 6379:6379 \
  -v redis-data:/data \
  --restart always \
  redis:alpine redis-server --appendonly yes --requirepass "knowengine"

分布式锁注解

我们的分布式锁的实现方案是基于Redisson来做的,用了他的分布式锁的lock和tryLock两个方案,同时为了方便实用,我们自定义了一个分布式锁的注解。(这块直接从我的数藏中复用的方案) 我们经常要在代码中使用分布式锁,有很多锁,是要在方法入口处加锁,结束后解锁的,于是为了方便,我们定义了一个通用的注解,来进行分布式锁的实现。

package cn.hollis.llm.mentor.know.engine.infra.lock;

import java.lang.annotation.ElementType;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
import java.lang.annotation.Target;

/**
 * 分布式锁注解
 *
 * @author Hollis
 */
@Target(ElementType.METHOD)
@Retention(RetentionPolicy.RUNTIME)
public @interface DistributeLock {

    /**
     * 锁的场景
     *
     * @return
     */
    public String scene();

    /**
     * 加锁的key,优先取key(),如果没有,则取keyExpression()
     *
     * @return
     */
    public String key() default DistributeLockConstant.NONE_KEY;

    /**
     * SPEL表达式:
     * <pre>
     *     #id
     *     #insertResult.id
     * </pre>
     *
     * @return
     */
    public String keyExpression() default DistributeLockConstant.NONE_KEY;

    /**
     * 超时时间,毫秒
     * 默认情况下不设置超时时间,会自动续期
     *
     * @return
     */
    public int expireTime() default DistributeLockConstant.DEFAULT_EXPIRE_TIME;

    /**
     * 加锁等待时长,毫秒
     * 默认情况下不设置等待时长,会一直等待直到获取到锁
     * @return
     */
    public int waitTime() default DistributeLockConstant.DEFAULT_WAIT_TIME;
}

其中定义了很多参数,其中 key 和keyExpression二选一,只能有一个,key 表示规定值当作锁的 key,keyExpression表示是一个SPEL 表达式。 使用方法如下:

@DistributeLock(keyExpression = "#payCreateRequest.bizNo", scene = "GENERATE_PAY_URL")
public PayCreateResponse generatePayUrl(PayCreateRequest payCreateRequest) {

}

在生成支付链接时,直接在方法上加上注解,实现自动加锁和解锁。 锁的实现逻辑如下:

package cn.hollis.llm.mentor.know.engine.infra.lock;

import org.aspectj.lang.ProceedingJoinPoint;
import org.aspectj.lang.annotation.Around;
import org.aspectj.lang.annotation.Aspect;
import org.aspectj.lang.reflect.MethodSignature;
import org.redisson.api.RLock;
import org.redisson.api.RedissonClient;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.core.StandardReflectionParameterNameDiscoverer;
import org.springframework.core.annotation.Order;
import org.springframework.expression.EvaluationContext;
import org.springframework.expression.Expression;
import org.springframework.expression.spel.standard.SpelExpressionParser;
import org.springframework.expression.spel.support.StandardEvaluationContext;
import org.springframework.stereotype.Component;

import java.lang.reflect.Method;
import java.util.concurrent.TimeUnit;

/**
 * 分布式锁切面
 *
 * @author hollis
 */
@Aspect
@Component
@Order(Integer.MIN_VALUE)
public class DistributeLockAspect {

    private RedissonClient redissonClient;

    public DistributeLockAspect(RedissonClient redissonClient) {
        this.redissonClient = redissonClient;
    }

    private static final Logger LOG = LoggerFactory.getLogger(DistributeLockAspect.class);

    @Around("@annotation(cn.hollis.llm.mentor.know.engine.infra.lock.DistributeLock)")
    public Object process(ProceedingJoinPoint pjp) throws Exception {
        Object response = null;
        Method method = ((MethodSignature) pjp.getSignature()).getMethod();
        DistributeLock distributeLock = method.getAnnotation(DistributeLock.class);

        String key = distributeLock.key();
        if (DistributeLockConstant.NONE_KEY.equals(key)) {
            if (DistributeLockConstant.NONE_KEY.equals(distributeLock.keyExpression())) {
                throw new DistributeLockException("no lock key found...");
            }
            SpelExpressionParser parser = new SpelExpressionParser();
            Expression expression = parser.parseExpression(distributeLock.keyExpression());

            EvaluationContext context = new StandardEvaluationContext();
            // 获取参数值
            Object[] args = pjp.getArgs();

            // 获取运行时参数的名称
            StandardReflectionParameterNameDiscoverer discoverer
                    = new StandardReflectionParameterNameDiscoverer();
            String[] parameterNames = discoverer.getParameterNames(method);

            // 将参数绑定到context中
            if (parameterNames != null) {
                for (int i = 0; i < parameterNames.length; i++) {
                    context.setVariable(parameterNames[i], args[i]);
                }
            }

            // 解析表达式,获取结果
            key = String.valueOf(expression.getValue(context));
        }

        String scene = distributeLock.scene();

        String lockKey = scene + "#" + key;

        int expireTime = distributeLock.expireTime();
        int waitTime = distributeLock.waitTime();
        RLock rLock= redissonClient.getLock(lockKey);
        try {
            boolean lockResult = false;
            if (waitTime == DistributeLockConstant.DEFAULT_WAIT_TIME) {
                if (expireTime == DistributeLockConstant.DEFAULT_EXPIRE_TIME) {
                    LOG.info(String.format("lock for key : %s", lockKey));
                    rLock.lock();
                } else {
                    LOG.info(String.format("lock for key : %s , expire : %s", lockKey, expireTime));
                    rLock.lock(expireTime, TimeUnit.MILLISECONDS);
                }
                lockResult = true;
            } else {
                if (expireTime == DistributeLockConstant.DEFAULT_EXPIRE_TIME) {
                    LOG.info(String.format("try lock for key : %s , wait : %s", lockKey, waitTime));
                    lockResult = rLock.tryLock(waitTime, TimeUnit.MILLISECONDS);
                } else {
                    LOG.info(String.format("try lock for key : %s , expire : %s , wait : %s", lockKey, expireTime, waitTime));
                    lockResult = rLock.tryLock(waitTime, expireTime, TimeUnit.MILLISECONDS);
                }
            }

            if (!lockResult) {
                LOG.warn(String.format("lock failed for key : %s , expire : %s", lockKey, expireTime));
                throw new DistributeLockException("acquire lock failed... key : " + lockKey);
            }


            LOG.info(String.format("lock success for key : %s , expire : %s", lockKey, expireTime));
            response = pjp.proceed();
        } catch (Throwable e) {
            throw new Exception(e);
        } finally {
            if (rLock.isHeldByCurrentThread()) {
                rLock.unlock();
                LOG.info(String.format("unlock for key : %s , expire : %s", lockKey, expireTime));
            }
        }
        return response;
    }
}

代码中给大家加了注释,就不展开介绍了。 后面还做了个优化,在这个切面上加了个@Order(Integer.MIN_VALUE),让它可以最早执行,避免并发时事务未提交但是锁释放导致的问题。

加锁

然后,我们在我们的DocumentProcessServiceImpl中增加分布式锁如下:

@DistributeLock(scene = "document-upload", keyExpression = "#uploadUser", waitTime = 0)
public KnowledgeDocument upload(MultipartFile file, String uploadUser, String accessibleBy) throws IOException {}


@Transactional
@DistributeLock(scene = "document-split", keyExpression = "#document.docId", waitTime = 0)
public int split(KnowledgeDocument document) {}

@Override
@DistributeLock(scene = "document-split", keyExpression = "#document.docId", waitTime = 0)
public boolean embedAndStore(KnowledgeDocument document) {}

upload方法针对用户加锁,限制同一个用户同时上传多份文件。 split和embedAndStore方法针对documentId加锁,即同一份文档,不能同时执行。 这里设置waitTime = 0,即尝试获取锁失败则返回,不阻塞。

版本提示

模型、框架与接口会持续变化。涉及版本号、参数与生产配置时,请在实践前对照对应官方文档。

LLMentor系统化学习大模型应用工程

内容来自个人课程知识库备份,并经过结构化整理。技术版本持续演进,生产使用前请结合官方文档验证。