From ff2179be67f0e071a1b744dc18d8b3eece97e00b Mon Sep 17 00:00:00 2001 From: ymq Date: Thu, 20 Aug 2026 15:44:51 +0800 Subject: [PATCH] =?UTF-8?q?feat(poller):=20=E6=AF=8F=E9=A1=B9=E7=9B=AE/?= =?UTF-8?q?=E7=A7=9F=E6=88=B7=E5=B9=B6=E5=8F=91=E4=B8=8A=E9=99=90=20max=5F?= =?UTF-8?q?agents=5Fper=5Fproject=EF=BC=8C=E9=98=B2=E5=8D=95=E9=A1=B9?= =?UTF-8?q?=E7=9B=AE=E5=8D=A0=E6=BB=A1=E5=85=A8=E5=B1=80=E5=B9=B6=E5=8F=91?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit agent_poller 项目轮询基础上,新增每项目并发门控: - 读 appbase params 表 max_agents_per_project(默认 4),改后下一轮(10s)立即生效 - 每轮查各项目 running 计数(GROUP BY tenant_id),达上限的项目跳过不派新任务 - 改小后不强制杀,超出的 running 任务自然跑完即「等其完成」语义 - 与全局 max_concurrent_agents(默认3→预置10) + 项目轮询配套,形成三级并发控制 --- pipeline_service/init.py | 26 ++++++++++++++++++++++++-- 1 file changed, 24 insertions(+), 2 deletions(-) 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]