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

异步代码中PyAudio同步音频流阻塞问题求助

问题分析与解决方案

核心问题定位

  1. 麦克风流实例被GC回收:主函数中临时创建的MicrophoneStreamer实例未被持有引用,Python垃圾回收机制会自动销毁该实例,导致音频流被关闭,后续无法获取新的音频数据。
  2. Websocket handler的终止逻辑错误:asyncio.wait设置了60秒超时,且return_when=FIRST_COMPLETED,会强制终止websocket连接和数据收发任务;同时任务取消逻辑会中断音频数据的持续发送。
  3. 主循环队列操作阻塞:当audio_queue满时,await audio_queue.get()会阻塞主循环,无法继续处理新的麦克风数据。
  4. 不必要的延迟操作:代码中的await asyncio.sleep()会导致音频数据处理和发送延迟,甚至阻塞事件循环调度。

修复步骤

  1. 持有麦克风流实例引用:在主函数中创建MicrophoneStreamer实例并赋值给变量,避免被GC回收。
  2. 修改Websocket handler逻辑:移除超时限制,使用asyncio.gather持续运行producer和consumer任务,直到websocket连接断开或手动终止;同时添加异常处理,确保连接异常时优雅关闭。
  3. 优化队列操作:当audio_queue满时,使用get_nowait()而非await get(),避免阻塞主循环;或者直接使用put_nowait()丢弃旧数据,优先保留最新音频数据。
  4. 移除不必要的sleep:删除producer和main中不必要的await asyncio.sleep(),确保事件循环高效调度。

修改后的完整代码

import asyncio
import msgpack
import os
import websockets
import pyaudio
from src.utils.constants import CHANNELS, CHUNK, FORMAT, RATE
from dotenv import load_dotenv

from .utils import websocket_data_packet  # 工具函数

load_dotenv()

MY_URL = os.environ.get("WEBSOCKETS_URL")

def make_iter():
    loop = asyncio.get_event_loop()
    queue = asyncio.Queue()
    def put(in_data, frame_count, time_info, status):
        # 使用put_nowait避免阻塞音频回调线程
        try:
            loop.call_soon_threadsafe(queue.put_nowait, in_data)
        except asyncio.QueueFull:
            # 队列满时丢弃旧数据,保留最新的
            pass
        return None, pyaudio.paContinue
    async def get():
        while True:
            yield await queue.get()
    return get(), put


class MicrophoneStreamer(object):
    chunk: int = CHUNK
    channels: int = CHANNELS
    format: int = FORMAT
    rate: int = RATE

    def __init__(self):
        self._pyaudio = pyaudio.PyAudio()
        self.is_stream_open: bool = True
        self.stream_get, stream_put = make_iter()
        self.stream = self._pyaudio.open(
            format=self.format,
            channels=self.channels,
            rate=self.rate,
            input=True,
            frames_per_buffer=self.chunk,
            stream_callback=stream_put,
        )
        self.stream.start_stream()
        
    def close(self):
        self.is_stream_open = False
        self.stream.stop_stream()
        self.stream.close()
        self._pyaudio.terminate()


async def consumer(websocket):
    async for message in websocket:
        print(f"Received message: {msgpack.unpackb(message)}")


async def producer(websocket, audio_queue):
    while True:
        chunk = await audio_queue.get()
        print(f"Sending message with audio data of size: {len(chunk)}")
        await websocket.send(msgpack.packb(websocket_data_packet(chunk)))
        audio_queue.task_done()


async def handler(audio_queue):
    print("Connecting to websocket...")
    try:
        async with websockets.connect(MY_URL) as websocket:
            print("Websocket connection established")
            producer_task = asyncio.create_task(producer(websocket, audio_queue))
            consumer_task = asyncio.create_task(consumer(websocket))
            # 持续运行直到某个任务抛出异常(如连接断开)
            await asyncio.gather(producer_task, consumer_task)
    except websockets.exceptions.WebSocketException as e:
        print(f"Websocket connection error: {e}")
    finally:
        print("Websocket connection closed")


async def main():
    audio_queue = asyncio.Queue(maxsize=5)
    # 持有麦克风流实例引用,避免被GC回收
    mic_streamer = MicrophoneStreamer()
    i = 0

    try:
        async for in_data in mic_streamer.stream_get:
            print(f"Processing audio chunk size: {len(in_data)}")
            # 触发websocket连接
            if i == 2:
                asyncio.create_task(handler(audio_queue))
            # 处理队列满的情况,丢弃旧数据
            if audio_queue.full():
                try:
                    audio_queue.get_nowait()
                except asyncio.QueueEmpty:
                    pass
            audio_queue.put_nowait(in_data)
            i += 1
    except KeyboardInterrupt:
        print("Stopping stream...")
    finally:
        mic_streamer.close()


if __name__ == "__main__":
    asyncio.run(main())

关键修改说明

  • 麦克风实例持有:将MicrophoneStreamer实例赋值给mic_streamer变量,确保GC不会回收它,音频流持续运行。
  • 队列操作优化:在音频回调和主循环中使用put_nowait(),队列满时丢弃旧数据,避免阻塞事件循环。
  • Websocket逻辑修复:使用asyncio.gather替代asyncio.wait,移除超时限制,确保websocket连接持续运行直到异常断开;添加异常处理,优雅处理连接错误。
  • 移除冗余sleep:删除不必要的延迟操作,让事件循环高效调度音频处理和websocket收发任务。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 01:27:04