Day 26 / 共 60 天 · 阶段4 Crew 与流程(收官)

事件系统:日志/遥测/流式/tracing 的统一底座

整个阶段4 你无数次见到 crewai_event_bus.emit(self, CrewKickoffStartedEvent(...))——kickoff 开始/结束/失败、训练开始/完成、记忆保存完成……crew 在关键节点发事件,但它不关心谁在听。而日志打印、遥测上报、流式输出、tracing 全都是事件的监听器。这就是经典的观察者模式(发布-订阅)。今天拆 events/ 目录:事件基类 BaseEvent、单例总线 CrewAIEventsBus@on 注册、emit 分发到同步/异步 handler、BaseEventListener 自动注册。这是阶段4 的收官,也是理解 CrewAI"可观测性"的钥匙。

📍 你在 60 天里的位置(阶段4 Crew 与流程 · 收官)
D19 Crew 全字段 D20 顺序流程 D21 层级+manager D22 kickoff 家族 D23 规划 D24 训练+replay D25 记忆开关 D26 事件系统 阶段5 工具
💡 先用一个类比兜住今天 事件系统就像公司的"广播站":crew 干活到关键节点就往广播里喊一声("kickoff 开始了!""任务完成!"),它不知道也不关心谁在收听。而各个部门自己去广播站登记"我要听哪类广播":日志部门听所有的打成日志,遥测部门听了上报云端,UI 部门听了刷新进度条。喊话的(crew)和收听的(监听器)完全解耦——crew 加一句 emit 不用改任何监听器,新增一个监听器也不用动 crew。这就是发布-订阅的威力。
L01

痛点:日志、遥测、UI……都要"插进执行流程"

🤔 痛点你希望 crew 执行时能:打印漂亮的日志、把耗时/token 上报监控、在 Web UI 上实时刷新进度、给 tracing 系统记 span……如果把这些逻辑全写进 kickoff/execute_tasks 里,那些核心方法会变成一坨"业务 + 日志 + 遥测 + UI"的意大利面,加一个新监控点就要改核心代码,改一处崩一片。怎么让"执行核心"只管执行,把这些旁路关注点干净地剥离出去?
💡 一句话本质 事件系统 = 发布-订阅(观察者模式)。执行核心只负责在关键节点 emit(事件)(发布),谁想响应就 @on(事件类型) 注册一个 handler(订阅)。中间是一个全局单例总线 CrewAIEventsBus 做撮合:它维护"事件类型 → handler 集合"的映射,emit 时查出对应 handler 挨个调用。核心代码里看不到任何日志/遥测逻辑,只有一行行 emit——旁路关注点全在监听器里,彻底解耦。
L02

BaseEvent:所有事件的共同底座

事件基类(events/base_events.py:66):

# events/base_events.py:66
class BaseEvent(BaseModel):
    """Base class for all events"""
    timestamp: datetime = Field(default_factory=lambda: datetime.now(timezone.utc))
    type: str                                     # 事件类型标识(字符串)
    source_fingerprint: str | None = None         # 谁发的(实体指纹)
    source_type: str | None = None                # "agent"/"task"/"crew"/"memory"...
    fingerprint_metadata: dict[str, Any] | None = None
    task_id: str | None = None
    task_name: str | None = None
    agent_id: str | None = None
    agent_role: str | None = None
继承 BaseModel★事件也是 Pydantic 模型——天生可序列化、带校验。这样事件能进 JSON、进 tracing、跨进程传。
timestamp 默认当前 UTC每个事件自动打时间戳。监听器据此排时间线(Day 24 的 replay 时间线、tracing 都靠它)。
type: str事件类型的字符串标识(如 "crew_kickoff_started")。子类会用 Literal[...] 固定它。
source_* / agent_* / task_*★一套溯源字段:这事件是哪个 crew/agent/task 发的。监听器不用猜,直接读这些字段就知道上下文。
💡 为什么所有事件共享这一组字段?因为监听器(尤其是通用的日志/遥测)需要不管什么事件都能拿到的公共信息——发生时间、谁发的、属于哪个 crew/agent。把这些收进基类,监听器就能写"对任意事件,先记时间戳和来源"的通用逻辑,不用为每种事件单独适配。基类承载共性,子类承载特性——这是继承体系设计的基本功。
L03

