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

train 与 replay:把团队"调教"好,出错能"从中间重跑"

前面几天都在讲"怎么跑"。今天讲两个偏"运维/调优"的能力:train(训练)——让 crew 反复跑同一批输入、每步收集你的人工反馈,最后为每个 agent 生成一份"以后该怎么做得更好"的改进建议存起来;replay(回放)——上次跑到第 5 个任务挂了,别从头重来,直接从第 5 个任务接着跑。两者都依赖同一个基础设施:_task_output_handler 把每个任务的产出存档到磁盘。核心在 crew.pytraincrew.py:914)和 replaycrew.py:1944)。

📍 你在 60 天里的位置(阶段4 Crew 与流程 · 共 8 天)
D19 Crew 全字段 D20 顺序流程 D21 层级+manager D22 kickoff 家族 D23 规划 D24 训练+replay D25 记忆开关 D26 事件系统
💡 先用一个类比兜住今天 train带新人做几遍演练:同一个任务反复做 N 遍,每做完你在旁边点评"这里该更简洁""数据要标来源",最后把你所有点评提炼成一份"给这位同事的成长建议"存进档案,以后他上岗就带着这份建议。replay视频从断点续播:一场三小时的会议录像,你不用从头看,直接拖到"上次看到的地方"接着看——因为前面每段都存了档。今天就看这两套机制怎么用同一个"存档"底座实现。
L01

痛点:agent 不稳定,且失败重跑太贵

🤔 痛点两个真实烦恼:① agent 每次输出风格/质量飘忽,你想"教"它按你的偏好来,但又不可能真去微调模型权重(太贵);② 一个 5 任务的 crew,跑到第 4 个任务因为网络错崩了——难道要把前 3 个(花了钱花了时间的)任务全部重跑?前 3 个的产出明明是好的。这两个诉求,分别对应 train 和 replay。
💡 一句话本质 train = "反复跑 + 收人工反馈 + 生成改进建议文件"(不改模型权重,改的是"喂给 agent 的建议",属于 prompt 级调教);replay = "读上次的任务产出存档,从指定任务往后重跑"(前面的任务用存档结果重建 context,不重新执行)。两者的公共底座是 _task_output_handler:每个任务执行完,它的产出+inputs 都被 _store_execution_log(Day 20 见过)存到磁盘
大白话注意 train 里的"训练"和你想的深度学习不是一回事。它不动模型参数——它做的是"让你当老师点评几遍,把点评整理成一段文字建议,以后拼进 agent 的提示里"。所以它便宜、快、你自己就能做。本质是"人类反馈 → 文字建议 → 注入 prompt"。
L02

train 入口 + 训练前的开关设置

训练前先把 crew 切到"训练模式"(crew.py:901):

# crew.py:901
def _setup_for_training(self, filename: str) -> None:
    """Sets up the crew for training."""
    self._train = True
    for task in self.tasks:
        task.human_input = True          # ★每个任务都强制要人工输入(点评)
    for agent in self.agents:
        agent.allow_delegation = False    # ★训练时禁止委派——要单独评估每个 agent
    CrewTrainingHandler(TRAINING_DATA_FILE).initialize_file()   # 建原始反馈文件
    CrewTrainingHandler(filename).initialize_file()             # 建最终建议文件

train 主入口(crew.py:914,节选头部):

# crew.py:914
def train(self, n_iterations: int, filename: str, inputs: dict | None = None) -> None:
    """Trains the crew for a given number of iterations."""
    inputs = inputs or {}
    try:
        crewai_event_bus.emit(self, CrewTrainStartedEvent(
            crew_name=self.name, n_iterations=n_iterations, filename=filename, inputs=inputs))
        train_crew = self.copy()                 # ★在副本上训练,不污染原 crew
        train_crew._setup_for_training(filename)
        ...
