openai-agents-python 语音流水线结果流:StreamedAudioResult 与 VoiceStreamEvent 全解析

📅 发布时间:2026/9/12 7:03:02
openai-agents-python 语音流水线结果流:StreamedAudioResult 与 VoiceStreamEvent 全解析
openai-agents-python 语音流水线结果流StreamedAudioResult 与 VoiceStreamEvent 全解析【免费下载链接】openai-agents-pythonA lightweight, powerful framework for multi-agent workflows项目地址: https://gitcode.com/GitHub_Trending/op/openai-agents-pythonStreamedAudioResult是 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.pyasync 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(fUnsupported 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 aVoicePipeline. Streams events and audio data as theyre generated.构造签名如下StreamedAudioResult( tts_model: TTSModel, tts_settings: TTSModelSettings, voice_pipeline_config: VoicePipelineConfig, )三个构造参数分别对应参数类型说明tts_modelTTSModel负责把文本合成为 PCM 音频字节流的文本转语音模型tts_settingsTTSModelSettingsTTS 设置如音色、语速、缓冲大小、输出 dtype 等voice_pipeline_configVoicePipelineConfig流水线级配置如追踪开关、敏感数据策略等公开属性tts_model底层 TTS 模型实例。tts_settings本次合成使用的 TTS 设置。total_output_text整个会话累计生成的完整文本。注意它随每个文本片段的到达而累积见 result.py适合做最终存档或字幕生成。instructionsTTS 模型的指令文本取自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.pyasync def stream(self) - AsyncIterator[VoiceStreamEvent]: Stream the events and audio data as theyre generated.其行为要点对应 result.py循环从内部asyncio.Queue取事件并yield遇到VoiceStreamEventError时记录异常并终止遇到session_ended生命周期事件时标记会话结束并终止流终止后统一检查后台任务异常若有则重新抛出。三、事件模型三种 VoiceStreamEvent消费stream()得到的每个元素都是VoiceStreamEvent类型别名的一种定义于 src/agents/voice/events.pyVoiceStreamEvent: 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_audiodata是一个 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 追踪 spanturn_ended某一轮次的全部音频已按序派发完毕后发出result.py、result.pysession_ended整个会话结束的终止事件由调度器在观察到会话完成后发出result.pyVoiceStreamEventError错误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[bytes]产出 PCM 格式字节见 model.py。StreamedAudioResult._stream_audio()负责消费这些字节并缓冲result.py。buffer_size最小流式块大小TTSModelSettings.buffer_size默认值为 120model.py表示每积累至少 120 个字节块才向外派发一个音频事件。源码逻辑result.pyif 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(dataaudio_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 queueresult.py同时把这些队列按顺序注册进_ordered_tasks。即便某段任务在协程启动前被取消其 done 回调也会向队列放入None哨兵确保调度器永远能前进result.py。调度器保证跨段顺序_dispatch_audio()result.py是唯一的排序入口它从_ordered_tasks按注册顺序弹出队列逐个消费其中的事件并转发到消费者可见的主队列。由于多段文本的合成是并发的先注册的段落一定先被派发从而在整体上保持文本出现顺序 音频播放顺序。最后一段以finish_turnTrue收尾时调度器在派发完turn_ended后结束本轮当观察到_completed_session后派发session_ended终止事件并退出。七、轮次与会话的终止语义_turn_done()result.py流水线在每轮 workflow 输出结束后调用。若缓冲区仍有残留文本则以finish_turnTrue合成最后一小段若没有文本但轮次已开始直接派发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_dataTrue是否在追踪中记录敏感文本如 TTS 输入、指令仅作用于语音流水线本身不影响 workflow 内部trace_include_sensitive_audio_dataTrue是否在追踪中上传/记录音频数据tracing_disabledFalse是否完全关闭流水线追踪workflow_nameVoice Agent追踪中显示的 workflow 名称group_id随机生成用于把同一对话的多条 trace 关联成组trace_metadataNone附加到 trace 的自定义元数据字典十、完整实践从流水线到播放综合以上内容一个完整的消费端写法如下融合 docs/voice/pipeline.md 的示例与本文的事件语义import numpy as np from agents.voice import AudioInput, VoicePipeline, VoicePipelineConfig, TTSModelSettings config VoicePipelineConfig( tts_settingsTTSModelSettings( voicealloy, speed1.0, dtypenp.int16, buffer_size120, ), ) pipeline VoicePipeline(workflowmy_workflow, configconfig) 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(fvoice pipeline failed: {e}) print(result.total_output_text) # 会话完整文本几点补充单轮场景预录音频、按键对讲使用AudioInput需要活动检测自动判断用户说完的多轮场景使用StreamedAudioInput见 pipeline.py 与 docs/voice/pipeline.mdVoicePipelineConfig支持直接传 dict 配置由coerce_dataclass_config自动转换pipeline.pyTTS 内置音色包括alloy、ash、ballad、coral、echo、fable、onyx、nova、sage、shimmer、verse、marin、cedar也支持自定义音色 IDTTSCustomVoice见 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【免费下载链接】openai-agents-pythonA lightweight, powerful framework for multi-agent workflows项目地址: https://gitcode.com/GitHub_Trending/op/openai-agents-python创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考