← 返回笔记

learning note

设计实时 ASR 链路,实现分段增量去重

梳理从前端音频分片、WebSocket 会话、讯飞 AST 推流,到稳定文本快照回推的完整实时转写链路。

背景:

在回答AI面试官的问题的时候,我们希望可以实现语音转写的功能。相关的演示视频在什么是实时ASR?

单纯的从接口角度来看,貌似我们只需要让前端发送音频文件,后端接受音频文件调用接口,之后再把结果推送给前端。而我们对ASR的调用并不是简单的调用接口,而是为了实现下面五个目标:

  • **实时性:**音频要边上传边转写,结果要边识别边回推,不能等用户说完才出完整文本。

  • **稳定性:**不能因为前端上传抖动、下游返回修正包或者连接关闭时机不完美,就让前端看到重复文本、前缀丢失或最终结果消失。

  • **解耦性:**WebSocket 接收音频、推流给 AST、处理 AST 回包、回推前端结果,这几个动作不应该死绑在一个线程里,也不应该互相污染职责。

  • **可消费性:**前端收到的最好不是一个模糊字符串,而是一份可以区分“稳定部分”和“活动部分”的结构化快照。

实时转写链路:

当前项目里,围绕这些问题形成了一套比较务实的实时转写实现:

前端页面 -> AudioTranscriptionWebSocketHandler -> TranscriptionSessionContext -> XunfeiAudioService.realTimeAudioToText -> 讯飞 AST WebSocket -> AstTranscriptionAssembler -> RealtimeTranscriptionUpdate -> AudioTranscriptionWebSocketHandler -> 前端页面

也就是说:

  • AudioTranscriptionWebSocketHandler 负责接住前端连接、控制消息和音频分片

  • TranscriptionSessionContext 负责保存会话级 Pipe 缓冲和最近快照

  • XunfeiAudioService 负责把音频按实时节奏推给讯飞 AST

  • AstTranscriptionAssembler 负责把供应商增量包装配成稳定文本快照

  • Handler 再把快照包装成前端协议里的 transcription / final

对应的核心代码主要在两个类:

  • admin/src/main/java/com/hewei/hzyjy/xunzhi/media/infrastructure/websocket/AudioTranscriptionWebSocketHandler.java

  • admin/src/main/java/com/hewei/hzyjy/xunzhi/media/infrastructure/integration/XunfeiAudioService.java

先看标准时序

一条标准实时转写链路,实际上有两条常见分支。

分支一:自然收口

这是最理想、也最符合“完整转写完成”语义的一条链路:

  1. 前端建立 WebSocket 连接

  2. 服务端返回 connected

  3. 前端发送 { "type": "start_transcription" }

  4. 服务端创建转写会话并返回 transcription_started

  5. 前端持续发送音频二进制分片

  6. 服务端边收音频、边推给讯飞、边回推 transcription

  7. 下游 AST 任务自然完成,服务端统一收口并返回 final

这种场景通常对应:

  • 用户正常说完了这一轮内容

  • 下游识别任务也顺利走到了最终完成态

  • 后端拿到了完整最终文本,并构造 final

分支二:用户主动停止

另一条常见链路是,用户还没等会话自然收口,就主动点击了“停止录音”:

  1. 前端建立 WebSocket 连接

  2. 服务端返回 connected

  3. 前端发送 { "type": "start_transcription" }

  4. 服务端创建转写会话并返回 transcription_started

  5. 前端持续发送音频二进制分片

  6. 服务端边收音频、边推给讯飞、边回推若干条 transcription

  7. 用户在前端点击停止录音,前端发送 { "type": "stop_transcription" }

  8. 服务端将 stopRequested = true,关闭当前会话的音频输出流,并返回 transcription_stopped

  9. 当前代码实现下,这次会话后续通常不会再回 final

这条链路的含义不是“请帮我整理最终结果”,而是:

  • 这次会话现在立刻停掉

  • 不再继续接收和推送后续音频

  • 不再把 final 当成协议上的必然返回

为什么要区分final和stop?

其实从整体场景来看,final 表示这次转写任务自然收口了,后端拿到了最终文本;stop_transcription 表示用户主动要求立刻结束这次会话。

当用户点击终止录音,发送stop_transcription的时候,后端就不应该再慢悠悠的等待ASR服务商的响应之后再收尾,而是立刻关 Pipe、停推流、尽快清理垃圾内容。

