先说结论:Voice Agent 接入流式 TTS,最容易低估的不是调用合成接口,而是接口前后的三段工程——什么时候把 LLM 文本切成句子、音频包怎样连续播放、用户打断后怎样同时清掉客户端队列和服务端任务

本文给出一套可以直接运行的最小实现:FastAPI 负责 WebSocket,服务端输出 PCM16 二进制帧,浏览器用 Web Audio API 排队播放,并用 turn_id 丢弃过期音频。为了让所有人不申请 API Key 也能复现,Demo 用本地生成的提示音替代真实 TTS。它验证链路和时序,不比较任何厂商音色。

适用范围:

  • 已经能拿到 LLM 流式文本,准备接 TTS;
  • 想定位“接口说是流式,用户为什么还是等很久”;
  • 正在处理分句、卡顿、爆音、音频堆积或打断后继续说的问题;
  • 前端是浏览器,后端是 Python/FastAPI;其他客户端也可复用同一协议。

不适用范围:

  • 本文不做 TTS 厂商自然度排名;
  • Mock 输出是提示音,不是语音模型;
  • 本地延迟数据不能外推为公网或生产环境性能。

一、流式 TTS 不是把 stream=True 打开就结束了

我最初把 TTS 接到 Voice Agent 时,链路看起来很顺:

LLM 文本 -> TTS 流式接口 -> 音频块 -> 播放器

真正跑起来后,问题集中出现在两头。

第一种情况是“假流式”。TTS 接口确实逐块返回音频,但程序一直等 LLM 把整段话生成完,才一次性把全文送去合成。TTS 首包可能只有 200 ms,用户却先等了 1~2 秒文本。

第二种情况是“包到了,声音仍然断”。如果每收到一个 PCM 包就立即播放一次,浏览器会创建许多互不相干的播放任务。网络或调度只要抖一下,包与包之间就会出现空洞,听感是咔哒声、断句和忽快忽慢。

第三种更隐蔽:用户已经打断,前端还存着一秒多待播音频。服务端即使停止生成,旧队列仍会继续说完,Voice Agent 就会像“没听见用户”。

所以完整链路应该画成这样:

LLM token 增量
      |
      v
文本缓冲器 -- 句末标点/最大长度 --> sentence_queue
                                      |
                                      v
                              Streaming TTS
                                      |
                         PCM16 + 16 字节包头
                                      |
                                      v
                              WebSocket 二进制帧
                                      |
                                      v
                      浏览器 AudioBuffer 排期队列
                                      |
                                      v
                                 扬声器播放

