Skip to content

实时智能体的会话管理

Supported in ADKPython v0.1.0

实时智能体是一种在用户说话、收听、打断和沉默期间始终保持连接的状态。

实时智能体使用与其他 ADK 智能体相同的 SessionSessionService 和状态模型,这些内容都在对话上下文中有介绍。实时会话新增的是一连接:它可能会断开、超时,或者比模型的上下文窗口存活得更久。关于该连接返回的内容,请参阅事件;关于影响连接的配置,请参阅配置

搭建实时应用

实时应用包含两种对象:一种在启动时创建一次并在所有会话中复用,另一种是每个会话新建的。

创建一次,全局复用:

  • Agent:你的模型、工具和指令。无状态且可复用。
  • SessionService:存储对话历史,使会话在重连和重启后得以保留。
  • Runner:驱动智能体并产出事件的运行时。
import os
from google.adk.agents import Agent
from google.adk.runners import Runner
from google.adk.sessions import InMemorySessionService
from google.adk.tools import google_search

APP_NAME = "live-agent"

agent = Agent(
    name="google_search_agent",
    model=os.getenv("DEMO_AGENT_MODEL", "gemini-live-2.5-flash-native-audio"),
    tools=[google_search],
    instruction="You are a helpful assistant that can search the web.",
)

runner = Runner(
    app_name=APP_NAME,
    agent=agent,
    session_service=InMemorySessionService(),
)

InMemorySessionService 在进程停止后会丢失状态。对于生产环境,请使用 DatabaseSessionService(SQLite、PostgreSQL 或 MySQL)或 VertexAiSessionService(在 Google Cloud 上托管)。参见会话服务

每个会话新建:

  • 一个 Session,在循环运行前获取或创建。
  • 一个 RunConfig,可以按用户设置不同(语音、转写、限制)。
  • 一个 LiveRequestQueue,你通过它发送用户输入的通道。
from google.adk.agents.live_request_queue import LiveRequestQueue
from google.adk.agents.run_config import RunConfig
from google.genai import types

# 获取或创建,同时处理新对话和重连。
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
    )

run_config = RunConfig(
    response_modalities=["AUDIO"],
    session_resumption=types.SessionResumptionConfig(),
)

live_request_queue = LiveRequestQueue()

user_idsession_id 是你定义的任意字符串;如果传入 session_id=None,ADK 会生成一个 UUID。在使用相同标识符调用 run_live() 之前,会话必须已存在,否则 run_live() 会抛出 ValueError: Session not found

每个会话一个队列

切勿跨会话复用 LiveRequestQueue。关闭信号会残留在队列中并被带到下一个会话,导致数据损坏。每次 run_live() 调用都要创建新的队列。

LiveRequestQueue

LiveRequestQueue 是你向智能体发送消息的通道。每条消息都是一个 LiveRequest,它是一个包含不同类型输入的单一容器:

参考:<a href="../api-reference/python/google-adk.html#google.adk.agents.LiveRequestQueue">LiveRequestQueue</a>
class LiveRequest(BaseModel):
    content: Optional[Content] = None            # 文本和结构化数据
    blob: Optional[Blob] = None                  # 音频/视频字节
    activity_start: Optional[ActivityStart] = None  # 手动轮次开始
    activity_end: Optional[ActivityEnd] = None      # 手动轮次结束
    close: bool = False                          # 优雅终止

contentblob 互斥。请使用便捷方法而不是自己构建 LiveRequest 对象;它们会设置正确的字段并确保你遵守这一约束。

方法 发送内容 模式
send_content(content) 文本,作为离散轮次 逐轮模式;触发响应
send_realtime(blob) 音频、图像或视频字节 持续流式传输
send_activity_start() / send_activity_end() 手动轮次边界 仅在禁用自动 VAD 时使用
close() 终止信号 结束会话
from google.genai import types

# 文本轮次。
live_request_queue.send_content(types.Content(parts=[types.Part(text=user_text)]))

# 音频块(持续流式传输)。
live_request_queue.send_realtime(
    types.Blob(mime_type="audio/pcm;rate=16000", data=audio_data)
)