一句话理解

  • final = 识别系统自己说:我这次已经完整做完了。

  • stop = 用户说:现在就停,别继续了。

因此无论是从语义区分上,还是代码逻辑区分上,我们都应该区分final和stop所以其实从一整个链路上来看,其实这个服务并不是简单的前端传输数据,后端请求ASR服务。

分清楚两条通路:

前端 -> 后端:产品侧 WebSocket 通道

这条通道的入口是:

  • AudioTranscriptionWebSocketHandler

它负责处理两类输入:

  • 文本控制消息,例如 start_transcriptionstop_transcriptionping

  • 二进制音频消息,也就是前端不断发来的 PCM 音频分片

它也负责往前端回推:

  • connected

  • transcription_started

  • transcription

  • final

  • transcription_stopped

  • heartbeat

  • error

后端 -> 讯飞:供应商 ASR 通道,这也是一个websocket链路

这条通道由:

  • XunfeiAudioService.realTimeAudioToText(...)

负责建立和维护。

它的职责是:

  • 读取后端会话缓冲中的音频字节

  • 按讯飞 AST 要求的节奏切块发送

  • 接收讯飞实时回包

  • 把供应商原始回包转成后端自己的结构化更新对象

所以整个系统并不是“前端直接把音频发给讯飞”,而是:

  • 前端把音频发给后端

  • 后端自己掌控推流节奏、会话状态、异常收尾和结果适配

全局有两条websocket,这一点一定要注意,不然看代码容易懵逼

  • 前端这条 WebSocket 负责:建连、start_transcription、stop_transcription、收音频帧、把结构化转写结果回推给前端:admin/src/main/java/com/hewei/hzyjy/xunzhi/media/infrastructure/websocket/AudioTranscriptionWebSocketHandler.java (line 69)

  • 后端到讯飞这条 WebSocket 负责:读取 Pipe 里的音频流、按节奏发送给讯飞、接收供应商回包并交给 AstTranscriptionAssembler:admin/src/main/java/com/hewei/hzyjy/xunzhi/media/infrastructure/integration/XunfeiAudioService.java (line 190)

中间怎么接起来

  • 前端发来的音频不会直接透传给讯飞,而是先进入后端会话上下文里的 Pipe,再由后端第二条 WebSocket 推给讯飞

  • 所以这不是“一条 WebSocket 直通到底”,而是:上游产品通道 + 下游供应商通道

从请求维度看服务

第一步:前端建立 WebSocket 连接

入口代码:

  • AudioTranscriptionWebSocketHandler.onOpen(...)

这一阶段后端会做几件事:

  • 校验当前 WebSocket 连接是否合法用户

  • 记录 userId -> sessionsessionId -> userId 的映射

  • 给前端回一条 connected

  • 启动心跳任务

这一步的意义是:

  • 建立“用户身份”和“WebSocket 会话”的绑定关系

  • 给后续的文本控制消息和二进制音频消息提供会话载体

****

第二步:前端发送 start_transcription

文本控制消息入口是:

  • AudioTranscriptionWebSocketHandler.onMessage(Session, String)

后端收到文本消息后,会先反序列化成 WebSocketMessage,再进入:

  • handleControlMessage(...)

type = start_transcription 时,会走到:

  • startTranscriptionSession(...)

这一步不是简单改个布尔值,而是会真正启动一条新的转写会话。

第三步:后端创建会话级转写上下文

真正创建会话运行时对象的地方是:

  • createAndStartTranscriptionSession(...)

这里会创建一套和当前 sessionId 绑定的会话资源:

  • PipedInputStream audioInputStream

  • PipedOutputStream audioOutputStream

  • AtomicBoolean active

  • AtomicBoolean stopRequested

  • AtomicReference<RealtimeTranscriptionUpdate> lastUpdate

这些资源会被装进:

你可以把它理解成:

  • 当前 WebSocket 会话对应的一条“后端内部实时转写管道”

其中最关键的是那对 Pipe:

  • audioOutputStream 是上游写端,专门给 WebSocket 音频入口写数据

  • audioInputStream 是下游读端,专门给 XunfeiAudioService 读数据

第四步:后端启动 XunfeiAudioService

