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

Azure Speech SDK+FastAPI WebSocket回调无法发送识别数据问题排查

问题

我正在用Azure Speech SDK和FastAPI构建实时语音识别服务,通过WebSocket接收Base64编码的二进制音频数据。Azure会话能正常启动,语音识别结果也能打印出来,但在回调函数里尝试把识别结果通过WebSocket发回客户端时,打印功能正常,send_text操作却不生效。日志里还出现了协程未被等待的警告。

相关代码

import asyncio
import io
import speechsdk
from fastapi import FastAPI, WebSocket
from azure.cognitiveservices.speech.audio import AudioStreamFormat, PushAudioInputStream

# 假设speech_config已提前配置完成

async def process_stream(stream, data, speech_recognizer, websocket):

    def recognized_callback(evt):
        recognized_text = evt.result.text
        print("I am in websocket callback : " + str(recognized_text)+" "+str(websocket))

        async def send_data():
            await websocket.send_text(recognized_text)

        asyncio.gather(send_data())

    # 每次推送的缓冲区字节数
    n_bytes = 4096
    data = io.BytesIO(data)

    speech_recognizer.recognized.connect(recognized_callback)

    # 持续推送音频数据直到读取完成
    try:
        speech_recognizer.start_continuous_recognition()
        while True:
            frames = data.read(n_bytes // 2)
            print('read {} bytes'.format(len(frames)))
            if not frames:
                speech_recognizer.stop_continuous_recognition()
                break
            stream.write(frames)

            await asyncio.sleep(0.03)
    finally:
        stream.close()

@app.websocket("/asr/en")
async def root(websocket: WebSocket):
    await websocket.accept()
    audio_format = AudioStreamFormat(
        channels=1,
        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)
    )

    try:
        while True:
            # 接收客户端音频数据
            data = await websocket.receive_bytes()
            break

        await process_stream(stream, data, speech_recognizer, websocket)

    except Exception as e:
        print(f"An error occurred: {e}")

日志输出

INFO:     Shutting down
INFO:     Waiting for application shutdown.
INFO:     Application shutdown complete.
INFO:     Finished server process [62395]
INFO:     Started server process [62529]
INFO:     Waiting for application startup.
INFO:     Application startup complete.
INFO:     ('127.0.0.1', 53888) - "WebSocket /asr/en" [accepted]
INFO:     connection open
SESSION STARTED: SessionEventArgs(session_id=02f165702334418a8635a40c4c16ea1d)
RECOGNIZING: SpeechRecognitionEventArgs(session_id=02f165702334418a8635a40c4c16ea1d, result=SpeechRecognitionResult(result_id=4011e2273ad742aa9e2df99eb3e8a854, text="thank you for", reason=ResultReason.RecognizingSpeech))
RECOGNIZING: SpeechRecognitionEventArgs(session_id=02f165702334418a8635a40c4c16ea1d, result=SpeechRecognitionResult(result_id=8bc91a19b5c2433d9fcc4dbe5125ea9c, text="thank you for contact", reason=ResultReason.RecognizingSpeech))
RECOGNIZING: SpeechRecognitionEventArgs(session_id=02f165702334418a8635a40c4c16ea1d, result=SpeechRecognitionResult(result_id=8ad3f7d5a8dc40e09fdab8dbaa13fc89, text="thank you for contacting", reason=ResultReason.RecognizingSpeech))
RECOGNIZING: SpeechRecognitionEventArgs(session_id=02f165702334418a8635a40c4c16ea1d, result=SpeechRecognitionResult(result_id=49d386798e3d4ebb8172a9b681d929fc, text="thank you for contacting us", reason=ResultReason.RecognizingSpeech))
RECOGNIZED: SpeechRecognitionEventArgs(session_id=02f165702334418a8635a40c4c16ea1d, result=SpeechRecognitionResult(result_id=f1d16bcf07da4b4ea2345e38b35c0300, text="Thank you for contacting us.", reason=ResultReason.RecognizedSpeech))
/Users/parikshit.mukherjee/PycharmProjects/pythonProject/./main.py:37: RuntimeWarning: coroutine 'root.<locals>.send_text_async' was never awaited
  send_text_async(evt.result.text)