用户打断
   |----------------> stop() 已排期音频
   |----------------> 清空 nextPlayTime
   |----------------> cancel 服务端 asyncio.Task
   `----------------> turn_id 丢弃迟到旧包

这篇文章就按这条链路往下搭。


二、运行环境与项目目录

本文实际运行环境:

项目 版本/配置
操作系统 macOS,Apple Silicon
Python 3.12.13
FastAPI 0.115.12
Uvicorn 0.34.3
websockets 15.0.1
音频 PCM16 little-endian、24 kHz、单声道
前端 Chrome + Web Audio API

项目目录:

streaming-tts-demo/
├── app.py                 # FastAPI、WebSocket、任务取消
├── segmenter.py           # LLM 增量文本分句
├── tts_provider.py        # 可复现的 Mock 流式 TTS
├── protocol.py            # 16 字节二进制包头
├── benchmark.py           # 首包与取消基准
├── test_demo.py           # 单元测试
├── requirements.txt
└── static/
    └── index.html         # PCM 播放队列与打断

requirements.txt

fastapi==0.115.12
uvicorn[standard]==0.34.3
websockets==15.0.1

安装并启动:

python3.12 -m venv .venv
source .venv/bin/activate
pip install -r requirements.txt
uvicorn app:app --host 127.0.0.1 --port 8000

浏览器打开:

http://127.0.0.1:8000

FastAPI 官方文档说明 WebSocket 端点可以收发文本、JSON 和二进制数据。本文把控制消息放在 JSON 帧里,把 PCM 放在二进制帧里,避免音频做 Base64 后额外膨胀。


三、先定义清楚“首包延迟”

很多 TTS 测试只报一个“首包 200 ms”,但 Voice Agent 至少有四个相关时间点:

t0  用户问题已经确定
t1  LLM 输出第一个 token
t2  分句器得到第一段可合成文本
t3  发起第一句 TTS 请求
t4  TTS 返回第一块 PCM
t5  客户端收到第一块 PCM
t6  播放器把第一块音频排入时间轴

本文重点记录三项:

首个可合成短句 = t2 - t0
TTS 自身首包     = t4 - t3
客户端首音频包   = t5 - t0

这三个数不能混用。

如果 LLM 到句末花了 700 ms,TTS 自身只花 180 ms,那么用户侧至少已经等了约 880 ms。此时继续压 TTS 的 20 ms,收益远不如修改分句策略。

本次 Demo 的一次完整运行中:

first_text_segment_ms = 116.8
tts_first_packet_ms   = 182.2
first_pcm_ms          = 299.0

关系非常直观:

约 116.8 ms 等到第一句 + 约 182.2 ms 等 TTS = 约 299.0 ms 服务端首 PCM

这里的 180 ms 是 Mock 中显式设置的等待参数,不是某个 TTS 产品的成绩。这样做的好处是,任何人运行代码都能验证埋点和队列,而不需要相信一组无法复现的“厂商实测”。


四、文本怎么切:不能每个 token 合成,也不能等完整段落

4.1 两个极端都不好

把每个 LLM token 都送去 TTS:

  • 请求数量暴涨;
  • 语音韵律被切碎;
  • 前后音色、音量和停顿可能不连续;
  • 很难撤销已经发出的几十个请求。

等 LLM 输出完整段落:

  • 首句话已经生成,却迟迟不能开口;
  • 答案越长,首音频等待越明显;
  • 用户会误以为系统卡死。

我在中文客服场景里采用的最小规则是:

  1. 遇到 。!?; 等强句末标点立即切;
  2. 没有句末标点但达到最大长度时,优先在逗号、冒号处切;
  3. 最后一段由 flush() 送出;
  4. 不把英文句点 . 直接当硬边界,避免把 URL、小数和版本号粗暴切开。

4.2 segmenter.py

from __future__ import annotations

from dataclasses import dataclass, field


HARD_BOUNDARIES = frozenset("。!?!?;;\n")
SOFT_BOUNDARIES = frozenset(",,::、")


@dataclass
class TextSegmenter:
    """把 LLM 的增量文本整理成适合送入 TTS 的短句。"""

    max_chars: int = 36
    min_soft_chars: int = 18
    _buffer: str = field(default="", init=False, repr=False)

    def feed(self, text_delta: str) -> list[str]:
        self._buffer += text_delta
        return self._drain()

    def flush(self) -> list[str]:
        segments = self._drain()
        tail = self._buffer.strip()
        self._buffer = ""
        if tail:
            segments.append(tail)
        return segments

    def _drain(self) -> list[str]:
        segments: list[str] = []

        while self._buffer:
            split_at = self._find_split()
            if split_at is None:
                break

            segment = self._buffer[:split_at].strip()
            self._buffer = self._buffer[split_at:].lstrip()
            if segment:
                segments.append(segment)

        return segments

    def _find_split(self) -> int | None:
        for index, char in enumerate(self._buffer):
            current_length = index + 1

            if char in HARD_BOUNDARIES:
                return current_length

            if current_length >= self.max_chars:
                soft_split = self._last_soft_boundary(
                    start=self.min_soft_chars - 1,
                    end=index,
                )
                return soft_split or current_length

        return None

    def _last_soft_boundary(self, start: int, end: int) -> int | None:
        for index in range(end, max(start - 1, -1), -1):
            if self._buffer[index] in SOFT_BOUNDARIES:
                return index + 1
        return None

36 不是标准答案。客服确认语可以更短,知识解释可以更长。生产环境建议按以下维度做 A/B 测试:

  • 第一段比后续段更短,优先让 Agent 尽快开口;
  • 数字、日期、金额、地址尽量保持完整;
  • Markdown 链接、代码和 URL 进入 TTS 前先清洗;
  • 标点规则之外,再加“最大等待时间”,防止 LLM 长时间不输出句末符号。

最后一条在真实 LLM 流里很重要。本文 Demo 用 4 字一组、每组 28 ms 的固定增量模拟 token,因此使用“句末 + 最大长度”即可稳定复现;生产环境应再加定时 flush。


五、为什么直接传 PCM,而不是每包都做成 MP3

浏览器拿到 MP3、AAC 等压缩音频后,通常还要先凑够可解码的数据,再交给解码器。不同编码器还可能带容器头、帧边界和额外缓冲。

本文使用:

PCM signed 16-bit
little-endian
24000 Hz
单声道

优点是收到字节就能按样本转换,延迟和边界都比较透明。缺点是带宽更大:

24000 samples/s × 2 bytes × 1 channel
= 48000 bytes/s
≈ 46.9 KiB/s

对浏览器 Demo 和内网调试很合适。公网生产环境是否使用 Opus,要结合带宽、解码延迟、客户端能力和电话网关格式再决定。

5.1 给每个音频包加 16 字节包头

只发裸 PCM 会遇到一个问题:用户打断后,网络里可能还有旧包在路上。客户端不知道这包属于哪一轮,就可能把它重新排进队列。

我给每个包加四个小端无符号整数:

字段 字节 作用
turn_id 4 属于哪一轮回答
sequence 4 包序号
sample_rate 4 采样率
flags 4 句末等标记

protocol.py

from __future__ import annotations

import struct


HEADER = struct.Struct("<IIII")
FLAG_SENTENCE_END = 1


def pack_audio_packet(
    *,
    turn_id: int,
    sequence: int,
    sample_rate: int,
    pcm: bytes,
    sentence_end: bool = False,
) -> bytes:
    flags = FLAG_SENTENCE_END if sentence_end else 0
    return HEADER.pack(turn_id, sequence, sample_rate, flags) + pcm


def unpack_audio_header(packet: bytes) -> tuple[int, int, int, int]:
    if len(packet) < HEADER.size:
        raise ValueError(f"audio packet is shorter than {HEADER.size} bytes")
    return HEADER.unpack_from(packet)

这里的 <IIII 中,< 表示 little-endian,四个 I 都是 32 位无符号整数。


六、做一个不依赖 API Key 的 Mock Streaming TTS

如果正文直接接某家云 TTS,会把大量篇幅花在账号、鉴权和 SDK 版本上,而且读者很难判断延迟到底来自网络、服务商还是我们的队列。

所以我先做一个可控替身:

  • 首包等待 180 ms;
  • 每块音频 40 ms;
  • 每隔 20 ms 产出一块,生成速度约为播放速度的 2 倍;
  • 输出能被播放器直接播放的 PCM16 提示音;
  • 协程在 sleep() 或下一次迭代时可以被取消。

tts_provider.py

from __future__ import annotations

import asyncio
import math
import sys
from array import array
from collections.abc import AsyncIterator
from dataclasses import dataclass


@dataclass(frozen=True)
class AudioChunk:
    pcm: bytes
    sample_rate: int
    sentence_end: bool


@dataclass(frozen=True)
class MockStreamingTTS:
    """
    本地可运行的流式 TTS 替身。

    它输出可播放的 PCM16 单声道提示音,不模拟任何厂商的音色质量。
    first_packet_delay_ms 和 packet_interval_ms 是可控实验参数,
    方便验证分句、二进制传输、播放队列和取消逻辑。
    """

    sample_rate: int = 24_000
    frame_ms: int = 40
    first_packet_delay_ms: int = 180
    packet_interval_ms: int = 20
    char_duration_ms: int = 55

    async def stream(self, text: str) -> AsyncIterator[AudioChunk]:
        visible_chars = max(1, sum(not char.isspace() for char in text))
        duration_ms = min(6_000, max(360, visible_chars * self.char_duration_ms))
        total_samples = int(self.sample_rate * duration_ms / 1_000)
        samples_per_frame = int(self.sample_rate * self.frame_ms / 1_000)
        emitted = 0
        phase = 0.0
        frequency = 330 + (sum(map(ord, text)) % 120)

        await asyncio.sleep(self.first_packet_delay_ms / 1_000)

        while emitted < total_samples:
            frame_samples = min(samples_per_frame, total_samples - emitted)
            pcm = array("h")

            for _ in range(frame_samples):
                envelope = min(1.0, emitted / max(1, self.sample_rate * 0.02))
                remaining = total_samples - emitted
                envelope *= min(1.0, remaining / max(1, self.sample_rate * 0.02))
                value = int(8_000 * envelope * math.sin(phase))
                pcm.append(value)
                phase += 2 * math.pi * frequency / self.sample_rate
                emitted += 1

            if pcm.itemsize != 2:
                raise RuntimeError("当前平台的 signed short 不是 16 bit")
            if sys.byteorder != "little":
                pcm.byteswap()

            sentence_end = emitted >= total_samples
            yield AudioChunk(
                pcm=pcm.tobytes(),
                sample_rate=self.sample_rate,
                sentence_end=sentence_end,
            )

            if not sentence_end:
                await asyncio.sleep(self.packet_interval_ms / 1_000)

替换真实厂商时,保留 stream(text) -> AsyncIterator[AudioChunk] 这层接口即可。厂商 SDK 返回什么对象并不重要,最后统一成 PCM、采样率和句末标志,后面的传输、播放、取消逻辑就不用跟着改。


七、FastAPI 后端:让 LLM 产句和 TTS 合成并行

如果代码是下面这种写法:

full_text = await llm.complete(question)
async for audio in tts.stream(full_text):
    ...

那仍然是假流式。正确做法是两个协程:

produce_sentences:持续读取 LLM delta,产出短句
consume_sentences:读取 sentence_queue,逐句合成音频

两者通过有界 asyncio.Queue 相连。本文设置 maxsize=4,防止文本生产无限领先。

7.1 app.py

from __future__ import annotations

import asyncio
import contextlib
import time
from pathlib import Path
from typing import Any

from fastapi import FastAPI, WebSocket, WebSocketDisconnect
from fastapi.responses import FileResponse

from protocol import pack_audio_packet
from segmenter import TextSegmenter
from tts_provider import MockStreamingTTS


BASE_DIR = Path(__file__).resolve().parent
INDEX_HTML = BASE_DIR / "static" / "index.html"

app = FastAPI(title="Voice Agent Streaming TTS Demo")
tts = MockStreamingTTS()


@app.get("/")
async def index() -> FileResponse:
    return FileResponse(INDEX_HTML)


async def fake_llm_stream(
    text: str,
    *,
    chars_per_delta: int = 4,
    delta_interval_ms: int = 28,
):
    """用固定节奏模拟 LLM 的 token 增量输出。"""
    for start in range(0, len(text), chars_per_delta):
        await asyncio.sleep(delta_interval_ms / 1_000)
        yield text[start : start + chars_per_delta]


async def stream_turn(websocket: WebSocket, turn_id: int, text: str) -> None:
    started_at = time.perf_counter()
    sentence_queue: asyncio.Queue[str | None] = asyncio.Queue(maxsize=4)
    metrics: dict[str, Any] = {
        "turn_id": turn_id,
        "sentence_count": 0,
        "audio_chunk_count": 0,
        "pcm_bytes": 0,
        "first_text_segment_ms": None,
        "first_tts_request_ms": None,
        "first_pcm_ms": None,
        "tts_first_packet_ms": None,
    }

    await websocket.send_json(
        {
            "type": "turn_start",
            "turn_id": turn_id,
            "audio_format": "pcm_s16le",
            "sample_rate": tts.sample_rate,
            "channels": 1,
        }
    )

    async def produce_sentences() -> None:
        segmenter = TextSegmenter(max_chars=36, min_soft_chars=18)

        async for delta in fake_llm_stream(text):
            for sentence in segmenter.feed(delta):
                if metrics["first_text_segment_ms"] is None:
                    metrics["first_text_segment_ms"] = elapsed_ms(started_at)
                await sentence_queue.put(sentence)

        for sentence in segmenter.flush():
            if metrics["first_text_segment_ms"] is None:
                metrics["first_text_segment_ms"] = elapsed_ms(started_at)
            await sentence_queue.put(sentence)

        await sentence_queue.put(None)

    async def consume_sentences() -> None:
        sequence = 0

        while True:
            sentence = await sentence_queue.get()
            if sentence is None:
                break

            sentence_index = metrics["sentence_count"]
            metrics["sentence_count"] += 1
            tts_started_at = time.perf_counter()
            if metrics["first_tts_request_ms"] is None:
                metrics["first_tts_request_ms"] = elapsed_ms(started_at)

            await websocket.send_json(
                {
                    "type": "sentence_start",
                    "turn_id": turn_id,
                    "sentence_index": sentence_index,
                    "text": sentence,
                }
            )

            async for audio_chunk in tts.stream(sentence):
                if metrics["first_pcm_ms"] is None:
                    metrics["first_pcm_ms"] = elapsed_ms(started_at)
                    metrics["tts_first_packet_ms"] = elapsed_ms(tts_started_at)

                packet = pack_audio_packet(
                    turn_id=turn_id,
                    sequence=sequence,
                    sample_rate=audio_chunk.sample_rate,
                    pcm=audio_chunk.pcm,
                    sentence_end=audio_chunk.sentence_end,
                )
                await websocket.send_bytes(packet)
                sequence += 1
                metrics["audio_chunk_count"] += 1
                metrics["pcm_bytes"] += len(audio_chunk.pcm)

            await websocket.send_json(
                {
                    "type": "sentence_end",
                    "turn_id": turn_id,
                    "sentence_index": sentence_index,
                }
            )

    async with asyncio.TaskGroup() as task_group:
        task_group.create_task(produce_sentences())
        task_group.create_task(consume_sentences())

    metrics["server_elapsed_ms"] = elapsed_ms(started_at)
    metrics["audio_duration_ms"] = round(
        metrics["pcm_bytes"] / 2 / tts.sample_rate * 1_000,
        1,
    )
    metrics["note"] = "Mock TTS:延迟为本地可控实验参数,不代表任何厂商性能"
    await websocket.send_json({"type": "turn_done", **metrics})


def elapsed_ms(started_at: float) -> float:
    return round((time.perf_counter() - started_at) * 1_000, 1)


async def cancel_task(task: asyncio.Task[None] | None) -> bool:
    if task is None or task.done():
        return False

    task.cancel()
    with contextlib.suppress(asyncio.CancelledError):
        await task
    return True


@app.websocket("/ws")
async def websocket_endpoint(websocket: WebSocket) -> None:
    await websocket.accept()
    current_task: asyncio.Task[None] | None = None
    current_turn_id = 0

    try:
        while True:
            message = await websocket.receive_json()
            message_type = message.get("type")

            if message_type == "speak":
                await cancel_task(current_task)
                current_turn_id += 1
                text = str(message.get("text", "")).strip()

                if not text:
                    await websocket.send_json(
                        {"type": "error", "message": "text 不能为空"}
                    )
                    continue

                current_task = asyncio.create_task(
                    stream_turn(websocket, current_turn_id, text),
                    name=f"tts-turn-{current_turn_id}",
                )

            elif message_type == "cancel":
                cancelled = await cancel_task(current_task)
                await websocket.send_json(
                    {
                        "type": "cancelled",
                        "turn_id": current_turn_id,
                        "cancelled": cancelled,
                    }
                )

            else:
                await websocket.send_json(
                    {"type": "error", "message": f"unknown type: {message_type}"}
                )

    except WebSocketDisconnect:
        await cancel_task(current_task)

这里有两个容易漏掉的点。

第一,取消任务后要 await task。Python 官方文档说明,调用 Task.cancel() 会在协程下一次获得执行机会时抛出 CancelledError;协程可以在 finally 中清理资源,然后通常应继续传播这个异常。只调用 cancel() 而不等待,清理可能尚未完成,新一轮任务就开始了。

第二,每次 speak 都递增 turn_id。取消是“阻止继续生产”,turn_id 是“阻止迟到旧包重新进入播放队列”,两者不是同一件事。


八、浏览器播放:不要用 setTimeout 一个包一个包播

8.1 用 AudioContext 时间轴连续排期

浏览器收到 PCM 后,把 Int16 样本转换成 [-1, 1] 范围的 Float32,再创建 AudioBufferSourceNode

关键不是调用 start(),而是给每个包计算准确的开始时间:

source.start(nextPlayTime);
nextPlayTime += audioBuffer.duration;

nextPlayTime 是音频时间轴上的游标。后一个包从前一个包结束的位置开始,而不是依赖 JavaScript 定时器“差不多到点再播”。

MDN 对 AudioBufferSourceNode.start(when) 的定义就是在 AudioContext 的时间坐标中安排播放。它也说明一个 AudioBufferSourceNode 只能启动一次,因此每个 PCM 块都要创建新的 source。

8.2 完整 static/index.html

下面代码包含连接、PCM 转换、排期、队列深度、旧包丢弃和立即打断。样式被压缩在顶部,不影响复制运行。

<!doctype html>
<html lang="zh-CN">
<head>
  <meta charset="utf-8">
  <meta name="viewport" content="width=device-width, initial-scale=1">
  <title>Voice Agent 流式 TTS Demo</title>
  <style>
    body{max-width:920px;margin:40px auto;padding:0 20px;background:#0b1020;color:#e8edff;font-family:system-ui}
    textarea{width:100%;min-height:120px;box-sizing:border-box;padding:14px;background:#151d33;color:inherit}
    button{margin:12px 8px 12px 0;padding:10px 18px}
    .metrics{display:grid;grid-template-columns:repeat(4,1fr);gap:10px}
    .card,pre{padding:14px;background:#121a2d}.value{display:block;color:#85a6ff}
  </style>
</head>
<body>
  <h1>Voice Agent 流式 TTS Demo</h1>
  <p>本页播放本地 PCM16 提示音,用来验证分句、首包、队列和打断。</p>
  <textarea id="text">您好,我已经查到您的订单。包裹目前在杭州转运中心,预计明天下午送达。如果地址需要修改,我可以继续帮您处理。</textarea>
  <div>
    <button id="speak">开始流式播放</button>
    <button id="cancel">立即打断</button>
  </div>
  <section class="metrics">
    <div class="card">服务端首包<span class="value" id="server-first">-</span></div>
    <div class="card">客户端收到首包<span class="value" id="client-first">-</span></div>
    <div class="card">待播队列<span class="value" id="queue-depth">0 ms</span></div>
    <div class="card">队列欠载<span class="value" id="underruns">0</span></div>
  </section>
  <pre id="log"></pre>

  <script>
    const HEADER_BYTES = 16;
    const START_BUFFER_SECONDS = 0.08;
    const textInput = document.querySelector("#text");
    const logElement = document.querySelector("#log");

    let socket;
    let audioContext;
    let activeTurnId = null;
    let activeSources = new Set();
    let nextPlayTime = 0;
    let clientRequestAt = 0;
    let firstPacketSeen = false;
    let underrunCount = 0;

    function log(message) {
      const time = new Date().toLocaleTimeString();
      logElement.textContent += `[${time}] ${message}\n`;
    }

    async function ensureAudioContext() {
      if (!audioContext) audioContext = new AudioContext();
      if (audioContext.state === "suspended") await audioContext.resume();
    }

    function ensureSocket() {
      if (socket?.readyState === WebSocket.OPEN) return Promise.resolve();

      return new Promise((resolve, reject) => {
        const protocol = location.protocol === "https:" ? "wss:" : "ws:";
        socket = new WebSocket(`${protocol}//${location.host}/ws`);
        socket.binaryType = "arraybuffer";
        socket.addEventListener("open", resolve, { once: true });
        socket.addEventListener("error", reject, { once: true });
        socket.addEventListener("message", handleMessage);
        socket.addEventListener("close", () => log("WebSocket 已断开"));
      });
    }

    function stopPlayback() {
      for (const source of activeSources) {
        try {
          source.stop();
        } catch (_) {
          // 节点可能已经自然结束。
        }
      }
      activeSources.clear();
      nextPlayTime = audioContext?.currentTime ?? 0;
      firstPacketSeen = false;
      document.querySelector("#queue-depth").textContent = "0 ms";
    }

    function handleMessage(event) {
      if (event.data instanceof ArrayBuffer) {
        handleAudioPacket(event.data);
        return;
      }

      const message = JSON.parse(event.data);
      if (message.type === "turn_start") {
        stopPlayback();
        activeTurnId = message.turn_id;
        log(`turn=${message.turn_id}${message.audio_format}/${message.sample_rate}Hz`);
      } else if (message.type === "sentence_start") {
        log(`开始合成第 ${message.sentence_index + 1} 句:${message.text}`);
      } else if (message.type === "turn_done") {
        document.querySelector("#server-first").textContent =
          `${message.first_pcm_ms} ms`;
        log(`完成:${message.audio_chunk_count}`);
      } else if (message.type === "cancelled") {
        log(`服务端取消结果:${message.cancelled}`);
      } else if (message.type === "error") {
        log(`错误:${message.message}`);
      }
    }

    function handleAudioPacket(packet) {
      const view = new DataView(packet);
      const turnId = view.getUint32(0, true);
      const sequence = view.getUint32(4, true);
      const sampleRate = view.getUint32(8, true);

      if (turnId !== activeTurnId) {
        log(`丢弃过期包:turn=${turnId}, seq=${sequence}`);
        return;
      }

      const isFirstPacket = !firstPacketSeen;
      if (isFirstPacket) {
        firstPacketSeen = true;
        const firstPacketMs = Math.round(performance.now() - clientRequestAt);
        document.querySelector("#client-first").textContent =
          `${firstPacketMs} ms`;
      }

      const pcmView = new DataView(packet, HEADER_BYTES);
      const sampleCount = pcmView.byteLength / 2;
      const audioBuffer = audioContext.createBuffer(1, sampleCount, sampleRate);
      const channel = audioBuffer.getChannelData(0);

      for (let index = 0; index < sampleCount; index++) {
        channel[index] = pcmView.getInt16(index * 2, true) / 32768;
      }

      const now = audioContext.currentTime;
      if (nextPlayTime < now) {
        if (!isFirstPacket && nextPlayTime > 0) {
          underrunCount += 1;
          document.querySelector("#underruns").textContent =
            String(underrunCount);
        }
        nextPlayTime = now + START_BUFFER_SECONDS;
      }

      const source = audioContext.createBufferSource();
      source.buffer = audioBuffer;
      source.connect(audioContext.destination);
      source.addEventListener(
        "ended",
        () => activeSources.delete(source),
        { once: true }
      );
      activeSources.add(source);
      source.start(nextPlayTime);
      nextPlayTime += audioBuffer.duration;

      const queueMs = Math.max(
        0,
        Math.round((nextPlayTime - now) * 1000)
      );
      document.querySelector("#queue-depth").textContent = `${queueMs} ms`;
    }

    document.querySelector("#speak").addEventListener("click", async () => {
      await ensureAudioContext();
      await ensureSocket();
      stopPlayback();
      activeTurnId = null;
      clientRequestAt = performance.now();
      socket.send(JSON.stringify({ type: "speak", text: textInput.value }));
      log("已发送 speak");
    });

    document.querySelector("#cancel").addEventListener("click", () => {
      stopPlayback();
      activeTurnId = null;
      if (socket?.readyState === WebSocket.OPEN) {
        socket.send(JSON.stringify({ type: "cancel" }));
      }
      log("客户端音频已立即停止,并发送 cancel");
    });
  </script>
