Spring Boot Kafka消费者约束违例后如何处理并规避会话异常
问题背景
我们的Java Spring Boot服务部署了多个跨集群的Kafka消费者(同消费组ID),负责将数据写入SQL Server数据库。代码逻辑是写入前检查ItemSet表是否存在对应记录,实际payload数据存在子表ItemValue中,表关系为ItemSet -> ItemName -> ItemValue的一对多层级结构,ItemSet表针对department id与season组合设置了唯一约束以避免重复添加。
现在需要在捕获约束违例异常后,将传入数据关联到已存在的ItemSet,但使用Spring Data JPA时,捕获异常后尝试查询现有记录会抛出:
org.hibernate.AssertionFailure: null id in ItemSet entry (don't flush the Session after an exception occurs).
catch块中的getItemSet()方法执行失败,同时还要解决多实例下的竞态条件问题。
原代码
ItemSet savedItemSet = null; try { String seasonName = itemSet.getSeasonName(); Long seasonYear = itemSet.getSeasonYear(); Long departmentId = itemSet.getDepartment().getId(); List<ItemSet> itemSets = attributeSetRepository.findBySeasonNameAndSeasonYearAndDepartmentId( seasonName, seasonYear, departmentId); LOGGER.info("Found {} item sets corresponding to season name : {}, season year : {}, " + "department id : {}", itemSets.size(), seasonName, seasonYear, departmentId); if(CollectionUtils.isEmpty(itemSets)) { savedItemSet = itemSetRepository.save(itemSet); } else { return new CreatedItemSet(itemSets.get(0).getId()); } } catch(PersistenceException | DataIntegrityViolationException e) { LOGGER.error("An exception occurred while saving itemSet set", e); if (e.getCause() instanceof ConstraintViolationException) { String seasonName = itemSet.getSeasonName(); Long seasonYear = itemSet.getSeasonYear(); Long deptId = itemSet.getDepartment().getId(); LOGGER.info("A duplicate item set found in the database corresponding " + "to season name : {}, season year : {} and department : {}", seasonName, seasonYear, deptId); ExistingItemSet existingItemSet = getItemSet(seasonName, seasonYear, deptId); if(existingItemSet == null) { LOGGER.info("No item set found"); return null; } return new CreatedItemSet(existingItemSet.getId()); } }
核心原因
- Session失效问题:Hibernate抛出
ConstraintViolationException这类持久化异常后,当前的Session会被标记为无效——异常已经破坏了Session的一致性状态,此时任何对该Session的操作(包括查询)都会触发断言失败。 - 竞态条件问题:原代码的"先查后存"逻辑在多实例并发场景下必然存在漏洞,查和存不是原子操作,两个实例可能同时查到空记录,然后同时执行保存,最终触发约束违例。
具体解决步骤
1. 异常处理时用新事务隔离失效Session
在catch块中查询现有记录时,不能复用当前已失效的Session,需要开启新事务执行查询:
给getItemSet()方法添加事务传播属性,强制它在新事务中运行,Spring会自动创建新的Session:
@Transactional(propagation = Propagation.REQUIRES_NEW) public ExistingItemSet getItemSet(String seasonName, Long seasonYear, Long deptId) { List<ItemSet> itemSets = attributeSetRepository.findBySeasonNameAndSeasonYearAndDepartmentId(seasonName, seasonYear, deptId); return CollectionUtils.isEmpty(itemSets) ? null : convertToExistingItemSet(itemSets.get(0)); }
2. 替换"先查后存"为原子操作,从根源解决竞态
原逻辑的竞态无法通过异常处理彻底解决,最好让数据库做原子性的"查存合并":
方式一:SQL Server原生MERGE语句(推荐)
通过自定义SQL实现原子的插入/查询操作,完全避免并发竞态:
在ItemSetRepository中定义自定义方法:
@Modifying @Query(value = "MERGE INTO ItemSet AS target " + "USING (VALUES (:seasonName, :seasonYear, :deptId)) AS source(seasonName, seasonYear, departmentId) " + "ON target.seasonName = source.seasonName AND target.seasonYear = source.seasonYear AND target.departmentId = source.departmentId " + "WHEN NOT MATCHED THEN " + "INSERT (seasonName, seasonYear, departmentId) VALUES (source.seasonName, source.seasonYear, source.departmentId) " + "OUTPUT inserted.id;", nativeQuery = true) Long insertOrGetItemSetId(@Param("seasonName") String seasonName, @Param("seasonYear") Long seasonYear, @Param("deptId") Long deptId);
调用该方法时,数据库会原子性完成"检查-插入-返回ID"操作,无需处理约束违例异常。
方式二:JPAmerge()方法
构造包含唯一约束字段的实体,通过merge()自动处理查询或插入:
// 构造仅含唯一约束字段的探针实体 ItemSet probe = new ItemSet(); probe.setSeasonName(itemSet.getSeasonName()); probe.setSeasonYear(itemSet.getSeasonYear()); probe.setDepartment(itemSet.getDepartment()); // merge会自动检查数据库中是否存在符合约束的记录,存在则返回现有实例,不存在则插入 ItemSet mergedItemSet = itemSetRepository.merge(probe); // 后续可将子表数据关联到mergedItemSet
3. 手动清理Session缓存(不推荐)
如果不想用新事务,也可以在catch块中手动清除Session缓存,但这种方式可靠性较低:
@Autowired private EntityManager entityManager; // 在catch块中执行 entityManager.clear(); // 清除Session缓存,尝试恢复可用状态 ExistingItemSet existingItemSet = getItemSet(seasonName, seasonYear, deptId);
最终推荐方案
优先选择SQL Server原生MERGE语句,它从根源解决竞态问题,性能和可靠性最高;如果无法使用自定义SQL,再选择merge()方法配合事务传播属性解决Session失效问题。
内容的提问来源于stack exchange,提问作者linuxNoob