有关音频、图像和视频格式,请参阅音频和视频。有关使用活动信号进行手动轮次控制,请参阅语音活动检测

每次调用只发送一个文本 Part

每次 send_content() 调用只发送一个文本 Part。某些实时模型会将多部分 Content 视为对话预填充(历史记录准备),而非需要响应的轮次,因此每次调用只用一个 Part 可以确保在不同模型间保持一致的行为。

并发和排序

LiveRequestQueue 封装了 asyncio.Queue,这带来三个影响:

  • 发送方法是同步的。 它们底层调用 put_nowait(),因此永远不会阻塞,也不需要 await
  • 投递是 FIFO 且不合并的。 请求按发送顺序到达模型,每次调用一个。
  • 队列是无界的。 发送速度快于模型消费速度会增加内存占用,而不是施加背压,因此对于高频率的音频或视频,请自行限制发送速率。

在异步上下文中创建队列,以便将其绑定到运行 run_live() 的事件循环。asyncio.Queue 在单个事件循环线程中是并发安全的;要从另一个线程写入数据,请使用 loop.call_soon_threadsafe()

run_live() 循环

run_live() 是一个异步生成器。它在事件生成时立即产出 Event 对象,没有缓冲或轮询,同时你可以通过队列并发发送新输入。这种并发正是实现打断功能的关键:智能体可以在说话的同时用户开始插话。

参考:<a href="../api-reference/python/google-adk.html#google.adk.runners.Runner.run_live">Runner.run_live()</a>
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,
):
    await websocket.send_text(event.model_dump_json(exclude_none=True, by_alias=True))

run_live() 在你调用时打开 Live API 连接,在循环运行期间双向流式传输,并在你调用 live_request_queue.close() 时关闭连接。关于它产出的事件类型及如何处理,请参阅事件

run_live() 的退出条件

退出条件 触发方式 优雅退出
手动关闭 live_request_queue.close()
工作流完成 实时工作流中最后一个智能体调用 task_completed()
会话超时 达到 Live API 时长限制(未启用压缩) 连接关闭
提前退出 工具或回调设置了 end_invocation
错误 连接失败或未处理的异常

会话结束时始终调用 close(),即使出错也不例外。跳过它会让 Live API 没有收到优雅终止信号,可能导致"僵尸"会话持续占用你的并发会话配额,直到超时。

try:
    await asyncio.gather(upstream_task(), downstream_task())
except WebSocketDisconnect:
    pass  # 客户端正常断开。
finally:
    live_request_queue.close()  # 始终关闭队列。

关于循环内的错误处理,请参阅错误事件。关于完整的上游/下游服务器模式,请参阅自定义服务器

保存到会话的内容

run_live() 退出时,只有部分事件会持久化到 ADK Session 中:

  • 已保存: 最终(非部分)转写、使用量元数据、函数调用和响应,以及大多数控制事件。音频文件仅在 save_live_blobTrue 时才会保存。
  • 临时的: 原始音频字节(inline_data)和部分转写,这些内容被产出用于实时播放和显示,但不会被存储。

ADK Session 与 Live API session 的区别

两个不同的东西共享了"会话"这个词:

  • ADK Session(由 SessionService 管理)是持久化的对话存储。它在多次 run_live() 调用和应用重启之间保持存在。
  • Live API session(由 Live API 后端管理)是一个临时的流式上下文,仅在循环运行期间存在。

run_live() 启动时,ADK 从 ADK Session 加载历史记录,用它初始化一个新的 Live API session,并在事件发生时更新 ADK Session。当循环结束时,Live API session 被销毁,而 ADK Session 持久保存。下一次调用会从存储的历史记录重建 Live API session。这种分离机制使得对话能够在网络断开和重启后继续。

在传输层,还有一个区分对可靠性很重要:

  • 连接是 ADK 与 Live API 之间的 WebSocket 链接。它可能会超时。
  • 会话是对话上下文,可以通过会话恢复跨越多个连接存在。

平台限制

两个后端都对连接时长、会话时长和并发会话数设置了上限。具体数值因后端而异且会随时间变化,因此支持的模型在一个地方追踪这些信息。

