跳到主要内容

结构与任务图:Agent / Pipeline / Workflow 怎么调度任务

30 秒导读: Griptape 把「一次 AI 编排」拆成两层——Structure(编排层)和 Task(执行单元)。Structure 把若干 Task 连成一张有向图(DAG),再决定「按什么顺序、要不要并行」地跑它们。Agent / Pipeline / Workflow 是三个 Structure 子类,共享同一套 run 骨架,只在调度策略上分道扬镳:单任务、一条链、一张图。本章只讲这条「编排 Task 的执行模型」骨架;单个 Task 里 LLM 怎么循环,是 02 章 的事。

本章属于 Griptape 系列,总览与阅读地图见 index.md


1. 这是什么(零基础也能懂)

一句话定义: Structure 是一个「任务调度器」——你把要做的活拆成一个个 Task,声明它们谁先谁后,Structure 负责按依赖关系把它们跑完,并收集结果。

它解决什么问题。 真实的 AI 应用很少是「问一句、答一句」。更常见的是一条流水:先检索资料 → 再让模型总结 → 再翻译成三种语言 → 最后汇总。这些步骤有先后、有分叉、有汇合。Structure 就是用来描述并执行这种「多步、有依赖」的编排的。

三种编排,对应三种形状。 你不必自己写调度循环,选一个现成的 Structure 即可:

Structure任务形状白话典型场景
Agent单个任务就一个活,直接干一个带工具的对话智能体
Pipeline一条链A→B→C 顺序做检索→总结→翻译 的流水线
Workflow一张图有分叉汇合,能并行一份资料 fan-out 成多路加工再汇总

用起来什么样。 下面这段感受一下「声明图、然后 run」的手感:

# 示意,非源码:用 >> 运算符声明「谁接谁」,再交给 Workflow 跑
from griptape.structures import Workflow
from griptape.tasks import PromptTask

research = PromptTask("查资料:{{ args[0] }}", id="research")
summary = PromptTask("总结:{{ parents_output_text }}", id="summary")
translate = PromptTask("翻译成中文:{{ parents_output_text }}", id="translate")

research >> summary # research 是 summary 的父任务
research >> translate # 也是 translate 的父任务(分叉)

flow = Workflow(tasks=[research, summary, translate])
flow.run("griptape 是什么") # summary 和 translate 会并行跑

一句话直觉: 把 Structure 当流程引擎,把 Task 当流程里的一个节点。Structure 不关心节点内部干了啥(调 LLM?查数据库?),它只关心「这个节点的父节点跑完了没、该不该轮到它」。


2. 顶层全景(它大概怎么转)

2.1 两层四类

整个编排层就四个主角:一个抽象基类 Structure、它的三个子类、以及被调度的 BaseTask

部件职责文件
Structure抽象编排基类:装任务、补全图关系、跑 run 生命周期structures/structure.py
Agent只容一个 Task 的最简 Structurestructures/agent.py
Pipeline顺序链式调度structures/pipeline.py
Workflow拓扑排序 + 线程池并行调度structures/workflow.py
BaseTask执行单元:状态机、父子图关系、run 契约tasks/base_task.py

2.2 一次 run 大致走的路

不进代码,先看主线。三个子类的前后段完全一样,只有中间「怎么跑」不同:

Structure.run(*args) ← 模板方法,基类里写死

├─ before_run(args) ← 重置所有 Task、补全父子关系、发 Start 事件

├─ try_run(*args) ★ 抽象方法 ★ ← 唯一由子类实现的一步
│ ├─ Agent: 就 task.run()
│ ├─ Pipeline: 从头顺着链递归
│ └─ Workflow: 拓扑排序 + 线程池并行

└─ after_run() ← 写对话记忆、发 Finish 事件

怎么读: 上下两头(before_run / after_run)是共用骨架,中间那颗星(try_run)是唯一的变量。这正是「模板方法」模式——基类定好流程骨架,把会变的那一步留成抽象方法给子类填。

