openai-agents-python 语音流水线结果流:StreamedAudioResult 与 VoiceStreamEvent 全解析
StreamedAudioResult 是 openai-agents-python 中 VoicePipeline 的最终产出对象,它以异步流的方式持续输出语音合成产生的音频块、轮次与会话生命周期事件以及错误信息。本文以 docs/ref/voice/result.md 所对应的 API 参考为主体,结合 src/agents/voice/result.py 的完整实现与 docs/voice/pipeline.md 的官方说明,深入讲解该结果对象的属性、事件模型、消费模式、底层音频缓冲与按序调度机制,以及错误传播与追踪集成语义,帮助你构建可实际运行的语音 Agent 应用。
一、结果对象在语音流水线中的位置
在 openai-agents-python 的语音体系中,VoicePipeline 是一个三步流水线:先将音频输入转写为文本(STT),再运行你提供的 workflow 生成文本回复序列,最后把文本转回流式音频输出(TTS)。该流程定义在 src/agents/voice/pipeline.py 的类 docstring 中。
StreamedAudioResult 正是这第三步的产出。调用 VoicePipeline.run(audio_input) 后会立即返回该对象(而非等待全部音频生成完毕),随后你可以通过 await result.stream() 以异步迭代的方式消费事件与音频数据。见 pipeline.py:
async def run(self, audio_input: AudioInput | StreamedAudioInput) -> StreamedAudioResult:
if isinstance(audio_input, AudioInput):
return await self._run_single_turn(audio_input)
elif isinstance(audio_input, StreamedAudioInput):
return await self._run_multi_turn(audio_input)
else:
raise UserError(f"Unsupported audio input type: {type(audio_input)}")
从源码看,StreamedAudioResult 实例由 pipeline.py 与 pipeline.py 创建,构造时注入 TTS 模型、TTS 设置与流水线配置,随后通过 output._set_task(asyncio.create_task(...)) 启动后台生产任务——生产与消费是解耦的:生产者把事件推入内部队列,消费者在 stream() 中逐个取出。
二、StreamedAudioResult 的公开接口
类定义位于 src/agents/voice/result.py,其 docstring 明确说明:"The output of a VoicePipeline. Streams events and audio data as they're generated." 构造签名如下:
StreamedAudioResult(
tts_model: TTSModel,
tts_settings: TTSModelSettings,
voice_pipeline_config: VoicePipelineConfig,
)
三个构造参数分别对应:
| 参数 | 类型 | 说明 |
|---|---|---|
tts_model |
TTSModel |
负责把文本合成为 PCM 音频字节流的文本转语音模型 |
tts_settings |
TTSModelSettings |
TTS 设置,如音色、语速、缓冲大小、输出 dtype 等 |
voice_pipeline_config |
VoicePipelineConfig |
流水线级配置,如追踪开关、敏感数据策略等 |
公开属性
tts_model:底层 TTS 模型实例。tts_settings:本次合成使用的 TTS 设置。total_output_text:整个会话累计生成的完整文本。注意它随每个文本片段的到达而累积(见 result.py),适合做最终存档或字幕生成。instructions:TTS 模型的指令文本,取自tts_settings.instructions,用于控制语音输出语气。text_generation_task:后台文本生成任务的asyncio.Task引用,由流水线通过_set_task()注入,供stream()在终止时等待生产者收尾。
其余字段(_queue、_tasks、_ordered_tasks、_dispatcher_task、_text_buffer、_turn_text_buffer 等)均为私有实现细节,驱动内部的事件队列、音频按序调度与缓冲逻辑。
核心方法:stream()
stream() 是唯一面向消费者的公开异步方法,签名与语义见 result.py:
async def stream(self) -> AsyncIterator[VoiceStreamEvent]:
"""Stream the events and audio data as they're generated."""
其行为要点(对应 result.py):
- 循环从内部
asyncio.Queue取事件并yield; - 遇到
VoiceStreamEventError时记录异常并终止; - 遇到
session_ended生命周期事件时标记会话结束并终止; - 流终止后统一检查后台任务异常,若有则重新抛出。
三、事件模型:三种 VoiceStreamEvent
消费 stream() 得到的每个元素都是 VoiceStreamEvent 类型别名的一种,定义于 src/agents/voice/events.py:
VoiceStreamEvent: TypeAlias = (
VoiceStreamEventAudio | VoiceStreamEventLifecycle | VoiceStreamEventError
)
VoiceStreamEventAudio(音频块)
events.py 定义:
@dataclass
class VoiceStreamEventAudio:
data: npt.NDArray[np.int16 | np.float32] | None
type: Literal["voice_stream_event_audio"] = "voice_stream_event_audio"
data 是一个 NumPy 数组,承载一段 PCM 音频。其 dtype 由 tts_settings.dtype 决定(默认 np.int16),可通过 transform_data 回调在产出前重整形(例如转为 float32 便于某些播放器或深度学习模型直接消费)。
VoiceStreamEventLifecycle(生命周期)
events.py 定义:
@dataclass
class VoiceStreamEventLifecycle:
event: Literal["turn_started", "turn_ended", "session_ended"]
type: Literal["voice_stream_event_lifecycle"] = "voice_stream_event_lifecycle"
三种事件语义:
| 事件 | 触发时机(对应源码) |
|---|---|
turn_started |
首个文本片段开始处理时,由 _start_turn() 发出(result.py),同时启动一个 speech group 追踪 span |
turn_ended |
某一轮次的全部音频已按序派发完毕后发出(result.py、result.py) |
session_ended |
整个会话结束的终止事件,由调度器在观察到会话完成后发出(result.py) |
VoiceStreamEventError(错误)
events.py 定义:
@dataclass
class VoiceStreamEventError:
error: Exception
type: Literal["voice_stream_event_error"] = "voice_stream_event_error"
错误事件携带原始 Exception 对象,消费端既可以在事件分支中处理,也可以依赖 stream() 在终止时重新抛出(见下文错误传播小节)。
四、消费模式:官方推荐写法
docs/voice/pipeline.md 给出了标准的消费循环,完整继承如下:
result = await pipeline.run(input)
async for event in result.stream():
if event.type == "voice_stream_event_audio":
# play audio
pass
elif event.type == "voice_stream_event_lifecycle":
# lifecycle
pass
elif event.type == "voice_stream_event_error":
# error
pass
三个分支分别处理音频块、生命周期事件与错误。结合源码可以进一步说明:
- 音频分支:
event.data是np.int16或np.float32数组。若你的播放库要求bytes,可按 dtype 自行转换:int16直接用data.tobytes(),float32通常需先还原为int16(乘以 32767 后转回)。若希望结果直接是某个特定形状,可在TTSModelSettings.transform_data中完成转换,见 model.py。 - 生命周期分支:可借此实现打断(interruption)处理——
turn_started表示新一轮开始,turn_ended表示该轮音频全部派发完毕。官方建议(docs/voice/pipeline.md):在模型开始输出时静音麦克风,在播放完该轮全部音频后再恢复收音。 - 错误分支:
event.error即底层异常;注意stream()迭代终止时还会把该异常重新抛出,因此更稳妥的做法是用try/except包裹整个async for循环,以VoiceStreamEventError分支做精细处理、外层except兜底。
五、音频数据是如何被组装与转换的
PCM 字节流到 NumPy 数组
TTS 模型的 run(text, settings) 返回 AsyncIterator<a href="https://link.gitcode.com/i/541e2133e4d22d0e4defc37891167530" target="_blank">bytes],产出 PCM 格式字节(见 [model.py)。StreamedAudioResult._stream_audio() 负责消费这些字节并缓冲(result.py)。
buffer_size:最小流式块大小
TTSModelSettings.buffer_size 默认值为 120(model.py),表示每积累至少 120 个字节块才向外派发一个音频事件。源码逻辑(result.py):
if len(buffer) >= self._buffer_size:
combined = pending_byte + b"".join(buffer)
if len(combined) % 2 != 0:
pending_byte = combined[-1:]
combined = combined[:-1]
else:
pending_byte = b""
if combined:
audio_np = self._transform_audio_buffer([combined], self.tts_settings.dtype)
if self.tts_settings.transform_data is not None:
audio_np = self.tts_settings.transform_data(audio_np)
await local_queue.put(VoiceStreamEventAudio(data=audio_np))
buffer = []
注意其中的奇数字节对齐处理:由于 np.int16 需要 2 字节对齐,若累计字节数为奇数,会把最后一个字节暂存为 pending_byte 留到下一轮拼接,避免产生错位的采样点。
dtype 转换规则
_transform_audio_buffer()(result.py)把字节数组解析为 np.int16 数组,再按目标 dtype 转换:
np.int16:原样返回;np.float32:先转为float32,再除以 32767.0,并reshape(-1, 1)成单声道列向量;- 其他 dtype:抛出
UserError("Invalid output dtype")。
源码通过 np.dtype(output_dtype) 解析配置值(result.py),因此 dtype 既可以是 np.int16 / np.float32 这类对象,也可以是 "int16" / "float32" 这类字符串;解析失败(TypeError/ValueError)时统一转换为 SDK 自有 UserError,并保留 NumPy 原始异常作为 cause,方便排查拼写问题。
transform_data:产出自定义形状
若设置了 transform_data 回调,每个派发块在入队前都会经过该函数(result.py),因此消费端拿到的 data 已经是你需要的形状与类型。
六、文本切分与按序调度:保证"先说的先播"
为了让 TTS 不必等待整段文本生成完毕,StreamedAudioResult 采用"按句切分 + 分段合成 + 顺序派发"的设计。
文本缓冲与切分
_add_text(text)(result.py)把流水线传入的文本追加到 _text_buffer 和 total_output_text,然后调用 tts_settings.text_splitter 把已积累的文本切分为"可合成的完整句子"与"残留缓冲区"两部分:
combined_sentences, self._text_buffer = self.tts_settings.text_splitter(self._text_buffer)
if combined_sentences:
local_queue = asyncio.Queue()
self._enqueue_audio_segment(local_queue)
self._create_audio_task(combined_sentences, local_queue)
默认切分器是 get_sentence_based_splitter()(model.py),按句子边界切分;你可以传入自定义 Callable[[str], tuple[str, str]] 替换(如按标点、按固定字数切分)。
每段文本一个独立任务
每产生一组完整句子,就创建一个独立 asyncio.Task(_stream_audio)把该段的音频合成结果推入该段专属的 local queue(result.py),同时把这些队列按顺序注册进 _ordered_tasks。即便某段任务在协程启动前被取消,其 done 回调也会向队列放入 None 哨兵,确保调度器永远能前进(result.py)。
调度器:保证跨段顺序
_dispatch_audio()(result.py)是唯一的排序入口:它从 _ordered_tasks 按注册顺序弹出队列,逐个消费其中的事件并转发到消费者可见的主队列。由于多段文本的合成是并发的,先注册的段落一定先被派发,从而在整体上保持"文本出现顺序 = 音频播放顺序"。最后一段以 finish_turn=True 收尾时,调度器在派发完 turn_ended 后结束本轮;当观察到 _completed_session 后,派发 session_ended 终止事件并退出。
七、轮次与会话的终止语义
_turn_done()(result.py):流水线在每轮 workflow 输出结束后调用。若缓冲区仍有残留文本,则以finish_turn=True合成最后一小段;若没有文本但轮次已开始,直接派发turn_ended。随后等待所有音频任务完成。_done()(result.py):整个会话结束时调用,标记_completed_session并唤醒调度器。源码特别处理了一种边界情况:如果会话从未产生任何音频,调度器从未启动,那么stream()将永远等不到session_ended,因此_done()会确保调度器任务被创建,让它观察到会话已完成并发出终止事件。_cancel()(result.py):取消合成同时保证终止事件按序送达——取消所有未完成的音频任务、等待调度器发布session_ended,最后关闭追踪 span。
八、错误传播:何时抛出、抛什么
官方文档(docs/voice/pipeline.md)明确指出:终端流水线错误在消费 StreamedAudioResult.stream() 时抛出,而非在 run() 时。实现细节如下:
- 后台任务出错时,先通过
_add_error()把VoiceStreamEventError放入队列(result.py),消费者遇到它即终止迭代; stream()在finally块中做三层收尾:先等待text_generation_task优雅结束(通过asyncio.shield包裹,避免取消信号错乱),再清理全部后台任务,最后按优先级选择要抛出的异常(result.py);- 异常优先级规则从源码可以归纳为:调用方取消(
CancelledError)优先于一切;否则保留消费者主异常;若消费者无异常,则优先传播生产者(TTS/转写)异常;再依次是消费收尾异常与清理异常; - 一个值得注意的语义:如果一轮对话本身已经失败,且随后关闭转写会话也失败,
stream()会保留原始的轮次错误作为主错误,而不会用会话关闭错误覆盖它(docs/voice/pipeline.md); - 被抑制的收尾异常会以
logger.warning("Voice stream finalization failed while preserving another exception")记录,但不会替换已选定的异常(result.py)。
九、追踪集成:语音产出的可观测性
每次音频合成都在一个 speech_span 内进行(result.py),而整个轮次则包在一个 speech_group_span 中(由 _start_turn() 启动、_finish_turn() 结束,见 result.py 与 result.py)。span 内记录:
- 模型名
model(来自tts_model.model_name); - 输入文本与
voice、instructions、speed等模型配置; - 首个音频字节到达时间
first_content_at; - 输出音频(PCM 经 base64 编码)——但仅当
VoicePipelineConfig.trace_include_sensitive_audio_data为True时才会保留整段音频数据(result.py、result.py),否则 span 输出为空字符串,避免无谓的内存占用与敏感数据暴露。
与之配套的配置项均来自 src/agents/voice/pipeline_config.py:
| 配置项 | 默认值 | 说明 |
|---|---|---|
trace_include_sensitive_data |
True |
是否在追踪中记录敏感文本(如 TTS 输入、指令),仅作用于语音流水线本身,不影响 workflow 内部 |
trace_include_sensitive_audio_data |
True |
是否在追踪中上传/记录音频数据 |
tracing_disabled |
False |
是否完全关闭流水线追踪 |
workflow_name |
"Voice Agent" |
追踪中显示的 workflow 名称 |
group_id |
随机生成 | 用于把同一对话的多条 trace 关联成组 |
trace_metadata |
None |
附加到 trace 的自定义元数据字典 |
十、完整实践:从流水线到播放
综合以上内容,一个完整的消费端写法如下(融合 docs/voice/pipeline.md 的示例与本文的事件语义):
import numpy as np
from agents.voice import AudioInput, VoicePipeline, VoicePipelineConfig, TTSModelSettings
config = VoicePipelineConfig(
tts_settings=TTSModelSettings(
voice="alloy",
speed=1.0,
dtype=np.int16,
buffer_size=120,
),
)
pipeline = VoicePipeline(workflow=my_workflow, config=config)
result = await pipeline.run(AudioInput(my_audio_bytes))
try:
async for event in result.stream():
if event.type == "voice_stream_event_audio":
audio_bytes = event.data.tobytes() # int16 PCM
await player.write(audio_bytes)
elif event.type == "voice_stream_event_lifecycle":
if event.event == "turn_started":
mic.mute()
elif event.event == "turn_ended":
mic.unmute()
elif event.event == "session_ended":
break
except Exception as e:
print(f"voice pipeline failed: {e}")
print(result.total_output_text) # 会话完整文本
几点补充:
- 单轮场景(预录音频、按键对讲)使用
AudioInput;需要活动检测(自动判断用户说完)的多轮场景使用StreamedAudioInput(见 pipeline.py 与 docs/voice/pipeline.md); VoicePipelineConfig支持直接传 dict 配置,由coerce_dataclass_config自动转换(pipeline.py);- TTS 内置音色包括
alloy、ash、ballad、coral、echo、fable、onyx、nova、sage、shimmer、verse、marin、cedar,也支持自定义音色 ID(TTSCustomVoice,见 model.py)。
十一、最佳实践与边界提醒
- 打断处理:SDK 不内置打断逻辑,每个检测到的轮次都会触发一次独立的 workflow 运行。请基于
turn_started/turn_ended自行实现麦克风静音与恢复(docs/voice/pipeline.md)。 - 偶数对齐:TTS 字节流可能出现奇数字节,
StreamedAudioResult内部已做对齐,但如果你自定义text_splitter或直接消费底层 TTS 字节流,需要注意 PCM 的 2 字节对齐要求。 - 错误必达:无论成功或失败,
stream()都会以session_ended(或错误)终结,不会无限挂起;消费端应始终用async for完整迭代,让finally清理逻辑(等待生产者、取消任务、关闭追踪 span)得以执行。 - 敏感数据策略:生产环境若涉及隐私音频,建议将
trace_include_sensitive_data与trace_include_sensitive_audio_data设为False,追踪记录中音频与文本将被置空。
延伸阅读
- API 参考入口:docs/ref/voice/result.md、docs/ref/voice/events.md、docs/ref/voice/pipeline.md
- 流水线与 workflow 官方指南:docs/voice/pipeline.md
- 核心实现:src/agents/voice/result.py、src/agents/voice/events.py、src/agents/voice/pipeline.py、src/agents/voice/pipeline_config.py、src/agents/voice/model.py
- 配套配置解析:src/agents/voice/pipeline_config.py
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust4.21 K638- DDeepSeek-V4.1-FlashDeepSeek-V4.1-Flash 是一个多模态混合专家(MoE)模型,拥有 5520 亿骨干参数,并支持最多一百万 token 的上下文长度。该模型原生支持图像和文本输入,并以自回归方式生成文本Python250
jforgamejforgame是一个一站式游戏服务器开发框架。包含游戏服务器开发所需要的各种组件,比如网关,socket服务端与客户端,自定义高效消息编解码,游戏热更新,游戏通用工具等等。包含游戏服,跨服,匹配服,后台管理系统等实现,同时提供大量业务案例以供学习。亦可用于其他socket应用,例如及时聊天等。Java301
fizz-gateway-nodeAn Aggregation API Gateway in Java . FizzGate 是一个基于 Java开发的微服务聚合网关,是拥有自主知识产权的应用网关国产化替代方案,能够实现热服务编排聚合、自动授权选择、线上服务脚本编码、在线测试、高性能路由、API审核管理、回调管理等目的,拥有强大的自定义插件系统可以自行扩展,并且提供友好的图形化配置界面,能够快速帮助企业进行API服务治理、减少中间层胶水代码以及降低编码投入、提高 API 服务的稳定性和安全性。Java210
certd开源SSL证书管理工具;全自动证书申请、更新、续期;通配符证书,泛域名证书申请;证书自动化部署到阿里云、腾讯云、主机、群晖、宝塔;https证书,pfx证书,der证书,TLS证书,nginx证书自动续签自动部署JavaScript190
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python300