会话上下文建好后,后端立即调用:

  • xunfeiAudioService.realTimeAudioToText(audioInputStream, callback)

这一步非常关键,因为它说明:

  • XunfeiAudioService 并不直接拿前端传来的 ByteBuffer

  • 它消费的是 TranscriptionSessionContext 里 Pipe 读端 audioInputStream

也就是说,从这一刻开始,后端内部形成了一个生产者-消费者模型:

  • 生产者:AudioTranscriptionWebSocketHandler.onMessage(Session, ByteBuffer)

  • 消费者:XunfeiAudioService.sendAudioStream(...)

  • 中间缓冲:PipedOutputStream -> PipedInputStream

第五步:前端持续发送二进制音频分片

二进制入口是:

  • AudioTranscriptionWebSocketHandler.onMessage(Session, ByteBuffer)

后端在这里会按顺序做三件事:

  1. ByteBuffer 转成 byte[]

  2. sessionId 找到对应的 TranscriptionSessionContext

  3. 把音频写入 context.audioOutputStream

所以这一刻的数据流是:

前端 ByteBuffer -> byte[] -> context.audioOutputStream

注意,这里还没有直接发给讯飞。它只是把当前这段音频写进了“当前会话专属的后端内部缓冲通道”。

为什么这里一定要先找到 TranscriptionSessionContext

  • 因为后端要知道这段音频属于哪一个 WebSocket 会话

  • 要知道这条转写是否已经通过 start_transcription 正式启动

  • 要拿到这条会话专属的 audioOutputStream

  • 要避免把音频写到错误的转写链路里

第六步:XunfeiAudioService从 Pipe 读端持续消费音频

XunfeiAudioService.realTimeAudioToText(...) 会先准备好讯飞 AST 的鉴权参数和 WebSocket URL,然后调用:

  • WS_CLIENT.newWebSocket(...)

当供应商 WebSocket 打开后,会异步启动:

  • sendAudioStream(webSocket, audioInputStream, sessionId, future)

这个方法就是后端真正把音频送给讯飞的地方。它会持续执行:

  • in.read(buffer)audioInputStream 读取音频

  • 每次按 CHUNK_SIZE_BYTES = 1280 的缓冲发送

  • 每发一块 sleep(40ms),按实时节奏推流

  • 读到流结束后,再发一条 {"end":true,"sessionId":"..."} 告诉讯飞音频结束

也就是说,前端音频并不是“来一块就立刻原样透传一块”,而是先进入后端 Pipe,再由后端自己的推流线程按统一节奏发送给讯飞

第七步:讯飞 AST 回实时增量包

供应商回包入口是:

  • WebSocketListener.onMessage(WebSocket, String)

后端在这里会先解析原始 JSON,然后提取出本轮增量需要的关键字段:

  • partialText

  • segmentId

  • pgs

  • rg

  • bg

  • ed

  • finalPacket

其中:

  • segmentId 表示当前分段标识

  • pgs 表示当前增量的修正模式,例如追加或替换

  • rg 表示替换范围

  • bg/ed 表示音频时间范围

  • finalPacket 表示这是否已经是供应商视角下的最终包

这些字段不是我们规定的,而是讯飞ASR服务返回的字段,对应的接口文档请看:讯飞ASR开放平台接口文档

第八步:AstTranscriptionAssembler 组装供应商增量

讯飞返回的不是一个永远单调追加的最终字符串,而是可能不断修正、替换、重算前面内容的增量包。

所以后端不会把 partialText 直接原样回给前端,而是先交给:

它内部维护:

  • TreeMap<Integer, SegmentState> segments

本质上就是一个有序分段池。

它会根据不同情况应用不同的合并策略:

  • 如果 pgs = rpl,按 rg 指定范围删除旧段,再插入新段

  • 如果 pgs = apd,按追加语义更新当前段

  • 如果没有 pgs,则用 bg/ed、文本相似性和重叠关系去判断这是不是同一句演化

这一步的目标不是“保留供应商原始包”,而是:

  • 在后端得到一份尽量稳定、尽量接近用户可见文本的快照

第九步:后端生成自己的 RealtimeTranscriptionUpdate

当组装器得到新的快照后,后端会构造一个标准化更新对象:

  • RealtimeTranscriptionUpdate