CrewAIEventsBus:全局唯一的总线

单例总线(events/event_bus.py:93)与模块级实例(events/event_bus.py:952):

# events/event_bus.py:93
class CrewAIEventsBus:
    """Singleton event bus for handling events in CrewAI.
    Synchronous handlers execute in a thread pool executor ...
    Asynchronous handlers execute in a dedicated event loop running in a daemon thread ..."""
    _instance: Self | None = None
    _instance_lock: threading.RLock = threading.RLock()
    _sync_handlers: dict[type[BaseEvent], SyncHandlerSet]       # 事件类型 → 同步 handler 集
    _async_handlers: dict[type[BaseEvent], AsyncHandlerSet]     # 事件类型 → 异步 handler 集

    def __new__(cls) -> Self:                     # ★单例:永远返回同一个实例
        if cls._instance is None:
            with cls._instance_lock:
                if cls._instance is None:
                    ...
        return cls._instance

# events/event_bus.py:952  模块级唯一实例——大家 import 的就是它
crewai_event_bus: Final[CrewAIEventsBus] = CrewAIEventsBus()
_sync_handlers / _async_handlers★核心数据结构:两个字典,键是事件类型,值是注册在该类型上的 handler 集合。emit 时按事件类型 O(1) 查出该通知谁。
__new__ 单例双重检查锁定的单例:不管谁 CrewAIEventsBus(),拿到的都是同一个实例。因为全局只能有一个"广播站"。
_instance_lock (RLock)单例初始化用锁保护,避免多线程同时创建出两个实例。
crewai_event_bus = ...(模块级)★整个框架import 的就是这一个模块级实例from ...event_bus import crewai_event_bus 到处都是它。
💡 设计取舍①:为什么事件总线要做成全局单例? 因为 emit 的地方(crew/agent/task/memory 深处)和监听的地方(用户在最外层注册)隔着十万八千里的调用栈——你不可能把一个 bus 实例层层当参数传下去。做成全局单例后,任何地方 import crewai_event_bus 就能 emit/on,无需传递。代价是单例是"隐式全局状态"(测试时要注意 handler 泄漏、要能 scoped_handlers 隔离)。但对"横切整个系统的事件总线"而言,单例是最自然的选择——这也是几乎所有框架的事件系统都用单例的原因。
L04

@on:把一个函数登记成监听器

注册用的装饰器(events/event_bus.py:245):

# events/event_bus.py:245
def on(self, event_type: type[BaseEvent],
       depends_on: Depends | list[Depends] | None = None):
    """Decorator to register an event handler for a specific event type.
    Handlers can accept 2 or 3 arguments:
        - (source, event) — standard handler
        - (source, event, state: RuntimeState) — handler with runtime state"""
    def decorator(handler):
        deps = None
        if depends_on is not None:
            deps = [depends_on] if isinstance(depends_on, Depends) else depends_on
        self._register_handler(event_type, handler, dependencies=deps)   # ★存进字典
        return handler                     # 原样返回,不包装
    return decorator

用起来是这样:

from crewai.events import crewai_event_bus
from crewai.events.types.crew_events import CrewKickoffStartedEvent

@crewai_event_bus.on(CrewKickoffStartedEvent)      # 订阅"kickoff 开始"事件
def my_logger(source, event):                       # (发送者, 事件对象)
    print(f"[{event.timestamp}] crew {event.crew_name} 开跑了!")
on(event_type)★参数是事件的类(如 CrewKickoffStartedEvent),不是字符串。类型安全,IDE 能补全。
_register_handler把你的函数存进 _sync_handlers[event_type](或异步字典)。注册 = 往字典对应键的集合里加一个函数
return handler(不包装)装饰器原样返回你的函数——它只做"登记"这个副作用,不改变函数本身。你还能正常单独调用它。
2 或 3 参数handler 签名可以是 (source, event) 或多一个 (source, event, state)。框架按参数个数自适应调用(很灵活)。
depends_on进阶:声明 handler 之间的依赖顺序("A 处理完才处理 B")。普通监听用不到。
L05

