Postgres乐观锁实现问题:单Executor超量获取Job
问题排查:Executor超额获取Job的原因及修复
需求与问题现象
- 核心需求:
- 每个Job必须分配给唯一服务实例
- 单个服务实例同时运行的Job不得超过2个
- 问题:测试中单个Executor最多获取了10个Job,远超预期上限
预期流程
- Job以
jobStatus='ACCEPTED'、executorId=null插入数据库 - Executor轮询数据库,提交查询请求
- 查询跳过已锁定行,避免实例间冲突
- 查询锁定第一条可用记录(按
createdAt倒序取1条) - JDBC客户端设置
executorId并更新记录 - 提交事务
- 若Executor已分配2个Job,查询应无返回结果
相关代码片段
SQL查询
select id -- some fields omitted for brevity from JOBS where ( jobStatus = 'ACCEPTED' -- job waits for execution, and executorId is null -- none of service instances got a job and ( select count(*) from JOBS where ( jobStatus in ( 'ACCEPTED', -- service took the job for execution 'PROCESSING' -- service executes job ) and executorId = 'service_instance_id' ) ) <= 2 -- each service instance should run not more than 2 jobs at a time ) order by createdAt desc limit 1 for update skip locked -- to avoid locks
Scala业务代码
val selectForUpdateQuery = generateQuery() // example above val ps = conn.prepareStatement(selectForUpdateQuery, ResultSet.TYPE_FORWARD_ONLY, ResultSet.CONCUR_UPDATABLE) val rs = ps.executeQuery() if (rs.next()) { val now = OffsetDateTimeFactory.nowOffsetDateTime() val jobRecord = resultSetRowToJobRecord(rs).copy(executorId = Some(executorId), updatedAt = now) val sqlTimestamp = java.sql.Timestamp.valueOf(now.atZoneSameInstant(ZoneOffset.UTC).toLocalDateTime()) rs.updateString("executorId", executorId) rs.updateObject("updatedAt", sqlTimestamp) rs.updateRow() logger.info(s"job ${rs.getLong("id")} was acquired by executorId: $executorId at: $now") Some(jobRecord) } else { None }
测试代码
val jobsWithExecutor = withClue(s"executor should acquire not more than 2 jobs") { Future .sequence( (1 to 100) .map { _ => jobService.acquireExecutorLockAtJob(executorId, concurrencyLimit = 2).map { case None => None case aJob @ Some(job) => logger.info(s"job: $job acquired by executorId: $executorId") aJob } } ) .futureValue .flatten }
问题根源分析
事务隔离导致的计数不一致
默认数据库隔离级别为读已提交,当同一个Executor的多个并发请求执行查询时,每个请求的子查询无法看到其他请求已经锁定并更新但未提交的Job记录。因此每个请求都认为当前实例的Job数量≤2,从而持续获取新Job,最终超出上限。Job状态未及时更新
业务代码仅更新了executorId,未将jobStatus从ACCEPTED改为PROCESSING。即使事务提交后,子查询仍会将已分配的Job统计为ACCEPTED状态,进一步干扰计数逻辑。子查询计数逻辑的边界错误
子查询使用<=2作为判断条件,意味着当实例已有2个Job时,仍允许获取第3个。正确的边界应该是<2,确保获取后总数不超过2。
修复方案
1. 修正SQL查询逻辑
通过锁定Executor已持有的Job记录,确保并发请求能正确感知当前Job数量,同时调整边界条件:
WITH executor_current_jobs AS ( SELECT COUNT(*) AS job_count FROM JOBS WHERE executorId = ? AND jobStatus IN ('ACCEPTED', 'PROCESSING') FOR UPDATE -- 锁定已分配的Job,防止并发计数不一致 ) SELECT j.id FROM JOBS j CROSS JOIN executor_current_jobs ecj WHERE j.jobStatus = 'ACCEPTED' AND j.executorId IS NULL AND ecj.job_count < 2 -- 确保获取后总数不超过2 ORDER BY j.createdAt DESC LIMIT 1 FOR UPDATE SKIP LOCKED;
2. 完善业务代码的状态更新
获取Job后,同步更新jobStatus为PROCESSING:
// 在rs.updateRow()前添加状态更新 rs.updateString("jobStatus", "PROCESSING") rs.updateString("executorId", executorId) rs.updateObject("updatedAt", sqlTimestamp) rs.updateRow()
3. 确保事务范围正确
将查询、更新操作包裹在同一个事务中,避免部分操作提交导致的状态不一致:
conn.setAutoCommit(false) try { // 执行查询、更新逻辑 conn.commit() } catch { case e: Exception => conn.rollback() throw e } finally { conn.setAutoCommit(true) }
内容的提问来源于stack exchange,提问作者Capacytron
相关产品推荐
相关产品推荐

