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

SpeechAsyncClient流式识别方法为何出现阻塞?

WebSocket音频流转写异步阻塞问题排查与修复

问题概述

基于websockets.asyncio实现的客户端-服务器音频流传输,结合Google Cloud SpeechAsyncClient做实时转写时,服务器端_build_requests方法阻塞,无法执行到print("Done"),且_read_audio异步生成器从未被调用。客户端确认持续发送音频,服务器缓冲区已被填充。

核心问题分析

  1. 服务器端错误调用异步流式方法:Google Cloud SpeechAsyncClient的streaming_recognize返回的是异步迭代器,而非可await的协程。原代码中错误使用await调用该方法,导致方法阻塞且无法触发异步生成器的迭代。
  2. 客户端同步队列阻塞事件循环:客户端使用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())

关键修复点说明

  1. 服务器端:移除streaming_recognize调用前的await,直接使用返回的异步迭代器进行async for遍历。同时在连接关闭时主动终止音频流,确保生成器正常结束。
  2. 客户端:将同步queue.Queue替换为asyncio.Queue,并将音频生成器改为异步生成器,避免阻塞事件循环;使用call_soon_threadsafe在音频回调线程中安全操作队列。
  3. 输出优化:所有print语句添加flush=True,确保日志及时输出,便于调试。

内容的提问来源于stack exchange,提问作者thomaoc

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 08:34:52