Skip to content

实时智能体的自定义服务器

Supported in ADKPython v0.1.0

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 连接状态,以防止连接关闭时出现错误
  • 认证和授权:为你的端点实现认证和授权
  • 速率限制和配额:添加速率限制和超时控制。有关并发会话和配额管理的指导,请参阅并发会话
  • 结构化日志:使用结构化日志进行调试。
  • 持久化会话服务:考虑使用持久化会话服务(DatabaseSessionServiceVertexAiSessionService)。有关更多详细信息,请参阅 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 webadk 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);
    }
};

你的客户端必须产生和消费的媒体格式(采样率、编码、分块大小)在音频和视频中。流式标志(partialturnCompleteinterrupted)以及转录如何分段的详细信息在事件中。

序列化事件

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))