RuntimeWarning: Enable tracemalloc to get the object allocation traceback
RECOGNIZING: SpeechRecognitionEventArgs(session_id=02f165702334418a8635a40c4c16ea1d, result=SpeechRecognitionResult(result_id=a16037fd2bec4790be32ec31b7430126, text="lands", reason=ResultReason.RecognizingSpeech))
RECOGNIZING: SpeechRecognitionEventArgs(session_id=02f165702334418a8635a40c4c16ea1d, result=SpeechRecognitionResult(result_id=508fa0fc387c43779c548dcc966133f7, text="lands had", reason=ResultReason.RecognizingSpeech))
RECOGNIZING: SpeechRecognitionEventArgs(session_id=02f165702334418a8635a40c4c16ea1d, result=SpeechRecognitionResult(result_id=8de79e5f5ac54b919cb8987dde425bae, text="yan's incorrectly", reason=ResultReason.RecognizingSpeech))
RECOGNIZING: SpeechRecognitionEventArgs(session_id=02f165702334418a8635a40c4c16ea1d, result=SpeechRecognitionResult(result_id=bfa9673b38be481d9a07263df3643794, text="lands had currently busy", reason=ResultReason.RecognizingSpeech))
RECOGNIZED: SpeechRecognitionEventArgs(session_id=02f165702334418a8635a40c4c16ea1d, result=SpeechRecognitionResult(result_id=382841509bba4ae68d5507cd594f2a9d, text="Yan's incorrectly busy.", reason=ResultReason.RecognizedSpeech))
RECOGNIZING: SpeechRecognitionEventArgs(session_id=02f165702334418a8635a40c4c16ea1d, result=SpeechRecognitionResult(result_id=5c05c9b20acc4d679a547bef5c97b473, text="how", reason=ResultReason.RecognizingSpeech))
RECOGNIZING: SpeechRecognitionEventArgs(session_id=02f165702334418a8635a40c4c16ea1d, result=SpeechRecognitionResult(result_id=c3b7c690ff61450ca2cb02f9a17d70e5, text="how pain", reason=ResultReason.RecognizingSpeech))
RECOGNIZING: SpeechRecognitionEventArgs(session_id=02f165702334418a8635a40c4c16ea1d, result=SpeechRecognitionResult(result_id=b1ee047ae4c143b78f17b75e6d306dca, text="how pain is", reason=ResultReason.RecognizingSpeech))
RECOGNIZING: SpeechRecognitionEventArgs(session_id=02f165702334418a8635a40c4c16ea1d, result=SpeechRecognitionResult(result_id=f396e7705c894f63a485378d0490d7f1, text="how pain is very", reason=ResultReason.RecognizingSpeech))
RECOGNIZING: SpeechRecognitionEventArgs(session_id=02f165702334418a8635a40c4c16ea1d, result=SpeechRecognitionResult(result_id=b502f32f7a1d47508a9af27b47abe24f, text="how pain is very important", reason=ResultReason.RecognizingSpeech))
RECOGNIZING: SpeechRecognitionEventArgs(session_id=02f165702334418a8635a40c4c16ea1d, result=SpeechRecognitionResult(result_id=618e549690534b7ca20e07d087de202f, text="how pain is very important to us", reason=ResultReason.RecognizingSpeech))
RECOGNIZED: SpeechRecognitionEventArgs(session_id=02f165702334418a8635a40c4c16ea1d, result=SpeechRecognitionResult(result_id=8af7a291290e45329b38bd172a1ddf65, text="How pain is very important to us.", reason=ResultReason.RecognizedSpeech))
/Users/parikshit.mukherjee/PycharmProjects/pythonProject/venv/lib/python3.9/site-packages/azure/cognitiveservices/speech/speech.py:652: RuntimeWarning: coroutine 'root.<locals>.session_stopped_cb' was never awaited
  cb(payload)
RuntimeWarning: Enable tracemalloc to get the object allocation traceback
INFO:     connection closed
解决方案

问题根源

  1. Azure Speech SDK的回调函数运行在独立的非异步线程中,并非FastAPI的异步事件循环线程,直接在回调里执行异步操作无法被正确调度。
  2. 代码中用asyncio.gather(send_data())创建了协程,但未等待其执行,且在同步线程中异步协程不会自动运行,导致send_text操作从未执行。

修复方案

  1. 在process_stream函数中获取FastAPI的异步事件循环。
  2. 使用asyncio.run_coroutine_threadsafe将发送WebSocket消息的协程提交到事件循环中,确保异步操作能被正确执行。
  3. 移除WebSocket接收逻辑中的break,保持持续接收音频数据,符合实时识别需求。

修改后的关键代码

async def process_stream(stream, data, speech_recognizer, websocket):
    # 获取当前运行的异步事件循环
    loop = asyncio.get_running_loop()

    def recognized_callback(evt):
        recognized_text = evt.result.text
        print("I am in websocket callback : " + str(recognized_text)+" "+str(websocket))

        # 定义异步发送函数
        async def send_data():
            await websocket.send_text(recognized_text)
        
        # 将协程提交到事件循环执行,跨线程安全
        asyncio.run_coroutine_threadsafe(send_data(), loop)

    # 剩余代码保持不变...

@app.websocket("/asr/en")
async def root(websocket: WebSocket):
    await websocket.accept()
    audio_format = AudioStreamFormat(
        channels=1,
        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)
    )

    try:
        while True:
            # 持续接收客户端的音频数据,去掉break实现实时识别
            data = await websocket.receive_bytes()
            await process_stream(stream, data, speech_recognizer, websocket)

    except Exception as e:
        print(f"An error occurred: {e}")

内容的提问来源于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 07:19:52