news 2026/9/14 12:02:29

使用 Pydantic AI 的 VercelAIAdapter 对接 Vercel AI Data Stream Protocol 构建流式聊天后端

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
使用 Pydantic AI 的 VercelAIAdapter 对接 Vercel AI Data Stream Protocol 构建流式聊天后端

使用 Pydantic AI 的 VercelAIAdapter 对接 Vercel AI Data Stream Protocol 构建流式聊天后端

【免费下载链接】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 官方提供的 Vercel AI 集成能力:通过VercelAIAdapter将前端(AI SDK UI 的useChat等 hooks)发来的 Vercel AI 请求输入转换为Agent.run_stream_events()的参数,运行 Agent 后把 Pydantic AI 事件流实时转换为 Vercel AI 的 SSE 事件流返回前端。读完本文,你将掌握 Starlette/FastAPI 下的一行式接入、无 Starlette 框架下的手动编排、取消与工具审批、自定义事件数据下发、客户端工具文件回传、消息元数据往返,以及系统提示词与信任模型的完整配置。

协议对接的总体架构

Vercel AI Data Stream Protocol 是 AI SDK UI 生态(useChatuseAssistant等 hooks 与 AI Elements 组件)默认使用的流式协议。Pydantic AI 在pydantic_ai.ui.vercel_ai模块中实现了该协议的双向转换,核心是两个类:

  • VercelAIAdapter:负责把前端请求体转换为Agent.run_stream_events()的输入、运行 Agent、再把 Pydantic AI 事件转换为 Vercel AI 事件。它继承自抽象的UIAdapter基类(与 AG-UI 协议 的AGUIAdapter平级,二者共用同一套适配器骨架)。
  • VercelAIEventStream:负责把 Pydantic AI 原生事件流转换为 Vercel AI 的 chunk 序列并编码为 SSE 字符串。常规请求路径下你无需直接使用它,只有 Agent 事件不是经由"服务于前端的那个请求"到达你时(例如 durable execution 工作流、消息队列或 WebSocket 扇出场景),才需要单独实例化它来转换事件,参见 UI Event Streams 文档 中的encode_events示例。

从源码看,VercelAIAdapter是一个 dataclass,声明了sdk_version(默认5)、server_message_id等字段,并复用了UIAdapter提供的dispatch_requestfrom_requestrun_streamrun_stream_nativeencode_streamstreaming_responsetransform_stream等完整方法族。整个请求处理链条为:build_run_input(请求体)VercelAIAdapter(agent, run_input, accept)run_stream(...)encode_stream(...)StreamingResponse

与 Starlette / FastAPI 一行式接入

如果后端基于 Starlette 系框架(FastAPI、Litestar 等),VercelAIAdapter.dispatch_request()类方法可以直接在端点函数中消费请求并返回流式响应,这是官方推荐的最简路径:

from fastapi import FastAPI from starlette.requests import Request from starlette.responses import Response from pydantic_ai import Agent from pydantic_ai.ui.vercel_ai import VercelAIAdapter agent = Agent('openai:gpt-5.2') app = FastAPI() @app.post('/chat') async def chat(request: Request) -> Response: return await VercelAIAdapter.dispatch_request(request, agent=agent)

除了requestagentdispatch_request还接受与Agent.run_stream_events()相同的可选参数(message_historydeferred_tool_resultsconversation_idrun_idmodelinstructionsdepsoutput_typemodel_settingsusage_limitsusagemetadatainfer_nametoolsetscapabilities等),以及两个回调:

  • on_complete:Agent 运行成功时触发,接收AgentRunResult,可额外产出 Vercel AI 事件;
  • on_cancel:Agent 因一等取消(first-party cancellation)结束时触发,接收RunCancelled,同样可额外产出事件。

调用链在 UIAdapter.dispatch_request 中实现:先from_request()解析请求,若请求体校验失败则返回 422(UNPROCESSABLE_ENTITY)JSON 错误,随后streaming_response(adapter.run_stream(...))生成流式响应。注意请求的Accept头会被用作流式响应的media_type

无 Starlette 框架时的手动编排(Advanced Usage)

