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

执行引擎 manager

整个平台的心脏。今天读 executor/manager.py:ExecutionManager 的进程架构(MQ + 线程池 + 事件总线)、一次图执行的生命周期、分布式锁去重。

📍 你在整门课的位置 · 第 3 周 执行引擎
D11 Graph 模型深入 D12 ExecutionManager D13 数据流/拓扑 D14 Scheduler
💡 今天的类比世界观:执行引擎 = 工厂的"调度指挥中心" manager 是整个平台的心脏,像一个调度中心:MQ + 线程池 = 任务排队 + 一群随时待命的工人事件总线 = 车间广播喇叭(谁有事喊一声,关心的人自己来听——发布 / 订阅);ClusterLock 分布式锁 = 一张任务只发给一个工人的"派工牌"(防两人抢同一活重复干);调度循环 = 调度员盯着看板不停派活。今天都用"调度中心"来想。

👶 小白:都用队列排队了,为什么还要一把"分布式锁"?队列不是已经保证顺序了吗?

👨‍🏫 老师:因为平台会开多个执行进程同时从队列抢活,可能两个进程几乎同时抓到同一个节点,导致同一步被跑两遍(既浪费钱又可能重复发邮件)。ClusterLock 就是"这活我领了"的派工牌——谁先抢到锁谁干,别人看到已锁就跳过。队列管"排队",锁管"同一件事只被做一次",两者解决的是不同问题。

L01

执行引擎全景

🤔 痛点:几千个用户同时点"运行",谁来跑、跑在哪、会不会互相拖垮? 一张图可能跑几分钟到几小时(多次 LLM 调用 + 外部 API)。如果 Web 服务器收到请求就地开跑,请求会一直挂着、服务器被长任务占满,多用户根本扛不住。得有个专门"干重活"的角色,把执行从 Web 层剥离出去。
💡 本质:ExecutionManager = 一个"长期待命的工头",从队列领活、丢给工人干 它是个常驻进程(manager.py:1353):盯着 RabbitMQ 队列,来一条"跑图请求"就从线程池executor 属性返回 ThreadPoolExecutormanager.py:1395-1397)派一个工人线程去跑;每个工人由 ExecutionProcessormanager.py:551)驱动,一张图从头到尾走 on_graph_execution:812),图里每个节点走 on_node_execution:578)。工头只管调度,不亲自干活。

ExecutionManagermanager.py:1353)是一个常驻进程。它由三样东西支撑:

RabbitMQ — 消费"图执行请求"和"取消请求"两个队列(跨进程/跨机器分发任务)
线程池 ThreadPoolExecutor — 每个图执行占一个 worker 线程(并发跑多个图)
Redis 事件总线 — 每次节点/图状态变化 publish 事件(推给前端/子 Agent)
RabbitMQ 跑图请求排队 ExecutionManager 工头 (:1353) 消费队列→派活 worker 1 · 图A ExecutionProcessor worker 2 · 图B on_graph_execution …共 pool_size 个 ThreadPoolExecutor (:1395) Redis 事件总线 ① 请求入队 → ② 工头领取 → ③ 派线程池 worker 跑图 → ④ 每步进度 publish 到 Redis
RabbitMQ 分发任务、线程池并发跑多张图(每图一个 worker 线程)、Redis 广播进度——工头 ExecutionManager 只负责"领活+派活"。
先建立全局图景 REST 把"图执行请求"发到 RabbitMQ(Day 05)→ ExecutionManager 消费 → 交给线程池的一个 worker → worker 跑这张图(拓扑执行各 Block)→ 每步状态变化 publish 到 Redis → WebSocket 转发给前端。三个基础设施各司其职:MQ 分发任务、线程池并发执行、事件总线广播进度。今天聚焦"一张图怎么被一个 worker 跑起来"。
L02

MQ + 线程池

ExecutionManager 用两个消费线程(manager.py:1379/1431)分别消费"运行"和"取消"队列,一个线程池(:1394)作为图执行工作池:

# 收到运行消息 → executor.submit(execute_graph, ...) 丢进线程池
# MQ 用 auto_ack=False(:1480):worker 挂了消息重回队列不丢
# prefetch_count=pool_size(:1476):一个 pod 不抢超过它能跑的任务
读法:auto_ack=Falseprefetch_count 是两个关键的 MQ 可靠性设置。前者保证"任务处理完才确认"——中途崩溃消息不丢、会被别的 worker 重新拿。后者保证"一个 worker 只预取它能处理的量"——不会一个 pod 贪心拿走所有任务却跑不过来。
为什么图执行要用线程池,一次一个 worker? 因为一张图可能跑很久(分钟到小时级——涉及多次 LLM 调用、外部 API)。用线程池能并发跑多张图(pool_size 个),每张图独占一个 worker 线程互不干扰。MQ + 线程池 = 可靠地、并发地处理长任务——这是任务型系统的标准架构。
L03

事件总线(发布-订阅)

基于 Redis 的 RedisExecutionEventBus——每次节点/图状态变化就 publish 事件到两个频道:exec/{graph_exec_id}(这次执行)和 {graph_id}/all(这个图的所有执行)。

谁在订阅这些事件?WebSocket 进程——订阅后转发给前端,你就看到画布上节点实时亮起(Day 15)。② 子 Agent——Agent-as-Block(Day 03)里,父块订阅子图执行事件,子图产出即父块产出。发布-订阅解耦:executor 只管"发生了什么就广播",不关心谁在听;订阅者各取所需。这和 OpenHands、CrewAI 的事件系统完全同构——事件驱动是 Agent 系统的通用可观测层。
L04

