1. 项目定位与核心问题
先说结论:Agent-Reach 是我最近在折腾多智能体协作时顺手搭起来的一套轻量级任务触达框架。名字里的 Reach 有“触达”的意思,核心就一句话——在海量任务和一群能力参差不齐的 Agent 之间,建立一条稳定、可度量、能容错的调度通道。很多做 Agent 项目的人一开始都会忽略这个问题:总是默认“Agent 收到指令就能执行”,但现实是 Agent 会挂起、会返回脏数据、会互相抢资源,甚至干脆失联。Agent-Reach 要解决的就是这些“看不见但致命”的协作问题。
这套东西适合谁?两类人最有必要往下看。一类是正在做 RPA 流程自动化、想让多个自动化脚本协同调度的人;另一类是研究 AI Agent 编排、想在个人项目里接入多智能体协作的同学。不是我吹,群里好几个朋友看完我的实现方案之后,第一反应都是“原来这里用了超时重试之后,整个链路稳了不是一点半点”。
Agent-Reach 不依赖任何重型中间件,核心实现不到一千行代码,你完全可以在自己机器上跑通整个流程。它做的事情抽象成一句话:把模糊的任务描述,转化成明确的、可路由的、有确认回执的执行请求,并在失败时自动补位。听起来不复杂,但真正落地时会踩出一堆坑,这篇文章把我踩过的坑和最终实现的稳定版本都写清楚了。
拆开来看,Agent-Reach 由三块拼成:任务解析层、路由调度层、触达确认层。任务解析层负责把自然语言或结构化输入拆成标准化的任务单元;路由调度层负责按能力标签、负载状态、历史信用分决定把任务交给谁;触达确认层负责跟踪任务是否真正被执行完毕、结果是否有效。这三层彼此独立、通过消息队列串起来,层与层之间不共享内部状态,只认数据契约。所以无论你底层用的是 OpenAI 的 Agent、国产模型还是简单的规则脚本,都能直接接入。
顺便说一下,这个项目最初是我处理本地繁重任务时的副产品。我有几十个自动化脚本散布在各台服务器上,彼此之间经常需要传递处理结果,但那时全靠一个共享目录加定时任务硬顶着。出了几次“文件写了一半”“上游还没生成下游就开始消费”的事故之后,我下定决心把整个触达过程重构成 Agent-Reach 的方案,效果可以说完美。后面全文都会围绕这三层展开,从设计动机、代码实现到问题排查,一步一步说清楚。
2. 整体架构与设计思路拆解
2.1 为什么不用现成的工作流引擎
动手之前,我特意花了两天时间比较了几条路线:直接用 Celery、AirFlow 这类老牌分布式任务框架,或者上 Temporal 这类号称“完美容错”的编排引擎,甚至考虑过给 LangGraph 写一套自定义的 Agent 协作节点。最终全部否掉了。
理由很简单:这些框架的抽象层级跟我想要的“Agent 触达”之间,隔着一层无法忽略的膜。Celery 解决的是“任务依次/并行执行”的问题,但它不关心某次执行到底发生了什么事、执行 Agent 的能力有没有变化、返回的结果是不是可信。AirFlow 适合的是有固定 DAG 形态的数据流,你要在里面动态决策“这个任务应该交给哪位 Agent”,就得写一堆侵入式插件,改起来恨不得把整个调度器拆了重装。Temporal 最强的重试和恢复能力,确实碾压我自己的实现,但它的部署复杂度对个人项目而言太重了。
Agent-Reach 需要的核心能力,是“触达”二字——确保任务被正确的 Agent 收到并执行,且结果被确认。这个目标不需要重型的分布式事务,也不需要对编排状态做持久化快照,它需要的是三层简单、边界清晰、用了最朴素的超时与重试机制就能自愈的管道。所以我写了这套自己的东西,状态只保存在本地 SQLite 以及普通文件里,部署时只要一台小机器就能跑。
2.2 三层各自负责什么
任务解析层是 Agent-Reach 唯一的入口。它接收两种输入:一段自然语言描述,或者一条 JSON 结构体。解析层要做的事有两件:第一,把输入拆成任务单元,每个任务单元必须包含三个字段——目标能力标签、输入数据部分、超时阈值(秒)。第二,判断任务单元之间是否存在依赖关系。比如“先抓取网页再分析情感”,就是两条任务单元,前者输出要喂给后者。在 Agent-Reach 里这种依赖叫“显式边”,解析层会在生成任务单元的同时生成一张依赖图。
这个设计背后有一个考量:不做隐式推断。市面上很多 Agent 编排框架试图从上下文里猜任务之间的顺序,只要猜错一次,就会产生连锁反应。Agent-Reach 的策略是宁可让用户显式标注依赖,也不要让框架替用户做模糊猜测。任务单元生成之后,解析层会把它们塞进本地队列,队列项带一个唯一 task_id,这个 ID 从生成到最终确认回执都在整个链路里流转。
路由调度层是整条流水线的核心。每个 Agent 在注册进系统时,会声明自己的能力标签、权重、最大并发数,以及心跳频率。路由层持有这些 Agent 的状态表,在任务队列触发时,先筛选出能力匹配的 Agent 候选集,再按“当前最少负载 + 本轮成功率最高”的综合评分选出最终执行者。这里的评分算法我踩过几次坑,最初只顾着成功率,结果某个高成功率 Agent 被塞了远超它处理能力的任务,直接拖垮了它的响应速度,导致它后续所有任务都超时。后来我把负载权重调高,才算稳定下来。
触达确认层关心的是结果本身。Agent 执行完任务后,会把结果通过确认端点回传。确认层会做三件事:校验结果结构合法性、记录执行耗时与返回码、把执行状态回写到 SQLite。如果 Agent 在超时时间内没回传结果,确认层会标记该 Agent 本轮失常并触发任务重新路由。为了不把“任务失败”误判成“Agent 死亡”,确认层还设计了一套双状态区分机制,细节放在后面参数调优那一节里详细说。
2.3 任务触达的关键路径
把三层拼起来之后,一条完整任务的生命周期是这样的:任务从入口提交,解析层把它拆成任务单元并写入队列;路由层从队列拉出任务,经过能力筛选和负载评分,确定执行者;任务通过 HTTP 或者本地 IPC 通道发给目标 Agent;Agent 执行完再把结果 POST 回确认端点。确认层写库、记录指标、触发依赖图中下游任务的释放。
整套链路里最容易被忽略的是结果回传本身的超时问题。任务发送成功不等于执行成功,执行成功不等于回传成功。如果不分清楚这两个环节,之后排查问题的时候会把时间浪费在没有意义的网络排查上。Agent-Reach 的设计中,发送动作和回传动作分别计时、分别记录,任务状态机里有两个独立的超时阈值:send_timeout 和 execute_timeout。前者是发送阶段等待连接的最大时间,后者是 Agent 处理任务到回传结果的最大时间。
3. 核心实现与实操配置
3.1 最小可运行版本的代码结构
Agent-Reach 的核心代码不依赖任何第三方库,只要 Python 3.10+ 就能跑。我把整个项目拆成五个文件:task_parser.py、router.py、dispatcher.py、confirmer.py 和 agent_runtime.py。前四个分别对应三层加一个入口,agent_runtime.py 是你要在你自己的 Agent 上嵌进去的客户端 SDK。为了让你直接跑起来,我把最小版本粘贴在下面。
先看最核心的路由调度逻辑:
# router.py import json import sqlite3 import time class AgentRegistry: def __init__(self, db_path="reach.db"): self.conn = sqlite3.connect(db_path) self._init_table() def _init_table(self): with self.conn: self.conn.execute(""" CREATE TABLE IF NOT EXISTS agents ( name TEXT PRIMARY KEY, tags TEXT NOT NULL, -- 能力标签 JSON max_concurrent INTEGER NOT NULL DEFAULT 1, active_tasks INTEGER NOT NULL DEFAULT 0, success_count INTEGER NOT NULL DEFAULT 0, fail_count INTEGER NOT NULL DEFAULT 0, last_heartbeat REAL NOT NULL DEFAULT 0 ) """) def register(self, name, tags, max_concurrent=1): with self.conn: self.conn.execute( "INSERT OR REPLACE INTO agents VALUES (?, ?, ?, 0, 0, 0, 0)", (name, json.dumps(tags), max_concurrent) ) def score(self, row, task_tags): tag_match = len(set(json.loads(row[1])) & set(task_tags)) load_score = row[3] / max(row[2], 1) success_rate = row[4] / max(row[4] + row[5], 1) return tag_match, load_score, success_rate def pick(self, task_tags): now = time.time() rows = self.conn.execute( "SELECT * FROM agents WHERE last_heartbeat > ?", (now - 60,) ).fetchall() candidates = [r for r in rows if self.score(r, task_tags)[0] > 0] if not candidates: return None # 评分策略:先比能力匹配度,再比负载,最后比成功率 candidates.sort(key=lambda r: (self.score(r, task_tags))) return candidates[-1][0]代码不复杂,但每个字段都是必须的。tags 字段里存的是该 Agent 的能力标签,比如["web_fetch", "html_parse"];路由时通过标签交集判断能力是否匹配。last_heartbeat 用时间戳判断 Agent 是否存活,任何超过 60 秒没上报心跳的 Agent 都不会参与本轮任务分配。这里有一个我现在认为很关键的设计:把 Agent 的心跳和任务的响应分开处理。很多框架把心跳断了就直接把任务全部转移到别的机器上,这是过度反应。心跳断了只能说明 Agent 进程可能挂了,但 Agent 进程在处理的任务结果也许下一秒就回传了,强行转移会让两个 Agent 处理同一份工作,结果就可能被覆盖。
路由算法只有三档权重:能力匹配度 > 负载状态 > 历史成功率。为什么要这样排序?因为“能不能做”永远比“做得好不好”更重要。一个成功率 99% 的 Agent 如果没有匹配的能力标签,把任务交给它只会等来一个解析失败的结果。
再看确认层的核心片段:
# confirmer.py def confirm(task_id, agent_name, result, status): now = time.time() with conn: cursor = conn.execute(""" UPDATE tasks SET result = ?, status = ?, done_at = ?, confirmed_by = ? WHERE task_id = ? AND status = 'RUNNING' """, (json.dumps(result), status, now, agent_name, task_id)) if cursor.rowcount == 0: # 任务不存在,或者任务已经被确认过了,属于重复回传 return "duplicate" return "ok"这段代码里最值钱的不是 UPDATE 语句本身,而是那个AND status = 'RUNNING'条件。它保证了幂等性——同一个 Agent 重复回传结果时,只有第一次回传会被接受,后续重复请求都会被识别成 duplicate。任务判重是分布式系统里最容易忽略的细节,没有这个条件的话,只要 Agent 有一次网络重传,就会把任务的执行结果覆盖成完全相同的第二份数据,虽然值一样,但确认时间会错位,后面做执行耗时统计的时候就全错了。
3.2 Agent 端接入的实际过程
我写了 agent_runtime.py 给接入方用,这个 SDK 本身是一个独立进程,通过 HTTP 长连接监听任务请求。你接入新 Agent 时只需要做三件事:实例化 runtime、注册能力标签、注册执行函数。
# agent_runtime_demo.py import time from agent_runtime import AgentRuntime # ========== 模拟一个能抓网页标题的 Agent ========== def fetch_title(url: str) -> dict: # 这里放你真正的爬虫逻辑 time.sleep(2) return {"title": f"mock-title: {url}"} rt = AgentRuntime(name="web-fetcher-01") rt.register_tags(["web_fetch", "html_parse"]) rt.register_handler("web_fetch", fetch_title) rt.run(blocking=True)接入方式是我故意设计成“注册制”的。Agent 自己声明能处理什么标签的任务,并且为每个标签绑定一个处理函数。这样设计的好处是:新增 Agent 不需要改动路由层代码,只要 run 起来、把能力标签上报上去、就能立刻参与调度。新 Agent 接入的成本,从“改代码重启服务”降低到“写一个几十行的注册脚本”。
AgentRuntime 内部其实维护了两个线程:一个是 HTTP 服务监听线程,负责接收调度器发来的任务并调用对应 handler;一个是心跳线程,每 15 秒向调度器 POST 一条心跳记录。心跳内容包含当前 active_tasks 数量,调度器会用这个数值计算负载量。有一点你必须注意,心跳频率不要太快,我最初调成 3 秒一次,导致调度器压力倒是不大,但因为心跳线程频繁唤醒,Agent 机器的 CPU 占用率莫名其妙高了不少。后来在 15 秒这个值上稳定了下来。
实际接入自己的 Agent 时,最常见的报错是“handler 抛异常了但调度器不知道”。我特意在 AgentRuntime 的 handler 外层包了一层 try-except,任何异常都会被打包成 status=failed 的结果回传给确认层,同时异常信息放在 result 字段里,调度器拿到失败结果后会自动重新路由这个任务给其他 Agent。这个细节帮我在一次线上事故里保住了整个流水线,后面具体讲。
3.3 触达确认机制是怎么工作的
任务从调度器发出后,状态机按这个顺序演变:PENDING→RUNNING→CONFIRMED 或者 FAILED。PENDING 表示任务已写入数据库但还没找到匹配的 Agent;RUNNING 表示任务已被某个 Agent 接收;CONFIRMED 表示 Agent 成功回传结果并通过结构校验;FAILED 则分两种子状态:EXEC_FAILED(执行报错)和 TIMEOUT(超时未回传)。
触达确认的核心,其实就是确认层在等 Agent 的 POST 请求。但 Agent 可能因为网络波动没收到任务、收到任务后执行到一半宕机、执行完了但回传请求丢包,这三种情况下调度器是不知道任务真实状态的。Agent-Reach 用“双超时 + 补偿查询”解决这个问题。主流程中调度器在 execute_timeout 到期后,会先发一条状态查询请求给 Agent(如果 Agent 还活着的话),Agent 返回当前任务的执行进度;只有确认 Agent 真的失联或者进度停滞时,调度器才把任务从 RUNNING 改判为 TIMEOUT,并重新路由。这套“先问再判”的机制,把误判率降到了几乎可以忽略。
4. 参数选型与调优实录
4.1 超时阈值怎么定
上面提到 agent_runtime.py 本身不涉及超时配置,超时阈值是在任务解析层指定、或在路由层用默认值兜底的。默认值我分别设成了 10 秒(send_timeout)和 120 秒(execute_timeout)。这两个数的选取不是拍脑袋,而是来自我对自己那几十个脚本的耗时分布的分析。
如果你任务的平均耗时在 30 秒以内,execute_timeout 建议设为 90 秒;如果任务里包含大文件下载或模型推理,建议直接拉到 300 秒以上。判断依据很简单:execute_timeout 必须大于 99.9% 的正常执行耗时,否则你就不是在管异常,而是在制造异常。我见过有同学把超时设成和平均耗时一样,结果一半的任务都因为波动而超时重试了,整个系统的吞吐量直接腰斩。
send_timeout 反而要设得相对小一些,因为发送环节本身不应该阻塞太久。如果连发送都超时,大概率是网络链路断了,再等也没意义。send_timeout 建议 2 到 10 秒之间,不要超过 15 秒。
4.2 负载因子的修正
之前提到过,路由评分时我一开始把成功率权重放得过高,导致高成功率 Agent 被打爆。后来引入了一个动态负载因子,公式是这样的:
score = tag_match * 100 - load_score * 50 + success_rate * 20load_score 的定义是active_tasks / max_concurrent。如果某个 Agent 最大并发数是 2,当前已经在处理 2 个任务,load_score 就是 1.0,会从评分里扣掉 50 分。这个扣分力度相当大,基本保证了一个满载的 Agent 不会继续被分配新任务。只有当所有 Agent 都满载时,才会轮到满载者继续接活,这是最极端情况下的兜底。
这里我要特别提醒一个边界问题:max_concurrent不要填太大。它不是越高越好,因为 Agent 进程的处理能力受 CPU 和内存限制,特别是那些调用了大模型的 Agent,并发数建议设置为单台机器的 CPU 核心数的一半。填高了只会让每个任务都在排队等资源,反而比低并发更慢。
4.3 心跳频率与失活判断之间的配合
Agent 的心跳周期默认 15 秒发送一次,路由层在挑选 Agent 时只看最近 60 秒内有过心跳的节点。这个“15 秒发送、60 秒判定”的策略组合里暗含一层缓冲:即使丢失了连续三个心跳包,Agent 仍然被视为存活,不会被踢出候选池。这耐受性很重要,因为你的网络链路偶尔抖动一秒钟就恢复了,没必要因此触发大规模任务迁移。
如果你希望 Agent 失活后能更快被感知,可以把心跳周期缩短到 5 秒、判定窗口缩短到 20 秒。但同理,代价是心跳线程对 Agent 进程的唤醒会更频繁,CPU 占用略高。多数场景下 15 秒心跳足够。
5. 真实场景下的坑与排查思路
5.1 案例一:任务被重复执行
这是我最早期碰到的诡异问题:某条任务明明已经执行完了,结果回传成功后,调度器又把同样的任务重新路由给了另一个 Agent,导致同一份任务被执行了两遍。刚开始我怎么也查不到原因,后来打开 SQLite 一看任务表才发现,确认层的 UPDATE 没加状态判断,第二次回传把第一次回传的 task_id 更新成了一条新的 RUNNING 状态,调度器一扫描发现 RUNNING 里的任务早就超时了,就又重新路由了。
这个问题用了 5.1 节开头提过的AND status = 'RUNNING'条件彻底解决。我后来在代码注释里写了一句“没有幂等就会出大事”,算是对自己的教训总结。
5.2 案例二:Agent 进程“假死”导致任务堆积
有次某个 Agent 进程本身还活着、心跳也正常,但实际上线程池已经全部卡死,任何一个新任务进去都会卡在资源等待里。调度器看到心跳正常就把任务发过去,结果每一个都超时,一连串任务的 execute_timeout 全部被打满,用户体验极差。
这个问题靠心跳根本发现不了,因为心跳线程并不关心业务线程的死活。我的解法是:给 Agent 心跳报文增加一个 execution_probe 指标,AgentRuntime 在心跳线程里顺带检测当前任务队列深度,如果队列深度大于 2 并且最近 30 秒没有任何任务完成,就把 status 标记为 degraded 上报给调度器。调度器收到 degraded 状态后会暂停向该 Agent 发送新任务,转到其他节点。这个机制本质上是把“Agent 健康”从进程层面细化到了工作能力层面,效果显著。
5.3 案例三:依赖任务中下游提前触发
Agent-Reach 的依赖图支持“下游任务等待上游确认后才释放”的能力。但我最开始实现时只考虑了任务在路由层被标记为 CONFIRMED 就释放下游,没有考虑到上游任务虽然确认了但结果结构不合法的情况。比如上游任务返回了一个空字符串,下游任务拿去做提取,直接抛异常。
现在的实现里,确认层在将任务标为 CONFIRMED 之前,多做了一个 result_schema 校验:每个任务单元可以附带预期返回结构,确认层用 jsonschema 轻量校验,不通过则按失败处理。这样一来,下游收到的永远是结构完整的数据对象。
5.4 常见问题速查表
| 现象 | 可能原因 | 排查思路 |
|---|---|---|
| 任务一直处于 PENDING | Agent 没有上报心跳、能力标签不匹配 | 检查 Agent 是否注册成功;确认标签拼写是否完全一致 |
| 出现重复执行 | 确认层缺少幂等保护 | 确认 UPDATE 是否带status=RUNNING条件 |
| 某个 Agent 持续超时但心跳正常 | Agent 线程池假死 | 检查 Agent 的队列深度指标,考虑加入探针检测 |
| 任务执行时间被莫名拉长 | execute_timeout设置过小导致大量重试 | 计算正常执行耗时的 99.9% 分位数,重新设定超时 |
| 结果偶尔丢失 | Agent 回传请求非幂等 | 确认端点是否支持重复 POST 不覆盖 |
6. 扩展场景:Agent-Reach 能怎么变着用
Agent-Reach 这套设计虽然是我为了解决本地脚本调度问题而写的,但把它抽象之后适配的领域其实挺广。最直接的是做 RPA 工作流:每个 RPA 机器人当成一个 Agent 注册进来,路由层根据“页面抓取”“Excel 填写”“邮件发送”等能力标签分配任务,能解决掉机器人之间互相顶替和重复劳动的问题。
另一个我认为很值得尝试的方向是把 Agent-Reach 接成 AI Agent 的调度中台。比如你有多套基于不同大模型构建的 Agent,分别擅长代码生成、文案改写、数据分析;Agent-Reach 作为统一入口,按任务的标签把请求路由到对应模型 Agent 上。加上触达确认层的校验逻辑,就能保证结果格式可控,不至于拿一个代码生成模型返回的长文本去直接填入数据表。
如果你有多个 Agent 跑在本地局域网内,部署 Agent-Reach 时不需要额外开放公网端口,只要调度器和 Agent 在一个子网里就能互通。跨机器通信时注意 AgentRuntime 监听地址要绑定到局域网 IP,不要绑 localhost。
Agent-Reach 后续可以扩展的方向也很多。目前我只做了基于 SQLite 的单机版,数据量如果大到一定程度,可以把 SQLite 换成 PostgreSQL,路由层的查询逻辑几乎不用改。触达确认层如果要支持消息回溯,可以引入一个简单的 append-only 日志文件,把每条任务的状态变化都追加进去。我在自己的版本里已经加了这一层,排查问题的时候方便得多。
如果打算自己动手改,我从经验上建议先动两个点:一是把 Agent 能力标签设计成层级结构,比如web_fetch.subpage,这样路由时可以做更精细的匹配;二是把心跳策略改成按耗时长短自动调节频率,对于长耗时任务可以降低心跳频率,省一点资源。
Agent-Reach 的整体思路说到底只是“清晰分层 + 明确超时 + 幂等确认”这三个原则的组合。很多分布式系统的问题,并不是缺一个重型框架,而是缺这些最基本的原则被严格执行。我把 Agent-Reach 放出来的时候,希望它至少能给你提供一个参考:当你的 Agent 开始协作的时候,别让它们失联,也别让它们重复干活——这两件事管好了,整个系统就稳了一大半。