使用R2DBC PostgreSQL与Webflux锁行时并发实体顺序异常问题
核心原因
Flow未被实际订阅执行
R2DBC的Flow是冷流,只有调用toList()、collect()这类终端操作时,才会触发实际的SQL执行。你的代码中existingEntities.size()如果没有先将Flow转换为集合(比如existingEntities.toList().size),这条加锁查询根本不会执行,自然不会对分区内的行加锁,导致第二个事务直接读到初始的1条数据。事务内数据库操作的绑定问题
R2DBC的事务基于Connection上下文,只有数据库操作被实际执行(即Flow被收集)时,才会绑定到当前事务。如果查询没执行,锁逻辑完全不生效,两个事务各自读取旧数据,生成相同的order值。FOR UPDATE锁的特性(次要)
PostgreSQL的FOR UPDATE只会锁定查询返回的现有行,虽然你预创建了1条数据,但如果查询没执行,锁就无从谈起。若分区下无数据时,FOR UPDATE不会创建锁,这种场景需要额外处理(比如插入占位行或用SELECT ... FOR UPDATE ON CONFLICT),但你的场景不属于这种情况。
修复步骤
1. 确保加锁查询被执行
修改事务内代码,将Flow转换为集合,触发查询并加锁:
operator.executeAndAwait { // 用toList()触发Flow执行,获取加锁后的实体列表 val existingEntities = repository.findByZoneWithLock(id).toList() val entity = TestEntity().apply { order = existingEntities.size + 1 } repository.save(entity) }
2. 优化加锁查询(可选但推荐)
查询所有行统计数量效率很低,改用SELECT COUNT(*) FOR UPDATE直接获取加锁后的计数:
interface EntityRepository<TestEntity> : CoroutineCrudRepository<TestEntity, String> { @Query("SELECT COUNT(*) FROM TestEntity WHERE zone = :zone FOR UPDATE") suspend fun countByZoneWithLock(zone: String): Long }
事务内代码改为:
operator.executeAndAwait { val count = repository.countByZoneWithLock(id) val entity = TestEntity().apply { order = count + 1 } repository.save(entity) }
3. 验证事务上下文绑定
确保所有数据库操作都在executeAndAwait的lambda内执行,且suspend函数调用不会脱离事务上下文。R2DBC的事务会自动绑定到当前协程上下文,只要不在lambda内启动新的非协程作用域(如launch(Dispatchers.IO)),就不会出现事务绑定问题。
其他排查点
- 确认
transactionManager是R2DBC的R2dbcTransactionManager,而非JDBC事务管理器; - 检查是否存在自定义缓存逻辑,Spring Data R2DBC默认不缓存查询结果,但如果有自定义缓存,需禁用该查询的缓存。
内容的提问来源于stack exchange,提问作者Михаил