emit:把事件分发给所有订阅者

分发的核心(events/event_bus.py:572,裁剪):

# events/event_bus.py:572
def emit(self, source: Any, event: BaseEvent) -> Future[None] | None:
    self._prepare_event(source, event)          # 填公共字段(时间戳/序号等)
    event_type = type(event)
    with self._rwlock.r_locked():                # 读锁:查 handler 时不被注册打断
        sync_handlers = self._sync_handlers.get(event_type, frozenset())
        async_handlers = self._async_handlers.get(event_type, frozenset())
    if not sync_handlers and not async_handlers:
        return None                              # ★没人订阅 → 直接返回,零开销

    if sync_handlers:
        if event_type is LLMStreamChunkEvent:    # 流式 chunk:必须同步、保证顺序
            self._call_handlers(source, event, sync_handlers, state)
        else:
            ctx = contextvars.copy_context()     # ★复制上下文给线程(Day 08 见过的手法)
            sync_future = self._sync_executor.submit(
                ctx.run, self._call_handlers, source, event, sync_handlers, state)
            ...
    if async_handlers:                            # 异步 handler 丢到专用事件循环线程
        return self._track_future(asyncio.run_coroutine_threadsafe(
            self._acall_handlers(source, event, async_handlers, state), self._loop))
    return None
_prepare_event发之前先填公共信息(发射序号、来源指纹等),保证监听器拿到完整事件。
get(event_type, frozenset())★按事件类型查订阅者。查不到就返回空集合——没人听的事件不会报错。
没订阅者 → return None★关键优化:没人订阅就零开销直接返回。所以框架里遍地 emit 也不心疼——没监听器时几乎免费。
同步 handler → 线程池同步 handler 丢进 _sync_executor 线程池跑,不阻塞主执行。用 contextvars.copy_context() 把上下文带给子线程(和 Day 08 并行工具同款手法)。
LLMStreamChunkEvent 特判★流式 token 事件必须同步、按序处理——否则 token 顺序乱了,流式输出就成乱码。所以它不走线程池。
异步 handler → 专用 loop异步 handler 丢到一个专用的后台事件循环线程跑。同步/异步两套 handler 各走各的执行器。
控制流:一次 emit 如何找到并调用订阅者 crew.emit(event) 发布者不关心谁听 总线:按 type 查字典 _sync_handlers[type] 同步 handler → 线程池 异步 handler → 专用 loop 订阅者(各自 @on 注册,互不知道对方): 日志 · 遥测 · 流式UI · tracing · 记忆保存追踪 没人订阅该事件 → emit 直接 return None(零开销)
图注:发布者只管 emit,总线按事件类型查字典,同步/异步 handler 分别在线程池和专用 loop 执行。
L06

BaseEventListener:把一组监听打包自动注册

监听器基类(events/base_event_listener.py):

# events/base_event_listener.py
class BaseEventListener(ABC):
    """Abstract base class for event listeners."""
    verbose: bool = False
    def __init__(self) -> None:
        super().__init__()
        self.setup_listeners(crewai_event_bus)      # ★构造时就自动注册所有监听
        crewai_event_bus.validate_dependencies()

    @abstractmethod
    def setup_listeners(self, crewai_event_bus: CrewAIEventsBus) -> None:
        """Setup event listeners on the event bus."""

框架自带的默认监听器就这么用(events/event_listener.py:161):

