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
相关产品推荐
相关产品推荐