其中两个上限会影响你的代码编写方式。上下文窗口压缩可以解除会话时长限制,而并发会话上限则是你在并发会话中需要设计应对的。

会话恢复

Live API 在大约 10 分钟后关闭每个 WebSocket 连接。会话恢复将对话迁移到新连接上,使其能够超过该限制继续进行。启用后,ADK 会为你处理所有重连逻辑,缓存恢复句柄、检测关闭并在后台重新连接。你的 run_live() 循环会不间断地继续产出事件。

from google.genai import types

run_config = RunConfig(session_resumption=types.SessionResumptionConfig())

ADK 仅管理 ADK 到 Live API 的连接。你的应用仍然负责自己的客户端连接(例如用户到你服务器的 WebSocket)以及任何客户端重连逻辑。

ADK 的重连方式:

  1. Live API 发送 session_resumption_update 消息;ADK 缓存最新的句柄。
  2. 在限制之前,Live API 可能会发送 go_away 警告;ADK 在断开之前重新连接,因此切换是无感的。
  3. 当连接优雅关闭时,ADK 的循环使用缓存的句柄重新连接,会话以完整上下文继续。
sequenceDiagram
    participant App as 你的应用
    participant ADK as ADK (run_live)
    participant API as Live API

    App->>ADK: run_live(run_config with session_resumption)
    ADK->>API: WebSocket connect()
    Note over ADK,API: 流式传输(0-10 分钟)
    API-->>ADK: session_resumption_update { handle }
    ADK->>ADK: 缓存句柄
    Note over API: ~10 分钟:连接优雅关闭
    ADK->>API: reconnect(handle)
    API-->>ADK: 会话恢复,完整上下文
    Note over App,API: 循环继续,无中断

重连尝试有上限

ADK 最多重试 5 次连续重连(DEFAULT_MAX_RECONNECT_ATTEMPTS)。计数器在每次成功重连后重置,因此长时间对话仅受限于连续五次失败,而非总共五次重连。ADK 仅在存在恢复句柄时才重试;如果未启用 session_resumption,第一次断开会直接从 run_live() 中抛出,你的应用必须自行处理。

仅在短会话(少于 10 分钟)、无状态请求-响应交互,或每次运行使用新会话有助于调试的开发场景中跳过恢复。

上下文窗口压缩

长对话会遇到两个限制:会话时长上限和模型的上下文窗口(因模型而异)。上下文窗口压缩可以同时解决这两个问题。当 token 数超过阈值时,它使用滑动窗口压缩较早的对话历史,同时保留最近的轮次为完整内容。启用后将移除会话时长限制。 权衡是:较早的上下文会变成摘要而非逐字历史记录。

from google.genai import types
from google.adk.agents.run_config import RunConfig

# 对于 128k 上下文的模型。
run_config = RunConfig(
    context_window_compression=types.ContextWindowCompressionConfig(
        trigger_tokens=100000,  # 在窗口约 78% 时开始压缩。
        sliding_window=types.SlidingWindow(
            target_tokens=80000,  # 压缩到约 62%,保留最近的轮次。
        ),
    )
)

trigger_tokens 设置为模型上下文窗口的约 70-80% 以留出余量,将 target_tokens 设置为 60-70% 以确保每次压缩能释放足够的空间容纳多个轮次。请根据你自己的对话模式进行测试。当会话必须运行超过平台限制或可能超出 token 限制时启用压缩;对于短会话或需要精确回忆早期轮次的场景,保持关闭。

并发会话

每个用户需要自己的 Live API session,且两个后端都对并发会话数设有上限。你的并发会话上限是同时在线用户的硬限制。关于当前上限和如何申请提升额度,请参阅支持的模型

针对上限进行设计:

  • 每用户一个会话是默认且正确的选择,只要峰值并发量在配额范围内。
  • 会话池(通过队列分配的固定会话集)可以在峰值并发超出配额时保持在限额内,代价是等待时间。释放时重置每个会话的状态,以避免对话在用户之间泄漏。

无论采用哪种方式,都需要自行统计活跃会话数,并在平台拒绝之前进行排队或拒绝新连接。配额拒绝表现为连接失败,这比可见的队列位置糟糕得多。