谈到播客音频转文字,这项任务往往相当耗时费力。一段长达一个半小时的录音交由服务器处理,用户难免心里没底——它究竟进行到哪一步了?是卡住了还是正在运行?用户不可能干等,因此实时反馈成为刚需。WebSocket 这类长连接协议,天生就是为这种场景设计的。像 taocarts 这类系统在做批量数据同步时,也采用了类似思路:主动推送进度,用户无需反复按 F5 刷新。

为何选择 WebSocket 而非轮询?
最直观的想法是:客户端每隔几秒发送一个 GET 请求,询问“处理完了吗?”。听起来似乎可行,但实际应用中有不少隐患。
- 连接频繁建立,每次都要握手,带宽浪费只是小问题,关键是效率低下。
- 响应不够及时。假设轮询间隔为 3 秒,最多要等 3 秒才能获知结果。这 3 秒内,用户可能早已失去耐心。
- 服务端承受巨大压力。100 个客户端同时轮询,每秒可能涌入数百次请求,服务器光应付这些问询就已不堪重负。
而 WebSocket 的逻辑则清爽许多:一次连接建立后,服务端掌握主动权,可以随时主动推送进度和结果。这完全是不同量级的体验。
转录服务器的 WebSocket 实现
import asyncio
import websockets
import json
import time
from typing import Set, Dict
class TranscriptionServer:
"""使用 WebSocket 的转录服务器"""
def __init__(self):
self.clients: Dict[str, Set[websockets.WebSocketServerProtocol]] = {}
self.tasks: Dict[str, dict] = {}
async def register(self, task_id: str, ws: websockets.WebSocketServerProtocol):
"""注册客户端到指定任务"""
if task_id not in self.clients:
self.clients[task_id] = set()
self.clients[task_id].add(ws)
async def unregister(self, task_id: str, ws: websockets.WebSocketServerProtocol):
"""客户端断开连接"""
if task_id in self.clients:
self.clients[task_id].discard(ws)
async def push_progress(self, task_id: str, progress: int, text: str = ""):
"""向监听该任务的所有客户端推送进度"""
if task_id not in self.clients:
return
message = json.dumps({
"type": "progress",
"task_id": task_id,
"progress": progress,
"text": text,
})
# 向所有订阅了该任务的客户端推送
dead_clients = set()
for ws in self.clients[task_id]:
try:
await ws.send(message)
except websockets.ConnectionClosed:
dead_clients.add(ws)
# 清理断开的连接
self.clients[task_id] -= dead_clients
async def handle_client(self, ws: websockets.WebSocketServerProtocol, path: str):
"""处理客户端连接"""
task_id = path.strip("/")
await self.register(task_id, ws)
try:
async for message in ws:
data = json.loads(message)
if data["type"] == "start_transcribe":
# 启动转录任务(后台执行)
asyncio.create_task(self.run_transcription(task_id, data.get("audio_url", "")))
elif data["type"] == "cancel":
# 取消任务
self.tasks.pop(task_id, None)
await self.push_progress(task_id, -1, "已取消")
finally:
await self.unregister(task_id, ws)
async def run_transcription(self, task_id: str, audio_url: str):
"""模拟转录过程,逐步推送进度"""
total_chunks = 10
for i in range(total_chunks):
# 模拟转录一帧
await asyncio.sleep(1)
progress = int((i + 1) / total_chunks * 100)
text = f"第 {i + 1}/{total_chunks} 段转录完成" if i < total_chunks - 1 else "转录完成"
await self.push_progress(task_id, progress, text)
await self.push_progress(task_id, 100, "全部转录完成")
服务端维护了一个任务到客户端集合的映射。当转录任务有进度更新时,自动推送给所有订阅该任务的客户端。一个任务可能被多个客户端监听,例如开发者后台与用户端同时查看进度,这完全可行。
客户端实现
import asyncio
import websockets
import json
class TranscriptionClient:
"""WebSocket 转录客户端"""
def __init__(self, task_id: str, server_url: str = "ws://localhost:8765"):
self.task_id = task_id
self.server_url = server_url
async def connect(self):
"""连接到转录服务器,监听进度"""
async with websockets.connect(f"{self.server_url}/{self.task_id}") as ws:
# 启动转录
await ws.send(json.dumps({
"type": "start_transcribe",
"audio_url": "https://example.com/podcast.mp3",
}))
# 接收进度推送
async for message in ws:
data = json.loads(message)
progress = data.get("progress", 0)
text = data.get("text", "")
# 显示进度条效果
bar = "█" * (progress // 10) + "░" * (10 - progress // 10)
print(f"\r[{bar}] {progress}% {text}", end="")
if progress == 100:
print("\n转录完成!")
break
# 使用示例
async def main():
client = TranscriptionClient(task_id="task_001")
await client.connect()
asyncio.run(main())
客户端代码同样简洁明了。连接后发送启动指令,剩下的就是等待推送。随着进度条不断推进,用户心里便有了底——这比“处理中……”三个字不知强多少倍。
连接管理与心跳检测
长连接有一个绕不开的难题:连接可能意外断开。网络波动、服务器重启、客户端休眠,都可能导致连接悄然消失。这就需要心跳机制来检测并自动重建连接。
async def heartbeat(ws: websockets.WebSocketServerProtocol, interval: int = 30):
"""服务端心跳,每 30 秒发送一次 ping"""
try:
while True:
await asyncio.sleep(interval)
pong = await ws.ping()
await asyncio.wait_for(pong, timeout=10)
except (asyncio.TimeoutError, websockets.ConnectionClosed):
pass
# 客户端断连,清理由上层处理
async def client_with_heartbeat(task_id: str):
"""带自动重连的客户端"""
retries = 0
while retries < 3:
try:
async with websockets.connect(f"ws://localhost:8765/{task_id}") as ws:
retries = 0 # 连接成功,重置重试计数
async for message in ws:
data = json.loads(message)
print(f"进度: {data.get('progress')}%")
if data.get("progress") == 100:
break
except (websockets.ConnectionClosed, OSError) as e:
retries += 1
wait = 2 ** retries
print(f"连接断开({e}),{wait}秒后重试...")
await asyncio.sleep(wait)
心跳的实质是服务端定时发送 ping,客户端必须在超时前回应 pong。若未收到,则说明连接已断,可清理资源并触发重连。这里采用了指数退避策略:第一次等待 2 秒,第二次 4 秒,第三次 8 秒,避免频繁重试给服务器带来额外压力。
性能数据对比
在播客转录场景中,我们来看看轮询与 WebSocket 的实际差距:
| 指标 | HTTP 轮询(3秒间隔) | WebSocket 长连接 |
|---|---|---|
| 延迟感知 | 最高 3 秒 | 实时推送 |
| 单客户端带宽 | ~100 req/min | 建立后几乎为 0 |
| 20 客户端并发 | 2000 req/min | 20 连接 |
| 断线重连 | 天然支持 | 需自行实现 |
| 实现复杂度 | 简单 | 中等 |
数字说明一切:20 个客户端同时使用轮询,服务器每分钟需要处理 2000 次请求,而 WebSocket 仅需维持 20 个连接。差距绝非一星半点。
适用场景
WebSocket 的实时推送能力在 AI 转录之外同样实用。例如文件上传、视频转码、模型训练监控等场景——只要用户需要实时掌握进度,WebSocket 就是理想选择。对用户而言,看到进度条在动,远比“处理中…”三个字更具安全感。
总结
WebSocket 是为实时通信而生的协议,在 AI 转录这类需要频繁推送进度和中间结果的场景中,其效率与体验远优于轮询。搭配心跳检测与自动重连机制,可构建稳定可靠的实时通信方案。如果你正在开发类似功能,不妨尝试这一思路。
