R2DBC中多条件插入新实体触发重复条目错误排查
问题分析与解决方案
你的条件插入逻辑看起来没问题,但实际执行重复实体时触发唯一键错误,核心原因是并发竞态条件加上Reactive场景下事务的特殊行为,我来详细拆解:
为什么会出现这个错误?
- flatMap的并发特性:
flatMap默认是并发处理Flux中的元素(默认并发数256),当你传入两个dbEntity1时,这两个请求会同时执行conditionalInsertEntity方法。 - 事务的独立性:
@Transactional在Reactive场景下,每个Mono<DBEntity>的订阅都会开启一个独立的事务。这意味着两个dbEntity1的请求会各自执行fetchObjectByPublicId,此时两者都查不到数据,接着同时进入插入分支,最终触发数据库的唯一键重复错误。 - 声明式事务无法解决跨事务的竞态:因为两个请求的事务是完全独立的,第一个事务的插入操作在提交前,第二个事务的查询是看不到未提交的数据(默认隔离级别下),所以第二个请求依然会执行插入。
解决方案
方案1:改用concatMap串行处理(最简单)
把flatMap换成concatMap,这样Flux中的元素会串行执行,前一个请求处理完成(事务提交)后,再处理下一个。这样第二个dbEntity1执行fetchObjectByPublicId时,第一个已经插入完成,就会直接返回已存在的实体,不会重复插入。
val result = Flux.fromIterable(listOf(dbEntity1, dbEntity1, dbEntity2)) .concatMap { conditionalInsertEntity(it) } // 串行处理,避免并发竞态 .collectList() .block()
- 优点:代码改动极小,逻辑简单易懂
- 缺点:牺牲了并发性能,适合数据量不大、并发要求不高的场景
方案2:数据库层面UPSERT(最可靠,推荐)
不管业务层怎么控制,数据库的唯一约束是最可靠的防线。首先给publicId字段添加唯一索引,然后把插入逻辑改成UPSERT(插入或更新),这样即使并发插入,数据库也会自动处理冲突,不会报错。
不同数据库的UPSERT语法略有不同,这里以MySQL和PostgreSQL为例:
MySQL 示例
// 修改conditionalInsertEntity中的插入分支 r2DatabaseClient.sql(""" INSERT INTO db_entity (public_id, col1, col2) VALUES (:publicId, :col1, :col2) ON DUPLICATE KEY UPDATE col1 = VALUES(col1), col2 = VALUES(col2) """) .bind("publicId", dbEntity.publicId) .bind("col1", dbEntity.col1) .bind("col2", dbEntity.col2) .fetch() .rowsUpdated() .flatMap { fetchObjectByPublicId(dbEntity.publicId) } // 插入/更新后返回最新实体
PostgreSQL 示例
r2DatabaseClient.sql(""" INSERT INTO db_entity (public_id, col1, col2) VALUES (:publicId, :col1, :col2) ON CONFLICT (public_id) DO UPDATE SET col1 = EXCLUDED.col1, col2 = EXCLUDED.col2 """) .bind("publicId", dbEntity.publicId) .bind("col1", dbEntity.col1) .bind("col2", dbEntity.col2) .fetch() .rowsUpdated() .flatMap { fetchObjectByPublicId(dbEntity.publicId) }
- 优点:从根源解决并发问题,不影响业务层并发性能,数据库层面保证数据一致性
- 缺点:需要依赖数据库特性,部分场景可能需要写原生SQL(如果R2DBC客户端没有封装UPSERT方法)
方案3:分布式锁控制并发(业务层兜底)
如果无法修改数据库,或者需要在业务层做更精细的控制,可以给每个publicId添加分布式锁,确保同一个publicId的请求只能串行执行。比如用Redis的Redisson实现:
val result = Flux.fromIterable(listOf(dbEntity1, dbEntity1, dbEntity2)) .flatMap { entity -> // 针对每个publicId创建独立锁 val lock = redissonClient.getLock("entity_insert_lock_${entity.publicId}") lock.lockAsync() // 异步获取锁 .then(conditionalInsertEntity(entity)) // 获取锁后执行插入逻辑 .doFinally { // 无论成功失败,都释放锁 if (lock.isHeldByCurrentThread) { lock.unlockAsync() } } } .collectList() .block()
- 优点:不依赖数据库特性,业务层灵活控制
- 缺点:引入分布式锁组件,增加系统复杂度,需要处理锁超时、死锁等问题
关于Reactive事务的额外提醒
在Reactive场景下,@Transactional的事务上下文是绑定在每个Mono/Flux的订阅上的,每个订阅都是一个独立的事务。所以这种跨请求的并发竞态条件,是无法通过声明式事务解决的,必须通过上面的并发控制或数据库约束来处理。
内容的提问来源于stack exchange,提问作者Akash Gupta
相关产品推荐
相关产品推荐

