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

kickoff 家族:一支团队的四种"喊 Action"方式

Day 19-21 把 crew 的"静态结构"和"两种执行流程"读完了。今天读启动入口——你天天写的那句 crew.kickoff()。但它其实是一整个家族:同步的 kickoff、批量的 kickoff_for_each、线程异步的 kickoff_async、原生异步的 akickoff,还有各自的 _for_each 变体。它们最终都汇聚到同一段"分流到流程 + 收尾"的核心。今天理清这个家族的关系,重点读 prepare_kickoff 开跑前那一长串准备、和 copy() 在批量里的关键作用。

📍 你在 60 天里的位置(阶段4 Crew 与流程 · 共 8 天)
D19 Crew 全字段 D20 顺序流程 D21 层级+manager D22 kickoff 家族 D23 规划 D24 训练+replay D25 记忆开关 D26 事件系统
💡 先用一个类比兜住今天 kickoff 家族就像餐厅的四种下单方式kickoff = 堂食(你坐着等菜上齐);kickoff_for_each = 一次点一桌人的餐(每人一份,各做各的);kickoff_async = 打包外卖(你先去干别的,好了叫你,但厨房还是那套老流程);akickoff = 专门的"异步厨房"(从头到尾都为并发优化)。不管哪种,后厨的"备料 → 炒菜 → 出餐 → 收拾"是同一套——这就是 prepare_kickoff → 流程 → 收尾
L01

痛点:为什么启动有这么多个方法?

🤔 痛点你只会写 crew.kickoff(inputs=...),但翻 API 发现还有 kickoff_for_each / kickoff_async / kickoff_for_each_async / akickoff / akickoff_for_each 一堆。它们到底啥区别?我想"用 100 条不同输入各跑一遍"该用哪个?想"在 FastAPI 里不阻塞事件循环"又该用哪个?kickoff_asyncakickoff 名字这么像,差在哪?选错了要么慢、要么把服务器卡死。
💡 一句话本质 整个家族 = 一个核心 kickoff(同步)+ 三个维度的包装。三个维度是:①批量_for_each:对每条输入 copy() 一份 crew 各跑一遍)、②线程异步kickoff_async:用 asyncio.to_thread 把同步 kickoff 丢到线程池)、③原生异步akickoff:全程 async/await,连任务执行都是异步的)。核心 kickoff 内部则永远是:prepare_kickoff(备战)→ 按 process 分流到流程 → after 回调 → _post_kickoff → finally 收尾。
数据结构:kickoff 家族六个方法的关系 kickoff (同步核心) kickoff_for_each kickoff_async akickoff (原生异步) akickoff_for_each 批量=循环调 批量=循环调 akickoff 核心内部:prepare_kickoff → process 分流 → after 回调 → 收尾 akickoff 走 _arun_* 异步版流程,其余走 _run_* 同步版
图注:批量方法循环调用单次方法;异步方法要么包线程、要么走原生异步流程。核心备战/收尾共用。
L02

kickoff 主干:备战 → 分流 → 收尾

剥掉流式/检查点分支后的主干(crew.py:966):

# crew.py:966
def kickoff(self, inputs=None, input_files=None, from_checkpoint=None) -> CrewOutput | ...:
    restored = apply_checkpoint(self, from_checkpoint)   # 检查点恢复(D24)
    if restored is not None:
        return restored.kickoff(inputs=inputs, input_files=input_files)
    get_env_context()
    if self.stream:                       # 流式:另起线程边跑边吐 chunk(略)
        ...
    baggage_ctx = baggage.set_baggage("crew_context", CrewContext(id=str(self.id), key=self.key))
    token = attach(baggage_ctx)           # 把 crew 身份挂到当前上下文(tracing 用)
    runtime_scope = crewai_event_bus._enter_runtime_scope()
    try:
        inputs = prepare_kickoff(self, inputs, input_files)   # ★备战(L03)
        if self.process == Process.sequential:
            result = self._run_sequential_process()           # 顺序(D20)
        elif self.process == Process.hierarchical:
            result = self._run_hierarchical_process()         # 层级(D21)
        else:
            raise NotImplementedError(f"The process '{self.process}' is not implemented yet.")
        for after_callback in self.after_kickoff_callbacks:   # after 回调改结果
            result = after_callback(result)
        result = self._post_kickoff(result)
        self.usage_metrics = self.calculate_usage_metrics()   # 统计 token
        return result
    except Exception as e:
        crewai_event_bus.emit(self, CrewKickoffFailedEvent(...))   # 失败发事件(D26)
        raise
    finally:
        self._drain_memory_writes()       # ★收尾(L07)
        clear_files(self.id)
        detach(token)
        crewai_event_bus._exit_runtime_scope(runtime_scope)
