数据流与拓扑执行
拓扑执行的最后一块拼图。一个节点跑完,输出怎么精确地"喂"给下游?多个上游的节点怎么"攒齐"再跑?今天读 _enqueue_next_nodes。
👶 小白:一个节点有两个上游,只到了其中一个数据,它会先跑起来吗?
👨🏫 老师:不会。就像做一道菜要"肉 + 菜"都备齐才下锅,只到了肉不能先炒。平台会先检查这个节点的必需输入是否全部到齐,齐了才把它入队执行——这就是 fan-in 汇聚。先到的数据会暂存等着,等最后一个上游也送到,才凑成完整一份料开工。这保证下游拿到的永远是完整输入。
核心问题
Day 12 的调度循环说"节点产出的输出决定下游哪些节点入队"。今天回答具体怎么做到:① 一个节点的输出怎么找到该连给谁;② 一个下游节点有多个输入口(来自多个上游),怎么等它们都到齐才跑。
单节点执行 execute_node
execute_node(manager.py:142)跑一个节点:校验输入 → 注入凭证(Day 08)→ 调 block.execute()(异步生成器)逐个产出:
# manager.py:341(概念)
async for output_name, output_data in block.execute(input_data, **extra_exec_kwargs):
yield output_name, output_data # Block 可流式产出多个 (口名, 值)
输出喂给下游 _enqueue_next_nodes
{"summary": "..."}。引擎不能瞎猜——它得知道 A 的这个输出口连到了哪些下游的哪些输入口,还得判断被喂到的下游节点"输入是不是齐了"。全靠图里那些 Link(Day 11)和一段专门的派发逻辑。execute_node(manager.py:142)像挤牙膏一样逐个 yield 出 (输出口名, 值);_enqueue_next_nodes(manager.py:389)拿着每个输出,顺着以它为源的每条 Link 找到下游输入口、把值写进去;写完检查下游"必需输入齐了没",齐了才把下游标 QUEUED 入队。整个拓扑执行,就是这一"产出→派发→凑齐→入队"的循环在推进。_enqueue_next_nodes(manager.py:389)——数据流的引擎。对当前节点的每条 output_link:
# 概念流程
# 1. parse_execution_output 判断输出值是否匹配这条连线的 source_name
# 2. 匹配 → upsert_execution_input 把值写给下游节点执行(攒输入)
# 3. validate_exec 检查下游输入齐了吗:
# 不齐 → 只落库,不入队
# 齐了 → 标 QUEUED,加进执行队列
攒齐输入才入队(核心规则)
关键机制 upsert_execution_input(data/execution.py:862):把上游的值"攒"到下游节点执行上。
# 找一个"还没有这个输入口、且 INCOMPLETE"的下游节点执行,把值插进去
# 找不到就新建一个 INCOMPLETE 执行
# —— 陆续到来的多个输入被"攒"到同一次执行上
research、data、title 三个输入:① 研究块先完成 →
upsert 找到一个 INCOMPLETE 执行,写入 {research: ...} → 检查:缺 data/title → 只落库,不入队② 标题块完成 → 写入
{title: ...}(攒到同一条执行)→ 仍缺 data → 不入队③ 数据块最后完成 → 写入
{data: ...} → 检查:三个全齐 → 标 QUEUED 入队,报告开跑。顺序无关:谁先谁后都行,只要三个都到,第三个到的那次触发入队。
upsert_execution_input 把陆续到来的值攒到同一个下游节点执行上(一个 INCOMPLETE 状态的执行记录),像往一个篮子里陆续放东西。每次放完就检查"篮子满了吗(所有必需输入都有了吗)"——满了才标 QUEUED 入队。这就实现了"等所有输入到齐才跑"。fan-in 汇聚 + fan-out 分叉
等两个输入齐
静态连线回填
Day 03/11 讲的静态连线(is_static)在这里有特殊处理(manager.py:458):补齐下游缺失的静态输入用"最近一次执行的输入值"回填;若连线是静态的,还会把之前"卡着等这个静态值"的 INCOMPLETE 执行重新补齐、重新入队(:501)。
原子锁防重复
_enqueue_next_nodes 用分布式锁 synchronized(f"upsert_input-{next_node_id}-{graph_exec_id}")(manager.py:446)保证"同一下游节点的输入攒集是原子的"。
今日小结 + 动手
🧠 今天你应该能回答
- 拓扑执行的核心难点是什么?(fan-in 多入汇聚)
- 输出怎么找到该喂给谁?(顺着 output_link 匹配 source_name)
- "攒齐输入才入队"怎么实现?(upsert_execution_input 攒 + validate 检查齐没齐)
- 局部规则怎么保证全局拓扑正确?
- 静态连线回填支持什么?攒输入为什么要加锁?
✋ 动手
P=autogpt_platform/backend/backend
sed -n '389,470p' $P/executor/manager.py | head -50 # _enqueue_next_nodes
sed -n '862,936p' $P/data/execution.py | head -40 # upsert_execution_input
grep -n 'synchronized\|is_static' $P/executor/manager.py | head