事件系统:日志/遥测/流式/tracing 的统一底座
整个阶段4 你无数次见到 crewai_event_bus.emit(self, CrewKickoffStartedEvent(...))——kickoff 开始/结束/失败、训练开始/完成、记忆保存完成……crew 在关键节点发事件,但它不关心谁在听。而日志打印、遥测上报、流式输出、tracing 全都是事件的监听器。这就是经典的观察者模式(发布-订阅)。今天拆 events/ 目录:事件基类 BaseEvent、单例总线 CrewAIEventsBus、@on 注册、emit 分发到同步/异步 handler、BaseEventListener 自动注册。这是阶段4 的收官,也是理解 CrewAI"可观测性"的钥匙。
痛点:日志、遥测、UI……都要"插进执行流程"
emit(事件)(发布),谁想响应就 @on(事件类型) 注册一个 handler(订阅)。中间是一个全局单例总线 CrewAIEventsBus 做撮合:它维护"事件类型 → handler 集合"的映射,emit 时查出对应 handler 挨个调用。核心代码里看不到任何日志/遥测逻辑,只有一行行 emit——旁路关注点全在监听器里,彻底解耦。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 发的。监听器不用猜,直接读这些字段就知道上下文。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 到处都是它。import crewai_event_bus 就能 emit/on,无需传递。代价是单例是"隐式全局状态"(测试时要注意 handler 泄漏、要能 scoped_handlers 隔离)。但对"横切整个系统的事件总线"而言,单例是最自然的选择——这也是几乎所有框架的事件系统都用单例的原因。@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")。普通监听用不到。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 各走各的执行器。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_attrs 里 event_listener = EventListener() 就实例化了它。一个类 = 一组相关监听把"和某个关注点相关的一堆监听"打包进一个 Listener 类,逻辑内聚。想加自定义监控,就写个自己的 BaseEventListener 子类。@on 也行(L04 那样),但当你有一组相关的监听(比如"完整的遥测"需要监听十几种事件)时,把它们收进一个类:① 能共享状态(self.formatter/self._telemetry);② 一次实例化全部生效,也能整体管理;③ 用户扩展时有明确的模板可继承。@on 是零件,BaseEventListener 是把零件组装成"一个完整监控模块"的外壳。事件类型体系 + 观察者模式的取舍
事件按领域分类,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——每种事件携带它特有的数据。监听器按需取。is_call_handler_safe 等做隔离);③ 顺序/时序敏感——比如 Day 22 讲的"完成事件必须在 drain 之后发",就是 emit 时序的坑。解耦不是免费的,它把"显式的调用"换成了"隐式的订阅",灵活性和可追踪性此消彼长。scoped_handlers()(event_bus.py:832)上下文管理器做临时隔离。② replay 重复副作用——Day 24 的 replay 会重放事件,但"写检查点""调外部 API"这类有副作用的 handler 不该重复执行。所以总线有 is_replaying() 标志(event_bus.py:72),这类 handler 要判断"如果在重放就早返回"。全局事件系统的两大经典陷阱,源码都专门留了应对机制。阶段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()
"
BaseTool(D27)、StructuredTool(D28)、工具调用协议(D29),到缓存复用、MCP 工具、自定义工具最佳实践。Agent 之所以能"干实事",全靠工具——下一阶段见。