self._train = True打开训练标志。执行时 agent 会知道"现在是训练模式,产出后要收反馈"。
task.human_input = True★强制每个任务结束后暂停等你输入点评(Day 14 的 human_input 机制)。这就是"人工反馈"的来源。
allow_delegation = False★训练时关掉委派——因为要独立评估每个 agent 自己的表现,委派会让功劳/问题混在一起。
两个 CrewTrainingHandler建两个文件:TRAINING_DATA_FILE原始反馈filename最终提炼的建议
train_crew = self.copy()★和 kickoff_for_each 一样:训练在副本上做。训练会改一堆开关(human_input/allow_delegation),不能污染你原来的 crew。
L03

训练主循环:跑 N 遍、攒反馈

train 的核心循环(crew.py:932):

# crew.py:932
            for n_iteration in range(n_iterations):        # ★跑 N 遍
                train_crew._train_iteration = n_iteration
                train_crew.kickoff(inputs=inputs)          # 每遍就是一次普通 kickoff

            training_data = CrewTrainingHandler(TRAINING_DATA_FILE).load()   # 读回攒的原始反馈
            ...
for n_iteration in range(n_iterations)★训练 = 用同一批 inputs 反复 kickoff N 次。每次每个任务都会停下来等你点评。
_train_iteration = n记下这是第几遍,反馈存档时按迭代号归类。
train_crew.kickoff(inputs)★训练里的每一遍就是一次标准 kickoff(Day 22)——没有独立的"训练执行引擎",完全复用。区别只在 human_input=True 会收集反馈。
CrewTrainingHandler(...).load()N 遍跑完,把这 N 遍里你给的所有原始反馈从磁盘读回来,准备提炼。
💡 为什么训练要"跑同一批输入好几遍"?因为一遍的反馈太单薄——你对某个 agent 可能只点评一次,样本不足。反复跑几遍,你会在不同迭代里针对同一个 agent 给出多角度的点评(这次嫌啰嗦、下次嫌漏了来源……)。多遍的反馈汇总起来,才能提炼出一份相对全面的改进建议n_iterations 就是你愿意花多少遍来调教的预算。
L04

评估:把反馈提炼成"给每个 agent 的建议"

训练的产出在这一步生成(crew.py:938):

# crew.py:938
            for agent in train_crew.agents:
                if training_data.get(str(agent.id)):
                    result = TaskEvaluator(agent).evaluate_training_data(   # ★LLM 提炼反馈
                        training_data=training_data, agent_id=str(agent.id))
                    CrewTrainingHandler(filename).save_trained_data(        # 存进最终文件
                        agent_id=str(agent.role),
                        trained_data=result.model_dump())
            crewai_event_bus.emit(self, CrewTrainCompletedEvent(
                crew_name=self.name, n_iterations=n_iterations, filename=filename))
    except Exception as e:
        crewai_event_bus.emit(self, CrewTrainFailedEvent(error=str(e), crew_name=self.name))
        self._logger.log("error", f"Training failed: {e}", color="red")
        CrewTrainingHandler(TRAINING_DATA_FILE).clear()   # ★失败清空,别留半截脏数据
        CrewTrainingHandler(filename).clear()
        raise
for agent in agents逐个 agent 提炼——每个 agent 得到属于自己的一份建议。
TaskEvaluator(agent).evaluate_training_data★用一次 LLM 调用,把你对这个 agent 的零散反馈提炼成结构化的改进建议("更简洁、标注来源、避免……")。
save_trained_data(agent_id=agent.role, ...)★存档,按 role 存(不是随机 id)——这样以后新建同 role 的 agent 也能匹配到建议。
CrewTrainCompletedEvent发训练完成事件(Day 26)。日志/UI 靠事件知道训练结束。
except: clear() + raise★训练出错:清空两个文件(避免留下半截无效数据误导后续)+ 发失败事件 + 抛出。事务性的清理。
💡 设计取舍①:建议按 role 存,而不是按 id 注意 save_trained_data(agent_id=str(agent.role), ...)——键用的是 role("资深分析师")而非随机 id。为什么?因为训练是离线一次性的,而使用建议是之后另起进程的——那时 agent 的随机 id 早变了(Day 19 讲过 id 每次新建都不同)。用 role 当键,才能让"下次跑的、同角色的 agent"通过 trained_agents_file(Day 19 字段)加载到这份建议。持久化的关联键必须选跨进程稳定的标识——role 稳,id 不稳。
L05

