异步代码中PyAudio同步音频流阻塞问题求助
问题分析与解决方案
核心问题定位
- 麦克风流实例被GC回收:主函数中临时创建的
MicrophoneStreamer实例未被持有引用,Python垃圾回收机制会自动销毁该实例,导致音频流被关闭,后续无法获取新的音频数据。 - Websocket handler的终止逻辑错误:
asyncio.wait设置了60秒超时,且return_when=FIRST_COMPLETED,会强制终止websocket连接和数据收发任务;同时任务取消逻辑会中断音频数据的持续发送。 - 主循环队列操作阻塞:当
audio_queue满时,await audio_queue.get()会阻塞主循环,无法继续处理新的麦克风数据。 - 不必要的延迟操作:代码中的
await asyncio.sleep()会导致音频数据处理和发送延迟,甚至阻塞事件循环调度。
修复步骤
- 持有麦克风流实例引用:在主函数中创建
MicrophoneStreamer实例并赋值给变量,避免被GC回收。 - 修改Websocket handler逻辑:移除超时限制,使用
asyncio.gather持续运行producer和consumer任务,直到websocket连接断开或手动终止;同时添加异常处理,确保连接异常时优雅关闭。 - 优化队列操作:当
audio_queue满时,使用get_nowait()而非await get(),避免阻塞主循环;或者直接使用put_nowait()丢弃旧数据,优先保留最新音频数据。 - 移除不必要的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
相关产品推荐
相关产品推荐

