Skip to content

第 14 章 · 回调、日志与事件监听

本章目标:分清 step_callbacktask_callback 的触发时机和签名,掌握 CrewAI 事件总线(Event Bus)与 BaseEventListener 体系,能实现监听 TaskCompletedEvent 的审计日志,并把结构化日志落盘为 JSONL。

14.1 两级回调:过程与产物

CrewAI 提供两个层级的内联回调,先明确分工再谈事件系统:

回调挂载位置触发时机典型用途
step_callbackAgent 或 Crew每个 Agent 的每个步骤结束后过程日志、实时 UI 刷新
callback / task_callbackTask(或 Crew 的 task_callback任务完成时产物审计、结果入库、指标统计
python
# 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 内部其实是一套事件总线架构

  1. CrewAIEventsBus——单例事件总线,负责事件的注册与发射;
  2. BaseEvent——所有事件的基类,携带 timestamptype
  3. BaseEventListener——自定义监听器的抽象基类。

Crew 启动、Agent 完成、工具调用、记忆读写、LLM 调用……整个生命周期都会在总线上发射事件。监听事件比回调强大得多:不需要侵入 Crew 定义代码,就能做全局监控、审计与第三方集成。

14.3 编写自定义事件监听器

标准四步:继承 BaseEventListener → 实现 setup_listeners → 用 @crewai_event_bus.on(事件类型) 注册处理函数 → 在 Crew/Flow 所在模块创建监听器实例。

python
# 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.outputevent.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 文件。

python
# 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()   # 模块级实例化,导入即生效
python
# 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.typeevent.timestamp,并按事件族分文件(llm.jsonl / tool.jsonl / task.jsonl)。

调试某个特定流程时,可以用官方提供的 scoped_handlers 上下文管理器注册临时监听器,出了作用域自动移除:

python
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. 只想在调试某段流程时临时监听事件,官方推荐的方式是?

🛠️ 动手实践

  1. 给第 13 章的 Crew 补一个 ToolUsageStartedEvent / ToolUsageFinishedEvent 监听器,统计每个工具的调用次数与总耗时,输出成 tool_stats.json
  2. 把本章的 TaskAuditListener 扩展为同时监听 TaskFailedEventCrewKickoffFailedEvent,失败记录额外写入 errors.jsonl 并带上事件类型字段。
  3. scoped_handlers 写一个调试脚本:临时监听 LLMCallStartedEvent,跑一个双任务 Crew,统计本次执行发起的 LLM 调用次数。

掌握了观测,下一章进入 CrewAI 的另一根支柱:第 15 章 · Flow 入门:事件驱动工作流