You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Hibernate乐观锁失效求助:并发场景下数据重复丢失

Hibernate乐观锁失效导致竞态条件问题排查与修复

问题现象

使用Hibernate乐观锁时,数据库版本号正常递增、异常场景能抛出OptimisticLockException,且已实现重试机制,但并发请求下仍出现ID重复、文件丢失的竞态问题。并发5个请求测试后,日志中出现重复ID(如[a5cc34ca-c276-419b-be90-b6e2e484c68c, 11159010-a303-4b59-8df4-81bf3cc0e1e8, 763bca97-d3ea-4786-ad8b-0880015792e6, 11159010-a303-4b59-8df4-81bf3cc0e1e8, 5c21bc3d-d278-47b9-b06c-693cbaede836])。

相关代码实现

重试与错误处理方法

public void updateIdsVariableInProcess(String processId, String processInstanceId, String objectKey, boolean insertFlag) {
    log.info("Adding ID variable to process. processId={}, processInstanceId={}, objectKey={}", new Object[]{processId, processInstanceId, objectKey});
    boolean lockAndVarsInContextUpdated = false;

    while(!lockAndVarsInContextUpdated) {
        try {
            this.processLockService.updateVariableWithLock2(processId, processInstanceId, objectKey, insertFlag);
            lockAndVarsInContextUpdated = true;
        } catch (ArcUndeclaredThrowableException | PersistenceException var8) {
            RuntimeException expected = var8;
            log.warn("Error while setting lock and updating process variable. processId={}, processInstanceId={}, objectKey={}, error={}", new Object[]{processId, processInstanceId, objectKey, expected.getMessage(), expected});
        } catch (Exception var9) {
            Exception unexpected = var9;
            log.error("Unexpected error while updating variable. processId={}, processInstanceId={}, cause={}", new Object[]{processId, processInstanceId, unexpected.getMessage(), unexpected});
            throw new RuntimeException("Unexpected error while updating variable", unexpected);
        }

        if (!lockAndVarsInContextUpdated) {
            try {
                Thread.sleep(50L);
            } catch (InterruptedException var7) {
                InterruptedException e = var7;
                Thread.currentThread().interrupt();
                throw new RuntimeException("Thread interrupted while waiting to retry", e);
            }
        }
    }

    log.info("SyncInteractionProcessService.addIdsVariableInProcess - Retrying ProcessId: {}, ProcessInstanceId: {}, ObjectKey: {}", new Object[]{processId, processInstanceId, objectKey});
}

核心事务方法

@Transactional
public Boolean updateVariableWithLock2(String processId, String processInstanceId, String fileId, boolean insertFlag) {
    if (!this.processLockRepository.existsByProcessIdAndProcessInstanceId(processId, processInstanceId)) {
        this.createNewLock(processId, processInstanceId);
    }

    int version;
    try {
        ProcessLockEntity processLockEntity = this.processLockRepository.getByProcessIdAndProcessInstanceId(processId, processInstanceId);
        version = processLockEntity.getVersion();
        log.info("Existing lock found. Updating processLockEntity. ProcessId: {}, ProcessInstanceId: {}, Version: {}", new Object[]{processId, processInstanceId, version});
        Map<String, Object> variables = this.processInteractionService.getVariables(processId, processInstanceId);
        S3ObjectIds s3ObjectIds = (S3ObjectIds)variables.get("s3ObjectIds");
        if (s3ObjectIds == null) {
            s3ObjectIds = new S3ObjectIds();
        }

        if (s3ObjectIds.getIdS() == null) {
            s3ObjectIds.setIdS(new ArrayList());
        }

        List<String> filesIds = new ArrayList(s3ObjectIds.getIdS());
        if (insertFlag) {
            filesIds.add(fileId);
        } else {
            filesIds.remove(fileId);
        }

        s3ObjectIds.setIdS(filesIds);
        // WE ARE UPDATING CONTEXT OF PROCESS HERE, WE DO NOT WRITE TO DB HERE, ONLY CONTEXT HERE
        this.patchIdsVariableInProcess(processId, processInstanceId, s3ObjectIds);
        this.updateVersionLock(processLockEntity);
    } catch (Exception var10) {
        Exception e = var10;
        log.error("Failed to update variable with lock", e);
        throw e;
    }

    log.info("Finished lock update. ProcessId: {}, ProcessInstanceId: {}, Version: {}", new Object[]{processId, processInstanceId, version});
    return true;
}

获取锁实体方法