replay:从存档里"接着跑"

replay 主流程(crew.py:1944):

# crew.py:1944
def replay(self, task_id: str, inputs: dict | None = None) -> CrewOutput:
    """Replay the crew execution from a specific task."""
    stored_outputs = self._task_output_handler.load()      # ① 读磁盘上的任务产出存档
    if not stored_outputs:
        raise ValueError(f"Task with id {task_id} not found in the crew's tasks.")

    start_index = self._find_task_index(task_id, stored_outputs)   # ② 找到从哪个任务开始
    if start_index is None:
        raise ValueError(f"Task with id {task_id} not found in the crew's tasks.")

    replay_inputs = inputs if inputs is not None else stored_outputs[start_index]["inputs"]
    self._inputs = replay_inputs
    if replay_inputs:
        self._interpolate_inputs(replay_inputs)             # 用当时/新的 inputs 插值

    if self.process == Process.hierarchical:
        self._create_manager_agent()                        # 层级流程要重新造经理

    for i in range(start_index):                            # ③ 重建断点之前的任务产出
        stored_output = stored_outputs[i]["output"]
        task_output = TaskOutput(
            description=stored_output["description"], agent=stored_output["agent"],
            raw=stored_output["raw"], pydantic=stored_output["pydantic"],
            json_dict=stored_output["json_dict"], output_format=stored_output["output_format"],
            messages=stored_output.get("messages", []))
        self.tasks[i].output = task_output                  # 把存档结果填回任务,不重跑
    self._logging_color = "bold_blue"
    return self._execute_tasks(self.tasks, start_index, True)   # ④ 从 start_index 往后真跑
_task_output_handler.load()★读回上次 kickoff 时 _store_execution_log 存的每个任务产出(Day 20 埋的伏笔)。没存档就没法 replay。
_find_task_index(task_id, ...)按你给的 task_id 在存档里找到它是第几个任务 → start_index
for i in range(start_index)关键:断点之前的任务不重跑,而是从存档里重建 TaskOutput 填回 task.output。这样后续任务的 context 拿得到它们。
_execute_tasks(tasks, start_index, True)★复用 Day 20 引擎,但从 start_index 开始,且 was_replayed=True。前面的靠 should_skip 跳过(用存档)。
L06

定位任务 + 存档格式

怎么按 task_id 找到位置(crew.py:1934):

# crew.py:1934
@staticmethod
def _find_task_index(task_id: str, stored_outputs: list[Any]) -> int | None:
    return next(
        (index for (index, d) in enumerate(stored_outputs)
         if d["task_id"] == str(task_id)),      # 在存档里逐个比对 task_id
        None)                                    # 没找到返回 None

存档的内容(回看 Day 20 的 _store_execution_logcrew.py:1445):

# crew.py:1445
log = {
    "task": task,
    "output": {"description": ..., "raw": output.raw, "pydantic": output.pydantic,
               "json_dict": output.json_dict, "agent": output.agent, "messages": ...},
    "task_index": task_index,
    "inputs": inputs,                # ★连当时的 inputs 也存了
    "was_replayed": was_replayed,
}
self._task_output_handler.update(task_index, log)   # 写进磁盘存档
next((...), None)Python 惯用法:找到第一个匹配的 index 就返回,找不到返回默认值 None。简洁的"查找或空"。
d["task_id"] == str(task_id)按 task_id 精确匹配。所以你 replay 时要传的是任务的 id(可以从上次日志里拿到)。
存档含 output + inputs★存档不只存产出,还存当时的 inputs——所以 replay 能复原"当时是用什么输入跑的",插值才对得上。
_task_output_handler.update每个同步任务完成后写一次。这是 train/replay 共用的持久化底座。
⚠️ 边界:replay 依赖"上一次真的存了档" replay 第一步 stored_outputs = load(),如果为空直接 raise。也就是说——你只能 replay "之前正常跑过、留下了存档"的执行。如果你换了台机器、或存档文件被清了(比如 kickoff_for_each 结尾会 _task_output_handler.reset(),Day 22),就没得 replay。而且如果你 replay 前改了任务列表的结构(加删任务),存档的 index 就对不上了,行为不可预期。replay 是"同一份任务定义 + 有存档"前提下的续跑,不是万能时间机器。
L07

