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

Postgres乐观锁实现问题:单Executor超量获取Job

问题排查:Executor超额获取Job的原因及修复

需求与问题现象

  • 核心需求:
    • 每个Job必须分配给唯一服务实例
    • 单个服务实例同时运行的Job不得超过2个
  • 问题:测试中单个Executor最多获取了10个Job,远超预期上限

预期流程

  1. Job以jobStatus='ACCEPTED'、executorId=null插入数据库
  2. Executor轮询数据库,提交查询请求
  3. 查询跳过已锁定行,避免实例间冲突
  4. 查询锁定第一条可用记录(按createdAt倒序取1条)
  5. JDBC客户端设置executorId并更新记录
  6. 提交事务
  7. 若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
}

问题根源分析

  1. 事务隔离导致的计数不一致
    默认数据库隔离级别为读已提交,当同一个Executor的多个并发请求执行查询时,每个请求的子查询无法看到其他请求已经锁定并更新但未提交的Job记录。因此每个请求都认为当前实例的Job数量≤2,从而持续获取新Job,最终超出上限。

  2. Job状态未及时更新
    业务代码仅更新了executorId,未将jobStatus从ACCEPTED改为PROCESSING。即使事务提交后,子查询仍会将已分配的Job统计为ACCEPTED状态,进一步干扰计数逻辑。

  3. 子查询计数逻辑的边界错误
    子查询使用<=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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 10:36:00