</body>
</html>

注意 socket.binaryType = "arraybuffer"。MDN 文档中,WebSocket 默认二进制类型是 Blob,设置为 arraybuffer 后才能直接用 DataView 按小端读取包头和 PCM。


九、打断时到底要清什么

一个可用的 Voice Agent 至少有四类“旧状态”:

1. LLM 还在生成的文本
2. TTS 正在生成的音频
3. WebSocket 途中尚未到达的音频包
4. 浏览器已经收到但还没播放的音频

本文 Demo 处理了后面三类:

9.1 服务端取消

current_task.cancel()
await current_task

这会终止等待中的 Mock TTS 协程。替换真实 SDK 后,还要检查 SDK 是否支持关闭流、取消请求或断开连接;如果底层请求不可取消,至少不要再把返回包向客户端转发。

9.2 客户端立即停掉已排期音频

for (const source of activeSources) {
  source.stop();
}
activeSources.clear();
nextPlayTime = audioContext.currentTime;

只清 JavaScript 数组不够。已经调用过 source.start(futureTime) 的节点仍在 AudioContext 时间轴上,必须逐个 stop()

9.3 丢弃迟到旧包

if (turnId !== activeTurnId) {
  return;
}

点击取消时把 activeTurnId 置空。即使旧包晚到,也进不了新队列。