对于 Django、Flask 等非 Starlette 框架,或需要细粒度控制输入输出的场景,可以直接实例化VercelAIAdapter并链式调用其方法,效果等价于dispatch_request

  1. 解析请求体VercelAIAdapter.build_run_input(await request.body())返回 Vercel AI 的RequestData对象。其实现是request_data_ta.validate_json(body)(见 _adapter.py),即用 Pydantic TypeAdapter 严格校验 JSON。如果是 Starlette/FastAPI 环境,也可直接用VercelAIAdapter.from_request()一步构建适配器实例。
  2. 运行 Agentadapter.run_stream(...)运行 Agent 并返回 Vercel AI 事件流,支持与Agent.run_stream_events()相同的参数及on_complete/on_cancel。它内部先走run_stream_native()(返回 Pydantic AI 原生事件),再经transform_stream()转换;你也可以自己先拿run_stream_native()的原生事件流,再调用transform_stream()手动转换。
  3. 编码 SSEadapter.encode_stream(event_stream)把 Vercel AI 事件流编码为 SSE 字符串;或者直接用adapter.streaming_response(...)生成 Starlette/FastAPI 的流式响应对象。

完整的可运行示例(含输入校验、422 返回、取消 token 注册表清理)参见官方文档 run_stream.py 中的对应示例与下方"取消"一节。该示例使用 FastAPI,但可改造适配任意 Web 框架。

取消机制:一等取消与外部取消

当一个运行以一等取消结束(来自工具内ctx.cancel()AgentRun.cancel()、或你服务器端在取消端点接线的CancellationToken),适配器会向流中发射一个 Vercelabortchunk。useChat会保留已产生的部分消息,并在onFinish中上报isAbort,而不是进入错误状态。此时on_cancel回调被触发,可在其中持久化可恢复的消息历史(RunCancelled.all_messages()返回可继续传递的消息列表)。

重要区别:客户端调用stop()会直接中断浏览器请求,服务器侧看到的是一次"断开"(disconnect),属于外部取消——运行以asyncio.CancelledError方式被撕裂,不会发射abortchunk,on_cancel也不会触发(客户端反正已经断开了)。若希望在用户点"停止"时也能拿到abortchunk 并执行on_cancel,应当保持流连接,改用一等取消:给运行传入CancellationToken,并暴露一个独立端点(如POST /chat/{id}/cancel)调用token.cancel()。完整实现如下:

import json from collections.abc import AsyncIterator from http import HTTPStatus from fastapi import FastAPI from fastapi.requests import Request from fastapi.responses import Response, StreamingResponse from pydantic import ValidationError from pydantic_ai import Agent, CancellationToken, RunCancelled from pydantic_ai.ui import SSE_CONTENT_TYPE from pydantic_ai.ui.vercel_ai import VercelAIAdapter agent = Agent('openai:gpt-5.2') app = FastAPI() cancellation_tokens: dict[str, CancellationToken] = {} async def on_cancel(cancelled: RunCancelled) -> None: messages = cancelled.all_messages() # (1)! print(f'cancelled after {len(messages)} messages') @app.post('/chat/{chat_id}') async def chat(chat_id: str, request: Request) -> Response: accept = request.headers.get('accept', SSE_CONTENT_TYPE) try: run_input = VercelAIAdapter.build_run_input(await request.body()) except ValidationError as e: return Response( content=json.dumps(e.json()), media_type='application/json', status_code=HTTPStatus.UNPROCESSABLE_ENTITY, ) adapter = VercelAIAdapter(agent=agent, run_input=run_input, accept=accept) cancellation_token = CancellationToken() cancellation_tokens[chat_id] = cancellation_token event_stream = adapter.run_stream( cancellation_token=cancellation_token, on_cancel=on_cancel ) async def encode_stream() -> AsyncIterator[str]: try: async for event in adapter.encode_stream(event_stream): yield event finally: if cancellation_tokens.get(chat_id) is cancellation_token: cancellation_tokens.pop(chat_id, None) return StreamingResponse(encode_stream(), media_type=accept) @app.post('/chat/{chat_id}/cancel', status_code=HTTPStatus.NO_CONTENT) async def cancel_chat(chat_id: str) -> None: if token := cancellation_tokens.get(chat_id): token.cancel()
  1. 这是需要持久化的可恢复历史——在后续运行中作为message_history传入即可续接对话。