# events/event_listener.py:161
class EventListener(BaseEventListener):
    def setup_listeners(self, crewai_event_bus):
        @crewai_event_bus.on(CrewKickoffStartedEvent)
        def on_crew_started(source, event):
            self.formatter.handle_crew_started(event.crew_name or "Crew", source.id)
            source._execution_span = self._telemetry.crew_execution_span(source, event.inputs)

        @crewai_event_bus.on(CrewKickoffCompletedEvent)
        def on_crew_completed(source, event):
            final = event.output.raw
            self._telemetry.end_crew(source, final)         # 遥测:结束 span
            self.formatter.handle_crew_status(event.crew_name or "Crew", source.id, "completed", final)
        # ... 还有 failed / train / task / agent / tool 等一大堆
__init__ 里 setup_listeners构造即注册:你只要 EventListener() 一实例化,它 __init__ 就把所有监听挂上总线了。省去手动逐个 on。
@abstractmethod setup_listeners抽象方法,子类必须实现——在里面用 @bus.on(...) 把这一组相关监听集中定义。
默认 EventListener框架内置的监听器,负责日志格式化(formatter)+ 遥测(telemetry)。Day 19 的 set_private_attrsevent_listener = EventListener() 就实例化了它。
一个类 = 一组相关监听把"和某个关注点相关的一堆监听"打包进一个 Listener 类,逻辑内聚。想加自定义监控,就写个自己的 BaseEventListener 子类。
💡 为什么要 BaseEventListener 这层,而不直接散着写 @on?散着写 @on 也行(L04 那样),但当你有一组相关的监听(比如"完整的遥测"需要监听十几种事件)时,把它们收进一个类:① 能共享状态(self.formatter/self._telemetry);② 一次实例化全部生效,也能整体管理;③ 用户扩展时有明确的模板可继承。@on 是零件,BaseEventListener 是把零件组装成"一个完整监控模块"的外壳。
L07

事件类型体系 + 观察者模式的取舍

事件按领域分类,crew 类事件如(events/types/crew_events.py:36):

# events/types/crew_events.py:12
class CrewBaseEvent(BaseEvent):
    """Base class for crew events with fingerprint handling"""
    crew_name: str | None
    crew: Crew | None = None
    def _set_crew_fingerprint(self):
        if self.crew is not None and self.crew.fingerprint:
            self.source_fingerprint = self.crew.fingerprint.uuid_str
            self.source_type = "crew"

# events/types/crew_events.py:36
class CrewKickoffStartedEvent(CrewBaseEvent):
    inputs: dict[str, Any] | None
    type: Literal["crew_kickoff_started"] = "crew_kickoff_started"   # ★固定 type

class CrewKickoffCompletedEvent(CrewBaseEvent):
    output: Any
    type: Literal["crew_kickoff_completed"] = "crew_kickoff_completed"
    total_tokens: int = 0
CrewBaseEventcrew 类事件的中间基类,统一处理"填 crew 指纹"。task/agent/llm/memory 各有自己的中间基类。三层继承:BaseEvent → 领域基类 → 具体事件。
type: Literal[...]★每个具体事件用 Literal 钉死自己的 type 字符串。序列化后靠它反序列化回正确的类型。
各带专属字段started 带 inputs、completed 带 output+total_tokens——每种事件携带它特有的数据。监听器按需取。
💡 设计取舍②:观察者模式的代价是什么? 好处很明显(解耦、可扩展)。但代价也要清楚:① 控制流变"隐式"——看 kickoff 代码你不知道一句 emit 背后到底触发了多少 handler、干了什么,调试时得去总线查订阅者;② handler 里出错要小心——一个监听器抛异常不该拖垮整个执行(框架用 is_call_handler_safe 等做隔离);③ 顺序/时序敏感——比如 Day 22 讲的"完成事件必须在 drain 之后发",就是 emit 时序的坑。解耦不是免费的,它把"显式的调用"换成了"隐式的订阅",灵活性和可追踪性此消彼长。
⚠️ 边界:监听器泄漏与 replay 重复副作用 两个真实坑:① 监听器泄漏——单例总线全局存在,测试里反复注册同名 handler 不清理,会导致一个事件被处理多次。框架提供 scoped_handlers()event_bus.py:832)上下文管理器做临时隔离。② replay 重复副作用——Day 24 的 replay 会重放事件,但"写检查点""调外部 API"这类有副作用的 handler 不该重复执行。所以总线有 is_replaying() 标志(event_bus.py:72),这类 handler 要判断"如果在重放就早返回"。全局事件系统的两大经典陷阱,源码都专门留了应对机制。
L08

