执行引擎 manager
整个平台的心脏。今天读 executor/manager.py:ExecutionManager 的进程架构(MQ + 线程池 + 事件总线)、一次图执行的生命周期、分布式锁去重。
👶 小白:都用队列排队了,为什么还要一把"分布式锁"?队列不是已经保证顺序了吗?
👨🏫 老师:因为平台会开多个执行进程同时从队列抢活,可能两个进程几乎同时抓到同一个节点,导致同一步被跑两遍(既浪费钱又可能重复发邮件)。ClusterLock 就是"这活我领了"的派工牌——谁先抢到锁谁干,别人看到已锁就跳过。队列管"排队",锁管"同一件事只被做一次",两者解决的是不同问题。
执行引擎全景
manager.py:1353):盯着 RabbitMQ 队列,来一条"跑图请求"就从线程池(executor 属性返回 ThreadPoolExecutor,manager.py:1395-1397)派一个工人线程去跑;每个工人由 ExecutionProcessor(manager.py:551)驱动,一张图从头到尾走 on_graph_execution(:812),图里每个节点走 on_node_execution(:578)。工头只管调度,不亲自干活。ExecutionManager(manager.py:1353)是一个常驻进程。它由三样东西支撑:
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=False 和 prefetch_count 是两个关键的 MQ 可靠性设置。前者保证"任务处理完才确认"——中途崩溃消息不丢、会被别的 worker 重新拿。后者保证"一个 worker 只预取它能处理的量"——不会一个 pod 贪心拿走所有任务却跑不过来。事件总线(发布-订阅)
基于 Redis 的 RedisExecutionEventBus——每次节点/图状态变化就 publish 事件到两个频道:exec/{graph_exec_id}(这次执行)和 {graph_id}/all(这个图的所有执行)。
分布式锁去重(ClusterLock)
收到运行消息后(_handle_run_message,manager.py:1529),先尝试获取一把 ClusterLock(Redis 分布式锁,:1670)——保证同一个 graph_exec_id 全集群只在一个 pod 上跑。
_on_run_done(:1717)ack 消息、释放锁。图执行调度循环(核心)
_on_graph_execution(manager.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)
取[读网页]→预扣费→异步执行→产出"网页文本"→喂给[总结],[总结]输入齐→入队取[总结]→执行→产出"摘要"→喂给[发邮件],入队取[发邮件]→执行→发送→无下游 → 队列空且无在跑节点 → 循环退出,图执行完成。循环本身不认识"读网页/总结/发邮件",它只做"取一个→跑→把产出喂下游→重复"。两个事件循环(并行度优化)
调度用了两个独立的 asyncio 事件循环线程(manager.py:794):
node_execution_loop:跑 Block 的业务逻辑(可能慢、I/O 密集——调 API、跑代码)。node_evaluation_loop:处理输出、算下游节点、入队(DB 写入)。
状态机与容错
执行状态用一个"合法跳转表" VALID_STATUS_TRANSITIONS(data/execution.py:131)约束——写进 SQL 的 where 子句做条件更新,实现并发安全的状态机。
# execution.py:131(概念)
VALID_STATUS_TRANSITIONS = {
QUEUED: [INCOMPLETE, TERMINATED, REVIEW],
RUNNING: [INCOMPLETE, QUEUED, ...],
COMPLETED: [RUNNING], ...
}
# 更新状态时 SQL where 加"当前状态必须∈允许来源"——数据库层原子保证
_cleanup_graph_execution(:1261)停掉所有在跑节点、清理队列残余、清临时文件。取消(:1495)只是 set 一个 cancel_event,主循环检测到就优雅 break。今日小结 + 动手
🧠 今天你应该能回答
- 执行引擎的三大支撑?(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