下一篇“用户打断以后,上下文怎么续”还会进一步区分:已经生成、已经播放、用户已经听到并确认的内容。本文先只解决音频层的停止。


十、实测结果:首包、队列与取消

10.1 单元测试

运行:

python -m unittest -v

结果:

test_stream_returns_pcm16_chunks ... ok
test_binary_header_round_trip ... ok
test_max_length_prefers_soft_boundary ... ok
test_streaming_chinese_sentences ... ok

Ran 4 tests in 0.023s
OK

10.2 完整链路连续运行 5 次

基准脚本在 WebSocket 建立后开始计时,避免把握手时间混进首音频包:

python benchmark.py --runs 5

本机结果:

first_binary_ms:
298.3, 297.8, 297.2, 298.8, 299.5

median = 298.3 ms

最后一次服务端日志:

{
  "first_text_segment_ms": 116.8,
  "first_tts_request_ms": 116.8,
  "first_pcm_ms": 299.0,
  "tts_first_packet_ms": 182.2,
  "server_elapsed_ms": 2195.5,
  "audio_duration_ms": 2915.0,
  "audio_chunk_count": 74
}

浏览器实测客户端收到首包约 299 ms,完整播放队列欠载为 0。

这里还有一个值得注意的现象:音频总时长是 2915 ms,服务端约 2195.5 ms 已经生成并发送完。因为 Mock 每 20 ms 生成 40 ms 音频,生成速度快于播放,客户端待播队列会逐渐变深。