run 的骨架在 structure.py:200-208:先 before_run(args),再 try_run(*args),最后 after_run()try_run 是抽象方法(structure.py:229-231),三个子类各给一份实现。


3. 核心原理之一:Structure 怎么管住这张图

3.1 tasks 是「扁平化」出来的

你传进去的 tasks 可以嵌套列表(为了并行时书写方便),但内部要的是一份扁平清单。_tasks 存原始结构,tasks 属性把嵌套列表拍平:

# structure.py:62-71 —— tasks 属性
@property
def tasks(self) -> list[BaseTask]:
tasks = []
for task in self._tasks:
if isinstance(task, list):
tasks.extend(task) # 嵌套列表 → 摊开
else:
tasks.append(task)
return tasks

在这份扁平清单上,input_task 取第 0 个、output_task 取最后一个(structure.py:77-83),output 就是 output_task.output(structure.py:85-91)。注意 Workflow 覆盖了这两个属性,改用拓扑序的首尾(下面 3.5 讲)。

3.2 父子关系是「双向补全」的

图的边用两份 id 列表表示:每个 Task 有 parent_idschild_ids(base_task.py:39-40)。问题是——你可能只声明了一半(只说了 A 有个孩子 B,没说 B 有个父亲 A)。resolve_relationships 在跑之前把两边补齐:

# structure.py:137-152 —— 遍历每个 task,把缺的反向边补上
for task in self.tasks:
for parent_id in task.parent_ids: # 我声明了父亲
parent = task_by_id[parent_id]
if task.id not in parent.child_ids: # 但父亲没把我登记成孩子
parent.child_ids.append(task.id) # → 补上
for child_id in task.child_ids: # 反方向同理
child = task_by_id[child_id]
if task.id not in child.parent_ids:
child.parent_ids.append(task.id)

它同时兼任校验:遇到重复 id(structure.py:133-134)或指向不存在的 task(structure.py:141149),直接抛错。这一步在 before_run 末尾调用(structure.py:170),保证「开跑那一刻,图是完整且合法的」。

3.3 >> / <<:声明边的糖

上面那个 research >> summary 的写法,靠的是 BaseTask 重载的两个运算符:

写法含义底层
a >> ba 是 b 的父(a 先跑)__rshift__add_child
a << ba 是 b 的子(b 先跑)__lshift__add_parent

__rshift__ / __lshift__ 定义在 base_task.py:52-66,内部调 add_child / add_parent(base_task.py:112-138)——后者同样是双向登记(既加自己的 child_ids,也加对方的 parent_ids),并在必要时把新任务自动 add_task 进 Structure。

3.4 run 生命周期:重置、发事件、写记忆

before_run(structure.py:154-170)做三件事:把 execution_args 存下、重置每个 Task([task.reset() for task in self.tasks],让上一轮的状态清零)、发 StartStructureRunEvent,最后 resolve_relationships()

after_run(structure.py:172-194)在结尾按 conversation_memory_strategy 决定要不要落对话记忆,再发 FinishStructureRunEvent:

  • per_structure(默认):整轮跑完,把 input_task.input → output_task.output 打包成一条 Run 存进 conversation_memory(structure.py:177-185)。
  • per_task:这里不做结构级落库,改由每个 Task 自己管(细节属于 05 章 记忆与制品)。

3.5 run_stream:后台线程 + 事件队列做流式

run 是阻塞的——跑完才返回。想「边跑边拿事件」(比如实时显示 LLM 吐字),用 run_stream(structure.py:210-227)。它的技巧是用一个后台线程跑 run,主线程从队列里捞事件:

# structure.py:217-227 —— 简化后的流式骨架
with EventListener(self._event_queue.put, event_types=event_types):
t = Thread(target=with_contextvars(self.run), args=args) # 后台跑 run
t.start()
while True:
event = self._event_queue.get() # 主线程阻塞取事件
if isinstance(event, FinishStructureRunEvent) and event.structure_id == self.id:
break # 收到收尾事件 → 停
else:
yield event # 其余事件流式吐出
t.join()

EventListener 把事件塞进 self._event_queue,while 循环不断 getyield,直到看见本 Structure 的 FinishStructureRunEvent 才收尾。with_contextvars 保证后台线程继承主线程的上下文变量(否则子线程里 EventBus 等 contextvar 会丢)。


4. 核心原理之二:BaseTask 的执行契约

Structure 负责「排班」,但「一个 Task 该不该跑、跑出错怎么办」由 BaseTask 自己定。这套契约是三种调度都依赖的地基。

4.1 四个状态

# base_task.py:31-35
class State(Enum):
PENDING = 1 # 还没轮到
RUNNING = 2 # 正在跑
FINISHED = 3 # 跑完了(成功或失败都算)
SKIPPED = 4 # 被跳过(父任务全被跳过时)

关键坑:FINISHED 不等于成功。下面会看到,即使 try_run 抛异常,状态照样进 FINISHED,只是 output 变成 ErrorArtifact

4.2 can_run:轮到我了吗

这是调度的核心裁决——一个 Task 满足什么条件才能跑:

# base_task.py:203-216
def can_run(self) -> bool:
if self.is_skipped() or not self.is_pending(): # 已跳过 / 不是 PENDING → 不跑
return False
if self.parents and all(parent.is_skipped() for parent in self.parents):
self.state = BaseTask.State.SKIPPED # 父任务全被跳过 → 我也跳过
return False
unskipped_parents = [p for p in self.parents if not p.is_skipped()]
return all(parent.is_finished() for parent in unskipped_parents) # 未跳过的父全 FINISHED → 可跑

三条判据,依次收紧:

  1. 自己得是 PENDING(没跑过、没被跳过)。
  2. 如果父任务全被跳过,自己也自动转 SKIPPED(跳过会沿图传染)。
  3. 否则,只要所有「没被跳过的」父任务都 FINISHED,就轮到我。

第 3 条用 unskipped_parents 而非全部父任务——这样图里某条分支被跳过时,汇合点不会被卡死。

4.3 run:把异常包成 ErrorArtifact

Task 的 run(base_task.py:170-188)是「不抛异常」的——它用 try/except/finally 把失败转成数据:

# base_task.py:170-188 —— 失败不炸,转成 ErrorArtifact
def run(self, *args) -> T:
try:
self._execution_args = args
self.state = BaseTask.State.RUNNING
self.before_run()
self.output = self.try_run() # ★ 真正干活,子类实现
self.after_run()
except Exception as e:
logger.exception(...)
self.output = cast("T", ErrorArtifact(str(e), exception=e)) # 异常 → 制品
finally:
self.state = BaseTask.State.FINISHED # 不论成败,一律 FINISHED
return self.output

这个设计是整套调度能容错的前提:一个 Task 崩了,不会把整个线程/递归炸掉,而是返回一个 ErrorArtifact。上层调度就靠检查返回值是不是 ErrorArtifact 来决定要不要中断(见第 5 节的 fail_fast)。

注意 try_run(base_task.py:225-227)在 BaseTask 里是抽象的——这跟 Structure 是同一个套路:run 是共用模板,try_run 是子类填的空。PromptTasktry_run 里那圈 LLM/ReAct 循环,是 02 章 的内容。

4.4 full_context:任务能看见什么

模板里 {{ parents_output_text }} 这类变量,数据来自 full_context(base_task.py:229-238):深拷贝自身 context,再并入 structure.context(self)。而 structure.context 被各子类覆盖,喂的料不同——这正是 Pipeline 和 Workflow 让下游拿到上游产物的通道:

Structurecontext 额外塞的键依据
Structure(基类)argsstructurestructure.py:127-128
Pipelineparent_outputparentchildtask_outputspipeline.py:57-69
Workflowparent_outputsparents_output_textparentschildrenworkflow.py:125-138

5. 核心原理之三:三种调度的差异

前四节讲的都是共用地基。真正让三个子类不同的,是各自那份 try_run

5.1 Agent:就一个任务

Agent 是退化情形——只容一个 Task。try_run 一行:

# agent.py:89-93
@observable
def try_run(self, *args) -> Agent:
self.task.run()
return self

它用一堆约束把「单任务」这件事焊死:add_tasks 收到 >1 个任务直接抛错(agent.py:84-87);add_taskclear() 再塞,保证永远只有一个(agent.py:75-82);fail_fast 被强制为 False,验证器里只要为 True 就抛「Agents cannot fail fast」(agent.py:42-45)。构造时若你没给任务,_init_task 会用默认 prompt driver 自动兜一个 PromptTask(agent.py:95-119)。

5.2 Pipeline:顺着链递归

Pipeline 把任务连成单链:add_task 时自动把新任务挂到当前 output_task 后面(pipeline.py:23-25)。跑的时候从头递归,一路取「第一个孩子」往下走:

# pipeline.py:71-74 —— 顺序递归
def __run_from_task(self, task: BaseTask | None) -> None:
if task is None or isinstance(task.run(), ErrorArtifact) and self.fail_fast:
return
self.__run_from_task(next(iter(task.children), None)) # 只取第一个孩子 → 线性

两个要点:

  1. 中断条件。因为 andor 结合更紧,判断实为 task is None or (返回 ErrorArtifact 且 fail_fast)——任务跑出 ErrorArtifact 且开了 fail_fast,链就此打住。
  2. 线性形状next(iter(task.children), None) 只取第一个孩子,所以 Pipeline 本质是一条线,不处理分叉。

5.3 Workflow:拓扑排序 + 线程池并行

Workflow 是唯一真正的 DAG 调度器。它不预设链形,而是每一轮重新算拓扑序、把当前能跑的任务一次性提交给线程池:

# workflow.py:102-123 —— 并行调度主循环
with self.create_futures_executor() as futures_executor:
while not self.is_finished() and not exit_loop:
futures_list = {}
ordered_tasks = self.order_tasks() # ① 拓扑排序
for task in ordered_tasks:
if task.can_run(): # ② 挑出就绪的
future = futures_executor.submit(with_contextvars(task.run)) # ③ 并行提交
futures_list[future] = task
for future in futures.as_completed(futures_list): # ④ 等这批完成
if isinstance(future.result(), ErrorArtifact) and self.fail_fast:
exit_loop = True # ⑤ 出错且 fail_fast → 中断
break

拓扑序怎么来的?靠 graphlib.TopologicalSorter:

# workflow.py:152-153
def order_tasks(self) -> list[BaseTask]:
return [self.find_task(tid) for tid in TopologicalSorter(self.to_graph()).static_order()]

to_graph(workflow.py:140-150)构造 {任务 id → 它的父任务 id 集合} 的映射(把每个把本任务当孩子的任务收成「前驱」),正好是 TopologicalSorter 要的「节点→前驱」格式。

为什么用 while 外层循环包着? 因为一批任务跑完后,原本被卡住的下游可能刚好就绪。每轮 while 重新 order_tasks + 挑 can_run 的,直到 is_finished()(所有任务都 can_run() 为假,structure.py:101-102)。fail_fast 命中时置 exit_loop 提前跳出。

5.4 一张图看懂 Workflow 的并行

以「一份资料 fan-out 成两路加工、再汇总」的菱形为例(A→B、A→C、B→D、C→D):