它包含的核心字段有:

  • fullText

  • committedText

  • liveText

  • displayText

  • revision

  • resultStatus

  • segmentId

  • segmentText

  • pgs

  • rg

  • bg

  • ed

  • finalPacket

其中几个字段的语义很关键:

  • fullText:当前全量文本快照

  • committedText:已经相对稳定的文本

  • liveText:仍在滚动修正的尾部文本

  • displayText:前端推荐直接展示的文本,通常等于 committedText + liveText

  • revision:后端维护的递增版本号

  • resultStatus:当前是 partial 还是 final

主要还是因为ASR包在不断的修正自己的结果。需要AstTranscriptionAssembler来做不断的修正汇总。这一步其实是在把供应商原始增量协议翻译成本系统稳定的产品语义协议(这句话听起来有点装逼。但其实意思就是把原始的ASR结果类 变成更符合我们前端展示需求的结果)

第十步:Handler 把 update 回推给前端

XunfeiAudioService 每得到一条新的 RealtimeTranscriptionUpdate,就会通过 callback 回调给 AudioTranscriptionWebSocketHandler

Handler 在 callback 里会做两件事:

  1. context.lastUpdate.set(update),保存最近一条实时快照

  2. sendMessage(session, createResponse("transcription", "Partial snapshot", update, true))

于是前端收到的就是一条后端自己封装过的 transcription 消息。这说明当前后端不是单纯透传供应商包,而是在做三层适配:

  • 供应商回包解析

  • 分段合并与快照组装

  • WebSocket 协议二次封装

第十一步:自然完成时构造 final

realTimeAudioToText(...) 返回的是一个 CompletableFuture<String>这个 future 在下面几种情况下会完成:

  • 收到供应商最终包,组装出最终文本

  • 供应商 WebSocket 正常关闭时,根据当前快照补出最终文本

Handler 在 future.whenComplete(...) 里统一收尾。如果满足下面条件:

  • 没有异常,或者异常属于主动停止导致的预期异常

  • 当前不是用户主动 stop_transcription

  • finalResult != null

后端就会发送:

  • createResponse("final", "Transcription completed", buildFinalUpdate(finalResult, context.lastUpdate.get()), true)

这里的关键点是:

  • finalResult 决定最终文本内容

  • lastUpdate 决定最终包继承哪些最后一帧的结构化元信息

也就是说,final 不是一个只有字符串的裸结果,而是会尽量继承:

  • revision

  • segmentId

  • segmentText

  • pgs

  • rg

  • bg

  • ed

同时后端还会把:

  • resultStatus 固定成 final

  • liveText 置空

  • finalPacket 置为 true

第十二步:收尾和资源清理

无论是自然完成、主动停止、连接关闭还是异常退出,最终都会走到一类清理动作:

  • 把 context 从 TRANSCRIPTION_CONTEXTS 里移除

  • active 设为 false

  • 关闭 audioOutputStream

  • 关闭 audioInputStream

  • 取消心跳任务

其中主动停止的特殊点在于:

  • stopTranscriptionSession(...) 会先把 stopRequested 设为 true

  • 然后关闭 audioOutputStream

这样做的效果是:

  • Pipe 写端关闭

  • XunfeiAudioService 读端最终读到流结束

  • 后端发出结束帧给讯飞

  • 整条链路有机会自然收口

但是当前实现里,如果是用户主动停止,后端不会再额外给前端发一条自动补出来的 final;它只会发:

  • transcription_stopped

用一段话记住整条链路

这条链路的本质是:

前端把音频流通过 WebSocket 持续送入 AudioTranscriptionWebSocketHandler,Handler 借助 TranscriptionSessionContext 做会话级缓冲与状态隔离,XunfeiAudioService 负责按实时节奏把音频推给讯飞 AST,并将返回的增量包交给 AstTranscriptionAssembler 做分段去重和有序重建,最终再以结构化快照和 final 结果的形式回推给前端。

如果你只记一件事,那就是:

  • 这不是一段音频进去,一段文本出来的同步接口

  • 它本质上是一个WebSocket 会话 + Pipe 缓冲 + 实时推流 + 分段装配 + 结构化快照回推的完整实时系统

ASR流程时序图

1.总体架构图

2.WebSocket 建连与启动转写时序图

3.实时音频流转写时序图

4.停止转写与资源清理时序图