跳到主要内容

数据截至 (上游 commit b21e54d6a845)

语音转文字:渐进转写与过期过滤

这一章讲什么: STT 阶段有两条并行的路——说话过程中每 500 ms 出一次「草稿」,说完后出一次「定稿」。这一章讲这两条路怎么互不干扰,以及为什么在算之前要花力气判断「这段音频还值不值得算」。


1. 两种模式:progressive 和 final

VAD 产出的 VADAudio 带一个 mode 字段(src/speech_to_speech/pipeline/messages.py:51),只有两个值:

mode什么时候产出音频内容下游做什么
progressive说话过程中,每隔一段时间当前累积的缓冲(不完整)PartialTranscription → 字幕增量
finalSilero 判定收尾时完整语音段(含前缀拼接)Transcription → 触发 LLM

渐进模式只在 --enable_live_transcription 开启时产出(默认开启,arguments_classes/module_arguments.py:53-58)。产出间隔不是固定的:_progressive_processing_pause(VAD/vad_handler.py:786-797)会随累积时长拉长——8 秒内是基准间隔,超过 30 秒变 6 倍,并硬顶在 2 秒。理由很实际:说得越久,每次重跑识别的成本越高,不能还按 500 ms 一次的频率烧算力。


2. 渐进转写:窗口怎么滑

它要解决的小问题

最朴素的做法是「每次拿到新音频就把从头到现在的全部重跑一遍」。说 30 秒就要重跑 30 秒的音频,而且每 500 ms 一次——算力爆炸。

思路:句子边界是天然的固化点

已经说完的完整句子几乎不会因为后面的话而改变。所以可以把它们「钉死」(fixed),之后只重跑最后那一小段(active)。

实际音频: [ 句子1 ][ 句子2 ][ 句子3 ][ 正在说的半句 ]
└────────────────┘└───────────────────────┘
钉死 fixed_text 重跑 active(末尾 2s 缓冲内的句子不钉死)

窗口 < 15s:fixed_text 为空,每次把整段重跑一遍
窗口 ≥ 15s:把 cutoff 之前的句子钉死,窗口起点前移,只重跑 active 段

原理演示

# 示意,非源码:固化-重跑的核心循环
fixed_sentences, fixed_end_time = [], 0.0 # fixed_end_time 是「绝对时间轴」上的窗口起点

def step(audio): # audio 是从头开始的完整缓冲
global fixed_end_time
base = fixed_end_time # 本轮窗口起点,循环期间不变
window = audio[int(base * SR):] # 只取「未固化」的部分
result = model.decode(window) # 识别这一段
if len(window) / SR >= 15.0 and len(result.sentences) > 1:
cutoff = len(window) / SR - 2.0 # 末尾留 2 秒不固化
new_end = base
for s in result.sentences: # 注意:s.end 是「窗口内的相对时间」
if s.end >= cutoff:
break
fixed_sentences.append(s.text) # 钉死这句
new_end = base + s.end # 绝对赋值,不是 += (累加会让时间轴越推越远)
fixed_end_time = new_end
result = model.decode(audio[int(fixed_end_time * SR):]) # 从新起点重跑
return " ".join(fixed_sentences), result.text # (固定部分, 活动部分)

两个重点:

  • s.end 是相对于当前窗口起点的秒数,不是绝对时间。所以换算成新的窗口起点必须写成 base + s.end 这样的绝对赋值;写成 fixed_end_time += s.end 会把每句的相对时长反复叠加,几轮之后窗口起点就跑到音频末尾之后,active 段变空。
  • 固化后要立刻用新起点再解一次,否则返回的 active_text 还包含刚被钉死的内容,拼起来会重复。

真实实现

SmartProgressiveStreamingHandler.transcribe_incremental(src/speech_to_speech/STT/smart_progressive_streaming.py:87-157)。参数默认:emission_interval=0.5max_window_size=15.0sentence_buffer=2.0(smart_progressive_streaming.py:40-46)。

源码里对应上面那段绝对赋值的是 sentence_abs_time = self.fixed_end_time + sentence.end,先攒进临时变量 new_fixed_end_time,整个循环跑完才一次性写回 self.fixed_end_time(smart_progressive_streaming.py:131-145)——循环期间 self.fixed_end_time 始终是本轮窗口的起点。

_decode_window(smart_progressive_streaming.py:69-85)兼容两种后端:MLX 版走 model.decode_chunk,nano-parakeet 版走 model.transcribe(timestamps=True) 再把 segment 时间戳包装成统一形状。

