SpeechAsyncClient流式识别方法为何出现阻塞?
WebSocket音频流转写异步阻塞问题排查与修复
问题概述
基于websockets.asyncio实现的客户端-服务器音频流传输,结合Google Cloud SpeechAsyncClient做实时转写时,服务器端_build_requests方法阻塞,无法执行到print("Done"),且_read_audio异步生成器从未被调用。客户端确认持续发送音频,服务器缓冲区已被填充。
核心问题分析
- 服务器端错误调用异步流式方法:Google Cloud SpeechAsyncClient的
streaming_recognize返回的是异步迭代器,而非可await的协程。原代码中错误使用await调用该方法,导致方法阻塞且无法触发异步生成器的迭代。 - 客户端同步队列阻塞事件循环:客户端使用
queue.Queue(同步阻塞队列)存储音频数据,其中的get()方法会阻塞事件循环线程,导致WebSocket发送操作无法及时执行,间接影响服务器端的音频消费流程。
修正后的代码
服务器端代码
import re from google.cloud import speech_v1p1beta1 as speech import google.api_core.retry_async as retries import google.api_core.exceptions as core_exceptions import asyncio from websockets.asyncio.server import serve streaming_config = speech.StreamingRecognitionConfig() streaming_config.interim_results = True streaming_config.config.encoding = speech.RecognitionConfig.AudioEncoding.LINEAR16 streaming_config.config.sample_rate_hertz = 16000 streaming_config.config.language_code = "en-US" streaming_config.config.audio_channel_count = 1 streaming_config.config.enable_automatic_punctuation = True streaming_config.config.profanity_filter = True retry = retries.AsyncRetry( initial=0.1, maximum=60.0, multiplier=1.3, predicate=retries.if_exception_type( core_exceptions.DeadlineExceeded, core_exceptions.ServiceUnavailable, ), deadline=5000.0, ) class Server: def __init__(self) -> None: self._is_streaming = False self._audio_queue = asyncio.Queue() self._speech_client = speech.SpeechAsyncClient() async def _read_audio(self): print("Reading audio", flush=True) config_request = speech.StreamingRecognizeRequest() config_request.streaming_config = streaming_config yield config_request while self._is_streaming: chunk = await self._audio_queue.get() if chunk is None: return data = [chunk] while True: try: chunk = await self._audio_queue.get_nowait() if chunk is None: return data.append(chunk) except asyncio.QueueEmpty: break request = speech.StreamingRecognizeRequest() request.audio_content = b"".join(data) yield request async def _build_requests(self): print("Building requests", flush=True) audio_generator = self._read_audio() # 移除await,直接获取异步迭代器 responses = self._speech_client.streaming_recognize( requests=audio_generator, retry=retry, ) print("Starting to listen for responses", flush=True) await self._listen_print_loop(responses) print("Done", flush=True) async def _handler(self, websocket): print("Connection", flush=True) task = asyncio.create_task(self._build_requests()) self._is_streaming = True try: async for message in websocket: await self._audio_queue.put(message) except Exception as e: print(f"Failed: {e}", flush=True) finally: # 连接关闭时终止流并结束生成器 self._is_streaming = False await self._audio_queue.put(None) await task async def launch_server(self): print("Waiting for connection...", end=" ", flush=True) server = await serve( self._handler, host="localhost", port=8080, ) await server.serve_forever() async def _listen_print_loop(self, responses) -> str: num_chars_printed = 0 transcript = "" async for response in responses: if not response.results: continue result = response.results[0] if not result.alternatives: continue transcript = result.alternatives[0].transcript overwrite_chars = " " * (num_chars_printed - len(transcript)) if not result.is_final: print(transcript + overwrite_chars, end="\r", flush=True) num_chars_printed = len(transcript) else: print(transcript + overwrite_chars, flush=True) if re.search(r"\b(exit|quit)\b", transcript, re.I): print("Exiting..", flush=True) break num_chars_printed = 0 return transcript if __name__ == "__main__": server = Server() asyncio.run(server.launch_server())
客户端代码
from websockets.asyncio.client import connect import asyncio import numpy as np import sounddevice as sd from asyncio import Queue # Audio recording parameters RATE = 16000 CHUNK = int(RATE / 10) # 100ms class MicrophoneStream: """Opens a recording stream as an async generator yielding the audio chunks.""" def __init__(self, rate: int = RATE, chunk: int = CHUNK): self._rate = rate self._chunk = chunk # 使用asyncio线程安全队列替代同步队列 self._buff = Queue() self.closed = True self._loop = asyncio.get_event_loop() def __enter__(self) -> "MicrophoneStream": self.closed = False # Start the audio stream self._stream = sd.InputStream( samplerate=self._rate, channels=1, dtype='int16', blocksize=self._chunk, callback=self._fill_buffer ) self._stream.start() return self def __exit__(self, type, value, traceback) -> None: """Closes the stream, regardless of whether the connection was lost or not.""" self._stream.stop() self._stream.close() self.closed = True # 在回调线程中安全终止队列 self._loop.call_soon_threadsafe(self._buff.put_nowait, None) def _fill_buffer( self, in_data: np.ndarray, frames: int, time, status ) -> None: """Continuously collect data from the audio stream, into the buffer.""" if not self.closed: # 线程安全地向asyncio队列添加数据 self._loop.call_soon_threadsafe(self._buff.put_nowait, in_data.tobytes()) async def generator(self): while not self.closed: # 使用asyncio队列的异步get方法,不阻塞事件循环 chunk = await self._buff.get() if chunk is None: return data = [chunk] # 批量获取队列中现有数据 while not self._buff.empty(): try: chunk = self._buff.get_nowait() if chunk is None: return data.append(chunk) except asyncio.QueueEmpty: break yield b"".join(data) async def start_client(): uri = "ws://localhost:8080" async with connect(uri) as websocket: with MicrophoneStream(RATE, CHUNK) as stream: # 异步迭代音频生成器 async for chunk in stream.generator(): await websocket.send(chunk) if __name__ == "__main__": asyncio.run(start_client())
关键修复点说明
- 服务器端:移除
streaming_recognize调用前的await,直接使用返回的异步迭代器进行async for遍历。同时在连接关闭时主动终止音频流,确保生成器正常结束。 - 客户端:将同步
queue.Queue替换为asyncio.Queue,并将音频生成器改为异步生成器,避免阻塞事件循环;使用call_soon_threadsafe在音频回调线程中安全操作队列。 - 输出优化:所有
print语句添加flush=True,确保日志及时输出,便于调试。
内容的提问来源于stack exchange,提问作者thomaoc
相关产品推荐
相关产品推荐

