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_async 和 akickoff 名字这么像,差在哪?选错了要么慢、要么把服务器卡死。💡 一句话本质
整个家族 = 一个核心
kickoff(同步)+ 三个维度的包装。三个维度是:①批量(_for_each:对每条输入 copy() 一份 crew 各跑一遍)、②线程异步(kickoff_async:用 asyncio.to_thread 把同步 kickoff 丢到线程池)、③原生异步(akickoff:全程 async/await,连任务执行都是异步的)。核心 kickoff 内部则永远是:prepare_kickoff(备战)→ 按 process 分流到流程 → after 回调 → _post_kickoff → finally 收尾。图注:批量方法循环调用单次方法;异步方法要么包线程、要么走原生异步流程。核心备战/收尾共用。
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_tasks(crew.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_writes(crew.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_async 用 to_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_async和akickoff的本质区别?各适合什么场景?- 为什么
_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.py 的 CrewPlanner。