Day 13 / 共 20 天 · 第 3 周 执行引擎

数据流与拓扑执行

拓扑执行的最后一块拼图。一个节点跑完,输出怎么精确地"喂"给下游?多个上游的节点怎么"攒齐"再跑?今天读 _enqueue_next_nodes

📍 你在整门课的位置 · 第 3 周 执行引擎
D12 ExecutionManager D13 数据流/拓扑 D14 Scheduler D15 WebSocket
💡 今天的类比世界观:拓扑执行 = 传送带上的"齐料才开工" 延续工厂比喻,看料怎么精确地流:execute_node = 某工位做完一道工序_enqueue_next_nodes = 把半成品放上传送带送给下游工位攒齐输入才入队(fan-in)= 下游工位要等所有必需的料都到齐才开工fan-out = 一份料复制着分给多个下游静态连线回填 = 常量料提前备在工位边原子锁 = 同一份料不会被误领两次。今天都用"传送带 / 齐料开工"来想。

👶 小白:一个节点有两个上游,只到了其中一个数据,它会先跑起来吗?

👨‍🏫 老师:不会。就像做一道菜要"肉 + 菜"都备齐才下锅,只到了肉不能先炒。平台会先检查这个节点的必需输入是否全部到齐,齐了才把它入队执行——这就是 fan-in 汇聚。先到的数据会暂存等着,等最后一个上游也送到,才凑成完整一份料开工。这保证下游拿到的永远是完整输入。

L01

核心问题

Day 12 的调度循环说"节点产出的输出决定下游哪些节点入队"。今天回答具体怎么做到:① 一个节点的输出怎么找到该连给谁;② 一个下游节点有多个输入口(来自多个上游),怎么等它们都到齐才跑。

难点在"多入"(fan-in) 如果图是一条直线(A→B→C),很简单:A 完了给 B,B 完了给 C。但真实图有汇聚——比如"生成报告"节点需要"研究结果"和"数据分析"两个上游都完成才能开始。这时不能研究一完就跑报告(数据分析还没好)。必须"攒齐所有输入才入队"——这是拓扑执行最核心也最微妙的地方。
L02

单节点执行 execute_node

execute_nodemanager.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 可流式产出多个 (口名, 值)
读法:Block 的 run 是异步生成器(Day 02),所以一个节点执行可以陆续产出多个 (输出口名, 值)。每产出一个,就交给下一步(L03)去看"这个输出该喂给哪些下游"。finally 里释放凭证锁、累加统计。
L03

输出喂给下游 _enqueue_next_nodes

🤔 痛点:节点跑出一个值,引擎凭什么知道"这值该给谁、下游能不能开跑"? 节点 A 产出一个 {"summary": "..."}。引擎不能瞎猜——它得知道 A 的这个输出口连到了哪些下游的哪些输入口,还得判断被喂到的下游节点"输入是不是齐了"。全靠图里那些 Link(Day 11)和一段专门的派发逻辑。
💡 本质:execute_node 产出 → _enqueue_next_nodes 顺着连线派发 → 齐了才入队 节点执行 execute_nodemanager.py:142)像挤牙膏一样逐个 yield 出 (输出口名, 值);_enqueue_next_nodesmanager.py:389)拿着每个输出,顺着以它为源的每条 Link 找到下游输入口、把值写进去;写完检查下游"必需输入齐了没",齐了才把下游标 QUEUED 入队。整个拓扑执行,就是这一"产出→派发→凑齐→入队"的循环在推进。

_enqueue_next_nodesmanager.py:389)——数据流的引擎。对当前节点的每条 output_link

# 概念流程
# 1. parse_execution_output 判断输出值是否匹配这条连线的 source_name
# 2. 匹配 → upsert_execution_input 把值写给下游节点执行(攒输入)
# 3. validate_exec 检查下游输入齐了吗:
#      不齐 → 只落库,不入队
#      齐了 → 标 QUEUED,加进执行队列
读法:一个输出口的值,顺着连它的每条 Link,写进对应下游节点的输入口。关键在第 3 步——写完后检查这个下游节点的所有必需输入是否都齐了,齐了才入队执行。
L04

攒齐输入才入队(核心规则)

关键机制 upsert_execution_inputdata/execution.py:862):把上游的值"攒"到下游节点执行上。

