跳到主要内容

实时语音管线:帧如何在处理器间流动

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+TTSOpenAI 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)给了完整解释,拆成三点:

  1. realtime LLM 既是 TranscriptionFrame 的源头(向下游广播),又是 LLMContextFrame 的终点(它消费掉、不再向下转发)。
  2. 把检测器放在 LLM 下游,下游的 TranscriptionFrame 才能到达分类器分支;而 UserStartedSpeaking/StoppedSpeaking 帧会被 LLM 透传下去。
  3. 主聚合器发出的 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:519-544 里,先判断是不是 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:524-530):realtime 服务不实现 run_inference,所以变量抽取、语音信箱判定这类「不出声的推理」得靠一个单独的文字 LLM。

四个工厂各自的分派表:

工厂函数位置覆盖供应商(部分)
create_stt_serviceservice_factory.py:138Deepgram、Deepgram Flux、OpenAI、Google、Cartesia、Sarvam、AssemblyAI、Gladia、Speechmatics、Azure、Smallest、Dograh
create_tts_serviceservice_factory.py:422Deepgram、OpenAI、Google、ElevenLabs、Cartesia、Inworld、Rime、Sarvam、MiniMax、Azure、Smallest、Dograh
create_llm_serviceservice_factory.py:1066OpenAI、Groq、OpenRouter、Google、Vertex、Azure、Bedrock、HuggingFace、MiniMax、Sarvam、Dograh
create_realtime_llm_serviceservice_factory.py:892OpenAI 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:761)做真正的分派。

4.2 Dograh 自带的托管栈

除了对接第三方,Dograh 还有自己托管的一套服务(MPS = Managed Provider Services),让用户不用自带 key 也能跑:

  • DograhSTTService / DograhFluxSTTService(service_factory.py:236-272):把 MPS_API_URL 的 http 改写成 ws/wss,连到 Dograh 的语音代理。语言若落在 Flux 多语种集合里,就走 Flux 变体(自带端点检测参数 eot_timeout_ms 等)。
  • DograhTTSService(service_factory.py:558-573):同样改写成 WebSocket,走 Dograh 托管合成。
  • DograhLLMService(service_factory.py:838-844):base_url 指向 {MPS_API_URL}/api/v1/llm,底层复用 OpenAI 兼容协议(OpenAILLMSettings)。

一条贯穿全局的细节:几乎每个 service 都带 text_filters=[XMLFunctionTagFilter()](service_factory.py:435),防止 TTS 把函数调用标签念出来;并统一带 skip_aggregator_types=["recording_router","recording"]silence_time_s=1.0

4.3 一个供应商的「外部回合」标记

stt_uses_external_turns(service_factory.py:110)是连接服务选择回合检测的关键钩子——它回答一个问题:这个 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:724-780,产出一个 UserTurnStrategies(start=[...], stop=[...]),连同静音策略、超时、VAD 一起塞进 LLMUserAggregatorParams

5.1 四种检测手段

先认识底层的四种「探测器」,后面的策略都是它们的组合:

手段靠什么判断类/参数
VAD纯声学:有没有人声能量SileroVADAnalyzer(VADParams(stop_secs=0.2))
transcriptionSTT 出没出新文字TranscriptionUserTurnStartStrategy / SpeechTimeoutUserTurnStopStrategy
smart-turn小模型判断「这句语义上说完没」LocalSmartTurnAnalyzerV3(SmartTurnParams)
externalSTT 自己发的端点信号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:736-767 按「STT 类型 + 工作流配置」三选一。注意 start 侧几乎都用 VAD + Transcription 双保险,差异主要在 stop 侧:

场景(条件)start 策略stop 策略适合
external(stt_uses_external_turns 为真)VAD + External(可打断)ExternalUserTurnStopStrategySTT 自带端点检测
turn_analyzer(配置选它)VAD + TranscriptionTurnAnalyzerUserTurnStopStrategy(smart-turn)带自然停顿的长回答
transcription(默认)VAD + TranscriptionSpeechTimeoutUserTurnStopStrategy短的 1-2 词回答

默认值是 turn_stop_strategy = "transcription"(run_pipeline.py:465),偏向短应答;想要不抢话的长应答体验,工作流配置里改成 "turn_analyzer"

5.3 realtime:把回合让给模型

realtime 服务往往自己就带服务端 VAD,本地再插一套会打架。_create_realtime_user_turn_config(run_pipeline.py:114)因此按供应商决定「本地 VAD 让到什么程度」:

realtime 供应商策略本地 VAD原因(源码注释)
OpenAI Realtime / Azure Realtime纯 external供应商已发 speaking-state 和打断事件,聚合器跟着走
Grok Realtime纯 externalGrok 服务端发 speech-start/stop 和打断信号
Google Live / Vertex Realtime本地 VAD,不启用打断Silero让 Gemini 用服务端 VAD 管 barge-in,本地 VAD 只做「回合开始」和状态跟踪
Ultravox本地 VAD,启用打断SileroUltravox 不发用户回合帧,靠本地 VAD 供生命周期信号