public ProcessLockEntity getByProcessIdAndProcessInstanceId(String processId, String processInstanceId) {
    CriteriaBuilder cb = this.entityManager.getCriteriaBuilder();
    CriteriaQuery<ProcessLockEntity> query = cb.createQuery(ProcessLockEntity.class);
    Root<ProcessLockEntity> root = query.from(ProcessLockEntity.class);
    query.select(root).where(cb.and(cb.equal(root.get("processId"), processId), cb.equal(root.get("processInstanceId"), processInstanceId)));
    TypedQuery<ProcessLockEntity> typedQuery = this.entityManager.createQuery(query).setLockMode(LockModeType.OPTIMISTIC);
    ProcessLockEntity processLockEntity = (ProcessLockEntity)typedQuery.getSingleResult();
    log.info("Retrieved ProcessLockEntity. ProcessId: {}, ProcessInstanceId: {}, Entity: {}", new Object[]{processId, processInstanceId, processLockEntity});
    return processLockEntity;
}

更新锁版本方法

private void updateVersionLock(ProcessLockEntity version) {
    this.processLockRepository.updateProcessInstanceVersion(version);
}
public void updateProcessInstanceVersion(ProcessLockEntity version) {
    log.info("ProcessLockCriteriaRepository class, updateProcessInstanceVersion: processId={}, processInstanceId={}, version={}", new Object[]{version.getProcessId(), version.getProcessInstanceId(), version.getVersion()});
    version.setLastUpdated(Instant.now());
}

实体类定义

@Entity
@Table(name = "process_lock", schema = "public", uniqueConstraints = {
        @UniqueConstraint(columnNames = {"process_id", "process_instance_id"})
})
@Getter
@Setter
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class ProcessLockEntity {

    /**
     * Unique identifier for the context version.
     */
    @Id
    @GeneratedValue
    private UUID id;

    /**
     * Identifier of the process. (for example ServiceDeskRequestRegistration)
     */
    @Column(name = "process_id")
    private String processId;

    /**
     * Identifier of the process instance.
     */
    @Column(name = "process_instance_id")
    private String processInstanceId;

    /**
     * TIMESTAMP of the last update.
     */
    @Column(name = "last_updated")
    private Instant lastUpdated;

    /**
     * Version number of the context.
     */
    @Version
    private Integer version;

}

问题根源分析

  1. 锁校验与业务操作不同步:乐观锁的生效依赖事务提交时的版本号校验,但当前业务数据的读取、修改、上下文更新操作完全脱离锁约束。多个线程可同时获取相同版本的锁实体,读取到相同初始filesIds列表,各自添加ID后提交,导致重复。
  2. 锁实体未被正确持久化:updateVersionLock仅修改内存中实体的lastUpdated字段,未触发EntityManager的脏检查或显式保存,事务提交时不会自动递增@Version字段,乐观锁未真正生效。
  3. 初始锁创建存在竞态:existsByProcessIdAndProcessInstanceId与createNewLock为两步无锁操作,并发场景下可能多个线程重复创建锁实体,导致后续流程读取不一致的锁数据。

修复方案

1. 绑定锁与业务操作的事务生命周期

将业务数据读写、修改与锁实体的生命周期绑定在同一事务中,确保锁校验覆盖完整业务逻辑:

@Transactional
public Boolean updateVariableWithLock2(String processId, String processInstanceId, String fileId, boolean insertFlag) {
    // 原子化创建或获取锁实体,避免竞态
    ProcessLockEntity processLockEntity = processLockRepository.findByProcessIdAndProcessInstanceId(processId, processInstanceId)
            .orElseGet(() -> {
                ProcessLockEntity newLock = ProcessLockEntity.builder()
                        .processId(processId)
                        .processInstanceId(processInstanceId)
                        .lastUpdated(Instant.now())
                        .build();
                return processLockRepository.save(newLock);
            });

    // 强制触发乐观锁校验
    entityManager.lock(processLockEntity, LockModeType.OPTIMISTIC_FORCE_INCREMENT);

    log.info("Existing lock found. Updating processLockEntity. ProcessId: {}, ProcessInstanceId: {}, Version: {}", 
            processId, processInstanceId, processLockEntity.getVersion());

    // 读取并修改业务数据
    Map<String, Object> variables = this.processInteractionService.getVariables(processId, processInstanceId);
    S3ObjectIds s3ObjectIds = (S3ObjectIds) variables.get("s3ObjectIds");
    if (s3ObjectIds == null) {
        s3ObjectIds = new S3ObjectIds();
    }
    if (s3ObjectIds.getIdS() == null) {
        s3ObjectIds.setIdS(new ArrayList<>());
    }

    List<String> filesIds = new ArrayList<>(s3ObjectIds.getIdS());
    if (insertFlag) {
        // 业务层去重兜底
        if (!filesIds.contains(fileId)) {
            filesIds.add(fileId);
        }
    } else {
        filesIds.remove(fileId);
    }
    s3ObjectIds.setIdS(filesIds);

    // 更新流程上下文
    this.patchIdsVariableInProcess(processId, processInstanceId, s3ObjectIds);

    // 修改锁实体触发脏检查,版本号自动递增
    processLockEntity.setLastUpdated(Instant.now());
    processLockRepository.save(processLockEntity);

    log.info("Finished lock update. ProcessId: {}, ProcessInstanceId: {}, Version: {}", 
            processId, processInstanceId, processLockEntity.getVersion());
    return true;
}