apply_checkpoint开头先看有没有检查点要恢复——有就用恢复出来的 crew 重新 kickoff(Day 24 会讲检查点/replay)。
attach(baggage_ctx)把 crew 的 id/key 挂进 OpenTelemetry 的上下文"行李箱",这样整条调用链的 trace 都知道属于哪个 crew。
prepare_kickoff★开跑前的一大串准备:回调、事件、插值、建 agent、规划……(L03 展开)。
if process ==★核心分流:顺序走 _run_sequential_process,层级走 _run_hierarchical_process(就是 Day 20/21)。
after_kickoff_callbacks结果拿到后,依次过用户挂的 after 回调,每个都能改结果(Day 19 的字段)。
finally: _drain_memory_writes★不管成功失败,收尾都跑:等后台记忆写完、清临时文件、解除上下文。
💡 为什么核心逻辑包在 try/finally 里?因为 kickoff 期间挂了不少需要成对清理的东西:上下文 token(attach 了要 detach)、临时文件、后台记忆写入、runtime scope。用 finally 保证哪怕中途抛异常也一定清理干净——否则一次失败的 kickoff 会泄漏上下文、留下垃圾文件,污染后续执行。这是"资源获取后必须释放"的标准结构。
L03

prepare_kickoff:开跑前的一长串备战

真正的"开机准备"在这(crews/utils.py:249,裁剪主干):

# crews/utils.py:249
def prepare_kickoff(crew, inputs, input_files=None) -> dict | None:
    resuming = crew.checkpoint_kickoff_event_id is not None
    if not resuming and get_current_parent_id() is None:
        reset_emission_counter()      # 事件序号归零(D26)
        reset_last_event_id()

    normalized = dict(inputs) if inputs is not None else None
    for before_callback in crew.before_kickoff_callbacks:   # ① before 回调改 inputs
        if normalized is None: normalized = {}
        normalized = before_callback(normalized)

    if not resuming:
        started_event = CrewKickoffStartedEvent(crew_name=crew.name, inputs=normalized)
        crew._kickoff_event_id = started_event.event_id
        crewai_event_bus.emit(crew, started_event)          # ② 发"开始"事件

    crew._task_output_handler.reset()                       # ③ 清上次的任务产出存档
    if normalized is not None:
        crew._inputs = normalized
        crew._interpolate_inputs(normalized)                # ④ 把 {topic} 替换成真实值
    crew._set_tasks_callbacks()

    agents_to_setup = list(crew.agents)                     # ⑤ 建齐所有 agent(含 task 里的)
    for task in crew.tasks:
        if task.agent is not None and id(task.agent) not in seen_agent_ids:
            agents_to_setup.append(task.agent)
    setup_agents(crew, agents_to_setup, crew.embedder, crew.function_calling_llm, crew.step_callback)

    if crew.planning:                                       # ⑥ 开了规划就先规划(D23)
        crew._handle_crew_planning()
    return normalized
before_kickoff_callbacks★链式过 before 回调:normalized = before_callback(normalized)——前一个的返回喂给后一个,可以改 inputs(Day 19 说过)。
emit(CrewKickoffStartedEvent)发"crew 开跑了"事件,并记下 _kickoff_event_id。日志/监控/tracing 从这里开始(Day 26)。
_task_output_handler.reset()清掉上一次 kickoff 存的任务产出,避免串味(也和 replay 的存档机制相关,D24)。
_interpolate_inputs★把任务/agent 里的占位符 {topic} 用 inputs 里的真实值替换。你传的 inputs 就是在这里"注入"进任务描述的。
setup_agents把所有要用的 agent(含只在 task 里出现的)初始化好:绑记忆、绑 embedder、绑 step 回调。
if crew.planning: 规划★如果开了 planning=True,在真正执行前先跑一轮规划给任务追加步骤(明天 Day 23)。
注意执行顺序:before 回调 → 发开始事件 → 插值 → 建 agent → 规划。这个顺序很讲究:before 回调能改 inputs,所以必须在插值之前跑;规划要看最终的任务描述,所以放在插值之后顺序错了功能就废了——读源码时留意这种"隐含的时序约束"。
L04

kickoff_for_each:为每条输入复制一份 crew

批量执行(crew.py:1061):

# crew.py:1061
def kickoff_for_each(self, inputs: list[dict], input_files=None) -> list[CrewOutput | ...]:
    results = []
    total_usage_metrics = UsageMetrics()
    for input_data in inputs:                 # 对每一条输入
        crew = self.copy()                    # ★关键:复制一份全新的 crew
        output = crew.kickoff(inputs=input_data, input_files=input_files)
        if not self.stream and crew.usage_metrics:
            total_usage_metrics.add_usage_metrics(crew.usage_metrics)   # 累加 token 用量
        results.append(output)
    if not self.stream:
        self.usage_metrics = total_usage_metrics
    self._task_output_handler.reset()
    return results