阶段4 收官

💡 一句话收束事件系统 发布者 emit(event) → 单例总线按 type(event) 查 handler 字典 → 同步 handler 走线程池、异步走专用 loop → 各监听器响应。监听器用 @bus.on(EventType) 单个注册,或继承 BaseEventListener 成组注册(构造即挂载)。事件都是 BaseEvent 的 Pydantic 子类,自带时间戳+溯源字段。整个系统让"执行核心"和"日志/遥测/UI/tracing"彻底解耦。

👶 小白:我想在 crew 每次开跑时发个钉钉通知,怎么做?

👨‍🏫 老师:不用改 CrewAI 任何源码——写一个 @crewai_event_bus.on(CrewKickoffStartedEvent) 装饰的函数,在里面调钉钉 API 就行。这就是发布-订阅的爽点:crew 早就在开跑时 emit 了那个事件,你只需去"订阅"它。想要成套监控就继承 BaseEventListener。核心执行代码一行都不用动。

🧠 今天你应该能回答

  • 事件系统解决什么问题?是什么设计模式?
  • BaseEvent 为什么继承 BaseModel?公共字段有哪些?
  • 事件总线为什么是单例?核心数据结构是什么?
  • @on 做了什么?为什么原样返回 handler?
  • emit 时同步/异步 handler 分别怎么执行?没订阅者会怎样?
  • BaseEventListener 相比散写 @on 好在哪?
  • 观察者模式的代价与两大经典陷阱是什么?

🎯 阶段4 知识地图(D19-26)

D19 Crew 是一张 Pydantic 字段表(六类字段+校验器)→ D20 顺序流程 = _execute_tasks 的 for 循环(传送带+异步 join)→ D21 层级流程 = 多造一个经理,靠委派工具把下属当工具调 → D22 kickoff 家族 = 一个同步核心 + 批量/线程异步/原生异步三维包装 → D23 规划 = 开工前用规划 agent 给任务贴计划 → D24 train/replay = 共享存档底座的调教与续跑 → D25 记忆开关 = 一个 memory 字段落地成共享档案柜 → D26 事件系统 = 让这一切的日志/遥测/UI 全部解耦的发布订阅底座。
一条主线:Crew(结构)→ 怎么跑(流程+kickoff)→ 跑得更好(规划/训练/记忆)→ 跑的过程可观测(事件)。

✋ 10 分钟动手

P=lib/crewai/src/crewai
sed -n '66,90p'   $P/events/base_events.py         # BaseEvent
sed -n '245,280p' $P/events/event_bus.py           # @on
sed -n '572,650p' $P/events/event_bus.py           # emit 分发
cat $P/events/base_event_listener.py               # BaseEventListener
# 亲手订阅一个事件
python -c "
from crewai import Agent, Task, Crew
from crewai.events.event_bus import crewai_event_bus
from crewai.events.types.crew_events import CrewKickoffStartedEvent
@crewai_event_bus.on(CrewKickoffStartedEvent)
def hi(source, event):
    print('>>> 我监听到 crew 开跑了:', event.crew_name)
a=Agent(role='助手', goal='帮忙', backstory='x')
t=Task(description='说句你好', expected_output='一句', agent=a)
Crew(name='测试组', agents=[a], tasks=[t]).kickoff()
"
下阶段预告 · Day 27 起(阶段5 工具系统):阶段4 我们把"团队怎么组织、怎么跑、怎么被观测"讲透了。接下来进入工具系统:从 BaseTool(D27)、StructuredTool(D28)、工具调用协议(D29),到缓存复用、MCP 工具、自定义工具最佳实践。Agent 之所以能"干实事",全靠工具——下一阶段见。
← Day 25 记忆开关 Day 27 · BaseTool →