注意:上述内存 token 注册表要求单进程部署或粘性路由。多 worker 部署时,需借助消息代理等共享协调机制,把取消请求路由到持有该运行的 worker。

从源码看,VercelAIEventStream.on_cancelled()发射的AbortChunk内容为reason='The agent run was cancelled.'(见 _event_stream.py);而on_error()会设置finish_reason='error'并发射ErrorChunk(_event_stream.py)。

向客户端下发数据:Data Chunks

运行过程中(例如长耗时工具的执行进度)想向客户端推送自定义数据,可以在工具内通过ctx.emit()发射一个CustomEvent

from dataclasses import dataclass from pydantic_ai import Agent, CustomEvent, RunContext agent = Agent('openai:gpt-5.2') @dataclass(kw_only=True) class FileUploadProgressEvent(CustomEvent): done: int total: int @agent.tool async def upload_files(ctx: RunContext, total: int) -> str: for done in range(1, total + 1): # 完成一个单位的工作,然后告诉前端进度 await ctx.emit(FileUploadProgressEvent(done=done, total=total)) return f'Uploaded {total} files'

每个事件会以DataChunk形式到达客户端,typedata-{name}datato_payload()的结果——上例即type='data-file_upload_progress'data={'done': 1, 'total': 3}。chunk 在事件发射的当下、工具仍在运行时即到达,天然支持进度条。

形状一致性:无论事件是否从工具调用内部发射,data形状都一致,因此前端针对某一形状编写的代码,不会因为同一事件类日后改在别处发射而失效。如需自定义形状(比如按前端期望命名字段、把工具归属信息放到线上),覆盖to_payload()即可:

from dataclasses import dataclass from typing import Any from pydantic_ai import CustomEvent @dataclass(kw_only=True) class FileUploadPhaseEvent(CustomEvent): done: int total: int def to_payload(self) -> dict[str, Any]: return { 'completed': self.done, 'total': self.total, 'toolCallId': self.tool_call_id, }

特殊规则:若to_payload()返回一个数据承载型 chunk(见下文),则该 chunk 原样透传;事件类声明为ui=False时永不下发(仅服务端消费的事件留在服务端);进程从未 import 过的事件类也不会被转发,因为其 opt-out 标记是挂在类上而非走线上。

工具返回时携带 chunk:工具可通过返回带metadata(单个或列表)的ToolReturn对象,把 Vercel AI data stream chunk 附加到工具结果上。支持四种类型:DataChunkSourceUrlChunkSourceDocumentChunkFileChunk

from pydantic_ai import Agent, ToolReturn from pydantic_ai.ui.vercel_ai.response_types import DataChunk, SourceUrlChunk agent = Agent('openai:gpt-5.2') @agent.tool_plain async def search_docs(query: str) -> ToolReturn: return ToolReturn( return_value=f'Found 2 results for "{query}"', metadata=[ SourceUrlChunk( source_id='doc-1', url='https://example.com/docs/intro', title='Introduction', ), DataChunk( type='data-search-results', data={'query': query, 'count': 2}, ), ], )

ctx.emit发射的事件不同,这些 chunk 属于消息的一部分,能随消息历史往返而幸存——这正是前端需要重建的数据(如答案背后的来源 URL)所期望的;代价是它们在工具返回时而非运行过程中下发。源码中iter_metadata_chunks()(见 _utils.py)只会转发上述四种>@app.post('/chat') async def chat(request: Request) -> Response: return await VercelAIAdapter.dispatch_request(request, agent=agent, sdk_version=6)

sdk_version=6时适配器会:

  1. 在调用requires_approval=True的工具时发射tool-approval-requestchunk;
  2. 自动从后续请求中提取审批响应;
  3. 为被拒绝的工具发射tool-output-deniedchunk。

前端方面,AI SDK UI 的useChathook 处理审批流程;可使用 AI Elements 的Confirmation组件做现成审批 UI,或用 hook 的addToolApprovalResponse自行构建。

