
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 通道发给目标 AgentAgent 执行完再把结果 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_pathreach.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_concurrent1): 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(keylambda 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: fmock-title: {url}} rt AgentRuntime(nameweb-fetcher-01) rt.register_tags([web_fetch, html_parse]) rt.register_handler(web_fetch, fetch_title) rt.run(blockingTrue)接入方式是我故意设计成“注册制”的。Agent 自己声明能处理什么标签的任务并且为每个标签绑定一个处理函数。这样设计的好处是新增 Agent 不需要改动路由层代码只要 run 起来、把能力标签上报上去、就能立刻参与调度。新 Agent 接入的成本从“改代码重启服务”降低到“写一个几十行的注册脚本”。AgentRuntime 内部其实维护了两个线程一个是 HTTP 服务监听线程负责接收调度器发来的任务并调用对应 handler一个是心跳线程每 15 秒向调度器 POST 一条心跳记录。心跳内容包含当前 active_tasks 数量调度器会用这个数值计算负载量。有一点你必须注意心跳频率不要太快我最初调成 3 秒一次导致调度器压力倒是不大但因为心跳线程频繁唤醒Agent 机器的 CPU 占用率莫名其妙高了不少。后来在 15 秒这个值上稳定了下来。实际接入自己的 Agent 时最常见的报错是“handler 抛异常了但调度器不知道”。我特意在 AgentRuntime 的 handler 外层包了一层 try-except任何异常都会被打包成 statusfailed 的结果回传给确认层同时异常信息放在 result 字段里调度器拿到失败结果后会自动重新路由这个任务给其他 Agent。这个细节帮我在一次线上事故里保住了整个流水线后面具体讲。3.3 触达确认机制是怎么工作的任务从调度器发出后状态机按这个顺序演变PENDING→RUNNING→CONFIRMED 或者 FAILED。PENDING 表示任务已写入数据库但还没找到匹配的 AgentRUNNING 表示任务已被某个 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 常见问题速查表现象可能原因排查思路任务一直处于 PENDINGAgent 没有上报心跳、能力标签不匹配检查 Agent 是否注册成功确认标签拼写是否完全一致出现重复执行确认层缺少幂等保护确认 UPDATE 是否带statusRUNNING条件某个 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 开始协作的时候别让它们失联也别让它们重复干活——这两件事管好了整个系统就稳了一大半。