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; }
问题根源分析
- 锁校验与业务操作不同步:乐观锁的生效依赖事务提交时的版本号校验,但当前业务数据的读取、修改、上下文更新操作完全脱离锁约束。多个线程可同时获取相同版本的锁实体,读取到相同初始
filesIds列表,各自添加ID后提交,导致重复。 - 锁实体未被正确持久化:
updateVersionLock仅修改内存中实体的lastUpdated字段,未触发EntityManager的脏检查或显式保存,事务提交时不会自动递增@Version字段,乐观锁未真正生效。 - 初始锁创建存在竞态:
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
相关产品推荐
相关产品推荐