这能避免句间断流,但队列越深,已经生成却来不及播放的音频越多。生产环境通常要设置高、低水位,例如:

低于 120 ms:容易欠载,优先补包
120~400 ms:正常缓冲区间
高于 400 ms:放慢预取或暂停继续合成

这组水位只是起始建议,不是通用标准。Web、App、SIP 网关和电话线路的抖动不同,需要用真实网络重新测。

10.3 400 ms 时取消,连续运行 5 次

python benchmark.py --cancel-after-ms 400 --runs 5

结果:

cancel_ack_ms:
0.6, 0.6, 0.6, 0.5, 0.5

median = 0.6 ms

这是本机进程内 Mock 协程的取消确认时间,不包含公网、真实厂商 SDK 和声卡输出延迟。浏览器点击“立即打断”时,会先同步停止已排期 source,再通知服务端,因此听感上的停止不必等待这条确认消息返回。


十一、基准脚本的关键写法

websockets 15.0.1 会读取代理配置。本机如果开着系统代理,本地 ws://127.0.0.1 也可能被尝试走代理,出现:

ImportError: python-socks is required to use a SOCKS proxy

本地基准明确关闭代理:

async with connect(
    "ws://127.0.0.1:8000/ws",
    max_size=4 * 1024 * 1024,
    proxy=None,
) as websocket:
    started_at = time.perf_counter()
    await websocket.send(json.dumps({"type": "speak", "text": text}))

