Pydantic AI ProcessEventStream 能力详解:把 Agent 事件流转发、改写与观测纳入 Capability 体系
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 中有精确定义EventStreamHandlerCallable[[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的 docstringprocess_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),仅供参考