这是一处很典型的「按供应商能力做取舍」:能信服务端就完全让位(external),半信半疑就保留本地 VAD 但关掉它的打断权(Google),完全不发信号的就本地全权接管(Ultravox)。

5.4 停止超时:external 为什么给 30 秒

_resolve_user_turn_stop_timeout(run_pipeline.py:104)决定「等多久没动静就强制收尾回合」:

# 示意,非源码
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:101),而默认只有 5 秒。原因:external 模式下「回合结束」由 STT 的端点信号来定,本地超时只是兜底保险,所以放得很宽,免得误伤。


6. 核心原理之四:静音策略与打断

回合检测决定「什么时候听用户」;静音(mute)和打断(interrupt)决定「什么时候听、以及用户能不能插话」。

6.1 用户静音:三条叠加规则

run_pipeline.py:717-721 把三条静音策略叠在一起,任一条命中就静音用户输入:

静音策略什么时候静音用户
MuteUntilFirstBotCompleteUserMuteStrategy机器人还没说完第一句开场白之前
FunctionCallUserMuteStrategy正在执行函数调用期间
CallbackUserMuteStrategy交给引擎回调 engine.should_mute_user 动态决定

前两条是固定规则(开场白期间、函数调用期间别让用户插嘴打乱状态),第三条把决定权交回 PipecatEngine(见 03 章),让图状态机能按节点动态决定要不要闭麦。

6.2 打断(barge-in)

「打断」= 用户在机器人说话时插话,机器人立刻停嘴。在本章范围内,打断权是通过 start 策略的 enable_interruptions 参数表达的:

  • external 场景显式 ExternalUserTurnStartStrategy(enable_interruptions=True)(run_pipeline.py:741)。
  • realtime 的 Google 分支特意 enable_interruptions=False,把 barge-in 让给模型服务端(§5.3)。

另外注意 §4 里各 STT 服务大多传了 should_interrupt=False(如 service_factory.py:171),注释说明「让 UserAggregator 去发 InterruptionFrame」——即打断的决定权集中在聚合器,而不是散落在每个 STT 里。至于工作流节点级的 allow_interrupt(run_pipeline.py:614),那属于通话编排,归 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:83)则按传输类型挑速率:电话供应商从注册表拿它的线路采样率,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:826-829)和 context compaction(run_pipeline.py:673-675),因为一体机自己在服务端管对话状态。
  • 回合检测没有银弹。 三种非 realtime 策略是「短应答 vs 长应答」的取舍,没有一种全场景最优——这也是为什么它做成可配置(turn_stop_strategy)。

9. 巧妙之处(可带走的设计)

  • 「按开关 append 处理器」而非写死链条(pipeline_builder.py:52-92):可选能力(voicemail、录音路由、gate)以「往列表里插」的方式表达,插入位置由注释解释清楚,新增能力不必重写整条链。
  • 一个布尔钩子连通服务选择与回合检测(stt_uses_external_turns, service_factory.py:110):STT 能力(自带端点检测)通过一个函数传递给回合策略装配,两个子系统解耦但对齐。
  • realtime 布局的不对称是被逼出来的正确(pipeline_builder.py:113-124):把 voicemail 检测器放到 LLM 下游,恰好利用「realtime LLM 既是转写源又是上下文汇」的双重身份,避免上下文泄漏——注释把这个非直觉决定讲透了。
  • 按供应商能力分级让位(run_pipeline.py:114-162):external / 本地 VAD 关打断 / 本地 VAD 全权,三档对应「完全信服务端 / 半信 / 不信」,是对接异构 realtime 供应商的干净模式。
  • 采样率单一真相源 + 16kHz 天花板(audio_config.py):把「VAD 只到 16kHz」这个物理约束集中到一个 dataclass 校验,避免各处理器各自为政。

10. 代码地图(导航索引)

主题文件符号
级联管线装配api/services/pipecat/pipeline_builder.pybuild_pipeline
语音到语音管线装配api/services/pipecat/pipeline_builder.pybuild_realtime_pipeline
共用组件(录音缓冲、上下文)api/services/pipecat/pipeline_builder.pycreate_pipeline_components
任务封装 / tracing / metricsapi/services/pipecat/pipeline_builder.pycreate_pipeline_task
STT 供应商分派api/services/pipecat/service_factory.pycreate_stt_service
TTS 供应商分派api/services/pipecat/service_factory.pycreate_tts_service
LLM 供应商分派api/services/pipecat/service_factory.pycreate_llm_service / create_llm_service_from_provider
realtime 一体机分派api/services/pipecat/service_factory.pycreate_realtime_llm_service
STT 是否自带回合边界api/services/pipecat/service_factory.pystt_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.pycreate_webrtc_transport
realtime 传输参数覆盖api/services/pipecat/transport_params.pyrealtime_param_overrides
采样率配置api/services/pipecat/audio_config.pyAudioConfig / create_audio_config