首包统计只在收到第一个二进制帧时记录:

if isinstance(message, bytes):
    if result["first_binary_ms"] is None:
        result["first_binary_ms"] = round(
            (time.perf_counter() - started_at) * 1_000,
            1,
        )

取消测试则在发出 cancel 后单独开始计时,直到收到 cancelled 控制消息。这样“首包”和“取消确认”不会混成同一个含糊的总耗时。


十二、常见报错与排查

12.1 浏览器没有声音,但能看到音频包

先检查:

console.log(audioContext.state);

浏览器自动播放策略可能让 AudioContext 处于 suspended。要在用户点击按钮后执行:

await audioContext.resume();

不要在页面加载时偷偷播放。

12.2 声音速度不对,像慢放或快放

最常见原因是采样率不一致。

服务端如果输出 24 kHz,客户端创建 AudioBuffer 时也必须传 24000:

audioContext.createBuffer(1, sampleCount, sampleRate);

不要把 audioContext.sampleRate 当成服务端 PCM 的采样率。浏览器会负责把 AudioBuffer 的原始采样率重采样到输出设备。

12.3 全是噪声

依次检查:

  • PCM 是有符号 16 位还是无符号 16 位;
  • little-endian 还是 big-endian;
  • 单声道还是双声道;
  • 是否把 16 字节自定义包头误当成音频;
  • Int16 转 Float32 时是否除以 32768。

