Day 20 / 共 60 天 · 阶段4 Crew 与流程

Process.sequential:任务如何一个接一个跑起来

Day 19 说 process 字段决定 kickoff 走哪条路。今天读最常用的那条——顺序流程。它的本质就是一句话:for task in tasks: 执行它。但真读源码你会发现里面藏着不少门道:上一个任务的输出怎么变成下一个的 context、异步任务为什么要"攒着最后一起 join"、条件任务怎么被跳过、执行日志何时写。核心方法是 _run_sequential_process → _execute_taskscrew.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 里就这一处分流。elseNotImplementedError——万一以后加了新枚举值却忘了写分支,会立刻炸出来提醒。
💡 为什么 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_contextcrew.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「调研」产出 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_taskscrew.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_taskcrew.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

取舍 + 边界 + 今日小结

控制流:一条流水线里同步/异步/条件三种任务 T1 同步·调研 T2 异步·查中文 T3 异步·查英文 join T2/T3 (遇同步前收割) T4 同步·写稿 CrewOutput =T4产出 task_outputs 传送带:T1产出 → 累积 → 每个下游任务的 context 从这里取 [T1.raw, T2.raw, T3.raw, T4.raw] → valid_outputs[-1] = 最终结果
图注:T2/T3 并行跑,遇到同步的 T4 前统一 join;最终 crew 结果 = 最后一个有效产出 T4。

👶 小白:顺序流程里 agent 之间能"互相商量"吗?

👨‍🏫 老师:不能直接商量。顺序流程是单向流水线——只能"上游产出 → 下游 context"这样往后传,下游改不了上游、agent 之间也不能来回讨论。想要"经理调度、来回委派",得用层级流程(明天 Day 21)。顺序流程胜在简单、可预测、便宜,绝大多数线性任务它就够了。

🧠 今天你应该能回答

  • 顺序流程和层级流程最后都调用哪个方法?区别在哪一步?
  • task_outputsfutures 两个列表各起什么作用?
  • 上一个任务的产出怎么变成下一个的 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。
← Day 19 Crew 全字段 Day 21 · 层级流程 →