for input_data in inputsinputs 是列表,每个元素是一套输入。挨个跑。
crew = self.copy()最关键的一行:每条输入都用一份全新复制的 crew 跑,而不是复用 self
crew.kickoff(inputs=input_data)在副本上跑标准 kickoff。副本各自持有独立的 messages、task_outputs、记忆写入队列。
add_usage_metrics把每个副本的 token 用量累加,最后汇总回 self.usage_metrics
💡 设计取舍①:为什么每条输入要 copy 一份 crew,而不是复用? 因为 kickoff 会往 crew 的实例状态里写东西_inputs_kickoff_event_id、任务的 output、agent 的对话历史……如果 100 条输入都用同一个 self 跑,第 2 条会污染第 1 条残留的状态,任务描述里的占位符还可能被上一条的值污染。copy()(Day 19 见过它 exclude 掉 id 等)保证每条输入都是干净的一次执行代价是复制开销,但换来的是执行之间的完全隔离——正确性远比这点开销重要。
⚠️ 边界:kickoff_for_each 默认是"串行"的 别被"批量"误导:kickoff_for_each 就是个 for 循环,一条跑完才跑下一条,并不并行。想真正并发跑多条输入,得用 akickoff_for_each(L06,原生异步版才会并发调度)。用错方法,100 条串行跑会慢得让你以为卡死了。
L05

kickoff_async:把同步 kickoff 丢进线程

最"偷懒"但实用的异步(crew.py:1097,非流式主干):

# crew.py:1097
async def kickoff_async(self, inputs=None, input_files=None, from_checkpoint=None):
    restored = apply_checkpoint(self, from_checkpoint)
    if restored is not None:
        return await restored.kickoff_async(inputs=inputs, input_files=input_files)
    inputs = inputs or {}
    if self.stream:
        ...                       # 流式异步(略)
    return await asyncio.to_thread(self.kickoff, inputs, input_files)   # ★核心就这一句
async def它是个协程,你要 await crew.kickoff_async(...)
asyncio.to_thread(self.kickoff, ...)★把同步的 kickoff 丢到线程池里跑,然后 await 它完成。同步代码原封不动,只是不再阻塞你的事件循环。
💡 为什么这么"偷懒"也管用?因为 crew 内部大量是同步阻塞代码(LLM 调用、工具执行),全改成原生 async 工程量巨大。asyncio.to_thread 是个务实的折中:不改一行同步逻辑,只是把它整个"搬到别的线程",你的主事件循环(比如 FastAPI 的)就不会被卡住。适合"我只想在 async 环境里不阻塞地调一次 crew"的场景。缺点:受线程池大小限制,不是真正为高并发优化。
L06

akickoff:从头到尾的原生异步

为并发而生的原生异步(crew.py:1177,主干):

# crew.py:1177
async def akickoff(self, inputs=None, input_files=None, from_checkpoint=None):
    """Native async kickoff method using async task execution throughout.
    Unlike kickoff_async which wraps sync kickoff in a thread, this method
    uses native async/await for all operations including task execution,
    memory operations, and knowledge queries."""
    ...
    try:
        inputs = prepare_kickoff(self, inputs, input_files)   # 同样的备战
        if self.process == Process.sequential:
            result = await self._arun_sequential_process()    # ★异步版流程
        elif self.process == Process.hierarchical:
            result = await self._arun_hierarchical_process()
        ...
    finally:
        self._drain_memory_writes()
        ...

# crew.py:1298  异步版顺序流程
async def _arun_sequential_process(self) -> CrewOutput:
    return await self._aexecute_tasks(self.tasks)

# crew.py:1151  批量原生异步(真并发)
async def kickoff_for_each_async(self, inputs, input_files=None):
    async def kickoff_fn(crew, input_data):
        return await crew.kickoff_async(inputs=input_data, input_files=input_files)
    return await run_for_each_async(self, inputs, kickoff_fn)
docstring 点题★源码自己讲清了区别:kickoff_async 是"包线程",akickoff 是"全程原生 async"——连任务执行、记忆、知识查询都是异步的。
await self._arun_sequential_process走的是异步版流程 _aexecute_taskscrew.py:1307),不是 Day 20 那个同步 _execute_tasks
run_for_each_async批量异步的真并发靠它调度——多条输入的 crew 副本可以真正同时跑,不是串行。
方法异步方式并发?适合
kickoff同步阻塞脚本、简单调用
kickoff_for_each同步串行少量批处理
kickoff_async线程包装受线程池限async 环境不阻塞地跑一次
akickoff原生 async高并发服务
akickoff_for_each原生 async 批量是(真并发)大批量输入并发
L07

finally 收尾:为什么记忆要"drain"两次

收尾里最有意思的是 _drain_memory_writescrew.py:1848):