本文读取方式:

pcmView.getInt16(index * 2, true) / 32768

true 代表 little-endian。

12.4 每个音频块之间有“咔哒”声

通常不是 WebSocket 丢包,而是播放方式不连续。

检查是否:

  • 每个包都 source.start() 立即播放;
  • 使用 setTimeout 猜测下一个包的开始时间;
  • 上一个包还没结束,下一个包就重叠;
  • 队列一度降到 0,发生 underrun;
  • TTS 按句独立合成,句首句尾都带明显淡入淡出或静音。

先改为 nextPlayTime 连续排期,再观察欠载计数。

12.5 点了打断,旧声音还在说

只取消服务端不够。打印这几个状态:

activeSources.size
nextPlayTime - audioContext.currentTime
activeTurnId
server current_task.done()

如果队列深度仍大于 0,说明前端没有 stop() 已排期节点;如果旧包继续进入,说明缺少 turn_id 过滤。

12.6 句间停顿很长

每一句都重新请求 TTS 时,第二句也会经历一次首包等待。解决方向有三个:

  1. 第一段播放时提前合成第二段;
  2. 保持 100~400 ms 左右的待播缓冲;
  3. 如果供应商支持单连接增量文本,复用同一个合成会话。

预取不能无限做。用户一旦打断,深队列里的音频全部变成浪费。