只有 Parakeet 后端接了这套渐进逻辑——注册表里只有 _create_parakeet 会把 enable_live_transcription 传进 handler(backend_registry.py:250-264)。其他 STT 后端只做 final。


3. 进/出双向闸门:算之前先问「还值得算吗」

它要解决的小问题

STT 推理可能要几百毫秒。这期间用户可能已经改口。三种浪费/错误:

  1. 队列里堆着的旧 progressive 请求,算完也没人要;
  2. 同一 revision 的 progressive 排在 final 后面,先算了纯属浪费;
  3. 算的时候还是最新的,算完就不是了

做法:一个基类装四道闸

BaseSTTHandler(src/speech_to_speech/STT/base_stt_handler.py)覆写基类的两个钩子,一共四道判定:

┌─ should_process_input ─────────────────────────────────────┐
│ ① 这个 (turn,rev) 的 final 已经出过结果了? → 丢 │
│ ② 我是 progressive,但队列里已经有同 rev 的 final? → 丢 │
│ ③ 不是最新 revision?(final 还要多等一个稳定窗口)→ 丢 │
└────────────────────────────────────────────────────────────┘
│ 通过

process() 跑推理(可能几百 ms)

┌─ should_emit_output ───────────────────────────────────────┐
│ ④ 算完了再查一次:还是最新 revision 吗? → 否则丢 │
└────────────────────────────────────────────────────────────┘

对应源码:should_process_input(base_stt_handler.py:24-61)、should_emit_output(base_stt_handler.py:63-71)。

第 ③ 道闸的特殊性:final 要等

# 真实源码,base_stt_handler.py:91-97
if wait_for_stability:
item_delay_s = max(0.0, getattr(item, "processing_delay_s", 0.0) - self._item_age_s(item))
is_latest = self.speculative_turns.is_latest_after_stability_window(
turn_id, turn_revision,
max(self.final_revision_settle_s, item_delay_s),
)

这里把 Smart Turn 塞进来的 processing_delay_s(见 02 章)减去消息在队列里已经等的时间,只补足剩余部分。也就是说:如果队列本来就堵着、消息已经躺了 600 ms,就不再额外等——延迟预算是「从产生算起」而不是「从取出算起」。

顺手清扫队列

发现一个过期输入时,_drop_stale_queued_inputs(base_stt_handler.py:104-128)会持锁遍历整条输入队列,把所有同样过期的一次性清掉,而不是一个一个慢慢丢。手法是 with self.queue_in.mutex: 直接操作底层 deque,过滤后写回。VAD 那边也有一个对称的 _drop_superseded_vad_audio(VAD/vad_handler.py:465-493),在入队前就把被取代的旧块清掉。

这是这个项目一个反复出现的手法:直接持 Queue.mutex 操作内部 deque,做队列内的批量筛选。 标准库没提供这个能力,但对实时系统很关键。

日志降噪

丢弃是高频事件。_log_stale_turn_item(base_stt_handler.py:130-156)按 (阶段, turn, rev) 计数,第一次用 INFO,后续降到 DEBUG。小技巧,但让生产日志能看。


4. Parakeet handler:两条路径的实际交汇

ParakeetTDTSTTHandler.process(src/speech_to_speech/STT/parakeet_tdt_handler.py:236-370)是把上面所有东西装配起来的地方:

process(vad_audio)

├─ mode == progressive 且开了实时转写
│ ├─ processing_final 为真? → 直接 return(final 优先)
│ ├─ 抢计算锁,超时 0.01s ← 抢不到就放弃,不排队
│ └─ 出 PartialTranscription

└─ mode == final
├─ processing_final = True ← 之后到的 progressive 全部作废
├─ 抢计算锁,超时 5.0s ← 这个必须成功
├─ MLX 走 _process_mlx_final,否则 _process_nano_parakeet
├─ 语言码校验(不在支持列表就沿用上一次)
└─ 出 Transcription(带 speech_stopped_at_s)

超时值的对比说明了优先级:progressive 只等 10 ms(parakeet_tdt_handler.py:266),final 等 5 秒(parakeet_tdt_handler.py:314)。草稿可以丢,定稿不能丢。

speech_stopped_at_s=vad_audio.created_at_s(parakeet_tdt_handler.py:369)——这个时间戳会一路传到 TTS,用来算「用户停止说话到第一声输出」的端到端延迟。

语言检测

_detect_language_from_text(parakeet_tdt_handler.py:379-403)用 lingua-py 从转写文本猜语言,并且短于 20 字符直接放弃(parakeet_tdt_handler.py:394)——短句的语言识别噪声太大。检测结果只有落在 SUPPORTED_LANGUAGES 里才采纳,否则沿用上一次(parakeet_tdt_handler.py:330-333)。


