- 人工智能
- AI Agent
- 多模态
- 语音
- AI 应用
【免费下载链接】ten-framework
Open-source framework for conversational voice AI agents
导读
AsyncAvatarBaseExtension是 TEN Framework 为「数字人 / 虚拟形象(Avatar / Digital Human)」类扩展提供的一个异步基类,它将生命周期管理、音频队列与处理循环、采样率校验、错误上报、flush/finalize消息处理等繁琐逻辑全部封装起来,开发者只需要实现7 个业务方法即可接入任意 Avatar 服务商。本文以仓库中的 AVATAR_BASE_README.md 为主线,结合 avatar_base.py 源码与 Spatius 参考实现,讲清基类的设计原理、每个回调方法的职责与调用时机,并给出可直接复制运行的完整示例。
一、基类能为你做什么
基类的核心价值在于把「与具体 Avatar 服务商无关」的通用逻辑全部接管,让子类只关注「服务商特有的业务」:
- ✅ 自动生命周期管理(
on_init/on_start/on_stop/on_deinit) - ✅ 内置音频队列(
asyncio.Queue)与异步处理循环 - ✅ 采样率校验(不支持的采样率直接拒绝并上报错误)
- ✅ 统一的错误处理与日志输出(支持标准化的
ModuleError错误载荷) - ✅ 消息处理(
flush命令、finalize数据) - ✅ 可选的音频落盘(audio dumping),便于排查音频问题
从源码看,基类继承自AsyncExtension与ABC,通过@abstractmethod强制子类实现业务方法,同时把on_init、on_start、on_stop等生命周期钩子标记为「由基类管理,不要覆写」,见 avatar_base.py。子类唯一要做的事,就是补上 7 个抽象方法。
二、快速上手:必须实现的 7 个方法
创建自己的 Avatar 扩展只需继承AsyncAvatarBaseExtension并实现以下 7 个方法。
1.validate_config(ten_env) -> bool:加载并校验配置
- 调用时机:在
on_init()阶段由基类自动调用(await self.validate_config(ten_env))。 - 返回值:配置有效返回
True,否则返回False。 - 失败后果:返回
False时,基类会记录错误日志、禁用音频处理,并且不会调用connect_to_avatar(),扩展处于「可运行但不工作」的安全状态。
async def validate_config(self, ten_env: AsyncTenEnv) -> bool: self.config = await MyAvatarConfig.create_async(ten_env) if not self.config.api_key: ten_env.log_error("api_key is required") return False return True2.get_target_sample_rate() -> list[int]:声明支持的采样率
- 调用时机:每当有音频帧到达时,基类会用该方法的返回值校验帧采样率。
- 返回值:支持的服务商采样率列表(单位 Hz)。
- 拒绝策略:不在列表中的采样率会被拒绝并上报一次错误(
code=1001),且错误只发送一次以避免刷屏。
def get_target_sample_rate(self) -> list[int]: return [24000] # Spatius 支持 24kHz # return [16000] # Sensetime 支持 16kHz # return [24000, 48000] # 支持多个采样率注意:基类不做重采样,音频数据会原样(as-is)交给服务商,因此你必须保证上游音频帧的采样率与get_target_sample_rate()声明一致,或由上游(如 TTS 扩展)负责转换。
3.connect_to_avatar(ten_env) -> None:建立与 Avatar 服务的连接
- 调用时机:配置校验通过后,在
on_start()阶段由基类调用。 - 失败行为:如果此方法抛出异常,基类会记录错误、发送标准错误载荷,然后重新抛出异常—— 扩展将启动失败。
async def connect_to_avatar(self, ten_env: AsyncTenEnv) -> None: self.client = MyAvatarClient(self.config) await self.client.connect() ten_env.log_info("Connected to avatar service")4.disconnect_from_avatar(ten_env) -> None:断开连接并释放资源
- 调用时机:在
on_stop()阶段由基类自动调用。 - 失败行为:与连接相反,这里抛出的异常只记录日志、不向上传播,保证清理流程继续执行。
async def disconnect_from_avatar(self, ten_env: AsyncTenEnv) -> None: if self.client: await self.client.disconnect() ten_env.log_info("Disconnected from avatar service")5.send_audio_to_avatar(audio_data: bytes) -> None:发送音频
- 调用时机:由音频处理循环自动调用,一帧一调。
- 注意:音频为原始 PCM 字节流,不做重采样;如果服务商要求 base64 等编码,在本方法内自行转换。
async def send_audio_to_avatar(self, audio_data: bytes) -> None: # 示例:如果服务商要求 base64 编码 base64_audio = base64.b64encode(audio_data).decode("utf-8") await self.client.send_audio(base64_audio)6.send_eof_to_avatar() -> None:通知音频流结束
- 调用时机:当收到
finalize数据时,基类会把 EOF 哨兵(队列中的None项)排到待发送音频之后,由处理循环自动调用本方法,确保「先发完已有音频,再发 EOF」。
async def send_eof_to_avatar(self) -> None: await self.client.send_eof()7.interrupt_avatar() -> None:立即打断当前播报
- 调用时机:收到
flush命令时调用,用于立即停止 Avatar 当前正在进行的语音播报。
async def interrupt_avatar(self) -> None: if self.client: await self.client.interrupt()三、可选方法:get_dump_config()音频落盘调试
除 7 个必选方法外,还有一个可选方法用于调试:
def get_dump_config(self) -> tuple[bool, str]: """返回值:(是否落盘, 落盘目录)""" if self.config.dump: return (True, self.config.dump_path) return (False, "") # 默认不落盘- 默认值:
(False, ""),即不落盘。 - 落盘规则:音频被保存为
{dump_path}/{扩展名}_in.pcm(追加写模式),目录不存在时会自动创建,见 avatar_base.py。这对排查「音频没发出去 / 波形异常 / 采样率不对」等问题非常有用。
四、完整示例:一个可直接运行的 Avatar 扩展
以下代码综合了文档示例与基类约定,是可复制运行的完整骨架:
from ten_runtime import AsyncTenEnv from ten_ai_base.config import BaseConfig from avatar_base import AsyncAvatarBaseExtension from dataclasses import dataclass import base64 @dataclass class MyAvatarConfig(BaseConfig): api_key: str = "" avatar_id: str = "default" sample_rate: int = 24000 dump: bool = False dump_path: str = "" class MyAvatarExtension(AsyncAvatarBaseExtension): def __init__(self, name: str): super().__init__(name) self.config: MyAvatarConfig | None = None self.client = None # 1. 校验配置 async def validate_config(self, ten_env: AsyncTenEnv) -> bool: self.config = await MyAvatarConfig.create_async(ten_env) if not self.config.api_key: ten_env.log_error("[MyAvatar] api_key is required") return False ten_env.log_info(f"[MyAvatar] Config validated (avatar={self.config.avatar_id})") return True # 2. 目标采样率 def get_target_sample_rate(self) -> list[int]: return [self.config.sample_rate] # 3. 连接服务 async def connect_to_avatar(self, ten_env: AsyncTenEnv) -> None: ten_env.log_info("[MyAvatar] Connecting...") self.client = MyAvatarClient(self.config) await self.client.connect() ten_env.log_info("[MyAvatar] Connected") # 4. 断开连接 async def disconnect_from_avatar(self, ten_env: AsyncTenEnv) -> None: if self.client: await self.client.disconnect() ten_env.log_info("[MyAvatar] Disconnected") # 5. 发送音频 async def send_audio_to_avatar(self, audio_data: bytes) -> None: if self.client: base64_audio = base64.b64encode(audio_data).decode("utf-8") await self.client.send_audio(base64_audio) # 6. 发送 EOF async def send_eof_to_avatar(self) -> None: if self.client: await self.client.send_eof() # 7. 打断播报 async def interrupt_avatar(self) -> None: if self.client: await self.client.interrupt() # 可选:音频落盘 def get_dump_config(self) -> tuple[bool, str]: if self.config: return (self.config.dump, self.config.dump_path) return (False, "")五、自动生命周期:你不需要覆写任何生命周期钩子
基类把完整生命周期封装成一条固定流水线:
1. on_init() └─> validate_config() 2. on_start() └─> connect_to_avatar() └─> 启动音频处理循环(asyncio.create_task) 3. 音频处理(自动) └─> on_audio_frame() 接收音频 └─> 校验采样率 └─> 入队(unbounded queue) └─> 处理循环调用 send_audio_to_avatar() 4. 消息处理(自动) └─> flush 命令 → interrupt_avatar() └─> finalize 数据 → send_eof_to_avatar() 5. on_stop() └─> 取消音频处理任务 └─> disconnect_from_avatar()你不需要覆写on_init()、on_start()或on_stop()!从 avatar_base.py 可以看到,基类在这些钩子内部完成了配置校验(on_init)、连接与任务启动(on_start)、任务取消与断开(on_stop),并且on_stop中还会用asyncio.CancelledError妥善收尾音频任务。
六、音频处理细节
采样率校验
- 每帧音频先取
source_rate = audio_frame.get_sample_rate(),再与get_target_sample_rate()返回的列表比对,见 avatar_base.py。 - 不支持的采样率会被拒绝,并通过
_send_error上报code=1001的错误数据。 _sample_rate_error_sent标志保证同一轮请求只报一次错,避免日志刷屏;该标志在_clear_request_context()中重置。
音频队列
- 音频帧被包装为
QueuedAudioFrame(audio=bytes)放入asyncio.Queue(无界队列),见 avatar_base.py。 - 处理循环
_process_audio_loop逐个取出并调用send_audio_to_avatar();循环内对CancelledError单独处理,其他异常记录日志并通过_send_error上报后继续处理下一帧。 - 收到
flush命令时队列会被清空(_clear_audio_queue会统计并记录清除了多少帧)。
音频落盘
- 通过
get_dump_config()返回(True, "/path/to/dump")开启。 - 音频写入
{dump_path}/{扩展名}_in.pcm(如spatius_avatar_python_in.pcm),适用于排查「上游是否真的发来了音频」「PCM 内容是否正确」等问题。
七、错误处理策略
基类对四类错误采用分层策略,这是它最值得借鉴的设计之一:
| 错误场景 | 处理方式 | 结果 |
|---|---|---|
配置校验失败(validate_config返回False) | 记录错误日志,禁用音频处理 | connect_to_avatar()不会被调用 |
连接失败(connect_to_avatar抛异常) | 记录错误 + 发送错误载荷 | 异常向上传播,扩展启动失败 |
断开失败(disconnect_from_avatar抛异常) | 只记录错误日志 | 异常不传播,清理流程继续 |
音频发送失败(send_audio_to_avatar抛异常) | 记录错误 + 发送错误载荷 | 处理循环继续处理下一帧 |
其中,连接与音频发送失败时,基类会调用_send_error构造一个标准的ModuleError载荷(module="avatar"、携带vendor/vendor_code/vendor_message等字段),通过名为error的Data消息发送出去,见 avatar_base.py。这套标准化错误协议在 tests/test_basic.py 中有对应的测试用例验证:测试断言错误载荷的module == "avatar"、vendor == "spatius"、且vendor_metadata中不包含空值。
八、消息处理:flush 与 finalize
flush 命令
当收到flush命令(CMD_IN_FLUSH)时,基类依次执行:
- 清空音频队列(丢弃未发送的积压音频);
- 调用
interrupt_avatar()打断当前播报; - 将
flush命令转发给下游(ten_env.send_cmd(Cmd.create(CMD_OUT_FLUSH))); - 返回
StatusCode.OK的CmdResult。
finalize 数据
当收到名为finalize的数据时,基类将 EOF 哨兵(None)放入音频队列尾部,排在所有待发送音频之后,由处理循环顺序消费到哨兵时调用send_eof_to_avatar()。这样既保证了「TTS 播报音频已全部发送完毕」的语义,又不会打断正在发送的音频流。
九、参考实现:Spatius Avatar 扩展
仓库中的spatius_avatar_python包是一个完整的参考实现,它通过 Spatius SDK 驱动真实数字人服务,并演示了基类的全部用法:
- avatar_base.py —— 基类实现(本文主体)。
- extension.py —— Spatius 参考实现,实现 7 个必选方法与
get_dump_config()等可选方法。 - addon.py —— 通过
@register_addon_as_extension("spatius_avatar_python")注册扩展。 - manifest.json —— 声明扩展的 API 契约(
audio_frame_in、cmd_in、data_in、data_out及全部配置属性)。 - property.json —— 默认配置,支持
${env:VAR|}环境变量注入(如SPATIUS_API_KEY、AGORA_APP_ID)。 - tests/test_basic.py —— 单元测试,覆盖基础命令往返与配置错误的标准载荷。
Spatius 实现中的几个亮点
1. 配置归一化与校验。SpatiusConfig通过update_params()把用户可见的params字典复制到规范化字段,再用validate_params()检查必填项、采样率范围(Ogg Opus 仅支持8000/12000/16000/24000/48000Hz)与 Agora token 二选一(agora_token或agora_appcert至少提供一个)。
2. Token 自动生成。若只配置了agora_appcert,resolve_agora_token()会用agora-token-builder的RtcTokenBuilder.buildTokenWithUid结合session_expire_minutes(默认 30 分钟)自动生成 RTC Token。
3. 敏感信息脱敏。日志与vendor_metadata中的 API Key、App Cert 等均通过encrypt()脱敏输出;get_vendor_metadata()还会剔除空值字段,避免把空串上报出去。
4. 连接流程。connect_to_avatar()用new_avatar_session(...)创建会话,随后await session.init()获取鉴权 token、await session.start()建立 WebSocket 连接并拿到connection_id;send_audio_to_avatar通过session.send_audio(bytes, end=False)发送,send_eof_to_avatar则用session.send_audio(b"", end=True)表示流结束。
配置参考(property.json)
{ "dump": false, "dump_path": "", "channel": "", "agora_uid": "", "agora_token": "", "agora_appid": "", "agora_appcert": "", "agora_channel": "", "params": { "spatius_api_key": "${env:SPATIUS_API_KEY|}", "spatius_app_id": "${env:SPATIUS_APP_ID|}", "spatius_avatar_id": "", "agora_uid": "", "agora_token": "", "agora_appid": "${env:AGORA_APP_ID|}", "agora_appcert": "${env:AGORA_APP_CERTIFICATE|}", "agora_channel": "", "region": "", "sample_rate": 24000, "session_expire_minutes": 30, "audio_format": "ogg_opus" } }其中agora_token与agora_appcert至少配置一个;sample_rate默认24000(与get_target_sample_rate()返回[self.config.sample_rate]保持一致);audio_format默认ogg_opus。
十、落地建议与总结
实现自己的 Avatar 扩展,遵循以下清单即可:
- ✅ 创建继承自
BaseConfig的配置类(dataclass),并在validate_config()中加载与校验; - ✅ 继承
AsyncAvatarBaseExtension; - ✅ 实现 7 个必选方法(外加可选的
get_dump_config()); - ✅ 使用统一的日志前缀(基类
LOG_PREFIX默认[Spatius],子类可按需覆盖),保证日志可检索; - ✅ 用不同采样率的音频测试
get_target_sample_rate()的校验逻辑; - ✅ 妥善处理错误:连接失败要让扩展启动失败,发送失败要保证循环继续。
其余一切 —— 生命周期、队列、循环、采样率校验、flush/finalize、错误上报 —— 都由AsyncAvatarBaseExtension自动完成。这种「基类兜底通用逻辑、子类只写业务差异」的模板方法设计,让接入新的数字人服务商从「理解整个框架消息流」简化为「实现 7 个方法」,是 TEN Framework 在语音 Agent 扩展开发上的一个值得复用的范式。
- 人工智能
- AI Agent
- 多模态
- 语音
- AI 应用
【免费下载链接】ten-framework
Open-source framework for conversational voice AI agents
相关推荐
TEN 框架数字人扩展开发指南:基于 AsyncAvatarBaseExtension 实现自定义虚拟形象扩展
TEN 框架数字人扩展开发指南:基于 AsyncAvatarBaseExtension 实现自定义虚拟形象扩展 导读 本文以 TEN 框架中 spatius_a
人工智能AI Agent多模态语音AI 应用昇腾C BatchNorm Tiling API
BatchNorm Tiling 功能说明 BatchNorm Tiling API用于获取BatchNorm kernel计算时所需的Tiling参数。获取T
人工智能深度学习算子库CANNAscend数字人Live2D终极指南:快速打造你的专属虚拟形象
数字人Live2D终极指南:快速打造你的专属虚拟形象 想要拥有一个会说话、会互动的数字人伙伴吗?《Awesome Digital Human Live2D》项目
人工智能AI 应用数字人语音AI Agent交互助手
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考