审批响应的严格性:审批响应按协议经useChataddToolApprovalResponse与参考 Next.js 后端往返,设计上被信任。但审批决定本身必须是真正的 JSON 布尔值——ToolApprovalResponded.approved字段是严格布尔(StrictBool),任何替代值(1"true"0"false")都会导致请求校验失败,而不会被强制转换成一个决定(见 request_types.py 的注释,防止{'approved': 1}被宽松模式强转为通过)。若需要把审批决定绑定到服务端状态而非请求本身,可拦截DeferredToolRequests,在服务端持久化审批 ID,并在续接时显式传入deferred_tool_results

补充说明:源码中sdk_version类型为Literal[5, 6, 7],默认5以保证向后兼容;7的线协议与6完全相同(v7 的>from fastapi import FastAPI from starlette.requests import Request from starlette.responses import Response from pydantic_ai import Agent from pydantic_ai.ui.vercel_ai import VercelAIAdapter agent = Agent('openai:gpt-5.2') app = FastAPI() @app.post('/chat') async def chat(request: Request) -> Response: return await VercelAIAdapter.dispatch_request( request, agent=agent, manage_system_prompt='client' )

协议细节的源码级补充

  • Finish reason 映射:Pydantic AI 的结束原因经_FINISH_REASON_MAP(_event_stream.py)映射为 Vercel 格式:stopstoplengthlengthcontent_filtercontent-filtertool_calltool-callserrorerror,未知值统一为other
  • SSE 编码与响应头:每个 chunk 编码为data: {json}\n\nencode_event),且流式响应带x-vercel-ai-ui-message-stream: v1响应头(VERCEL_AI_DSP_HEADERS),这是 AI SDK UI 识别 data stream 协议的标准信号。
  • 消息 ID 生成_generate_message_id按优先级生成确定性消息 ID:有provider_response_id{provider_response_id}-{index},有run_id{run_id}-{index},否则用uuid5(timestamp-kind-role-index)(_adapter.py)。
  • conversation_id关联:适配器把请求体顶层的id(chat ID)作为conversation_id,用于跨多次运行关联 OpenTelemetry span 与消息历史。
  • 往返的已知损耗dump_messages → load_messages对工具结果并非完全无损——RetryPromptPart重新加载后会变成outcome='failed'ToolReturnPart(协议没有独立 retry 概念);无法解析为 JSON 对象的ToolCallPart.args会被重写为{'INVALID_JSON': '<raw args>'}。需要 retry 语义跨往返存活时,应在进程内维护对话而非经由 Vercel AI 线格式持久化。

验证与测试

仓库的测试套件(tests/test_ui.py)覆盖了适配器的请求解析、事件流编码与消息往返;tests/test_vercel_ai.pytests/cassettes/test_vercel_ai/中的 VCR 磁带用于验证真实模型调用下的流式行为。测试同时覆盖了sdk_version=5/6两种线格式下的 chunk 输出差异,可作为理解协议行为的补充参考。

总结

VercelAIAdapter把 Pydantic AI 的 Agent 运行无缝接入 Vercel AI Data Stream Protocol 生态:dispatch_request一行式接入 FastAPI,手动编排适配任意框架,取消、自定义数据、工具审批、消息元数据与压缩均开箱即用。安全方面,系统提示词归属、文件 URL scheme 白名单与上传文件开关共同构成默认信任模型;在对外暴露该端点时,请务必把它当作内部后端服务,放在你自己已认证的路由处理器之内。

【免费下载链接】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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/14 12:00:51

光储系统双层优化模型与改进粒子群算法实践

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/14 11:55:05

音频在线预览工具开发实战:URL校验、跨域与防盗链排查指南

拿一条音频链接想快速听一下效果&#xff0c;最常见的动作是下载到本地再打开播放器&#xff0c;遇到大文件或者临时地址过期&#xff0c;一折腾就是好几分钟。我做了一个音频在线预览工具&#xff0c;输入URL即刻播放远程音频&#xff0c;同时能看出这个链接到底能不能用、失败…

作者头像 李华