2. 修复锁实体查询与更新逻辑

修改锁实体查询方法,强制递增版本号确保锁生效:

public ProcessLockEntity getByProcessIdAndProcessInstanceId(String processId, String processInstanceId) {
    CriteriaBuilder cb = this.entityManager.getCriteriaBuilder();
    CriteriaQuery<ProcessLockEntity> query = cb.createQuery(ProcessLockEntity.class);
    Root<ProcessLockEntity> root = query.from(ProcessLockEntity.class);
    query.select(root).where(cb.and(
            cb.equal(root.get("processId"), processId),
            cb.equal(root.get("processInstanceId"), processInstanceId)
    ));
    TypedQuery<ProcessLockEntity> typedQuery = this.entityManager.createQuery(query)
            .setLockMode(LockModeType.OPTIMISTIC_FORCE_INCREMENT);
    return typedQuery.getSingleResult();
}

修改锁实体更新方法,确保持久化触发版本递增:

public void updateProcessInstanceVersion(ProcessLockEntity version) {
    log.info("ProcessLockCriteriaRepository class, updateProcessInstanceVersion: processId={}, processInstanceId={}, version={}", 
            version.getProcessId(), version.getProcessInstanceId(), version.getVersion());
    version.setLastUpdated(Instant.now());
    this.entityManager.merge(version);
}

3. 优化重试机制

仅针对乐观锁异常重试,添加次数限制与指数退避:

public void updateIdsVariableInProcess(String processId, String processInstanceId, String objectKey, boolean insertFlag) {
    log.info("Adding ID variable to process. processId={}, processInstanceId={}, objectKey={}", processId, processInstanceId, objectKey);
    boolean lockAndVarsInContextUpdated = false;
    int retryCount = 0;
    final int MAX_RETRIES = 5;

    while(!lockAndVarsInContextUpdated && retryCount < MAX_RETRIES) {
        try {
            this.processLockService.updateVariableWithLock2(processId, processInstanceId, objectKey, insertFlag);
            lockAndVarsInContextUpdated = true;
        } catch (OptimisticLockException e) {
            retryCount++;
            log.warn("Optimistic lock conflict, retrying. Attempt {} of {}, processId={}, processInstanceId={}", 
                    retryCount, MAX_RETRIES, processId, processInstanceId);
        } catch (ArcUndeclaredThrowableException | PersistenceException var8) {
            RuntimeException expected = var8;
            log.warn("Error while setting lock and updating process variable. processId={}, processInstanceId={}, objectKey={}, error={}", 
                    processId, processInstanceId, objectKey, expected.getMessage(), expected);
            throw expected;
        } catch (Exception var9) {
            Exception unexpected = var9;
            log.error("Unexpected error while updating variable. processId={}, processInstanceId={}, cause={}", 
                    processId, processInstanceId, unexpected.getMessage(), unexpected);
            throw new RuntimeException("Unexpected error while updating variable", unexpected);
        }

        if (!lockAndVarsInContextUpdated) {
            try {
                Thread.sleep(50L * retryCount);
            } catch (InterruptedException var7) {
                Thread.currentThread().interrupt();
                throw new RuntimeException("Thread interrupted while waiting to retry", var7);
            }
        }
    }

    if (!lockAndVarsInContextUpdated) {
        throw new RuntimeException("Max retries reached for updating variable, processId=" + processId + ", processInstanceId=" + processInstanceId);
    }

    log.info("Successfully updated ID variable. ProcessId: {}, ProcessInstanceId: {}, ObjectKey: {}", processId, processInstanceId, objectKey);
}

内容的提问来源于stack exchange,提问作者Arseniy Porfirev

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.13 06:09:51