FastAPI WebSocket实时ASR(Azure Speech SDK)AsyncIO异常排查
问题:Azure Speech SDK实时ASR的FastAPI WebSocket服务后续请求异常
首次启动服务后运行正常,但后续请求会抛出AsyncIO异常,仅能写入音频流却无法持续进行识别处理。
相关代码
# Create a speech config speech_config = speechsdk.SpeechConfig(subscription=AZURE_SPEECH_KEY, region=AZURE_SPEECH_REGION) # auto_detect_source_language_config = \ # speechsdk.languageconfig.AutoDetectSourceLanguageConfig(languages=["en-US", "ae-AR"]) speech_config.speech_recognition_language = "en-US" # Set the language audio_format = AudioStreamFormat( # channels=1, samples_per_second=16000, bits_per_sample=16 ) stream = speechsdk.audio.PushAudioInputStream(audio_format) speech_recognizer = speechsdk.SpeechRecognizer(speech_config=speech_config, audio_config=speechsdk.audio.AudioConfig(stream=stream)) speech_config = speechsdk.SpeechConfig(subscription=AZURE_SPEECH_KEY, region=AZURE_SPEECH_REGION) speech_config.speech_recognition_language = "en-US" # Set the language audio_format = AudioStreamFormat(samples_per_second=16000, bits_per_sample=16) stream = speechsdk.audio.PushAudioInputStream(audio_format) speech_recognizer = speechsdk.SpeechRecognizer(speech_config=speech_config, audio_config=speechsdk.audio.AudioConfig(stream=stream)) queue = [] done=asyncio.Event() stream_lock = asyncio.Lock() recognizing_alive,recognized_alive=datetime.datetime.now(),datetime.datetime.now() def recognized_callback(evt): recognized_text = evt.result.text queue.append(recognized_text) global recognized_alive recognized_alive = datetime.datetime.now() print("Recognized ALive" + str(recognized_alive)) def recognizing_callback(evt): global recognizing_alive recognizing_alive=datetime.datetime.now() print("Recognizing ALive"+str(recognizing_alive)) speech_recognizer.recognized.connect(recognized_callback) speech_recognizer.recognizing.connect(recognizing_callback) async def validate_stream_processed(): while True: global recognizing_alive global recognized_alive diff=(datetime.datetime.now()-recognized_alive).total_seconds() print("The diff is "+str(diff)) if diff>10: print("I am breaking of validate stream process") done.set() break await asyncio.sleep(1) async def read_queue(websocket): print("About to read") while True: await asyncio.sleep(1) if queue: async with stream_lock: if queue: item = queue.pop(0) print("About to send: ", item) await websocket.send_text(item) if done.is_set(): break async def write_stream(stream,data): # The number of bytes to push per buffer n_bytes = 4096 data = io.BytesIO(data) try: while True: frames = data.read(n_bytes // 2) print('read {} bytes'.format(len(frames))) if not frames: print("I am breaking from read stream") break async with stream_lock: stream.write(frames) print("Wrote one") print("Entring Wait in write") await asyncio.sleep(0) print("Cameout in write") except Exception as e: print(f"An error occurred: {e}") async def write_stream_and_process(speech_recognizer, stream, data): await write_stream(stream, data) # After writing to the stream, start continuous recognition @app.websocket("/asr/en") async def root(websocket: WebSocket): global speech_recognizer global stream await websocket.accept() try: data = await websocket.receive_bytes() # speech_recognizer.stop_continuous_recognition() speech_recognizer.start_continuous_recognition() task1 = asyncio.create_task(write_stream_and_process(speech_recognizer, stream, data)) task2 = asyncio.create_task(read_queue(websocket)) task3 = asyncio.create_task(validate_stream_processed()) await asyncio.gather(task1, task2, task3) except asyncio.CancelledError: print("WebSocket task cancelled") except Exception as e: print(f"An error occurred: {e}") finally: global queue print("Closing WebSocket Connection") speech_recognizer.stop_continuous_recognition() print("Stopped Recognition") stream.close() print("Clearing Queue") queue.clear() print("Items in Queue at Finish", queue) print("Stream closed") await websocket.close()
异常信息
Exception in callback StreamReaderProtocol.connection_made.<locals>.callback(<Task cancell...erver.py:102>>) at /opt/homebrew/Cellar/python@3.11/3.11.7/Frameworks/Python.framework/Versions/3.11/lib/python3.11/asyncio/streams.py:248 handle: <Handle StreamReaderProtocol.connection_made.<locals>.callback(<Task cancell...erver.py:102>>) at /opt/homebrew/Cellar/python@3.11/3.11.7/Frameworks/Python.framework/Versions/3.11/lib/python3.11/asyncio/streams.py:248> Traceback (most recent call last): File "/opt/homebrew/Cellar/python@3.11/3.11.7/Frameworks/Python.framework/Versions/3.11/lib/python3.11/asyncio/events.py", line 80, in _run self._context.run(self._callback, *self._args) File "/opt/homebrew/Cellar/python@3.11/3.11.7/Frameworks/Python.framework/Versions/3.11/lib/python3.11/asyncio/streams.py", line 249, in callback exc = task.exception() ^^^^^^^^^^^^^^^^ File "/opt/homebrew/lib/python3.11/site-packages/uvicorn/server.py", line 105, in handler await handle_http( File "/opt/homebrew/lib/python3.11/site-packages/uvicorn/_handlers/http.py", line 80, in handle_http data = await reader_read ^^^^^^^^^^^^^^^^^ File "/opt/homebrew/Cellar/python@3.11/3.11.7/Frameworks/Python.framework/Versions/3.11/lib/python3.11/asyncio/streams.py", line 708, in read await self._wait_for_data('read') File "/opt/homebrew/Cellar/python@3.11/3.11.7/Frameworks/Python.framework/Versions/3.11/lib/python3.11/asyncio/streams.py", line 540, in _wait_for_data await self._waiter asyncio.exceptions.CancelledError
问题原因及修复方案
核心问题
- 全局资源复用冲突:
speech_recognizer和stream是全局变量,首次请求关闭后资源已失效,后续请求复用直接导致异步IO异常。 - 初始化冗余无效:代码开头重复创建两次
speech_config、stream和speech_recognizer,属于无效冗余代码。 - 识别启动时机错误:
write_stream_and_process中注释了启动识别的代码,且首次请求中启动识别的时机在写入音频流之前,易导致部分音频丢失。 - 全局状态未隔离:
done事件、queue和时间变量为全局变量,未在每次请求时重置,导致后续请求逻辑混乱。
修复后的代码
import asyncio import io import datetime import speechsdk from speechsdk.audio import AudioStreamFormat, PushAudioInputStream from fastapi import FastAPI, WebSocket AZURE_SPEECH_KEY = "你的Azure密钥" AZURE_SPEECH_REGION = "你的Azure区域" app = FastAPI() @app.websocket("/asr/en") async def asr_websocket(websocket: WebSocket): await websocket.accept() # 每个请求独立创建Speech相关资源,避免全局复用冲突 speech_config = speechsdk.SpeechConfig(subscription=AZURE_SPEECH_KEY, region=AZURE_SPEECH_REGION) speech_config.speech_recognition_language = "en-US" audio_format = AudioStreamFormat(samples_per_second=16000, bits_per_sample=16) stream = PushAudioInputStream(audio_format) speech_recognizer = speechsdk.SpeechRecognizer( speech_config=speech_config, audio_config=speechsdk.audio.AudioConfig(stream=stream) ) # 每个请求独立维护状态,避免全局变量污染 queue = [] done = asyncio.Event() stream_lock = asyncio.Lock() recognizing_alive = recognized_alive = datetime.datetime.now() def recognized_callback(evt): nonlocal recognized_alive recognized_text = evt.result.text queue.append(recognized_text) recognized_alive = datetime.datetime.now() print(f"Recognized ALive {recognized_alive}") def recognizing_callback(evt): nonlocal recognizing_alive recognizing_alive = datetime.datetime.now() print(f"Recognizing ALive {recognizing_alive}") speech_recognizer.recognized.connect(recognized_callback) speech_recognizer.recognizing.connect(recognizing_callback) async def validate_stream_processed(): nonlocal recognized_alive while not done.is_set(): diff = (datetime.datetime.now() - recognized_alive).total_seconds() print(f"The diff is {diff}") if diff > 10: print("结束流处理校验") done.set() break await asyncio.sleep(1) async def read_queue(): print("准备读取识别结果") while not done.is_set(): if queue: async with stream_lock: if queue: item = queue.pop(0) print(f"准备发送: {item}") await websocket.send_text(item) await asyncio.sleep(0.1) # 缩短轮询间隔,提升实时性 async def write_stream(data): n_bytes = 4096 data_buffer = io.BytesIO(data) try: while not done.is_set(): frames = data_buffer.read(n_bytes // 2) print(f"读取 {len(frames)} 字节") if not frames: print("音频流读取完毕") stream.close() # 主动关闭流,通知ASR服务结束识别 break async with stream_lock: stream.write(frames) await asyncio.sleep(0) except Exception as e: print(f"写入流时出错: {e}") done.set() try: # 提前启动连续识别,避免音频丢失 speech_recognizer.start_continuous_recognition() data = await websocket.receive_bytes() # 创建异步任务并等待完成 task_write = asyncio.create_task(write_stream(data)) task_read = asyncio.create_task(read_queue()) task_validate = asyncio.create_task(validate_stream_processed()) await asyncio.gather(task_write, task_read, task_validate) except asyncio.CancelledError: print("WebSocket任务被取消") except Exception as e: print(f"发生错误: {e}") finally: print("关闭WebSocket连接") speech_recognizer.stop_continuous_recognition() stream.close() queue.clear() await websocket.close()
关键修复点
- 每个WebSocket请求独立创建
speech_recognizer、stream等资源,彻底解决全局复用的状态冲突问题。 - 使用
nonlocal替代全局变量,确保每个请求的状态完全隔离。 - 调整识别启动时机,在接收音频前启动连续识别,避免音频数据丢失。
- 优化异步任务轮询间隔,提升识别结果的实时推送效率。
- 音频流读取完毕后主动关闭
stream,明确通知ASR服务结束识别流程。 - 移除冗余的重复初始化代码,精简逻辑。
内容的提问来源于stack exchange,提问作者Parikshit Mukherjee
相关产品推荐
相关产品推荐

