Fun-ASR-Realtime 这类实时识别服务,识别端一般只维持当前连接的会话状态,不会替你记住“上次念到哪一句”。所以重连后从哪一段音频续上,答案在推流侧:先看本地环形缓存里有没有那段时间的音频,再看日志里记的“最后确认时刻”落在缓存的哪个位置。缓存没留、日志没记,重连就只能重推整段,或者接受一段缺口。
断线续传的起点由推流侧决定,而不是识别端。建议推流进程维护一个带时间戳的音频环形缓存,断开时记下最后发送时刻与最后返回时刻,重连后以“最后返回时刻往前多取一小段”作为续传区间,再由文本侧按时间区间或相似度去重。适用场景是分钟级以上的会议转写;边界在于缓存时长受内存限制,往前多取必然带来重复文本,去重规则没定好,记录里就会出现两遍同一句话。
在推流侧维护一个带时间戳的音频环形缓存
缓存容量可以直接按音频格式算出来:容量字节 = 采样率 × 每采样字节数 × 声道数 × 缓存秒数。16 kHz、16 bit、单声道的会议音频大约 32 KB/s,缓存上限设 900 秒时,常驻内存约 27 MB。如果一台机器上同时跑多路转写,要按路数乘一遍再决定缓存上限。
import time
from collections import deque
SAMPLE_RATE = 16000
BYTES_PER_SAMPLE = 2
CHANNELS = 1
BPS = SAMPLE_RATE * BYTES_PER_SAMPLE * CHANNELS # 约 32000 字节/秒
MAX_SECONDS = 900 # 缓存时长上限
class AudioRing:
def __init__(self, max_seconds=MAX_SECONDS):
self.q = deque() # 元素为 (帧起始时间戳秒, 音频字节)
self.size = 0
self.max_bytes = max_seconds * BPS
def append(self, chunk, ts=None):
ts = time.time() if ts is None else ts
self.q.append((ts, chunk))
self.size += len(chunk)
while self.size > self.max_bytes and len(self.q) > 1:
_, old = self.q.popleft()
self.size -= len(old)
def slice_from(self, start_ts, end_ts=None):
out = bytearray()
for ts, chunk in self.q:
if ts >= start_ts and (end_ts is None or ts <= end_ts):
out.extend(chunk)
return bytes(out)
append 进来的 chunk 最好是对齐的固定时长帧,比如 20 ms 或 100 ms。如果帧长不固定,按“帧起始时间戳”切片会在边界上多切或少切几个采样点,续传接缝处可能出现极短的杂音,识别结果里表现为一个莫名其妙的字。
缓存时长上限是一笔取舍。设得太短,客户端休眠、网络切换这类几十秒级的断线就接不上;设得太长,内存占用和音频在内存中的停留时间都会上去,涉及会议的音频还多一层隐私考虑。可以先取 5 到 15 分钟,跑一段时间后看日志里实际出现过的最长断线时间,再往上或往下调。
捕获连接断开事件并记录最后确认时刻
断开事件的监听点通常放两处。一处是发送侧的异常与关闭回调,例如 WebSocket 的 onclose/onerror 或者识别 SDK 自带的 disconnect 回调;另一处是自己的心跳超时定时器。只靠 onclose 有时会晚很多,对端静默断开、链路被中间设备清掉时甚至不触发,心跳超时是兜底。
断开时至少要往日志里落这几个字段,缺一个都会让续传判断变得含糊:
last_send_ts:最后一帧成功写入连接的本地时刻,代表“我推出去到哪了”。last_result_ts:最后收到识别结果的时刻,代表“对端确认处理到哪了”。ring_base_ts:环形缓存里最早一帧的时间戳,用来判断续传点是否已被覆盖。last_seq或已发送字节数:和识别端返回的序号对账时用得上。reason:断开原因分类,区分本地异常、对端关闭、心跳超时。
这两个时刻通常不相等:发送成功不等于被识别。续传起点一般取 last_result_ts,因为从那之后的内容对端大概率没有产出文本;用 last_send_ts 会丢掉一段已经推出去但还没出结果的音频。
重连后按最后确认时刻从缓存取后续音频
取区间的代码骨架很短,关键参数只有一个“往前多取多少”。
OVERLAP_SEC = 0.5
ring_earliest_ts = ring.q[0][0] if ring.q else None
resume_ts = last_result_ts - OVERLAP_SEC
if ring_earliest_ts is not None:
resume_ts = max(resume_ts, ring_earliest_ts)
audio = ring.slice_from(resume_ts)
if not audio:
# 缓存已被覆盖,只能从当前时刻重新开始,并在记录里显式标注缺口
log.warning('ring exhausted, session=%s', session_id)
else:
stream.send(audio)
log.info('resume from %.3f, bytes=%d', resume_ts, len(audio))
往前多取一小段(可以先按 200 到 500 ms 试)的理由是:识别结果的返回本身就带延迟,日志里 last_result_ts 记录的是“你收到结果的时刻”,而那段结果对应音频的尾部往往比它更晚。不多取的话,接缝处容易掉半句话,表现为重连后第一句缺几个字。
代价同样明确:多取的这段会被二次提交,一定会产生重复文本,必须配合下一节的去重;overlap 设得越大,去重压力越大,重连后第一批结果的延迟也越明显。如果 resume_ts 早于 ring_base_ts,说明断线时间已经超过缓存上限,此时应当明确记一条缺口日志,而不是悄悄从当前时刻开始,否则事后排查会找不到丢失区间。
对重叠区间产生的重复文本做去重
如果识别结果里带着时间区间,优先用时间区间判重,规则简单且不容易误杀:
- 每条结果记录
[start_ms, end_ms],落库前查这个区间与已落库片段的重叠比例。 - 重叠比例超过阈值(例如 50%)时判定为重复片段,只保留一条。
- 冲突时保留哪条:优先保留时间区间更长、以句末标点收尾、文本更完整的那条;两者几乎等价时保留先落库的那条并丢弃新到的,保证重复执行时结果一样。
from difflib import SequenceMatcher
def similar(a, b, thr=0.85):
return SequenceMatcher(None, a, b).ratio() >= thr
def pick(old, new):
if (new['end_ms'] - new['start_ms']) > (old['end_ms'] - old['start_ms']):
return new
return old
识别端不返回时间戳时只能退到文本相似度兜底。阈值不要设得太低,会议里“好,那我们看下一个议题”这类短句本来就会反复出现,纯文本判重容易把正常内容删掉。这类场景可以加一层限制:只对重连后前 N 条结果做相似度判重,因为重复只可能出现在接缝附近,而不是整场。
用断网开关做一次可复现的断线演练
重连逻辑写完不演练,等于没写。建议用可快速撤销的方式制造断网,别直接拔网线:
- 起一路正常转写,确认结果在持续落库。
- 制造断线:
sudo iptables -A OUTPUT -p tcp `--dport` 443 -j DROP,演练结束后立刻用sudo iptables -D OUTPUT -p tcp `--dport` 443 -j DROP撤销。端口按实际服务端口替换。 - 断线期间持续说话,念一段带编号的短句,方便事后逐句核对。
- 恢复网络,看重连日志里的续传起点和发送字节数。
- 把本地录音与落库文本逐句对比,检查接缝处是否缺字、是否出现重复句。
断线时长要分多档测,而不是只测一次。1 秒、3 秒、10 秒、30 秒,再加一档故意超过缓存上限(比如缓存 900 秒时设成远大于它的时长不现实,可以临时把 MAX_SECONDS 调小到 10 秒再断 30 秒),目的是确认两件事:当前配置下还能续上的断线上限大致在哪,超过之后是明确记缺口日志还是静默丢内容。
| 断线时长 | 是否自动重连 | 重连耗时 | 续传起点相对断点的偏移 | 是否缺句 | 是否重复 | 备注 |
| 1s | ||||||
| 3s | ||||||
| 10s | ||||||
| 30s | ||||||
| 超过缓存上限一档 |
记录表填完之后再回去看第 1 节的缓存时长和第 3 节的 overlap 取值:如果多档都出现重复,overlap 可以调小;如果小断线也缺句,overlap 需要调大,或者检查 last_result_ts 是不是记在了错误的时机。