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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 00:04:59