feat(worker): 认领加同项目串行(NOT EXISTS 无同 tenant 活跃任务),避免多 worker 并发写同一 workspace

This commit is contained in:
ymq 2026-08-18 17:33:17 +08:00
parent 1eba1c3f87
commit 8c6e49992c

View File

@ -525,16 +525,26 @@ async def _claim_task(sor, tenant_id, role, state='submitted', match_role=True,
from appPublic.uniqueID import getID
claim_token = getID()
# claimed_by IS NULL 保证原子认领(并发 poller / start_agents 不会双重认领);
# NOT EXISTS 保证同项目(tenant_id)串行:同一项目同一时刻只允许一个任务处于
# 活跃状态running / review·已认领 / qc_review·已认领避免多 worker 并发写
# 同一 workspace。跨项目完全并发。derived table 包裹规避 MySQL 1093同表子查询
# updated_at=NOW() 作为心跳,供 stale 回收判断。
await sor.sqlExe(
"UPDATE pipeline_tasks SET state=${setstate}$, claimed_by=${cb}$, updated_at=NOW() "
"WHERE id=${tid}$ AND state=${state}$ AND claimed_by IS NULL",
{"setstate": set_state, "cb": claim_token, "tid": task_id, "state": state})
"WHERE id=${tid}$ AND state=${state}$ AND claimed_by IS NULL "
"AND NOT EXISTS ("
" SELECT 1 FROM (SELECT 1 FROM pipeline_tasks t2 "
" WHERE t2.tenant_id=${tidx}$ AND t2.claimed_by IS NOT NULL "
" AND t2.state IN ('running','review','qc_review')) AS z"
")",
{"setstate": set_state, "cb": claim_token, "tid": task_id,
"state": state, "tidx": tenant_id})
check = await sor.sqlExe(
"SELECT id FROM pipeline_tasks WHERE id=${tid}$ AND state=${setstate}$ AND claimed_by=${cb}$",
{"tid": task_id, "setstate": set_state, "cb": claim_token})
if not check:
logger.info(f"claim lost race: task={task_id}")
# 并发抢输task 已被别人认领)或同项目已有活跃任务被串行阻塞,均属正常,下轮再试
logger.info(f"claim skipped: task={task_id}")
return None
return task