一次 kickoff 的完整旅程(阶段1 收官)
前 5 天我们分别拆了 Agent、Task、Crew、Process。今天把它们串成一条完整的执行链:你按下 crew.kickoff() 后,控制流从 Crew 一层层下沉到 Task、到 Agent、到执行引擎的思考循环,再把结果层层上卷成 CrewOutput。看完这一天,你脑子里就有了一张"数据在五大主角间怎么流"的地图——阶段2 之后所有深水区,都是在这张地图的某一站里放大。
先看全景:五层调用栈
result.raw?不把这条链打通,读阶段2 的执行器源码时会"只见树木不见森林"。第1站:Crew.kickoff 主干
入口在 crew.py:966,主干在 crew.py:1021(D02 见过,这里从"旅程"角度重看):
# crew.py:1021
inputs = prepare_kickoff(self, inputs, input_files) # ① 准备:回调/事件/插值/重置
if self.process == Process.sequential: # ② 分流
result = self._run_sequential_process()
elif self.process == Process.hierarchical:
result = self._run_hierarchical_process()
for after_callback in self.after_kickoff_callbacks: # ③ 跑完回调
result = after_callback(result)
result = self._post_kickoff(result)
self.usage_metrics = self.calculate_usage_metrics() # ④ 统计用量
return result # ⑤ 上卷的终点:CrewOutput
prepare_kickoffD02 讲过:before 回调 → 发 CrewKickoffStartedEvent → inputs 插值(把 {sector} 换成真实值)→ 重置输出缓存。按 process 分流D05 讲过:sequential 直接 _execute_tasks;hierarchical 先 _create_manager_agent 再 _execute_tasks。after_kickoff_callbacks结果出来后还能被钩子改写(比如统一后处理)。calculate_usage_metrics把这趟花掉的 token 汇总,最终写进 CrewOutput.token_usage。第2站:_execute_tasks 把活派到具体 task
主循环在 crew.py:1519(D05 详解过)。旅程视角只看关键三行(crew.py:1575):
# crew.py:1575
context = self._get_context(task, task_outputs) # ★把前面任务的产出组装成上下文
task_output = task.execute_sync( # ★下沉到第3站:Task 执行
agent=exec_data.agent,
context=context,
tools=exec_data.tools,
)
task_outputs.append(task_output) # 收集这一站的 TaskOutput
_get_context★流水线的血脉:任务 N 执行前,把任务 1..N-1 的输出拼成上下文传进去。这就是"上一道工序喂下一道"的落点。exec_data.agent由 prepare_task_execution 定好:sequential 用 task 自带的 agent,hierarchical 由 manager 委派。两种流程的唯一分叉点(D05 取舍②)。task.execute_sync(...)★把控制权交给第3站——Task 自己的执行方法(D04 讲过)。append(task_output)每完成一个任务,就把它的 TaskOutput 收进列表,为最后汇总 CrewOutput 攒料。第3站:Task → Agent 的交接
execute_sync 转 _execute_core(task.py:790),核心是把活交给 agent:
# task.py:787
crewai_event_bus.emit(self, TaskStartedEvent(context=context, task=self)) # 发"任务开始"事件
# task.py:790
result = agent.execute_task( # ★下沉到第4站:Agent 执行
task=self,
context=context,
tools=tools,
)
进入 Agent.execute_task(agent/core.py:740)——它负责"拼最终 prompt + 决定要不要超时保护":
# agent/core.py:761
task_prompt = self._prepare_task_execution(task, context) # role/goal/backstory + task 拼 prompt
task_prompt = handle_knowledge_retrieval(self, task, task_prompt, ...) # 若有知识库,检索注入
task_prompt = self._finalize_task_prompt(task_prompt, tools, task)
crewai_event_bus.emit(self, AgentExecutionStartedEvent(agent=self, ...)) # 发"agent 开始"事件
validate_max_execution_time(self.max_execution_time)
if self.max_execution_time is not None:
result = self._execute_with_timeout(task_prompt, task, self.max_execution_time) # 带超时
else:
result = self._execute_without_timeout(task_prompt, task) # 不带超时
return self._finalize_task_execution(task, result)
_prepare_task_execution★把 D03 的人设三件套、D04 的 task 描述/期望输出,拼成真正发给模型的 prompt。handle_knowledge_retrieval如果 agent 挂了知识库(knowledge),先检索相关内容拼进 prompt(阶段6)。max_execution_time 分支★D03 的 max_execution_time 字段在这里生效:设了就用带超时的线程池执行,超时抛 TimeoutError;没设就直接跑。两条都通向 executor带不带超时,最终都会调用第4站的 agent_executor.invoke——只是外面套没套超时保护。execute_task 这一层用 ThreadPoolExecutor 把整段执行包起来,靠 future.result(timeout=...) 统一兜底(见 _execute_with_timeout,agent/core.py:811)。好处:超时逻辑只写一处、与执行细节解耦,无论里面循环多复杂都能被整体掐断。代价是多一个线程、超时后被掐断的任务清理要小心。这是"横切关注点(超时)用一层包装统一处理"的常见做法。第4/5站:executor 的思考循环(这才是"智能"发生的地方)
无论带不带超时,最终都进 CrewAgentExecutor.invoke(agents/crew_agent_executor.py:208):
# agents/crew_agent_executor.py:219
self.messages = []
self.iterations = 0 # ★迭代计数从 0 开始
self._setup_messages(inputs) # 组装初始消息(system + user)
# ...
formatted_answer = self._invoke_loop() # ★核心循环
# ...
return {"output": formatted_answer.output}
_invoke_loop(agents/crew_agent_executor.py:309)先判断模型支不支持"原生函数调用",决定走哪种循环:
# agents/crew_agent_executor.py:318
use_native_tools = (
hasattr(self.llm, "supports_function_calling")
and self.llm.supports_function_calling()
and self.original_tools
)
if use_native_tools:
return self._invoke_loop_native_tools() # 用模型原生 tool-calling
return self._invoke_loop_react() # 否则用 ReAct 文本模式
# agents/crew_agent_executor.py:341(ReAct 循环骨架)
while not isinstance(formatted_answer, AgentFinish): # ★循环直到拿到"最终答案"
if has_reached_max_iterations(self.iterations, self.max_iter): # 到上限就收尾
formatted_answer = handle_max_iterations_exceeded(...); break
enforce_rpm_limit(self.request_within_rpm_limit) # 限速(D03 的 max_rpm)
answer = get_llm_response(llm=self.llm, messages=self.messages, ...) # 问模型
# ...解析出 是"要用工具(Action)" 还是 "最终答案(Final Answer)"...
iterations / max_iter★D03 的 max_iter 在这里兜底:循环每转一圈 +1,到上限强制收尾,防止死循环烧钱。while not AgentFinish★这就是 agent 的"思考":反复"问模型→模型说要用某工具→执行工具→把结果喂回去再问",直到模型给出 Final Answer(AgentFinish)才跳出。native vs ReAct 两条路模型支持原生 function-calling 就用结构化的工具调用;不支持就退回 ReAct——把工具说明写进 prompt、靠解析模型输出的文本来触发工具。阶段2 D08/D09 逐行拆。enforce_rpm_limitD03 的 max_rpm 在每轮开头生效,保护你的 API 额度。返回 output循环结束,invoke 返回 {"output": ...}——这就是要往上卷的原始结果。max_iter 就是"最多折腾几轮,别没完没了"。第5站往回:结果层层上卷成 CrewOutput
模型给出最终答案后,结果原路返回,每一层各自"再加工一次":
| 层 | 拿到什么 | 加工成什么 / 交给谁 |
|---|---|---|
| ④ Executor | 模型 Final Answer | {"output": ...} 返回给 Agent |
| ③ Agent.execute_task | output 文本 | _finalize_task_execution 处理后返回给 Task |
| ② Task._execute_core | 原始结果 | 按 output_pydantic/json 包成 TaskOutput(D04 L06) |
| ② _execute_tasks | 一串 TaskOutput | task_outputs.append(...) 收集 |
| ① Crew | 全部 TaskOutput | _create_crew_output 汇成 CrewOutput |
最终 CrewOutput 的结构 D02 见过(crews/crew_output.py:13):raw / pydantic / json_dict / tasks_output[] / token_usage。你在 L02 看到的 return result 返回的就是它。
result = crew.kickoff(inputs={"sector":"tech"}):① Crew 插值把
{sector}→"tech"、分流到 sequential;② 主循环取第一个 task,组装 context,调
execute_sync;③ Task 发事件、交给 analyst 这个 agent 的
execute_task;④ agent 拼好 prompt、进 executor 的
_invoke_loop 反复思考;⑤ 模型给出分析(Final Answer);
⑥ 结果上卷:包成 TaskOutput → 汇成 CrewOutput → 就是你手里的
result,result.raw 就是那段分析。CrewKickoffStartedEvent,第3站发 TaskStartedEvent 和 AgentExecutionStartedEvent……CrewAI 在每个关键节点都 emit 事件。关键认知:这些事件是旁路观测(给日志、遥测、监听器用),不参与主流程的控制——就算没有任何监听器,kickoff 照样跑通。所以 debug 时你可以挂监听器"偷看"每一站发生了什么,而不用改动执行代码。但也有个坑:事件是同步 emit 的,如果你写的监听器很慢或抛异常,可能拖慢甚至干扰主流程(源码在某些 emit 处用 try/except 兜底正是防这个)。事件系统是阶段4(D26)的专题。阶段1 收官 + 今日小结
👶 小白:这五层里,哪一层是"AI 真正在思考"的?
👨🏫 老师:第4/5站——executor._invoke_loop 里那个 while not AgentFinish 循环。前三层(Crew/Task/Agent.execute_task)都是"准备与调度",把 prompt 拼好、把上下文备好;真正"问模型、看它要不要用工具、反复迭代"发生在 executor。所以阶段2 一上来就深挖它,是有道理的。
🧠 阶段1 通关自测(能答上来就毕业)
- 画出 kickoff 的五层调用栈:Crew → Task → Agent → Executor → LLM,各管什么?
- inputs 的
{占位符}在哪一站被替换?(第1站 prepare_kickoff 的插值) - "上一个任务的输出喂给下一个"发生在哪?(第2站
_get_context) max_iter/max_rpm/max_execution_time分别在哪一站生效?(前两个在 executor 循环,后者在 Agent.execute_task 的超时包装)- 结果怎么从模型输出变成
result.raw?(上卷:output→TaskOutput→tasks_output[]→CrewOutput) - 事件系统在这条链里扮演什么角色?(旁路观测,不控制主流程)
- 为什么要分五层?(职责单一、可独立替换/测试/深挖)
✋ 10 分钟动手:给旅程装个"监听器"看每一站
# 顺着五层各读一段真代码
sed -n '1021,1038p' lib/crewai/src/crewai/crew.py # ① Crew
sed -n '1575,1582p' lib/crewai/src/crewai/crew.py # ② _execute_tasks
sed -n '787,795p' lib/crewai/src/crewai/task.py # ③ Task→Agent 交接
sed -n '761,809p' lib/crewai/src/crewai/agent/core.py # ④ Agent.execute_task
sed -n '309,329p' lib/crewai/src/crewai/agents/crew_agent_executor.py # ⑤ 循环分流
# 用事件监听"偷看"旅程(体会 L06 的旁路观测)
from crewai import Agent, Task, Crew
from crewai.events import crewai_event_bus
from crewai.events.types.task_events import TaskStartedEvent
@crewai_event_bus.on(TaskStartedEvent)
def _peek(source, event):
print("👀 一个任务开始了:", event.task.description[:40])
a = Agent(role="Analyst", goal="analyze", backstory="veteran", llm="gpt-4o")
t = Task(description="Analyze {sector} sector", expected_output="a summary", agent=a)
Crew(agents=[a], tasks=[t]).kickoff(inputs={"sector": "tech"})
# 控制台会打出 👀 —— 你没改任何执行代码,只是旁路观测
crew_agent_executor 的执行循环、D09 看输出解析 parser.py(Thought/Action/Final Answer 怎么解析出来的)。地图已在手,开始深潜。