Day 20 / 共 60 天 · 阶段4 Crew 与流程
Process.sequential:任务如何一个接一个跑起来
Day 19 说 process 字段决定 kickoff 走哪条路。今天读最常用的那条——顺序流程。它的本质就是一句话:for task in tasks: 执行它。但真读源码你会发现里面藏着不少门道:上一个任务的输出怎么变成下一个的 context、异步任务为什么要"攒着最后一起 join"、条件任务怎么被跳过、执行日志何时写。核心方法是 _run_sequential_process → _execute_tasks(crew.py:1475/1519),我们一行行拆。
📍 你在 60 天里的位置(阶段4 Crew 与流程 · 共 8 天)
D19 Crew 全字段→
D20 顺序流程→
D21 层级+manager→
D22 kickoff 家族→
D23 规划→
D24 训练+replay→
D25 记忆开关→
D26 事件系统
💡 先用一个类比兜住今天
顺序流程就像工厂的一条流水线:工位 1 做完,半成品放到传送带上,工位 2 拿起来接着做,再传给工位 3……每个工位(task)拿到的"半成品"就是前面工位攒下来的产出(context)。唯一的花样是:有些工序可以并行开工(异步任务),流水线会先把它们都启动、最后到需要用的时候再统一"取货"(join)。今天就是读这条流水线的传送带逻辑。
L01
痛点:任务之间怎么"传接力棒"?
🤔 痛点你给 crew 三个任务:调研 → 写稿 → 校对。写稿的 agent 怎么知道调研 agent 查到了什么?它俩又不共享脑子。而且你还想让"查中文资料"和"查英文资料"同时进行省时间——框架怎么在"顺序"里塞进"并行"?如果第二个任务是"仅当第一个失败才执行"的条件任务,又怎么跳过?这些全是顺序流程主循环要解决的。
💡 一句话本质
顺序流程 =
_execute_tasks 里的一个 for task in tasks 循环。每一圈:① 算出这个 task 的 context(把之前的产出拼进去);② 同步任务就当场执行、等结果;异步任务就启动后先攒着(放进 futures),等碰到下一个同步任务或循环结束再统一 join;③ 把产出存进 task_outputs 供后面用。最后 _create_crew_output 拿最后一个任务的产出当 crew 的最终结果。入口极简(crew.py:1475):
# crew.py:1475
def _run_sequential_process(self) -> CrewOutput:
"""Executes tasks sequentially and returns the final output."""
return self._execute_tasks(self.tasks)
# crew.py:1479 对比:层级流程只是多了一步"造经理"(D21)
def _run_hierarchical_process(self) -> CrewOutput:
self._create_manager_agent()
return self._execute_tasks(self.tasks)
大白话两种流程最后都调同一个
_execute_tasks!区别只在层级流程会先 _create_manager_agent() 造个经理。所以 _execute_tasks 才是真正的"流水线引擎",今天的主角。顺序流程只是"不造经理,直接开跑"。L02
Process 枚举 + kickoff 里的分流
process.py 整个文件就 11 行(process.py:4):
# process.py:1
from enum import Enum
class Process(str, Enum):
"""Class representing the different processes that can be used to tackle tasks"""
sequential = "sequential" # 顺序(今天)
hierarchical = "hierarchical" # 层级(D21)
# TODO: consensual = 'consensual' # 共识——还没实现
kickoff 里靠它二选一(crew.py:1023,节选):
# crew.py:1023 在 kickoff 内部
if self.process == Process.sequential:
result = self._run_sequential_process()
elif self.process == Process.hierarchical:
result = self._run_hierarchical_process()
else:
raise NotImplementedError(
f"The process '{self.process}' is not implemented yet.")
class Process(str, Enum)★同时继承 str 和 Enum:所以 Process.sequential == "sequential" 成立,YAML/JSON 里直接写字符串就能反序列化成枚举。很实用的技巧。consensual 被注释掉共识流程还是 TODO。读源码看到 TODO 就知道"这功能官方在规划但没做"。if/elif/elsekickoff 里就这一处分流。else 抛 NotImplementedError——万一以后加了新枚举值却忘了写分支,会立刻炸出来提醒。💡 为什么 Process 要继承 str?用户在 YAML 配置里写
process: sequential(一个字符串),Pydantic 要能把它变成 Process.sequential 枚举。继承 str 后,枚举成员本身就是字符串,比较、序列化、日志打印都天然兼容。这是 Python 里"配置友好枚举"的标准写法。L03
_execute_tasks:流水线主循环
核心引擎(crew.py:1519,裁剪到主干):
# crew.py:1519
def _execute_tasks(self, tasks, start_index=0, was_replayed=False) -> CrewOutput:
custom_start = self._get_execution_start_index(tasks) # 断点续跑用(D24)
if custom_start is not None:
start_index = custom_start
task_outputs: list[TaskOutput] = []
futures: list[tuple[Task, Future[TaskOutput], int]] = [] # 攒异步任务
last_sync_output: TaskOutput | None = None
for task_index, task in enumerate(tasks): # ★流水线主循环
exec_data, task_outputs, last_sync_output = prepare_task_execution(
self, task, task_index, start_index, task_outputs, last_sync_output)
if exec_data.should_skip: # 断点前的任务:跳过(用存档)
continue
if isinstance(task, ConditionalTask): # 条件任务:判断要不要跳(L06)
skipped = self._handle_conditional_task(task, task_outputs, futures, task_index, was_replayed)
if skipped:
task_outputs.append(skipped)
continue
if task.async_execution: # 异步任务:启动后攒着
context = self._get_context(task, [last_sync_output] if last_sync_output else [])
future = task.execute_async(agent=exec_data.agent, context=context, tools=exec_data.tools)
futures.append((task, future, task_index))
else: # 同步任务:先 join 之前的异步,再当场执行
if futures:
task_outputs.extend(self._process_async_tasks(futures, was_replayed))
futures.clear()
context = self._get_context(task, task_outputs)
task_output = task.execute_sync(agent=exec_data.agent, context=context, tools=exec_data.tools)
task_outputs.append(task_output)
self._process_task_result(task, task_output)
self._store_execution_log(task, task_output, task_index, was_replayed)
if futures: # 收尾:还有没 join 的异步任务
task_outputs.extend(self._process_async_tasks(futures, was_replayed))
return self._create_crew_output(task_outputs)
task_outputs 列表★"传送带"本体:每完成一个任务就 append 进去。后面任务的 context 从这里取。futures 列表★存已启动但没等结果的异步任务(Future 对象)。攒起来延迟 join。for task_index, task in enumerate★主循环:严格按 tasks 列表顺序(下标)一个个来。这就是"sequential"的含义。prepare_task_execution每个任务开跑前的准备:算出该用哪个 agent、哪些工具、是否该跳过(断点续跑)。task.execute_sync同步:当场执行并阻塞等结果,拿到 task_output 立刻用。task.execute_async异步:启动就返回 Future,不等结果,塞进 futures 继续下一个任务。L04
接力棒:上游输出如何变成下游 context
传接力棒的关键在 _get_context(crew.py:1827):
# crew.py:1827
@staticmethod
def _get_context(task: Task, task_outputs: list[TaskOutput]) -> str:
if not task.context:
return "" # 任务没声明要 context → 空
return (
aggregate_raw_outputs_from_task_outputs(task_outputs) # 默认:把之前所有产出拼起来
if task.context is NOT_SPECIFIED
else aggregate_raw_outputs_from_tasks(task.context) # 显式指定:只拼指定的那几个任务
)
if not task.context任务把 context 设成空/None → 不给它喂任何上文(独立任务)。NOT_SPECIFIED(默认)★用户没显式设 context → 默认行为:把到目前为止所有任务的产出拼成一段文本喂进去。这就是"流水线自动传递"。else task.context用户显式指定 context=[task_a, task_b] → 只拼这两个任务的产出。精确控制依赖(Day 16 讲过)。aggregate_raw_outputs_*两个辅助函数负责把多个 TaskOutput 的 .raw 文本拼接成一段 context 字符串。📝 例子:调研 → 写稿 的接力
任务1「调研」产出
轮到任务2「写稿」,它没显式设 context(=NOT_SPECIFIED)→
raw="2024 年 AI 市场规模 $1.2 万亿",append 进 task_outputs。轮到任务2「写稿」,它没显式设 context(=NOT_SPECIFIED)→
_get_context 把任务1 的 raw 拼成 context 字符串 → 作为上文喂给写稿 agent。写稿 agent 于是"看得到"调研结果,能基于 $1.2 万亿这个数字写。接力棒就这么传过去了。💡 设计取舍①:默认"拼全部" vs 显式"指定依赖"
默认把之前所有产出都喂给下游,好处是零配置就能用——小 demo 不用操心依赖关系。但坏处是任务一多,context 越滚越长、烧 token、还可能塞进无关信息干扰模型。所以源码同时支持显式
context=[...]:认真做项目时,你精确声明"这个任务只依赖那两个",既省 token 又更聚焦。默认求方便,进阶求精确——两头都照顾到。L05
同步 vs 异步:为什么异步要"攒着最后 join"
异步任务的收割在 _process_async_tasks(crew.py:1918):
# crew.py:1918
def _process_async_tasks(self, futures, was_replayed=False) -> list[TaskOutput]:
task_outputs = []
for future_task, future, task_index in futures:
task_output = future.result() # ★阻塞:等这个异步任务真正跑完
task_outputs.append(task_output)
self._process_task_result(future_task, task_output)
self._store_execution_log(future_task, task_output, task_index, was_replayed)
return task_outputs
回看主循环里的关键分流(crew.py:1558):
# crew.py:1558
if task.async_execution:
context = self._get_context(task, [last_sync_output] if last_sync_output else [])
future = task.execute_async(...) # 启动,不等,塞 futures
futures.append((task, future, task_index))
else:
if futures: # ★碰到同步任务前,先把攒的异步 join 掉
task_outputs.extend(self._process_async_tasks(futures, was_replayed))
futures.clear()
... # 再执行这个同步任务
execute_async → future异步任务立刻返回一个 Future("取货凭证"),任务在后台线程跑,主循环继续往下。futures.append把凭证攒起来。此时还没有结果。else: if futures: join★关键:一旦碰到同步任务,说明"下面这步要用前面的结果了",必须先 _process_async_tasks 把攒的异步全部 future.result() 等回来,才能继续。future.result()阻塞直到该异步任务完成。多个异步任务已经并行跑过了,这里只是收割,所以省了总时间。循环末尾再 join 一次如果最后几个任务都是异步的(后面没同步任务触发 join),循环结束时补一次 join。💡 设计取舍②:为什么不立刻 join 每个异步任务?
如果异步任务一启动就
future.result() 等它,那和同步没区别、白瞎了并行。源码的策略是"能攒就攒,用时才等":连续的多个异步任务会同时在后台跑,直到遇到一个"需要它们结果"的同步任务,才统一收割。这样 N 个独立异步任务的耗时从"N 个相加"变成"最慢那一个"。延迟 join 是压榨并行度的经典手法。⚠️ 边界:末尾最多一个异步任务
Day 19 提过校验器
validate_end_with_at_most_one_async_task(crew.py:754):crew 的最后一个任务不能是"多个异步任务扎堆"。因为 _create_crew_output 要拿"最后一个有效产出"当 crew 结果——如果结尾一堆异步任务并行、谁先谁后不定,"最终结果是哪个"就说不清了。所以框架在构造期就限制:结尾的异步任务最多一个。L06
条件任务:主循环里的"要不要跳过"
条件任务在主循环里被特判(crew.py:1590):
# crew.py:1590
def _handle_conditional_task(self, task: ConditionalTask, task_outputs,
futures, task_index, was_replayed) -> TaskOutput | None:
if futures: # 先 join 掉攒着的异步——因为
task_outputs.extend(self._process_async_tasks(futures, was_replayed))
futures.clear() # 条件判断要看前面的真实产出
return check_conditional_skip(
self, task, task_outputs, task_index, was_replayed)
先 join futures★条件任务的 condition(前一个输出) 要基于真实产出判断,所以必须先把攒着的异步任务等回来,不能拿"半成品"做条件判断。check_conditional_skip调用条件函数。返回 None → 不跳过,继续正常执行;返回一个 TaskOutput(空产出)→ 跳过,把这个占位产出记进列表。主循环里 if skipped: continue拿到跳过标记就 continue,不执行这个任务本体。(Day 18 详细讲过 ConditionalTask 本身)💡 为什么条件判断也要先 join 异步?因为条件任务的"要不要跑"是数据驱动的——它看前一个任务的产出内容来决定。如果前一个是异步任务还没跑完,你拿到的是个空 Future,条件判断就成了瞎猜。所以任何"需要看真实产出"的节点(同步任务、条件任务)都会先强制 join。这是一条贯穿主循环的一致规则。
L07
收尾:产出汇总与执行日志
循环跑完,_create_crew_output 收尾(crew.py:1880,节选):
# crew.py:1880
def _create_crew_output(self, task_outputs: list[TaskOutput]) -> CrewOutput:
if not task_outputs:
raise ValueError("No task outputs available to create crew output.")
valid_outputs = [t for t in task_outputs if t.raw] # 过滤掉空产出(被跳过的条件任务)
if not valid_outputs:
raise ValueError("No valid task outputs available to create crew output.")
final_task_output = valid_outputs[-1] # ★最后一个有效产出 = crew 结果
...
self._drain_memory_writes() # 等后台记忆写完(D25)
crewai_event_bus.emit(self, CrewKickoffCompletedEvent(...)) # 发"完成"事件(D26)
return CrewOutput(
raw=final_task_output.raw,
pydantic=final_task_output.pydantic,
json_dict=final_task_output.json_dict,
tasks_output=task_outputs, # 全部任务产出也一并带上
token_usage=self.token_usage)
每个同步任务完成还会写执行日志(crew.py:1445):
# crew.py:1445 _store_execution_log(供 replay 用,D24)
log = {
"task": task,
"output": {"description": ..., "raw": output.raw, "agent": output.agent, ...},
"task_index": task_index,
"inputs": inputs,
"was_replayed": was_replayed,
}
self._task_output_handler.update(task_index, log) # 存进磁盘,replay 时读回
valid_outputs[-1]★crew 的最终结果 = 最后一个非空产出。所以流水线末尾那个任务的产出最重要。过滤 t.raw被跳过的条件任务产出是空的,用 if t.raw 过滤掉,避免拿空产出当最终结果。tasks_outputCrewOutput 里也保留全部任务产出,方便你事后逐个查看每步结果。_store_execution_log把每个任务的产出+inputs 存磁盘。Day 24 的 replay 就靠这些日志"从第 N 个任务重跑"。L08
取舍 + 边界 + 今日小结
图注:T2/T3 并行跑,遇到同步的 T4 前统一 join;最终 crew 结果 = 最后一个有效产出 T4。
👶 小白:顺序流程里 agent 之间能"互相商量"吗?
👨🏫 老师:不能直接商量。顺序流程是单向流水线——只能"上游产出 → 下游 context"这样往后传,下游改不了上游、agent 之间也不能来回讨论。想要"经理调度、来回委派",得用层级流程(明天 Day 21)。顺序流程胜在简单、可预测、便宜,绝大多数线性任务它就够了。
🧠 今天你应该能回答
- 顺序流程和层级流程最后都调用哪个方法?区别在哪一步?
task_outputs和futures两个列表各起什么作用?- 上一个任务的产出怎么变成下一个的 context?默认行为 vs 显式指定?
- 异步任务为什么"攒着最后 join"?遇到什么会触发提前 join?
- crew 的最终结果是哪个任务的产出?为什么要过滤
t.raw? - 为什么条件任务判断前也要先 join 异步?
✋ 10 分钟动手
P=lib/crewai/src/crewai
cat $P/process.py # 11 行的枚举
sed -n '1475,1590p' $P/crew.py # _execute_tasks 主循环
sed -n '1827,1836p' $P/crew.py # _get_context 接力棒
sed -n '1918,1934p' $P/crew.py # _process_async_tasks 收割
# 亲眼看顺序执行:两个任务,第二个用第一个的产出
python -c "
from crewai import Agent, Task, Crew, Process
a=Agent(role='研究员', goal='查资料', backstory='x', verbose=True)
b=Agent(role='作家', goal='写稿', backstory='x', verbose=True)
t1=Task(description='用一句话总结 Python 是什么', expected_output='一句话', agent=a)
t2=Task(description='基于上文写一句宣传语', expected_output='一句话', agent=b)
print(Crew(agents=[a,b], tasks=[t1,t2], process=Process.sequential).kickoff())
"
明日预告 · Day 21:今天说层级流程只比顺序多一步
_create_manager_agent()。明天就读这一步:经理 agent 是怎么被自动造出来的、为什么"经理不能带工具"、AgentTools 怎么把每个下属包装成一个"委派工具"让经理能调、层级流程为什么不需要给 task 指定 agent。