第 14 章 · 回调、日志与事件监听
本章目标:分清
step_callback与task_callback的触发时机和签名,掌握 CrewAI 事件总线(Event Bus)与BaseEventListener体系,能实现监听TaskCompletedEvent的审计日志,并把结构化日志落盘为 JSONL。
14.1 两级回调:过程与产物
CrewAI 提供两个层级的内联回调,先明确分工再谈事件系统:
| 回调 | 挂载位置 | 触发时机 | 典型用途 |
|---|---|---|---|
step_callback | Agent 或 Crew | 每个 Agent 的每个步骤结束后 | 过程日志、实时 UI 刷新 |
callback / task_callback | Task(或 Crew 的 task_callback) | 任务完成时 | 产物审计、结果入库、指标统计 |
# callbacks_demo.py —— 两级回调的最小对照示例
import os
from crewai import Agent, Task, Crew, LLM
llm = LLM(
model="openai/deepseek-chat",
base_url="https://api.deepseek.com/v1",
api_key=os.getenv("DEEPSEEK_API_KEY"),
)
def on_step(step):
"""步骤级:每步触发一次,参数是该步骤的输出/动作对象"""
print(f"[STEP] {str(step)[:120]}")
def on_task_done(task_output):
"""任务级:任务完成触发一次,参数是 TaskOutput 对象"""
print(f"[TASK] {task_output.description[:40]} -> "
f"{len(task_output.raw)} 字符, 来自 {task_output.agent}")
analyst = Agent(role="分析师", goal="给出结论", backstory="资深分析师", llm=llm)
t1 = Task(
description="分析远程办公对团队效率的影响。",
expected_output="150 字以内的结论段。",
agent=analyst,
callback=on_task_done, # Task 级回调
)
crew = Crew(
agents=[analyst],
tasks=[t1],
step_callback=on_step, # Crew 级步骤回调
)
crew.kickoff()三个容易踩的坑:
step_callback定义在 Agent 上会覆盖 Crew 上的同名回调,不要指望两层都触发;TaskOutput.raw是字符串正文,pydantic/json_dict属性则承载结构化结果(第 12 章);- 回调里抛出的异常会干扰执行流程——审计代码务必自带 try/except,别让日志问题弄崩业务。
14.2 事件总线:比回调更完整的观测面
回调只能覆盖"步骤"和"任务完成"两个点。CrewAI 内部其实是一套事件总线架构:
CrewAIEventsBus——单例事件总线,负责事件的注册与发射;BaseEvent——所有事件的基类,携带timestamp与type;BaseEventListener——自定义监听器的抽象基类。
Crew 启动、Agent 完成、工具调用、记忆读写、LLM 调用……整个生命周期都会在总线上发射事件。监听事件比回调强大得多:不需要侵入 Crew 定义代码,就能做全局监控、审计与第三方集成。
14.3 编写自定义事件监听器
标准四步:继承 BaseEventListener → 实现 setup_listeners → 用 @crewai_event_bus.on(事件类型) 注册处理函数 → 在 Crew/Flow 所在模块创建监听器实例。
# my_listeners.py —— 监听 Crew 启动/结束与 Agent 完成
from crewai.events import (
BaseEventListener,
CrewKickoffStartedEvent,
CrewKickoffCompletedEvent,
AgentExecutionCompletedEvent,
)
class MyCustomListener(BaseEventListener):
def __init__(self):
super().__init__()
def setup_listeners(self, crewai_event_bus):
@crewai_event_bus.on(CrewKickoffStartedEvent)
def on_crew_started(source, event):
print(f"Crew '{event.crew_name}' 已开始执行!")
@crewai_event_bus.on(CrewKickoffCompletedEvent)
def on_crew_completed(source, event):
print(f"Crew '{event.crew_name}' 执行完成!")
print(f"输出: {event.output}")
@crewai_event_bus.on(AgentExecutionCompletedEvent)
def on_agent_completed(source, event):
print(f"Agent '{event.agent.role}' 完成了任务")
# 关键:必须创建实例,否则处理函数不会注册
my_listener = MyCustomListener()实例必须存活且被导入
只定义类是不够的。官方文档强调三点:处理函数要注册到事件总线;监听器实例要保持在内存中(不能被垃圾回收);实例所在的模块要被你的应用导入。最简单的做法就是在定义 Crew 的文件顶部 from my_listeners import MyCustomListener 并实例化。
处理函数的统一签名是 (source, event):source 是发射事件的对象,event 是事件实例——除了基类的 timestamp/type,每种事件还有自己的字段(如 CrewKickoffCompletedEvent.output、event.agent.role)。
14.4 监听 TaskCompletedEvent 写审计日志
事件类型覆盖面很广,常用的几族:
- Crew 事件:
CrewKickoffStartedEvent/CrewKickoffCompletedEvent/CrewKickoffFailedEvent - Agent 事件:
AgentExecutionStartedEvent/AgentExecutionCompletedEvent/AgentExecutionErrorEvent - Task 事件:
TaskStartedEvent/TaskCompletedEvent/TaskFailedEvent - 工具事件:
ToolUsageStartedEvent/ToolUsageFinishedEvent/ToolUsageErrorEvent - LLM 事件:
LLMCallStartedEvent/LLMCallCompletedEvent/LLMCallFailedEvent - Flow 事件:
FlowStartedEvent/MethodExecutionStartedEvent/MethodExecutionFinishedEvent等
下面是一个生产可用的审计监听器:把每个任务的完成情况追加写入 JSONL 文件。
# audit_listener.py —— TaskCompletedEvent 审计落盘
import json
import os
from datetime import datetime
from crewai.events import BaseEventListener, TaskCompletedEvent
class TaskAuditListener(BaseEventListener):
def __init__(self, path="audit_log.jsonl"):
super().__init__()
self.path = path
def setup_listeners(self, crewai_event_bus):
@crewai_event_bus.on(TaskCompletedEvent)
def on_task_completed(source, event):
try: # 审计代码绝不抛异常影响主流程
record = {
"ts": datetime.now().isoformat(timespec="seconds"),
"task": getattr(event, "task", None) is not None
and str(event.task.description)[:80],
"summary": str(getattr(event, "output", ""))[:2000],
}
with open(self.path, "a", encoding="utf-8") as f:
f.write(json.dumps(record, ensure_ascii=False) + "\n")
except Exception as e: # noqa: BLE001
print(f"[audit] 写日志失败: {e}")
task_audit = TaskAuditListener() # 模块级实例化,导入即生效# main.py —— 在业务代码里导入监听器即可
from crewai import Agent, Task, Crew, LLM
import audit_listener # noqa: F401 导入即注册,无需显式使用
llm = LLM(
model="openai/deepseek-chat",
base_url="https://api.deepseek.com/v1",
api_key=os.getenv("DEEPSEEK_API_KEY"),
)
writer = Agent(role="写手", goal="写摘要", backstory="专业写手", llm=llm)
crew = Crew(
agents=[writer],
tasks=[Task(description="给一段产品文案写 50 字摘要。",
expected_output="50 字摘要。", agent=writer)],
)
crew.kickoff() # 执行后查看 audit_log.jsonl 即有审计记录多监听器时,官方建议做成 listeners/ 包:每个监听器模块在底部创建实例,__init__.py 里统一 from .xxx import xxx 导出,业务文件只需 import my_project.listeners 一行。
14.5 结构化日志与临时监听
生产环境的日志要结构化(JSON 行)而不是 print:字段固定、机器可解析、能直接喂给 ELK/Loki。上面 JSONL 的写法就是最小实现;更进一步可以在记录里带上 event.type、event.timestamp,并按事件族分文件(llm.jsonl / tool.jsonl / task.jsonl)。
调试某个特定流程时,可以用官方提供的 scoped_handlers 上下文管理器注册临时监听器,出了作用域自动移除:
from crewai.events import crewai_event_bus, CrewKickoffStartedEvent
with crewai_event_bus.scoped_handlers():
@crewai_event_bus.on(CrewKickoffStartedEvent)
def temp_handler(source, event):
print("这个处理器只在这个 with 块内存在")
crew.kickoff() # 期间的事件会被临时处理器捕获
# 出了 with 块,temp_handler 已被移除14.6 回调 vs 事件监听:怎么选
- 一次性、局部的观测(某个 Task 的结果入库)→ 用 Task 回调,代码就近、直观;
- 全局、跨 Crew 的横切关注(审计、监控、计费、通知)→ 用事件监听,零侵入、可统一开关;
- 两者不互斥:成熟项目通常回调管"业务动作",事件管"平台观测"。
14.7 本章小结
step_callback步骤级触发(Agent 级覆盖 Crew 级),Task(callback=...)任务完成级触发;- 事件系统三件套:
CrewAIEventsBus单例总线、BaseEvent基类、BaseEventListener监听器基类; - 监听器必须实例化且被导入才生效;处理函数统一签名
(source, event); TaskCompletedEvent+ JSONL 落盘是最小可用审计方案;scoped_handlers适合临时调试;- 审计/日志代码必须自带异常保护,避免观测系统拖垮业务执行。
🧪 随堂测验
点击你认为正确的选项。答错时会展示正确答案与原因解析。
1. 关于 step_callback 与 Task 回调(callback/task_callback)的分工,正确的是?
2. 定义了 BaseEventListener 子类但运行时监听不生效,最可能的原因是?
3. 事件处理函数的统一签名是?
4. 只想在调试某段流程时临时监听事件,官方推荐的方式是?
🛠️ 动手实践
- 给第 13 章的 Crew 补一个
ToolUsageStartedEvent/ToolUsageFinishedEvent监听器,统计每个工具的调用次数与总耗时,输出成tool_stats.json。 - 把本章的
TaskAuditListener扩展为同时监听TaskFailedEvent与CrewKickoffFailedEvent,失败记录额外写入errors.jsonl并带上事件类型字段。 - 用
scoped_handlers写一个调试脚本:临时监听LLMCallStartedEvent,跑一个双任务 Crew,统计本次执行发起的 LLM 调用次数。
掌握了观测,下一章进入 CrewAI 的另一根支柱:第 15 章 · Flow 入门:事件驱动工作流。