# crew.py:1848
def _drain_memory_writes(self) -> None:
    """Block until all pending background memory saves have completed.
    ... Must run before CrewKickoffCompletedEvent is emitted: listeners
    (e.g. telemetry sessions) tear down on that event, and any
    MemorySaveCompletedEvent emitted after teardown is lost ..."""
    seen: set[int] = set()
    candidates = [self._memory, self.memory,
                  getattr(self.manager_agent, "memory", None),
                  *(getattr(agent, "memory", None) for agent in self.agents)]
    for mem in candidates:
        if mem is None or isinstance(mem, bool):
            continue
        backing = getattr(mem, "_memory", None) or mem
        if id(backing) in seen:      # 同一个底层 memory 只 drain 一次
            continue
        seen.add(id(backing))
        drain = getattr(backing, "drain_writes", None)
        if callable(drain):
            drain()                   # ★阻塞等这个记忆池的后台写入全部落盘
candidates 收集所有记忆池crew 的、每个 agent 的、经理的记忆各自独立,都可能有后台写入没完成。全收集起来。
seen 去重多个 agent 可能共享同一个底层 memory,用 id(backing) 去重,避免重复 drain。
drain_writes()★阻塞,直到这个记忆池的异步写入全部落盘。记忆是后台异步存的(不阻塞主流程),但结束前必须等它们完成。
在成功路径和 finally 各调一次★成功时在 _create_crew_output 里 drain(crew.py:1897),异常路径在 finally 里再兜一次。
⚠️ 边界:为什么必须在"完成事件"之前 drain? docstring 讲透了这个坑:CrewKickoffCompletedEvent 一发出,一些监听器(如遥测会话)会拆除自己。如果记忆的 MemorySaveCompletedEvent 在这之后才发,就没人接收了——那个"保存追踪"就变成孤儿、span 永远不闭合。所以顺序必须是:先 drain(等记忆写完、事件都发完)→ 再发完成事件(触发拆除)。这是异步事件系统里典型的"生命周期顺序"陷阱。
L08

取舍 + 今日小结

💡 设计取舍②:为什么同时提供"包线程"和"原生异步"两套异步? kickoff_async(包线程)和 akickoff(原生)看似重复,其实是迁移期的务实并存。原生异步要把整条链路(任务执行、记忆、工具)都改成 async,工程量大、周期长。在完全改完之前,kickoff_asyncto_thread 让用户立刻能在 async 环境里用(哪怕不是最高效);akickoff 则是"正确答案",为真高并发准备。先给一个能用的,再给一个更好的——这是成熟框架照顾存量用户的常见做法。

👶 小白:我在 FastAPI 接口里调 crew,该用哪个?

👨‍🏫 老师:至少用 kickoff_async——直接用同步 kickoff阻塞整个事件循环,一个请求把所有并发请求都卡住。如果你的服务要同时扛很多路 crew 执行、追求吞吐,就上 akickoff。批量并发跑就 akickoff_for_each。记住:在 async 框架里绝不要裸调同步 kickoff。

🧠 今天你应该能回答

  • kickoff 家族六个方法的关系?三个包装维度是什么?
  • kickoff 主干的固定四步是什么?为什么包在 try/finally?
  • prepare_kickoff 里 before 回调、插值、规划的先后为什么不能乱?
  • kickoff_for_each 为什么要 copy()?它并发吗?
  • kickoff_asyncakickoff 的本质区别?各适合什么场景?
  • 为什么 _drain_memory_writes 必须在完成事件之前跑?

✋ 10 分钟动手

P=lib/crewai/src/crewai
sed -n '966,1057p'  $P/crew.py          # kickoff 主干
sed -n '249,362p'   $P/crews/utils.py   # prepare_kickoff 备战全流程
sed -n '1061,1096p' $P/crew.py          # kickoff_for_each(看 copy())
sed -n '1177,1270p' $P/crew.py          # akickoff 原生异步
# 对比串行 for_each 和普通 kickoff
python -c "
from crewai import Agent, Task, Crew
a=Agent(role='翻译', goal='翻译成英文', backstory='资深译者')
t=Task(description='把这句话翻成英文:{sentence}', expected_output='英文', agent=a)
crew=Crew(agents=[a], tasks=[t])
outs=crew.kickoff_for_each(inputs=[{'sentence':'你好'},{'sentence':'再见'}])
for o in outs: print(o.raw)
"
明日预告 · Day 23:今天 prepare_kickoff 里那句 if crew.planning: crew._handle_crew_planning() 一带而过。明天专门拆它:一个"规划 agent"怎么在正式执行前,读遍所有任务和工具,生成"每个任务该怎么一步步做"的计划,再把计划追加进任务描述。源码在 utilities/planning_handler.pyCrewPlanner
← Day 21 层级流程 Day 23 · 规划 →