分布式锁去重(ClusterLock)

收到运行消息后(_handle_run_messagemanager.py:1529),先尝试获取一把 ClusterLock(Redis 分布式锁,:1670)——保证同一个 graph_exec_id 全集群只在一个 pod 上跑

为什么需要分布式锁? 生产环境有多个 executor pod(多台机器/容器)。如果一条 MQ 消息因为网络抖动被投递两次,或两个 pod 同时抢到——同一次执行可能被跑两遍(重复扣费、重复副作用如发两封邮件)。ClusterLock 保证"同一执行全集群只有一个 pod 在跑"——拿到锁的才跑,没拿到的跳过。分布式系统里,"防止重复执行"必须靠分布式锁(单机锁不够,因为跨机器)。完成时 _on_run_done:1717)ack 消息、释放锁。
L05

图执行调度循环(核心)

_on_graph_executionmanager.py:942)是真正的调度循环:

# 1. 建进程内 ExecutionQueue(节点级队列)
# 2. 预填队列:把 QUEUED/RUNNING 的节点执行捞出来入队(起始节点已建为 QUEUED)
# 3. 主调度循环:
while 队列有节点 或 有在跑的节点:
    node = 队列.取一个
    billing.charge_usage(node)          # 预扣费(Day 09)
    异步执行 on_node_execution(node)    # 丢到专门的事件循环跑
    # 内层:取在跑节点的输出 → _process_node_output → 入队下游(Day 13)
📝 举个例子:一张"读网页→总结→发邮件"的图怎么被这个循环跑完 起始节点[读网页]已是 QUEUED,进循环:
取[读网页]→预扣费→异步执行→产出"网页文本"→喂给[总结],[总结]输入齐→入队
取[总结]→执行→产出"摘要"→喂给[发邮件],入队
取[发邮件]→执行→发送→无下游 → 队列空且无在跑节点 → 循环退出,图执行完成。循环本身不认识"读网页/总结/发邮件",它只做"取一个→跑→把产出喂下游→重复"。
读法:调度循环就是"从队列取节点 → 预扣费 → 异步执行 → 产出的输出决定下游哪些节点入队 → 重复,直到没有节点可跑"。这就是 Day 05 讲的拓扑执行的代码核心。起始节点(Day 11 的 starting_nodes)在创建执行时已入队。
L06

两个事件循环(并行度优化)

调度用了两个独立的 asyncio 事件循环线程manager.py:794):

  • node_execution_loop:跑 Block 的业务逻辑(可能慢、I/O 密集——调 API、跑代码)。
  • node_evaluation_loop:处理输出、算下游节点、入队(DB 写入)。
为什么分两个循环? 如果"跑 Block"和"算下游/入队"在同一个循环里排队,那么一个慢 Block(比如调 LLM 等 10 秒)会阻塞"算下游"的工作——即使别的 Block 已经产出了输出、下游本可以开跑,也得干等。分成两个循环:一边埋头跑 Block、一边同时处理已产出的输出并入队下游——两者并行,不互相阻塞,并行度更高。这是对"执行"和"调度"两种不同性质工作的解耦(跑 Block 是计算/IO,算下游是数据处理)。
L07

状态机与容错

执行状态用一个"合法跳转表" VALID_STATUS_TRANSITIONSdata/execution.py:131)约束——写进 SQL 的 where 子句做条件更新,实现并发安全的状态机。

# execution.py:131(概念)
VALID_STATUS_TRANSITIONS = {
    QUEUED: [INCOMPLETE, TERMINATED, REVIEW],
    RUNNING: [INCOMPLETE, QUEUED, ...],
    COMPLETED: [RUNNING], ...
}
# 更新状态时 SQL where 加"当前状态必须∈允许来源"——数据库层原子保证
为什么把状态机做进 SQL? 多个线程/进程可能同时想改同一条执行记录的状态。如果先读状态、判断、再写,中间可能被别人插队(竞态)——比如已 COMPLETED 又被改回 RUNNING。把"当前状态必须是合法来源"写进 SQL 的 where 条件,用数据库的原子更新保证——非法跳转的更新直接"落空"(影响 0 行)。这是"用数据库的原子性做并发状态机"的漂亮技巧。容错:图执行结束(含异常)时 _cleanup_graph_execution:1261)停掉所有在跑节点、清理队列残余、清临时文件。取消(:1495)只是 set 一个 cancel_event,主循环检测到就优雅 break。
L08

今日小结 + 动手

🧠 今天你应该能回答

  • 执行引擎的三大支撑?(MQ / 线程池 / 事件总线)
  • auto_ack=False 和 prefetch_count 保证什么?
  • ClusterLock 分布式锁防什么?
  • 调度循环的核心逻辑?为什么用两个事件循环?
  • 状态机为什么做进 SQL?

✋ 动手

P=autogpt_platform/backend/backend
grep -n 'class ExecutionManager\|ThreadPoolExecutor\|auto_ack\|prefetch_count\|ClusterLock' $P/executor/manager.py | head
sed -n '942,1030p' $P/executor/manager.py | head -50    # 调度循环
grep -n 'VALID_STATUS_TRANSITIONS' $P/data/execution.py
明天预告 · Day 13数据流与拓扑执行——输出如何精确地"喂"给下游输入(_enqueue_next_nodes)、"攒齐输入才入队"的 fan-in、静态连线回填。拓扑执行的最后一块拼图。
← Day 11 Graph 模型 Day 13 · 数据流 →