diff --git a/pipeline_service/init.py b/pipeline_service/init.py index cd77c55..cac8dba 100644 --- a/pipeline_service/init.py +++ b/pipeline_service/init.py @@ -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]