十三、生产环境还要补哪些东西

这个 Demo 已经能说明核心时序,但上线前至少还要补:

13.1 真实 TTS 适配器

把厂商返回统一成:

AudioChunk(
    pcm=...,
    sample_rate=...,
    sentence_end=...,
)

同时记录厂商请求 ID、错误码和重试次数,不要只记一个总耗时。

13.2 队列水位回传

客户端定期回传:

{
  "type": "playback_state",
  "turn_id": 12,
  "queued_ms": 260,
  "underruns": 0
}

服务端根据 queued_ms 决定是否继续预取,避免两秒音频压在客户端。

13.3 超时与降级

建议分别设置:

  • TTS 建连超时;
  • TTS 首包超时;
  • 相邻音频包超时;
  • 整句合成超时。

首包超时和整句超时不是一回事。已经收到前几个包后突然断流,要能停止当前句、提示重试或切换备用 TTS。

13.4 端到端埋点

至少记录:

turn_id
sentence_index
text_chars
t_text_ready
t_tts_request
t_first_pcm
t_client_receive
t_play_schedule
queued_ms
underrun_count
cancel_requested
cancel_ack

下一篇我会把 ASR、LLM、TTS 和播放统一放到一条 t0-t8 时间线上,专门拆 Voice Agent 到底慢在哪。


十四、FAQ

Q1:流式 TTS 的首包多少才算合格?

没有脱离场景的统一数字。先把“供应商首包”“服务端首 PCM”“客户端首包”和“真正开始播放”分开,再看用户能否感知。电话、网页、App 的网络和播放缓冲完全不同。

Q2:分句越短,首包是不是越快?

通常第一句能更早发出,但过短会破坏韵律、增加请求次数,并让句间首包等待更频繁。优化目标不是句子最短,而是“第一段尽快开口,后续语音连续”。

Q3:PCM 包多大合适?

本文每包 40 ms,便于观察。包太小会增加 WebSocket 帧、调度和对象创建开销;包太大则增加等待和打断粒度。可以从 20~60 ms 起测。

Q4:为什么不直接播放 Blob?

Blob 更适合完整文件或浏览器已有解码路径。本文要读自定义包头、按 Int16 样本转换并精确排期,所以设置 binaryType = "arraybuffer"

Q5:为什么不用 AudioWorklet?

这个 Demo 先用 AudioBufferSourceNode 把队列和打断讲清楚。对于长时间网络流、环形缓冲、混音或更严格的低延迟控制,AudioWorklet 更适合,但代码和跨线程通信也更复杂。

Q6:服务端取消很快,为什么用户仍听到尾音?

可能是声卡/系统缓冲、浏览器已经排期的音频,或真实 TTS SDK 不支持立即中止。先确认前端 stop() 执行时间,再分别检查 SDK 和输出设备,不要只盯服务端日志。

Q7:能不能把 Base64 音频放进 JSON?

能,但 Base64 会增加体积,也会多一次编码和解码。控制消息用 JSON、音频用二进制帧,协议更清楚。


十五、这次实现的判断标准

这条链路是否可用,我会看五件事:

  1. 第一段文本是否在完整答案生成前进入 TTS;
  2. 客户端能否区分控制帧与 PCM 二进制帧;
  3. 音频能否沿 AudioContext 时间轴连续排期;
  4. 打断时能否同时停客户端队列、取消服务端任务;
  5. 日志能否区分分句耗时、TTS 首包和客户端首包。

如果这五项没有拆开,“流式 TTS”很容易只剩一个接口参数。


系列阅读

上一篇:

Voice Agent 的 RAG 为什么会答非所问?从混合检索到引用校验

下一篇计划:

Voice Agent 到底慢在哪?拆解 ASR、LLM、TTS 全链路延迟

会把 t0-t8 埋点、首字时间、首音频包、播放排期和用户实际等待统一起来。


官方资料

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