feat(poller): 每项目/租户并发上限 max_agents_per_project,防单项目占满全局并发
agent_poller 项目轮询基础上,新增每项目并发门控: - 读 appbase params 表 max_agents_per_project(默认 4),改后下一轮(10s)立即生效 - 每轮查各项目 running 计数(GROUP BY tenant_id),达上限的项目跳过不派新任务 - 改小后不强制杀,超出的 running 任务自然跑完即「等其完成」语义 - 与全局 max_concurrent_agents(默认3→预置10) + 项目轮询配套,形成三级并发控制
This commit is contained in:
parent
7a2175f363
commit
ff2179be67
@ -610,6 +610,25 @@ def load_pipeline_service():
|
||||
# 达到并发上限,本轮不派发(退出 context 后末尾 sleep 再 poll)
|
||||
return
|
||||
|
||||
# 每项目/租户并发上限(max_agents_per_project,appbase params 表,默认 4)。
|
||||
# 改小后 poller 不再给已达上限的项目派新任务,超出的 running 任务自然跑完(不强制杀)——
|
||||
# 即「改小后超出任务等其完成」。
|
||||
from .workspace import get_param
|
||||
per_max = 4
|
||||
try:
|
||||
per_max = max(1, int(float(str(await get_param(sor, 'max_agents_per_project', '4')))))
|
||||
except (ValueError, TypeError):
|
||||
per_max = 4
|
||||
|
||||
# 每项目 running 计数(GROUP BY tenant_id),供每项目上限门控用
|
||||
per_running = {}
|
||||
_per_rows = await sor.sqlExe(
|
||||
"SELECT tenant_id, COUNT(*) as c FROM pipeline_tasks "
|
||||
"WHERE state='running' AND pipeline_id='role_task' GROUP BY tenant_id", {})
|
||||
for _r in (_per_rows or []):
|
||||
per_running[getattr(_r, 'tenant_id', '')] = getattr(_r, 'c', 0) or 0
|
||||
await sor.sqlExe("COMMIT", {})
|
||||
|
||||
def _deps_of(depends_on_raw):
|
||||
"""解析 depends_on(JSON 数组字符串或 list)→ 依赖任务 ID 列表。空/非法 → []"""
|
||||
if not depends_on_raw:
|
||||
@ -648,8 +667,8 @@ def load_pipeline_service():
|
||||
if all(dep_state.get(d, '') in ('completed', 'approved') for d in deps):
|
||||
ready.append(rec)
|
||||
|
||||
# 项目轮询(round-robin):每个项目轮流取一个就绪任务,保证项目间公平,
|
||||
# 避免单项目大量就绪任务占满全局并发名额、饿死其他项目。(dict 保持插入序)
|
||||
# 项目轮询(round-robin + 每项目并发上限):每个项目轮流取一个就绪任务,
|
||||
# 保证项目间公平 + 单项目不占满全局并发名额、饿死其他项目。(dict 保持插入序)
|
||||
by_project = {}
|
||||
for rec in ready:
|
||||
by_project.setdefault(getattr(rec, 'tenant_id', ''), []).append(rec)
|
||||
@ -659,9 +678,12 @@ def load_pipeline_service():
|
||||
for pid in list(by_project.keys()):
|
||||
if len(selected) >= avail:
|
||||
break
|
||||
if per_running.get(pid, 0) >= per_max:
|
||||
continue # 该项目已达每项目上限,跳过(等其 running 任务完成再派)
|
||||
tasks = by_project[pid]
|
||||
if tasks:
|
||||
selected.append(tasks.pop(0))
|
||||
per_running[pid] = per_running.get(pid, 0) + 1 # 计入本轮即将派发的名额
|
||||
progressed = True
|
||||
if not tasks:
|
||||
del by_project[pid]
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user