取舍与边界

💡 设计取舍②:train/replay 为什么都"复用 kickoff + _execute_tasks"? 你会发现一个反复出现的模式:train 的每一遍是 kickoff,replay 的执行是 _execute_tasks——没有一套独立的"训练引擎"或"回放引擎"。这是刻意的:如果训练/回放走另一套代码,就会和正常执行行为不一致(正常跑对了、训练时却因为引擎不同出别的 bug)。复用同一个执行核心,保证"训练时看到的行为 = 真实上线的行为"。代价是要在核心引擎里塞一些 was_replayedshould_skip 之类的分支参数,但换来行为一致性,非常值。
数据结构:一个存档底座,喂养 train 与 replay _task_output_handler 磁盘存档:每任务 output+inputs train:跑 N 遍 kickoff 收人工反馈→TaskEvaluator→建议文件 replay:读存档定位断点 重建前置 output→从断点 _execute_tasks
图注:train 靠反复 kickoff 攒反馈,replay 靠读存档续跑,两者共享 _task_output_handler 存档。

👶 小白:train 之后 agent 真的"变聪明"了吗?下次自动生效?

👨‍🏫 老师:不是自动的。train 只是生成了一个建议文件。要让它生效,你得在建 crew 时设 trained_agents_file="那个文件"(Day 19 的字段),框架才会在推理时把建议拼进 agent 的提示。而且它没改模型权重——所谓"变聪明"是"提示里多了针对性建议",属于 prompt 工程,不是真微调。效果有限但零成本、可解释。

L08

今日小结

💡 一句话收束 traincopy() 副本 → _setup_for_training(human_input=True + 禁委派)→ 跑 N 遍收反馈 → TaskEvaluatorrole 提炼成建议文件。replay_task_output_handler.load() 读存档 → _find_task_index 定位断点 → 重建前置 task.output → 从断点 _execute_tasks。二者共享"每任务产出存磁盘"的底座,且都复用正常执行引擎保证行为一致。

🧠 今天你应该能回答

  • train 的"训练"改的是什么?是模型权重吗?
  • 训练前 _setup_for_training 为什么要开 human_input、关委派?
  • 为什么训练要跑 N 遍?建议为什么按 role 存而非 id?
  • replay 从哪读数据?断点之前的任务重跑吗?
  • replay 依赖什么前提?什么情况下会失败?
  • train/replay 为什么都复用 kickoff / _execute_tasks?

✋ 10 分钟动手

P=lib/crewai/src/crewai
sed -n '901,965p'   $P/crew.py     # _setup_for_training + train 全流程
sed -n '1934,1982p' $P/crew.py     # _find_task_index + replay
sed -n '1445,1475p' $P/crew.py     # _store_execution_log 存档格式
cat $P/utilities/training_handler.py   # CrewTrainingHandler
# 体验 replay(需先正常跑一次留下存档,再用某个 task_id replay)
python -c "
from crewai import Agent, Task, Crew
a=Agent(role='作家', goal='写作', backstory='x')
t1=Task(description='写一句关于春天的话', expected_output='一句', agent=a)
t2=Task(description='基于上文再写一句', expected_output='一句', agent=a)
c=Crew(agents=[a], tasks=[t1,t2])
c.kickoff()                          # 先正常跑,留存档
print('task2 id =', t2.id)          # 记下这个 id
# c.replay(task_id=str(t2.id))       # 再从 task2 续跑
"
明日预告 · Day 25:今天多次提到 agent/crew 的 memory。明天专门拆记忆开关与 crew 级配置memory=Truecreate_crew_memory 怎么造出统一记忆、root_scope 命名空间怎么定、_memory_llm 怎么兜底选模型、后台写入 drain_writes 的机制。记忆的"总开关"这一层。
← Day 23 规划 Day 25 · 记忆开关 →