实时智能体的自定义服务器¶
adk web 工具用于开发目的运行实时智能体。它提供了一个浏览器客户端,可以捕获麦克风和摄像头、播放模型音频并渲染转录文本,这样你无需编写任何代码就能与智能体对话。部署到生产环境意味着替换它:运行你自己的服务器,将客户端桥接到 run_live(),在启动时初始化一次 runner 和会话服务,每个连接用户对应一个 LiveRequestQueue。
接下来是该桥接的完整 FastAPI 实现,以及客户端与之通信所需了解的内容。本文假设你已阅读 Sessions,其中涵盖了本示例所实践的生命周期。
FastAPI 应用示例¶
这个 FastAPI 应用实现了桥接。它运行两个并发任务:一个上游任务将 WebSocket 消息转发到 LiveRequestQueue,一个下游任务将 run_live() 事件转发回去。
import asyncio
from fastapi import FastAPI, WebSocket, WebSocketDisconnect
from google.adk.runners import Runner
from google.adk.agents.run_config import RunConfig
from google.adk.agents.live_request_queue import LiveRequestQueue
from google.adk.sessions import InMemorySessionService
from google.genai import types
from google_search_agent.agent import agent
# 应用设置(启动时执行一次)
APP_NAME = "live-agent"
app = FastAPI()
# 定义会话服务
session_service = InMemorySessionService()
# 定义 runner
runner = Runner(
app_name=APP_NAME,
agent=agent,
session_service=session_service
)
@app.websocket("/ws/{user_id}/{session_id}")
async def websocket_endpoint(websocket: WebSocket, user_id: str, session_id: str) -> None:
await websocket.accept()
# 每个会话的设置:RunConfig、会话、队列
response_modalities = ["AUDIO"]
run_config = RunConfig(
response_modalities=response_modalities,
input_audio_transcription=types.AudioTranscriptionConfig(),
output_audio_transcription=types.AudioTranscriptionConfig(),
session_resumption=types.SessionResumptionConfig()
)
session = await session_service.get_session(
app_name=APP_NAME,
user_id=user_id,
session_id=session_id
)
if not session:
await session_service.create_session(
app_name=APP_NAME,
user_id=user_id,
session_id=session_id
)
live_request_queue = LiveRequestQueue()
async def upstream_task() -> None:
"""从 WebSocket 接收消息并发送到 LiveRequestQueue。"""
try:
while True:
# 从 WebSocket 接收文本消息
data: str = await websocket.receive_text()
# 发送到 LiveRequestQueue
content = types.Content(parts=[types.Part(text=data)])
live_request_queue.send_content(content)
except WebSocketDisconnect:
# 客户端断开连接 - 通知队列关闭
pass
async def downstream_task() -> None:
"""从 run_live() 接收事件并发送到 WebSocket。"""
async for event in runner.run_live(
user_id=user_id,
session_id=session_id,
live_request_queue=live_request_queue,
run_config=run_config
):
# 将事件以 JSON 格式发送到 WebSocket
await websocket.send_text(
event.model_dump_json(exclude_none=True, by_alias=True)
)
# 并发运行两个任务
try:
await asyncio.gather(
upstream_task(),
downstream_task(),
return_exceptions=True
)
finally:
live_request_queue.close() # 始终关闭,即使出错也要关闭。
需要异步上下文
所有 ADK 双向流式应用必须在异步上下文中运行。这一要求来自多个组件:
run_live():ADK 的流式方法是一个异步生成器,没有同步包装器(与run()不同)- 会话操作:
get_session()和create_session()是异步方法 - WebSocket 操作:FastAPI 的
websocket.accept()、receive_text()和send_text()都是异步的 - 并发任务:上游/下游模式需要
asyncio.gather()来实现并发执行
所有代码示例都假设在异步上下文中(在 async def 或协程内)。它们展示了核心逻辑,省略了样板包装函数。
为什么需要两个任务¶
桥接是同时运行的两个循环,这正是实现双向通信的关键:
- 上游从 WebSocket 读取并推送到
LiveRequestQueue,这样用户可以在任何时刻发送输入,包括在智能体说话的过程中。 - 下游从
run_live()读取事件并写入 WebSocket,实时流式传输响应、转录和工具活动。
如果顺序运行它们,你将失去中断能力:当用户试图打断时,服务器会被阻塞在读取智能体的输出上。asyncio.gather() 正是让两个方向同时保持活跃的关键。
live_request_queue.close() 必须在每个退出路径上运行,包括异常情况。未关闭的队列会使 Live API 缺少终止信号,并可能导致会话卡在你的并发会话配额上直到超时,这就是 try/finally 的作用。
gather(..., return_exceptions=True) 会收集异常而不是抛出它们,因此如果需要区分正常断开和失败,请检查返回值。
生产环境注意事项¶
本示例展示了核心模式。对于生产应用,需要考虑:
- 错误处理(ADK):为 ADK 流式事件添加适当的错误处理。有关错误事件处理的详细信息,请参阅错误事件。
- 在关闭期间通过捕获
asyncio.CancelledError来优雅地处理任务取消 - 使用
return_exceptions=True检查asyncio.gather()的异常 - 异常不会自动传播
- 在关闭期间通过捕获
- 错误处理(Web):处理上游/下游任务中的 Web 应用特定错误。例如,使用 FastAPI 时你需要:
- 捕获
WebSocketDisconnect(客户端断开连接)、ConnectionClosedError(连接丢失)和RuntimeError(向已关闭的连接发送数据) - 在发送前使用
websocket.client_state验证 WebSocket 连接状态,以防止连接关闭时出现错误
- 捕获
- 认证和授权:为你的端点实现认证和授权
- 速率限制和配额:添加速率限制和超时控制。有关并发会话和配额管理的指导,请参阅并发会话。
- 结构化日志:使用结构化日志进行调试。
- 持久化会话服务:考虑使用持久化会话服务(
DatabaseSessionService或VertexAiSessionService)。有关更多详细信息,请参阅 ADK 会话服务文档。
连接客户端¶
你的服务器暴露了一个 WebSocket;需要有东西与之通信。在开发阶段,那是 adk web。在生产环境中,那是你编写的客户端:浏览器应用、移动应用,或电话/WebRTC 桥接。无论你构建什么,都继承相同的契约,因此值得了解 adk web 具体做了什么以及到哪里为止。
adk web 为你处理的功能:
| 功能 | 内置客户端的行为 |
|---|---|
| 麦克风 | 捕获并重采样为 16 kHz 单声道 PCM,以 audio/pcm;rate=16000 格式流式传输 |
| 播放 | 以 24 kHz 单声道 PCM 播放模型音频,无缝播放 |
| 摄像头 | 以约 1 fps 发送 JPEG 帧,格式为 image/jpeg |
| 转录 | 渲染用户和模型的转录文本,合并部分片段 |
| 打断 | 当事件到达且 interrupted 被设置时停止播放 |
它不做的事情,而生产客户端可能需要:
- 不支持屏幕共享,且没有活跃音频通话时不支持视频。
- 不支持模态选择;响应始终为
AUDIO。 - 没有用于主动性、情感对话、会话恢复、
save_live_blob或显式 VAD 信号的 UI。这些通过服务器端的RunConfig设置。 - 不支持手动 VAD;它依赖默认启用的服务器端自动检测。
adk web 和 adk api_server 都提供相同的 /run_live WebSocket;adk api_server 除非传入 --with_ui,否则不提供浏览器客户端。因此你可以针对 adk web 进行开发,然后将自定义客户端指向其中任何一个。
线路协议和数据格式¶
/run_live 端点仅使用 JSON 文本帧。你的客户端发送序列化的 LiveRequest 对象并接收序列化的 Event 对象。二进制数据(音频和图像字节)在 JSON 内部进行 base64 编码,而不是作为二进制 WebSocket 帧发送。
在客户端上,按照与 Python 中相同的事件字段进行分支处理,使用驼峰命名法:
websocket.onmessage = (message) => {
const adkEvent = JSON.parse(message.data);
if (adkEvent.interrupted) {
stopAudioPlayback(); // 用户打断;丢弃排队的音频
finishCurrentBubble();
return;
}
if (adkEvent.turnComplete) {
finishCurrentBubble();
return;
}
for (const part of adkEvent.content?.parts ?? []) {
if (part.text) appendText(part.text);
if (part.inlineData) enqueueAudio(part.inlineData.data);
}
};
你的客户端必须产生和消费的媒体格式(采样率、编码、分块大小)在音频和视频中。流式标志(partial、turnComplete、interrupted)以及转录如何分段的详细信息在事件中。
序列化事件¶
ADK 与 Live API 之间的 /run_live 端点仅使用 JSON 文本,但你的服务器和你的客户端之间的传输方式由你设计,在那里你可以发送二进制帧的音频以避免 base64 开销。
Event 是一个 Pydantic 模型,因此 model_dump_json() 可以将其转换为 JSON 字符串用于 WebSocket 或 SSE 传输。使用 by_alias=True 在客户端获得驼峰命名的字段名,使用 exclude_none=True 丢弃空字段:
async for event in runner.run_live(...):
await websocket.send_text(event.model_dump_json(exclude_none=True, by_alias=True))
inline_data 中的二进制音频在 JSON 中进行 base64 编码,这会使有效负载膨胀约 33%。对于音频密集型流,可以将音频作为二进制帧发送,元数据作为 JSON 发送:
async for event in runner.run_live(...):
parts = event.content.parts if event.content else []
audio_parts = [p for p in parts if p.inline_data]
if audio_parts:
for part in audio_parts:
await websocket.send_bytes(part.inline_data.data)
# 不包含音频字节的元数据。
await websocket.send_text(event.model_dump_json(
exclude={"content": {"parts": {"__all__": {"inline_data"}}}},
by_alias=True,
))
else:
await websocket.send_text(event.model_dump_json(exclude_none=True, by_alias=True))