# 找一个"还没有这个输入口、且 INCOMPLETE"的下游节点执行,把值插进去
# 找不到就新建一个 INCOMPLETE 执行
# —— 陆续到来的多个输入被"攒"到同一次执行上
📝 举个例子:报告节点的 3 个输入陆续到齐的过程 报告节点需要 researchdatatitle 三个输入:
① 研究块先完成 → upsert 找到一个 INCOMPLETE 执行,写入 {research: ...} → 检查:缺 data/title → 只落库,不入队
② 标题块完成 → 写入 {title: ...}(攒到同一条执行)→ 仍缺 data → 不入队
③ 数据块最后完成 → 写入 {data: ...} → 检查:三个全齐 → 标 QUEUED 入队,报告开跑。
顺序无关:谁先谁后都行,只要三个都到,第三个到的那次触发入队。
"攒输入"是怎么回事? 一个下游节点有 3 个输入口,分别来自 3 个上游。这 3 个上游不会同时完成——A 先完成给了口1(此时口2、口3还空)、然后 C 完成给了口3、最后 B 完成给了口2。upsert_execution_input 把陆续到来的值攒到同一个下游节点执行上(一个 INCOMPLETE 状态的执行记录),像往一个篮子里陆续放东西。每次放完就检查"篮子满了吗(所有必需输入都有了吗)"——满了才标 QUEUED 入队。这就实现了"等所有输入到齐才跑"。
L05

fan-in 汇聚 + fan-out 分叉

研究 Block分析 Block
都完成才 → 报告 Block
等两个输入齐
fan-out(一出多入)与 fan-in(多入汇聚) 触发/输入 研究 Block 分析 Block 报告 Block 攒齐 2 个输入才 QUEUED fan-out:一个输出喂多个下游 fan-in:两个上游都到齐,报告才开跑
同一套"攒齐输入才入队"机制同时支持 fan-out(触发同时喂研究+分析)和 fan-in(报告耐心等两条支线都完成)——无需为复杂图另写逻辑。
读法:"攒齐输入才入队"优雅地处理了 fan-in(多入汇聚)——报告节点耐心等研究和分析都完成。同理 fan-out(一出多入)也自然支持:一个节点的一个输出口可以连多条 Link,每条都触发一个下游(一个输出喂给多个下游)。
为什么这套机制能处理任意复杂的图? 无论图多复杂(分支、汇聚、并行、菱形),本质都是"节点等输入齐了就跑、跑完把输出喂给下游"。只要每个节点都遵守"输入齐了才跑",整张图的执行顺序就自动正确——不需要预先算好一个全局执行顺序。这是数据流驱动的美妙之处:局部规则(等输入齐)保证全局正确(拓扑顺序)。和 eino 的 channel"多上游合并"、Pregel 引擎异曲同工。
L06

静态连线回填

Day 03/11 讲的静态连线(is_static)在这里有特殊处理(manager.py:458):补齐下游缺失的静态输入用"最近一次执行的输入值"回填;若连线是静态的,还会把之前"卡着等这个静态值"的 INCOMPLETE 执行重新补齐、重新入队(:501)。

静态连线为什么要特殊处理? 普通输入"用一次就消费掉",但静态输入(如配置值、循环里的常量)要能被反复用。想象一个循环:每轮都要用同一个"配置"值。如果配置值用一次就没了,第二轮就卡住。静态连线回填 = "这个值一直可用,谁等它就随时给"——支持循环、支持一值多用。这让 AutoGPT 的图能表达迭代循环(不只是一次性 DAG)。
L07

原子锁防重复

_enqueue_next_nodes 用分布式锁 synchronized(f"upsert_input-{next_node_id}-{graph_exec_id}")manager.py:446)保证"同一下游节点的输入攒集是原子的"。

为什么攒输入要加锁? 想象两个上游节点几乎同时完成,都要给下游节点写输入。如果不加锁,可能两个都读到"篮子里还差 1 个"、都以为"我写完还差、不该入队"——结果输入其实齐了却没人入队(下游卡死);或者都判断"齐了"、重复入队两次(下游跑两遍)。加锁让"写输入 + 检查是否齐 + 入队"这一串操作原子化——同一时刻只有一个上游在处理,不会错判。并发编程里,"检查后操作"(check-then-act)必须加锁,否则有竞态。
L08

今日小结 + 动手

🧠 今天你应该能回答

  • 拓扑执行的核心难点是什么?(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
明天预告 · Day 14Scheduler 定时执行——让 Agent 定时/周期性自动跑("每天早上 8 点")。基于 APScheduler、持久化 JobStore、cron 兼容处理、孤儿任务自愈。
← Day 12 执行引擎 Day 14 · Scheduler →