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

使用R2DBC PostgreSQL与Webflux锁行时并发实体顺序异常问题

问题原因及解决方案

核心原因

  1. Flow未被实际订阅执行
    R2DBC的Flow是冷流,只有调用toList()、collect()这类终端操作时,才会触发实际的SQL执行。你的代码中existingEntities.size()如果没有先将Flow转换为集合(比如existingEntities.toList().size),这条加锁查询根本不会执行,自然不会对分区内的行加锁,导致第二个事务直接读到初始的1条数据。

  2. 事务内数据库操作的绑定问题
    R2DBC的事务基于Connection上下文,只有数据库操作被实际执行(即Flow被收集)时,才会绑定到当前事务。如果查询没执行,锁逻辑完全不生效,两个事务各自读取旧数据,生成相同的order值。

  3. 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,提问作者Михаил

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 12:25:09