Pydantic AI ProcessEventStream 能力详解:把 Agent 事件流转发、改写与观测纳入 Capability 体系
【免费下载链接】pydantic-aiHow Python does AI. Agents, realtime voice, image generation, embeddings. Every model, every interface, typed end to end.项目地址: https://gitcode.com/GitHub_Trending/py/pydantic-ai
本文围绕 Pydantic AI 的ProcessEventStream能力(capability)展开:它把一次 Agent 运行中的全部AgentStreamEvent(模型流式事件与工具执行事件)转发给你提供的处理器,注册该能力后agent.run()会自动启用流式,无需再显式传event_stream_handler参数。读完后你能掌握两种处理器形态(观察者 / 处理器)的语义差异、事件流在能力链中的传递规则,以及该能力与 durable execution、realtime 会话的组合边界,并能直接在项目里落地事件转发、审计日志与流式过滤等场景。
ProcessEventStream 是什么
[ProcessEventStream][pydantic_ai.capabilities.ProcessEventStream] 是 Pydantic AI 能力体系 中用于「处理事件流」的内置能力。它将 Agent 运行过程中产生的AgentStreamEvent流——包括模型流式响应事件和工具执行事件——转发给一个用户提供的异步处理器。注册该能力后,agent.run()会自动启用流式(streaming),因此即使不传显式的event_stream_handler参数,处理器也会照常触发。
在 realtime(实时语音)会话期间,同一事件流中还会出现 realtime 专属的RealtimeEvent成员,处理器同样会收到。
最简用法如下(与仓库文档 process-event-stream.md 中的示例一致):
from collections.abc import AsyncIterable from pydantic_ai import Agent, AgentStreamEvent, RunContext from pydantic_ai.capabilities import ProcessEventStream async def log_events(ctx: RunContext, events: AsyncIterable[AgentStreamEvent]) -> None: async for event in events: print(event) # (1)! agent = Agent('openai:gpt-5.2', capabilities=[ProcessEventStream(log_events)])- 典型去向:websocket 转发、进度条、审计日志等。
由于能力可以直接挂在capabilities=[...]上,事件观测逻辑与具体某次调用解耦——同一个 Agent 实例的所有运行都会经过这个处理器,适合做成统一的事件出口(如日志、遥测、UI 推送)。
两种处理器形态:观察者与处理器
ProcessEventStream的handler字段接受两种形态(见 process_event_stream.py 中handler: EventStreamHandlerFunc | EventStreamProcessorFunc的类型声明):
EventStreamHandler(观察者形态)——一个返回None的async def,如上文示例。事件被转发给该处理器的同时原样透传给下游,因此多个ProcessEventStream处理器(以及顶层event_stream_handler参数)可以互不干扰地观察同一条流。事件是同步送达的:一个慢处理器会形成反压(back-pressure),拖慢整条流。EventStreamProcessor(处理器形态)——一个 yield 事件的异步生成器。它 yield 出来的事件会替换下游消费者看到的流,因此可以修改、丢弃或注入事件。
这两种类型别名在 abstract.py 中有精确定义:
EventStreamHandler:Callable[[RunContext, AsyncIterable[AgentStreamEvent]], Awaitable[None]];EventStreamProcessor:接收同样参数、返回AsyncIterator[AgentStreamEvent]的异步生成器,官方注释明确说明其用于「通过ProcessEventStream能力修改、丢弃或添加能力链其余部分可见的事件」。
一个处理器形态的完整示例——丢弃所有PartStartEvent,只放行其余事件:
from collections.abc import AsyncIterable, AsyncIterator from pydantic_ai import Agent, AgentStreamEvent, RunContext from pydantic_ai.capabilities import ProcessEventStream from pydantic_ai.messages import PartStartEvent async def drop_part_starts( ctx: RunContext, stream: AsyncIterable[AgentStreamEvent] ) -> AsyncIterator[AgentStreamEvent]: async for event in stream: if isinstance(event, PartStartEvent): continue yield event agent = Agent('openai:gpt-5.2', capabilities=[ProcessEventStream(drop_part_starts)])仓库测试 test_processor_replaces_stream 系列用例 验证了该语义:处理器丢掉的PartStartEvent不会出现在下游(包括显式event_stream_handler参数注册的观察者)看到的流中;而 test_multiple_handlers_and_param_all_observe 则证明观察者形态下,两个ProcessEventStream加顶层参数三者收到的事件序列完全一致。
处理器形态的“全局替换”边界
源码 docstring 特别强调了几点容易被误解的边界,值得逐条对照:
- 替换是全局的,不是私有视图。一次运行只有一条事件流,处理器塑造的是整条流:丢弃或改写
PartDeltaEvent会同时改变run_stream()调用方stream_text()拿到的内容。 - 部分事件是控制信号。例如
FinalResultEvent告诉agent.run_stream()最终输出已经开始;把它丢掉会让run_stream()退化为等整个模型响应完成后再交付结果,而不是流式交付。过滤要谨慎。 - 不影响运行的最终输出。
ModelResponse在处理器看到事件之前就已从原始模型流累积完成,所以stream_output()和最终校验过的输出不受影响——丢事件只能改变「部分快照何时发出」,不能改变其内容。只想观察事件就用观察者形态。 - realtime 会话中同样只是消费者侧视图:转换或丢弃事件不会影响会话历史与工具执行。
测试 test_processor_shapes_streamed_text_but_not_the_output 把这条边界钉死了:改写所有文本 delta 为'XXX'后,stream_text()得到'hello XXX',而result.get_output()仍是原始内容'hello world'。
注册能力即自动启用流式
ProcessEventStream的 docstring(process_event_stream.py)说明:注册该能力后,agent.run()与AgentRun.next()会自动启用流式,处理器无需显式event_stream_handler参数即可触发;无论运行如何被驱动——包括agent.iter()以及手动对节点调用node.stream()——处理器都能看到相同的事件。
测试 test_handler_fires_under_every_drive_mode 用参数化的四种驱动方式(run、agent.iter()裸async for、next()逐步推进、手动node.stream())验证了这一点,并固化了包含工具调用时处理器实际看到的事件序列:
['PartStartEvent', 'PartEndEvent', 'FunctionToolCallEvent', 'FunctionToolResultEvent', 'PartStartEvent', 'FinalResultEvent', 'PartEndEvent']另一条相关保障是 test_next_does_not_force_streaming_without_event_hooks:没有任何能力注册事件流钩子时,next()不会强制走流式请求,运行继续使用非流式模型请求——自动启用流式是「能力注册」的副作用,而不是框架的默认行为。
同时要注意前提:模型必须支持流式。test_non_streaming_model_raises_a_clear_error 表明,对一个只实现了request()的模型,能力注册后运行会抛出带does not support streamed requests提示的UserError;而如果请求被wrap_model_request短路(缓存响应、SkipModelRequest)或经由模型选择能力替换成了支持流式的模型,则不需要配置模型本身支持流式(见 tests 中对应的三个用例)。
源码实现:观察者如何与主流并行
wrap_run_event_stream()是 AbstractCapability 提供的生命周期钩子之一,ProcessEventStream正是靠它接入事件链。其实现(process_event_stream.py)有几个值得细看的设计点:
形态探测:处理器被调用一次得到probe。若返回的是AsyncIterator,说明是处理器形态,直接async for消费并 yield;否则说明返回的是尚未 await 的协程(观察者形态),源码会先close()掉这个探测协程(此时什么都没执行),再以分叉出的接收流重新调用一次处理器。这种「靠返回值类型探测」的方式对普通函数和 callable 实例都稳健——测试 test_callable_instance_processor 专门验证了以类实例作为处理器也能被正确识别。
观察者分叉用内存对象流:观察者形态下,框架用anyio.create_memory_object_stream()建一对发送/接收流,主循环每取出一个事件就send(event)给观察者,然后原样yield event向下传递。观察者的await send_stream.send(event)是同步的,因此慢观察者天然形成反压,与文档描述一致。
观察者跑在独立asyncio任务里:源码注释(L118-L124)解释了为何不用 anyio 任务组——任务组绑定进入它的任务,而节点流会被记忆化(memoized)、可能在其他任务中恢复(比如在另一个任务里消费StreamedRunResult),跨任务退出 cancel scope 会抛 anyio 的 "cancel scope in a different task" 错误;而普通任务没有这种亲和性。测试 test_streamed_result_can_be_consumed_in_another_task 正锁定了这个场景:asyncio.create_task(result.get_output())跨任务消费不会崩溃。
优雅退出语义:
- 观察者提前结束迭代(如读到第一个事件就
return)只停止它自己的投递,下游仍能看到全部事件——见 test_observer_bailout_does_not_break_downstream。 - 观察者抛异常则把异常传播给整个运行;源码还会取消在途的上游拉取并尽力关闭源迭代器,见 test_failing_observer_interrupts_stalled_stream。
- 被包装的流失败、消费者提前退出、或节点流被关闭时,观察者任务会被
cancel_and_drain拆掉,不会滞留在receive上。一系列 teardown 用例(test_failing_stream_tears_down_the_handler、test_abandoned_model_request_stream_tears_down_the_handler 等)覆盖了这些路径。
另外两个能力元信息:get_serialization_name()返回None,即该能力持有 callable,不能参与 YAML/JSON agent spec 的声明式构建(测试 test_not_spec_serializable 予以确认);_emits_app_events属性为True,表示该能力会让运行产生应用侧事件,这也是框架判定「需要启用流式」的依据之一。
与其他流式机制的组合
ProcessEventStream与 Pydantic AI 的既有流式机制是叠加关系而非替代关系:
- 顶层
event_stream_handler=参数、agent.iter()+node.stream()手动流式、run_stream()的结果流,都可以与能力注册并存,且观察者形态下彼此看到的是同一条未改写的流(多个处理器 + 参数处理器三者事件序列相等的测试见上文)。 - 事件词汇表与更多处理器示例见 agent.md 的 "Streaming All Events" 一节。
- 在 hooks 能力 的
on_event监听场景下,能力事件监听会触发相同的「自动启用流式」行为,两者语义一致。
Durable Execution 下的确定性约束
文档特别提示(与源码 docstring 一致):在 durable execution 能力(TemporalDurability、DBOSDurability、PrefectDurability,总览见 durable_execution/overview.md)下,ProcessEventStream的处理器运行在 workflow/flow 代码里,必须保持确定性,因为 workflow 重放(replay)时它会再次执行。具体时序为:工具调用与最终输出事件实时送达,而模型事件是每次模型请求 activity/step/task 完成后重放的真实捕获事件。若处理器内部有必须「恰好执行一次」的 I/O(写库、发通知等),应改把event_stream_handler=传给 durability 能力本身,而不是依赖这个 capability。
仓库的 durable 测试(test_durability.py、test_dbos.py、test_prefect.py)中均有ProcessEventStream的用法,可作为该场景下的参考实现。
小结
ProcessEventStream把「看事件」从每次调用的参数提升为 Agent 的一等配置:
- 观察者形态适合旁路消费(websocket、审计、遥测),不影响下游,多观察者互不干扰,代价是慢处理器会反压全流;
- 处理器形态可以重塑下游看到的流(过滤、改写、注入),但要清楚它同时影响
stream_text()与控制信号(如FinalResultEvent),且永远不改变运行的最终输出; - 注册即自动流式,四种运行驱动方式下行为一致;
- 持有 callable 因而不可 spec 序列化,durable 场景下需保证处理器确定性,或将一次性 I/O 移交给 durability 能力的
event_stream_handler=参数。
深入阅读路径:能力总览 docs/capabilities/overview.md、事件词汇表 docs/agent.md#streaming-all-events、实现 pydantic_ai_slim/pydantic_ai/capabilities/process_event_stream.py、行为测试 tests/test_capability_process_event_stream.py、持久化运行 docs/durable_execution/overview.md。
【免费下载链接】pydantic-aiHow Python does AI. Agents, realtime voice, image generation, embeddings. Every model, every interface, typed end to end.项目地址: https://gitcode.com/GitHub_Trending/py/pydantic-ai
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考