5. MLX 全局锁:Apple Silicon 上的隐形瓶颈

问题

模块 docstring 说得很直接:MLX 模型(STT、LLM、TTS)在 Apple Silicon 上不能被多线程并发使用,因为 Metal command buffer 的限制(src/speech_to_speech/utils/mlx_lock.py:1-15)。

做法

一把进程级全局 RLock(mlx_lock.py:27),所有 MLX handler 用前必须拿。用 RLock 是为了同线程可重入。

附带一层可观测性:_record_lock_acquired / _record_lock_released(mlx_lock.py:45-85)记录持有者线程名、handler 名、持有时长和重入深度,抢锁失败时能打出「谁拿着、拿了多久」。

连锁后果

这把锁把「STT / LLM / TTS 可以并行」变成了「Mac 上必须串行」,进而导致:

  • 渐进转写抢不到锁就直接放弃(所以字幕会跳过若干次更新,但不影响最终结果);
  • --num_pipelines > 1 时,启动阶段直接关掉实时转写,免得日志被抢锁失败刷屏(s2s_pipeline.py:634-640)。

非 MLX 后端走的是 handler 自己的 compute_lock,不是全局锁(parakeet_tdt_handler.py:405-415_compute_lock_context 分支)。


6. TranscriptionNotifier:一个只发事件的 handler

它坐在 STT 和 LLM 之间,但从不往下游队列放业务数据——process 的最后一行是 yield from ()(src/speech_to_speech/STT/transcription_notifier.py:105)。

它做三件事:

  1. PartialTranscriptionPartialTranscriptionEvent 放旁路队列;
  2. TranscriptionTranscriptionCompletedEvent 放旁路队列,即使转写为空也发(transcription_notifier.py:79-91),因为客户端可能已经收到增量,需要一个终结事件来收尾;
  3. 转写为空时把 should_listen 重新置位(transcription_notifier.py:95-97),因为不会有响应来触发正常的「恢复收听」路径。

转写增量的协议语义

项目自带的英文设计文档说明了一个重要细节(src/speech_to_speech/api/openai_realtime/README.md 的 “Input transcription semantics” 段):内部的 partial 是累积假设而非增量;协议层在发 conversation.item.input_audio_transcription.delta 前,会比对相邻两次假设、在词边界上扣住最新的那个词,只把确认增长的部分发出去。因为 Realtime 协议没有「撤回转写」的事件——发出去就收不回来了。

注意 PartialTranscriptionEvent 的字段名仍叫 delta,但注释明说它的值是累积假设,保留这个名字只为 API 兼容(pipeline/events.py:56-67)。这是一个容易踩的命名陷阱。


7. 直接音频输入:跳过 STT

--stt none 时,STT 位置换成 AudioInputNotifier(src/speech_to_speech/LLM/audio_input_notifier.py),它把 VADAudio 转成 AudioInputCompletedEvent 走旁路,由协议层组装成带音频的 GenerateResponseRequest(service.py:680-713)。

约束(s2s_pipeline.py:321-323):必须搭配 capabilities.supports_audio_input 为真的 LLM 后端,当前只有 chat-completions(backend_registry.py:418)。README 解释了原因:有些模型支持 /v1/chat/completions 的音频输入但不支持 /v1/responses


8. 代码地图

主题文件路径符号名
STT 双向闸门基类src/speech_to_speech/STT/base_stt_handler.pyBaseSTTHandler.should_process_input, should_emit_output, _drop_stale_queued_inputs
渐进窗口滑动src/speech_to_speech/STT/smart_progressive_streaming.pySmartProgressiveStreamingHandler.transcribe_incremental, _decode_window
Parakeet 主逻辑src/speech_to_speech/STT/parakeet_tdt_handler.pyParakeetTDTSTTHandler.process, _compute_lock_context, _process_mlx_final
语言检测src/speech_to_speech/STT/parakeet_tdt_handler.py_detect_language_from_text, SUPPORTED_LANGUAGES
转写事件桥src/speech_to_speech/STT/transcription_notifier.pyTranscriptionNotifier.process
直接音频输入src/speech_to_speech/LLM/audio_input_notifier.pyAudioInputNotifier
MLX 全局锁src/speech_to_speech/utils/mlx_lock.pyacquire_mlx_lock, MLXLockContext, _owner_snapshot
其他 STT 后端src/speech_to_speech/STT/WhisperSTTHandler, FasterWhisperSTTHandler, MLXAudioWhisperSTTHandler, ParaformerSTTHandler