┌──────────────┐
│ A 查资料 │ 第 1 轮:只有 A 无父,can_run
└──────┬───────┘ → 提交 1 个 future
┌─────┴─────┐
▼ ▼
┌───────────┐ ┌───────────┐
│ B 总结 │ │ C 翻译 │ 第 2 轮:A 已 FINISHED
└─────┬─────┘ └─────┬─────┘ → B、C 同时就绪
│ │ → 一次提交 2 个 future 并行跑
└──────┬──────┘

┌──────────────┐
│ D 汇总 │ 第 3 轮:B、C 都 FINISHED
└──────────────┘ → D 就绪,提交 1 个 future

怎么读: 竖向是数据流向(上游→下游);同一横排 = 同一轮里被并行提交can_run 决定谁进得了这一轮,futures.as_completed 等这一轮全部落地,再进下一轮。同样这张菱形图,Pipeline 只能取「第一个孩子」,走不出并行,也接不住 D 的双父汇合。

5.5 三者对照

维度AgentPipelineWorkflow
任务数恰好 1多个,单链多个,任意 DAG
调度方式直接 task.run()递归取第一个孩子拓扑排序 + 线程池
并行有(每轮就绪任务并发)
fail_fast 默认强制 FalseTrueTrue
中断判据返回 ErrorArtifactfail_fast同左,置 exit_loop
依据agent.py:89-93pipeline.py:71-74workflow.py:102-123

6. 边界与局限

  • 本章不碰 Task 内部。 PromptTask.try_run 里的提示词组装、ReAct 子任务循环,是 02 章;具体 Task 子类的业务只在这里被当作「实现了 try_run 契约的黑盒」。
  • Pipeline 不处理分叉。 __run_from_task 只顺第一个孩子,一个任务挂了多个孩子时,只有第一个会被走到——多路并行/汇合请用 Workflow
  • fail_fast 只看返回值类型。 中断靠「try_run 返回的是不是 ErrorArtifact」判断(base_task.py:184)。若某 Task 内部吞掉了异常、返回了正常制品,fail_fast 就拦不住——容错的边界取决于 Task 自己怎么处理异常。
  • run_stream 起后台线程。 事件顺序依赖队列,且 try_run 必须靠 with_contextvars 才能在子线程里继承 contextvar,否则 EventBus 等上下文会丢。

7. 横向对比

同为 agent 框架,「编排层怎么表达任务依赖」各有取舍。Griptape 的特点是:把「图的定义」和「图的执行」彻底分开——BaseTask 只持有 parent_ids/child_ids 两份 id 列表当边,三个 Structure 各写一份遍历策略。相比「用装饰器/回调隐式串联」的框架,Griptape 的依赖是显式、可序列化的(边就是数据),代价是要手动声明 >>/<<。跨库原理对比见总库 doc 的「任务编排」条目。


8. 代码地图(导航索引)

主题文件路径符号名
编排基类 / run 模板方法structures/structure.pyStructure.runStructure.try_run(abstract)
tasks 扁平化structures/structure.pyStructure.tasks
双向补全父子关系structures/structure.pyStructure.resolve_relationships
run 生命周期structures/structure.pyStructure.before_runStructure.after_run
流式(后台线程+队列)structures/structure.pyStructure.run_stream
对话记忆策略structures/structure.pyStructure.conversation_memory_strategy
单任务调度structures/agent.pyAgent.try_runAgent.add_task
顺序递归调度structures/pipeline.pyPipeline.__run_from_task
拓扑并行调度structures/workflow.pyWorkflow.try_runWorkflow.order_tasksWorkflow.to_graph
任务状态机tasks/base_task.pyBaseTask.StateBaseTask.can_run
任务执行契约tasks/base_task.pyBaseTask.runBaseTask.try_run(abstract)
图关系运算符tasks/base_task.pyBaseTask.__rshift__BaseTask.__lshift__add_childadd_parent
任务可见上下文tasks/base_task.pyBaseTask.full_context
模板方法钩子mixins/runnable_mixin.pyRunnableMixin.before_runRunnableMixin.after_run