数据截至 (上游 commit 689ca048bb0a)
实时语音管线:帧如何在处理器间流动
30 秒导读: 一通语音通话,本 质是把用户的声音变成文字、喂给 LLM、再把回答变成声音播回去。Dograh 用 pipecat 把这一串步骤搭成一条处理器流水线:每个环节是一个处理器,数据以「帧(Frame)」的形式从上一个处理器流到下一个。本章只讲语音这一半——管线怎么装、供应商怎么选、以及语音 agent 最难的那件事:判断用户什么时候说完了。
本章聚焦三件事:管线结构、服务选择、回合/打断机制。至于「整通通话怎么被编排起来」留给 04-call-orchestration,「工作流图怎么驱动 LLM」留给 03-pipecat-engine。想先看全景请回 index。
1. 这是什么(零基础也能懂)
一句话定义: 语音管线就是一条单向传送带——用户说的话从一头进来,机器人的回答从另一头出去,中间每个工位(处理器)做一件事。
它要解决的问题: 电话/网页里的实时对话有个硬约束——延迟要低、要能被打断。你不能等用户说完一整段、转成文字、算完答案、合成完语音再一次性播出去,那样机器人会「慢半拍」。所以管线里每个环节都是流式的:边听边转写、边生成边合成、用户一插话就得能停下来。
两种搭法: 同一条传送带,Dograh 提供两种装配方式。
| 形态 | 中间怎么处理 | 典型供应商 |
|---|---|---|
| 级联管线(cascaded) | STT → LLM → TTS 三个独立环节串起来 | Deepgram + OpenAI + ElevenLabs 等自由组合 |
| 语音到语音(realtime) | 一个 realtime 模型内部同时干了 STT+LLM+TTS | OpenAI Realtime、Gemini Live、Grok Voice 等 |
一句话直觉: 把管线想成一条工厂流水线。级联管线是「三台不同的机器排队加工」;realtime 管线是「一台一体机把三道工序都包了」——所以后者的传送带上少了两个工位,但布局要跟着改(见 §3.2)。
2. 顶层全景(一条帧的旅程)
怎么读下面这张图: 从上到下就是数据流向。左边一列是级联管线(声音→文字→文字→声音),右边是语音到语音管线(声音→声音)。带 ? 的框是可选处理器(按配置插入)。
级联管线 build_pipeline 语音到语音 build_realtime_pipeline
(pipeline_builder.py:28) (pipeline_builder.py:97)
transport.input() ← 用户音频进 transport.input()
│ │
STT 语音转文字 user_context_aggregator 先聚合用户回合
│ │
voicemail? 语音信箱检测(可选) realtime_llm 一体机:STT+LLM+TTS 全在里面
│ │
user_context_aggregator 攒够一个用户回合 voicemail? ← 注意:检测器放在 LLM 之后
│ │
llm_gate? 分类完再放行(可选) engine_callback_processor 引擎回调
│ │
LLM 生成回答文字 transport.output() ← 机 器人音频出
│ │
engine_callback_processor 引擎回调 audio_buffer 录音
│ │
recording_router? 播录音 or 走 TTS assistant_context_aggregator 聚合机器人回合
│ │
TTS 文字转语音 metrics 指标
│
transport.output() ← 机器人音频出
│
audio_buffer 录音(输入+输出合并)
│
assistant_context_aggregator 聚合机器人说了啥
│
metrics 用量/延迟指标
部件一句话职责:
| 处理器 | 干什么 | 在哪 |
|---|---|---|
transport.input/output | 收发用户音频(WebRTC / 电话) | transport_setup.py:create_webrtc_transport |
stt | 语音 → 文字(TranscriptionFrame) | service_factory.py:create_stt_service |
user_context_aggregator | 把碎片转写攒成一个完整「用户回合」 | pipecat LLMContextAggregatorPair.user() |
llm | 文字 → 回答文字 | service_factory.py:create_llm_service |
engine_callback_processor | 把帧事件回调给 PipecatEngine(见 03 章) | PipelineEngineCallbacksProcessor |
tts | 回答文字 → 语音 | service_factory.py:create_tts_service |
audio_buffer | 录制输入+输出的合并音频 | pipeline_builder.py:create_pipeline_components |
assistant_context_aggregator | 记录机器人实际说出的话 | pipecat LLMContextAggregatorPair.assistant() |
metrics | 汇总用量/延迟指标 | PipelineMetricsAggregator |
三类可选处理器(voicemail 检测、llm_gate、recording_router)属于「引擎旁路能力」,本章只标出它们在管线里的位置,能力细节留给 04-call-orchestration。
3. 核心原理之一:管线怎么装
这节讲处理器链是怎么被组装出来的——两个 build_* 函数、组件工厂、以及任务封装。
3.1 级联管线:build_pipeline
装配逻辑不是写死的一条链,而是按开关往一个列表里追加处理器,最后 Pipeline(processors) 把列表变成传送带。
先放固定的头两个,再按需插入可选项:
# 示意,非源码 —— 重点看「按开关 append」这个模式
processors = [transport.input(), stt]
if voicemail_detector: # 语音信箱检测紧跟在 STT 后
processors.append(voicemail_detector.detector())
processors.append(user_context_aggregator) # 用户回合聚合器必须在检测器之后
if voicemail_detector:
processors.append(voicemail_detector.llm_gate()) # 分类完再放行主 LLM
processors.extend([llm, engine_callback_processor, tts,
transport.output(), audio_buffer,
assistant_context_aggregator, metrics])
真实实现见 api/services/pipecat/pipeline_builder.py:28(build_pipeline)。几个次序上的讲究,代码注释里写明了原因:
user_context_aggregator必须在voicemail_detector之后(pipeline_builder.py:62-64):否则聚合器发出的LLMContextFrame会触发「语音信箱分类器」去跑 LLM 补全,污染主流程。recording_router插在引擎回调和 TTS 之间(pipeline_builder.py:71-72):它负责在「播放预录音频」和「走动态 TTS」之间做路由。audio_buffer放在transport.output()之后(pipeline_builder.py:88):这样它能同时录到输入和输出,合并成一条录音。
3.2 语音到语音管线:build_realtime_pipeline
realtime 服务(OpenAI Realtime、Gemini Live)内部就把 STT+LLM+TTS 全包了,所以传送带上没有独立的 STT 和 TTS 工位(pipeline_builder.py:107-110)。链子短了,但布局出现一处不对称,值得单独讲:
语音信箱检测器被放到了 realtime LLM 的下游,而级联管线里它在 STT 和用户聚合 器之间。为什么反过来?源码注释(pipeline_builder.py:113-124)给了完整解释,拆成三点:
- realtime LLM 既是
TranscriptionFrame的源头(向下游广播),又是LLMContextFrame的终点(它消费掉、不再向下转发)。 - 把检测器放在 LLM 下游,下游的
TranscriptionFrame才能到达分类器分支;而UserStartedSpeaking/StoppedSpeaking帧会被 LLM 透传下去。 - 主聚合器发出的
LLMContextFrame会被 realtime LLM 吸收,不会泄漏到分类器——否则分类器会拿主上下文去跑一次语音信箱补全。
另外,realtime 模式下不用 TTS gate 和 LLM gate:realtime LLM 直接对音频反应,而不是对 LLMContextFrame 反应;检测到语音信箱就直接用 end_call_with_reason 挂断(pipeline_builder.py:126-130)。
3.3 组件工厂与任务封装
两个 build_* 之外,还有两个辅助函数:
create_pipeline_components(pipeline_builder.py:13):造出跨两种形态共用的东西——AudioBufferProcessor(录音,采样率取自audio_config.pipeline_sample_rate)和LLMContext(对话上下文容器)。create_pipeline_task(pipeline_builder.py:155):把Pipeline包成一个可运行的PipelineWorker,并挂上PipelineParams(开启 metrics、usage metrics、heartbeats)和 tracing。它还会读环境变量ENABLE_TURN_LOGGING,若开启就注册on_turn_started回调,把回合号写进日志上下文(pipeline_builder.py:209-227)。
4. 核心原理之二:按用户配置选供应商
这节讲同一条管线,怎么塞进不同厂商的 STT/TTS/LLM。答案是一组 create_*_service 工厂函数——它们读 user_config,if/elif 分派到具体供应商,返回一个 pipecat service 实例。
4.1 四个工厂入口
run_pipeline.py:672-697 里,先判断是不是 realtime,再决定造哪些服务:
# 示意,非源码 —— 重点看 realtime 分支「stt/tts 为 None」
if is_realtime:
llm = create_realtime_llm_service(user_config, audio_config) # 一体机
stt = tts = None
inference_llm = create_llm_service(...) # 另配一个文字 LLM 做变量抽取等旁路推理
else:
stt = create_stt_service(user_config, audio_config, keyterms=...)
tts = create_tts_service(user_config, audio_config)
llm = create_llm_service(user_config)
注意 realtime 模式额外造了一个 inference_llm(run_pipeline.py:677-683):realtime 服务不实现 run_inference,所以变量抽取、语音信箱判定这类「不出声的推理」得靠一个单独的文字 LLM。
四个工厂各自的分派表:
| 工厂函数 | 位置 | 覆盖供应商(部分) |
|---|---|---|
create_stt_service | service_factory.py:251 | Deepgram、Deepgram Flux、OpenAI、Google、Cartesia、Sarvam、AssemblyAI、Gladia、Speechmatics、Azure、Smallest、Dograh |
create_tts_service | service_factory.py:554 | Deepgram、OpenAI、Google、ElevenLabs、Cartesia、Inworld、Rime、Sarvam、MiniMax、Azure、Smallest、Dograh |
create_llm_service | service_factory.py:1287 | OpenAI、Groq、OpenRouter、Google、Vertex、Azure、Bedrock、HuggingFace、MiniMax、Sarvam、Dograh |
create_realtime_llm_service | service_factory.py:1078 | OpenAI Realtime、Grok、Ultravox、Gemini Live、Vertex Realtime、Azure Realtime |
create_llm_service 是一层薄壳:它按 provider 从 user_config.llm 里挑出该供应商需要的 kwargs(base_url / endpoint / aws 密钥 / project_id 等),再转调 create_llm_service_from_provider(service_factory.py:936)做真正的分派。
4.2 Dograh 自带的托管栈
除了对接第三方,Dograh 还有自己托管的一套服务(MPS = Managed Provider Services),让用户不用自带 key 也能跑:
DograhSTTService/DograhFluxSTTService(service_factory.py:349-385):把MPS_API_URL的 http 改写成 ws/wss,连到 Dograh 的语音代理。语言若落在 Flux 多语种集合里,就走 Flux 变体(自带端点检测参数eot_timeout_ms等)。DograhTTSService(service_factory.py:688-703):同样改写成 WebSocket,走 Dograh 托管合成。DograhLLMService(service_factory.py:1022-1029):base_url 指向{MPS_API_URL}/api/v1/llm,底层复用 OpenAI 兼容协议(OpenAILLMSettings)。
一条贯穿全局的细节:几乎每个 service 都带 text_filters=[XMLFunctionTagFilter()](service_factory.py:567),防止 TTS 把函数调用标签念出来;并统一带 skip_aggregator_types=["recording_router","recording"] 和 silence_time_s=1.0。
4.3 一个供应商的「外部回合」标记
stt_uses_external_turns(service_factory.py:222)是连接服务选择和回合检测的关键钩子——它回答一个问题:这个 STT 自己会不会告诉我们「用户说完了」?
# 示意,非源码 —— 判断 STT 是否自带回合边界
def stt_uses_external_turns(user_config) -> bool:
if provider == DEEPGRAM: return model in DEEPGRAM_FLUX_MODELS # Flux 自带端点检测
if provider == DOGRAH: return 用的是 Flux 多语种
if provider == CARTESIA: return model == "ink-2"
return False
Deepgram Flux、Cartesia ink-2、Dograh Flux 这几个模型内建了端点检测(endpointing),会自己发出回合结束信号。这直接决定了下一节要用哪套回合策略。
5. 核心原理之三:回合检测(语音 agent 的核心难点)
它要解决的小问题: 打字聊天里,「用户说完了」有个明确信号——回车。语音里没有回车。机器人得自己猜:用户是真的说完了,还是只是句子中间喘了口气?猜早了会抢话,猜晚了会冷场。这就是回合检测(turn-taking)。
Dograh 把「一个用户回合」拆成两个独立问题,各由一组策略回答:
- 回合何时开始(start): 用户开始说话了吗?
- 回合何时结束(stop): 用户说完了吗?
组装点在 run_pipeline.py:892-938,产出一个 UserTurnStrategies(start=[...], stop=[...]),连同静音策略、超时、VAD 一起塞进 LLMUserAggregatorParams。
5.1 四种检测手段
先认识底层的四种「探测器」,后面的策略都是它们的组合:
| 手段 | 靠什么判断 | 类/参数 |
|---|---|---|
| VAD | 纯声学:有没有人声能量 | SileroVADAnalyzer(VADParams(stop_secs=0.2)) |
| transcription | STT 出没出新文字 | TranscriptionUserTurnStartStrategy / SpeechTimeoutUserTurnStopStrategy |
| smart-turn | 小模型判断「这句语义上说完没」 | LocalSmartTurnAnalyzerV3(SmartTurnParams) |
| external | STT 自己发的端点信号 | ExternalUserTurnStart/StopStrategy |
- VAD(Voice Activity Detection,语音活动检测):最底层,只回答「现在有没有人在出声」。
stop_secs=0.2意思是静音 0.2 秒就认为声学上停了。它快但「笨」——分不清「说完了」和「思考中的停顿」。 - smart-turn:用一个本地小模型(
LocalSmartTurnAnalyzerV3)判断这句话语义上是否完整,专治「longer responses with natural pauses」(带自然停顿的长回答),不会因为中途停顿就抢话。 - external:当 STT 本身内建端点检测(§4.3 的 Flux/ink-2),就直接信它的信号,本地不用再猜。
5.2 非 realtime:三种策略的取舍
run_pipeline.py:904 起按「STT 类型 + 工作流配置」(_create_non_realtime_user_turn_stop_strategies,run_pipeline.py:183)三选一。注意 start 侧几乎都用 VAD + Transcription 双保险,差异主要在 stop 侧:
| 场景(条件) | start 策略 | stop 策略 | 适合 |
|---|---|---|---|
external(stt_uses_external_turns 为真) | VAD + External(可打断) | ExternalUserTurnStopStrategy | STT 自带端点检测 |
| turn_analyzer(配置选它) | VAD + Transcription | TurnAnalyzerUserTurnStopStrategy(smart-turn) | 带自然停顿的长回答 |
| transcription(默认) | VAD + Transcription | SpeechTimeoutUserTurnStopStrategy | 短的 1-2 词回答 |
默认策略是 transcription(未显式配置时落到 SpeechTimeoutUserTurnStopStrategy,run_pipeline.py:191),偏向短应答;想要不抢话的长应答体验,工作流配置里改成 "turn_analyzer"。
5.3 realtime:把回合让给模型
realtime 服务往往自己就带服务端 VAD,本地再插一套会打架。_create_realtime_user_turn_config(run_pipeline.py:206)因此按供应商决定「本地 VAD 让到什么程度」:
| realtime 供应商 | 策略 | 本地 VAD | 原因(源码注释) |
|---|---|---|---|
| OpenAI Realtime / Azure Realtime | 纯 external | 无 | 供应商已发 speaking-state 和打断事件,聚合器跟着走 |
| Grok Realtime | 纯 external | 无 | Grok 服务端发 speech-start/stop 和打断信号 |
| Google Live / Vertex Realtime | 本地 VAD,不启用打断 | Silero | 让 Gemini 用服务端 VAD 管 barge-in,本地 VAD 只做「回合开始」和状态跟踪 |
| Ultravox | 本地 VAD,启用打断 | Silero | Ultravox 不发用户回合帧,靠本地 VAD 供生命周期信号 |
这是一处很典型的「按供应商能力做取舍」:能信服务端就完全让位(external),半信半疑就保留本地 VAD 但关掉它的打断权(Google),完全不发信号的就本地全权接管(Ultravox)。
5.4 停止超时:external 为什么给 30 秒
_resolve_user_turn_stop_timeout(run_pipeline.py:122)决定「等多久没动静就强制收尾回合」:
# 示意,非源码
if "user_turn_stop_timeout" in run_configs: return 用户配置的值
if uses_external_turns: return 30.0 # EXTERNAL_TURN_USER_STOP_TIMEOUT
return 5.0 # DEFAULT_USER_TURN_STOP_TIMEOUT
external 场景给到 30 秒(常量 EXTERNAL_TURN_USER_STOP_TIMEOUT,run_pipeline.py:119),而默认只有 5 秒。原因:external 模式下「回合结束」由 STT 的端点信号来定,本地超时只是兜底保险,所以放得很宽,免得误伤。
6. 核心原理之四:静音策略与打断
回合检测决定「什么时候听用户」;静音(mute)和打断(interrupt)决定「什么时候不听、以及用户能不能插话」。
6.1 用户静音:三条叠加规则
run_pipeline.py:885-889 把三条静音策略叠在一起,任一条命中就静音用户输入:
| 静音策略 | 什么时候静音用户 |
|---|---|
MuteUntilFirstBotCompleteUserMuteStrategy | 机器人还没说完第一句开场白之前 |
FunctionCallUserMuteStrategy | 正在执行函数调用期间 |
CallbackUserMuteStrategy | 交给引擎回调 engine.should_mute_user 动态决定 |
前两条是固定规则(开场白期间、函数调用期间别让用户插嘴打乱状态),第三条把决定权交回 PipecatEngine(见 03 章),让图状态机能按节点动态决定要不要闭麦。
6.2 打断(barge-in)
「打断」= 用户在机器人说话时插话,机器人立刻停嘴。在本章范围内,打断权是通过 start 策略的 enable_interruptions 参数表达的:
- external 场景显式
ExternalUserTurnStartStrategy(enable_interruptions=True)(run_pipeline.py:178)。 - realtime 的 Google 分支特意
enable_interruptions=False,把 barge-in 让给模型服务端(§5.3)。
另外注意 §4 里各 STT 服务大多传了 should_interrupt=False(如 service_factory.py:284),注释说明「让 UserAggregator 去发 InterruptionFrame」——即打断的决定权集中在聚合器,而不是散落在每个 STT 里。至于工作流节点级的 allow_interrupt(run_pipeline.py:781),那属于通话编排,归 04 章。
7. 传输层与音频配置
管线两端的 transport.input/output 抽象了「音频从哪来、到哪去」。 本章只覆盖非电话的 WebRTC 传输(电话传输在 services/telephony/providers/<name>/transport.py,归 04 章)。
7.1 WebRTC 传输
create_webrtc_transport(transport_setup.py:15)造一个 SmallWebRTCTransport,关键在 TransportParams 里对齐进出采样率,并挂上一个环境音混音器(build_audio_out_mixer,可播放背景白噪声让通话更自然):
# 示意,非源码 —— 重点看进出采样率对齐 + realtime 覆盖
SmallWebRTCTransport(params=TransportParams(
audio_in_enabled=True, audio_out_enabled=True,
audio_in_sample_rate=audio_config.transport_in_sample_rate,
audio_out_sample_rate=audio_config.transport_out_sample_rate,
audio_out_mixer=mixer,
**realtime_param_overrides(is_realtime), # realtime 下把 bot_vad_stop_secs 调到 0.5s
))
realtime_param_overrides(transport_params.py:17)是个小而关键的补丁:realtime LLM 不发 TTSStoppedFrame,「机器人说完了」只能靠「输出队列排空」兜底;默认 3 秒尾巴太长,realtime 下压到 0.5 秒(REALTIME_BOT_VAD_STOP_SECS),让对话不拖泥带水。
7.2 采样率:16kHz 天花板
AudioConfig(audio_config.py:14)是全管线采样率的唯一真相源,确保 VAD、音频缓冲、传输序列化器口径一致。最重要的一条约束:
管线内部采样率上限 16kHz,因为 VAD 只支持到这个档。 传输层负责在更高的外部速率(24kHz/48kHz)之间做重采样。
这条约束写死在 __post_init__ 里(audio_config.py:44-53):pipeline_sample_rate 若没指定就取 min(transport_out_sample_rate, 16000),超过 16kHz 会告警并强制封顶。VAD 采样率也校验只能是 8000 或 16000(audio_config.py:38)。create_audio_config(audio_config.py:73)则按传输类型挑速率:电话供应商从注册表拿它的线路采样率,WebRTC 固定 16kHz。
8. 边界与局限(本章范围内)
- 本章不讲通话生命周期。
run_pipeline.py里连接建立、DB 取配置、引擎初始化、post-call 处理都属于编排,归 04-call-orchestration。本章只摘了其中「服务/回合/静音装配」的片段。 - 本章不讲图如何驱动 LLM。
PipecatEngine、节点转移、should_mute_user的内部逻辑归 03-pipecat-engine;本章只在管线里标出它的挂载点。 - realtime 与部分能力互斥。 realtime 模式下关掉了 voicemail 检测(
run_pipeline.py:984-987)和 context compaction(run_pipeline.py:840-842),因为一体机自己在服务端管对话状态。 - 回合检测没有银弹。 三种非 realtime 策略是「短应答 vs 长应答」的取舍,没有一种全场景最优——这也是为什么它做成可配置(
turn_stop_strategy)。
9. 巧妙之处(可带走的设计)
- 「按开关 append 处理器」而非写死链条(
pipeline_builder.py:52-92):可选能力(voicemail、录音路由、gate)以「往列表里插」的方式表达,插入位置由注释解释清楚,新增能力不必重写整条链。 - 一个布尔钩子连通服务选择与回合检测(
stt_uses_external_turns,service_factory.py:222):STT 能力(自带端点检测)通过一个函数传递给回合策略装配,两个子系统解耦但对齐。 - realtime 布局的不对称是被逼出来的正确(
pipeline_builder.py:113-124):把 voicemail 检测器放到 LLM 下游,恰好利用「realtime LLM 既是转写源又是上下文汇」的双重身份,避免上下文泄漏——注释把这个非直觉决定讲透了。 - 按供应商能力分级让位(
run_pipeline.py:206-254):external / 本地 VAD 关打断 / 本地 VAD 全权,三档对应「完全信服务端 / 半信 / 不信」,是对接异构 realtime 供应商的干净模式。 - 采样率单一真相源 + 16kHz 天花板(
audio_config.py):把「VAD 只到 16kHz」这个物理约束集中到一个 dataclass 校验,避免各处理器各自为政。
10. 代码地图(导航索引)
| 主题 | 文件 | 符号 |
|---|---|---|
| 级联管线装配 | api/services/pipecat/pipeline_builder.py | build_pipeline |
| 语音到语音管线装配 | api/services/pipecat/pipeline_builder.py | build_realtime_pipeline |
| 共用组件(录音缓冲、上下文) | api/services/pipecat/pipeline_builder.py | create_pipeline_components |
| 任务封装 / tracing / metrics | api/services/pipecat/pipeline_builder.py | create_pipeline_task |
| STT 供应商分派 | api/services/pipecat/service_factory.py | create_stt_service |
| TTS 供应商分派 | api/services/pipecat/service_factory.py | create_tts_service |
| LLM 供应商分派 | api/services/pipecat/service_factory.py | create_llm_service / create_llm_service_from_provider |
| realtime 一体机分派 | api/services/pipecat/service_factory.py | create_realtime_llm_service |
| STT 是否自带回合边界 | api/services/pipecat/service_factory.py | stt_uses_external_turns |
| 非 realtime 回合/静音装配 | api/services/pipecat/run_pipeline.py | _run_pipeline_impl(724-780 行) |
| realtime 回合策略选择 | api/services/pipecat/run_pipeline.py | _create_realtime_user_turn_config |
| 回合停止超时 | api/services/pipecat/run_pipeline.py | _resolve_user_turn_stop_timeout |
| WebRTC 传输 | api/services/pipecat/transport_setup.py | create_webrtc_transport |
| realtime 传输参数覆盖 | api/services/pipecat/transport_params.py | realtime_param_overrides |
| 采样率配置 | api/services/pipecat/audio_config.py | AudioConfig / create_audio_config |