You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Azure实时语音转写GStreamer报错:运行8分45秒后转录中断求助

解决Azure实时语音转写8分45秒后GStreamer流式终止问题

问题背景

使用GStreamer将流式音频转换为Azure支持的PCM格式,通过Python WebSocket部署实时语音转写服务时,转录可正常运行约8分45秒,随后抛出错误:

Azure Gstreamer error - basesrc gstbasesrc.c:3072: error: streaming stopped, reason error (-5)

转录突然终止,目标是实现Azure语音转写支持的最长240分钟连续转录。

调整方案

1. 恢复WebSocket连接心跳机制

当前代码禁用了WebSocket心跳(ping_interval=None),长时间无数据传输时,负载均衡或防火墙会主动断开长连接,导致音频流中断。修改WebSocket服务启动代码:

async def run_websocket_server():
    # 设置30秒心跳间隔,60秒超时,保持连接活跃
    start_server = websockets.serve(on_connect, "0.0.0.0", 8000, ping_interval=30, ping_timeout=60)
    await start_server

2. 替换同步队列为异步队列

代码中使用的queue.Queue是同步队列,在异步环境下会阻塞事件循环,导致音频流推送不及时,触发GStreamer报错。全部替换为asyncio.Queue:

  • 修改receive_audio中的队列初始化和操作:
async def receive_audio(uuid, path):
    # 替换为异步队列
    audio_queue = asyncio.Queue()

    try:
        conversation_transcriber, push_stream = create_conversation_transcriber(
            CONNECTIONS.connections[uuid]
        )
        logger.info("Conversation transcriber initialized - socket_service.py")
        conversation_transcriber.start_transcribing_async().get()

        while True:
            websocket = CONNECTIONS.connections[uuid]["websocket"]
            data = await websocket.recv()
            if data:
                await audio_queue.put(data)
                # 异步获取队列数据
                while not audio_queue.empty():
                    chunk = await audio_queue.get()
                    CONNECTIONS.connections[uuid]["audio_buffer"] += chunk
                    push_stream.write(chunk)
  • 修改CONNECTIONS中的队列类型:
CONNECTIONS.connections[UID] = {
    # ... 其他字段
    "transcribed_queue": asyncio.Queue(),
    "transcribing_queue": asyncio.Queue(),
    # ... 其他字段
}
  • 修改消息发送函数的队列操作:
async def send_transcribed_messages_async(uuid):
    try:
        while True:
            msg_queue = CONNECTIONS.connections[uuid]["transcribed_queue"]
            while not msg_queue.empty():
                # 异步获取消息
                msg = await msg_queue.get()
                websocket = CONNECTIONS.connections[uuid]["websocket"]
                if websocket.open:
                    await websocket.send(json.dumps(msg))
                    logger.info(f"Message sent to: {uuid}")
                else:
                    logger.warning("WebSocket is closed, unable to send message.")
                    break
            await asyncio.sleep(0.01)
    except Exception as e:
        logger.error(f"Error sending transcribed message: {e}")

3. 配置Azure Speech SDK长会话参数

默认情况下,Speech SDK的会话超时远小于240分钟,需手动调整超时参数并开启自动重连:
在generate_speech_config函数中添加以下配置:

def generate_speech_config(self):
    try:
        speech_config = speechsdk.SpeechConfig(
            subscription=self.SPEECH_KEY,
            region=self.SERVICE_REGION,
        )
        # ... 现有配置
        # 设置最大会话超时为240分钟(单位:毫秒)
        speech_config.set_property(
            speechsdk.PropertyId.SpeechServiceConnection_InitialSilenceTimeoutMs,
            "14400000"
        )
        speech_config.set_property(
            speechsdk.PropertyId.SpeechServiceConnection_EndSilenceTimeoutMs,
            "14400000"
        )
        speech_config.set_property(
            speechsdk.PropertyId.SpeechServiceConnection_ConnectionTimeoutMs,
            "14400000"
        )
        # 开启自动重连
        speech_config.set_property(
            speechsdk.PropertyId.SpeechServiceConnection_EnableAutoReconnect,
            "true"
        )
        # ... 现有配置
        return speech_config
    except Exception as e:
        print(f"Error generating speech config: {e}")

4. 明确指定音频流格式

当前使用AudioStreamContainerFormat.ANY自动检测格式,长时间运行中可能出现格式识别漂移,触发GStreamer解码错误。建议明确指定输入音频格式(例如MP3):

def create_conversation_transcriber(self):
    try:
        # 替换为实际输入格式,如MP3/WAV等
        audio_stream_format = speechsdk.audio.AudioStreamFormat(
            compressed_stream_format=speechsdk.audio.AudioStreamContainerFormat.MP3
        )
        push_stream = speechsdk.audio.PushAudioInputStream(
            stream_format=audio_stream_format
        )
        # ... 现有代码
    except Exception as e:
        print(f"Error in create_conversation_transcriber: {e}")

若无法提前确定格式,建议在客户端完成音频格式统一,或在服务端用GStreamer明确解码后再推送给Speech SDK。

5. 完善异常日志与资源清理

添加详细的错误日志排查终止原因,同时确保资源正确释放:

  • 修改canceled事件回调,打印错误详情:
def on_transcription_canceled(evt):
    logger.error(f"Transcription canceled: {evt.reason}")
    if evt.reason == speechsdk.CancellationReason.Error:
        logger.error(f"Error details: {evt.error_details}")

conversation_transcriber.canceled.connect(on_transcription_canceled)
  • 优化finally块中的资源清理逻辑:
finally:
    try:
        conversation_transcriber.stop_transcribing_async().get()
        push_stream.close()
    except Exception as e:
        logger.error(f"Error cleaning up resources: {e}")
    # 延迟删除目录,确保日志文件写入完成
    await asyncio.sleep(1)
    shutil.rmtree(f"results/{uuid}")

内容的提问来源于stack exchange,提问作者Googler Thiru

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.17 19:27:03