OpenAI实时语音Bot音频录制与播放并发问题求助
问题描述
我正使用OpenAI实时API结合Semantic Kernel开发异步实时语音Bot,初始运行正常,但遇到任务同步与事件循环相关问题:
- 无法中断Bot或音频播放,必须等待播放完全结束才能进行下一步操作
- 首次交互正常,但后续只能在Bot完成响应后才能再次提问
- Bot响应时,麦克风录制的音频变得断断续续,仅在播放结束后恢复正常录制
排查后发现问题出在生成器同步上:当OpenAI发送事件时,接收生成器处理事件并播放响应,输入生成器会停止收集数据,仅偶尔短暂回到麦克风生成器,导致音频录制不完整。期望实现麦克风输入流持续录制与模型输出播放并行,考虑过多线程但受Python GIL限制,寻求解决方案或变通方法,以及适用于该场景的并发管理技术/库。
解决方案与分析
核心问题定位
当前代码存在两个关键阻塞点:
VoiceBotService.run()中,麦克风输入的异步生成器循环是主阻塞循环,当receive_task中的音频播放操作占用事件循环时,麦克风读取协程无法被及时调度LocalAudioPlayer.send_audio_output中的stream.write()是同步阻塞操作,会占用事件循环时间片,挤压麦克风录制的执行机会
具体解决方案
1. 用线程分离音频I/O操作
音频录制/播放属于I/O密集型任务,可将其放到独立线程中执行,避开asyncio事件循环的阻塞限制:
- 麦克风录制:将音频读取逻辑放到线程,通过队列传递音频块给asyncio协程
- 音频播放:将同步
write操作包装到线程中执行,避免阻塞事件循环
麦克风录制修改示例:
import asyncio from queue import Queue from threading import Thread class LocalAudioRecorder(AudioInputPort): def __init__(self, device, sample_rate, channels, dtype, frame_size) -> None: super().__init__() self.device = device self.sample_rate = sample_rate self.channels = channels self.dtype = dtype self.frame_size = frame_size self.audio_queue = Queue() self.running = False def _recording_thread(self): try: with InputStream( samplerate=self.sample_rate, channels=self.channels, device=self.device, dtype=np.int16, ) as stream: with wave.open('recorded_audio.wav', 'wb') as wav_file: wav_file.setnchannels(self.channels) wav_file.setsampwidth(np.dtype(np.int16).itemsize) wav_file.setframerate(self.sample_rate) while self.running: if self._is_key_pressed(): input() print("Stopping recording...") break if stream.read_available >= self.frame_size: audio_chunk, _ = stream.read(self.frame_size) wav_file.writeframes(audio_chunk.tobytes()) self.audio_queue.put(audio_chunk) except Exception as e: print(f"An error occurred: {e}") async def get_input_audio_frames(self) -> AsyncGenerator[np.ndarray, None]: self.running = True thread = Thread(target=self._recording_thread, daemon=True) thread.start() while self.running: if not self.audio_queue.empty(): yield self.audio_queue.get() await asyncio.sleep(0.001)
音频播放修改示例:
async def send_audio_output(self, audio_frame: ndarray | RealtimeAudioEvent) -> None: if isinstance(audio_frame, RealtimeAudioEvent): audio_frame = np.frombuffer(audio_frame.audio.data, dtype=np.int16) if self.stream is None: self.stream = OutputStream( channels=self.channels, samplerate=self.sample_rate, dtype="float32" ) self.stream.start() audio_chunk = audio_frame.astype(np.float32) / np.iinfo(np.int16).max # 将同步write操作放到线程执行,不阻塞事件循环 await asyncio.to_thread(self.stream.write, audio_chunk)
2. 优化asyncio任务调度
- 移除
receive_task中不必要的await asyncio.sleep(0.01),减少事件循环切换延迟 - 重构
VoiceBotService.run(),将输入和接收任务作为独立asyncio任务并行运行:
async def run(self): async def input_task(): async for audio_frame in self.audio_in.get_input_audio_frames(): print("Sending audio") await self.ai_service.send(audio_frame) async def receive_task(): async for event in self.ai_service.receive(): if isinstance(event, (RealtimeAudioEvent, np.ndarray)): await self.audio_out.send_audio_output(event) if isinstance(event, CallInterruptedEvent): break else: await self.event_handler(event) # 并行运行两个任务,避免主循环阻塞 await asyncio.gather(input_task(), receive_task())
3. 使用异步音频库替代同步库
优先选择支持异步I/O的音频库,比如sounddevice的异步回调模式、aioaudio,这类库能直接融入asyncio事件循环,无需额外线程处理:
- 利用
sounddevice.InputStream的callback参数,在回调中直接将音频块放入队列供协程读取 - 异步音频库从根源上避免同步阻塞,提升事件循环的调度效率
相关代码
VoiceBotService
import asyncio import numpy as np from semantic_kernel.contents.realtime_events import RealtimeAudioEvent, RealtimeEvents from acev_realtime_voice_bot.service.ports import ( AIStreamingServicePort, AudioInputPort, AudioOutputPort, ) from tests.events import CallInterruptedEvent class VoiceBotService: def __init__( self, audio_in: AudioInputPort, audio_out: AudioOutputPort, ai_service: AIStreamingServicePort, ) -> None: self.audio_in = audio_in self.audio_out = audio_out self.ai_service = ai_service async def run(self): async def receive_task(): async for event in self.ai_service.receive(): await asyncio.sleep(0.01) if ( isinstance(event, RealtimeAudioEvent) or isinstance(event, np.ndarray) ): await self.audio_out.send_audio_output(event) if isinstance(event, CallInterruptedEvent): break else: await self.event_handler(event) receive_task_future = asyncio.create_task(receive_task()) async for audio_frame in self.audio_in.get_input_audio_frames(): print("Sending audio") await self.ai_service.send(audio_frame) await receive_task_future async def event_handler(self, event: RealtimeEvents): print(event.service_type)
LocalAudioRecorder
import sys import wave import select import numpy as np from sounddevice import InputStream from collections.abc import AsyncGenerator class LocalAudioRecorder(AudioInputPort): def __init__(self, device, sample_rate, channels, dtype, frame_size) -> None: super().__init__() self.device = device self.sample_rate = sample_rate self.channels = channels self.dtype = dtype self.frame_size = frame_size async def get_input_audio_frames(self) -> AsyncGenerator[np.ndarray, None]: """Generator function to yield audio data chunks and save them to a WAV file.""" try: with InputStream( samplerate=self.sample_rate, channels=self.channels, device=self.device, dtype=np.int16, ) as stream: with wave.open('recorded_audio.wav', 'wb') as wav_file: wav_file.setnchannels(self.channels) wav_file.setsampwidth(np.dtype(np.int16).itemsize) wav_file.setframerate(self.sample_rate) while True: if self._is_key_pressed(): input() print("Stopping recording...") break if stream.read_available < self.frame_size: await asyncio.sleep(0) continue audio_chunk, _ = stream.read(self.frame_size) print("Read audio chunk.") print(f"Content of audio: {audio_chunk}") wav_file.writeframes(audio_chunk.tobytes()) await asyncio.sleep(0) yield audio_chunk except Exception as e: print(f"An error occurred: {e}") def _is_key_pressed(self): return select.select([sys.stdin], [], [], 0) == ([sys.stdin], [], [])
LocalAudioPlayer
import numpy as np from numpy import ndarray from sounddevice import OutputStream from semantic_kernel.contents.realtime_events import RealtimeAudioEvent class LocalAudioPlayer(AudioOutputPort): def __init__(self, channels, sample_rate) -> None: self.channels = channels self.sample_rate = sample_rate self.stream = None async def send_audio_output( self, audio_frame: ndarray | RealtimeAudioEvent ) -> None: if isinstance(audio_frame, RealtimeAudioEvent): audio_frame = np.frombuffer(audio_frame.audio.data, dtype=np.int16) if self.stream is None: self.stream = OutputStream( channels=self.channels, samplerate=self.sample_rate, dtype="float32" ) self.stream.start() audio_chunk = audio_frame.astype(np.float32) / np.iinfo(np.int16).max self.stream.write(audio_chunk)
OpenAI实时连接适配器
import base64 from collections.abc import AsyncGenerator, Callable, Coroutine from typing import Any, cast import numpy as np from numpy import ndarray from semantic_kernel.connectors.ai import FunctionChoiceBehavior from semantic_kernel.connectors.ai.open_ai import ( AzureRealtimeExecutionSettings, AzureRealtimeWebsocket, ) from semantic_kernel.contents.audio_content import AudioContent from semantic_kernel.contents.realtime_events import RealtimeAudioEvent, RealtimeEvents from acev_realtime_voice_bot.service.ports import AIStreamingServicePort class OpenAIRealtimeAdapter(AIStreamingServicePort): def __init__( self, system_prompt: str, endpoint: str | None = None, api_version: str | None = None, deployment_name: str | None = None, ) -> None: super().__init__() self.client = AzureRealtimeWebsocket( endpoint=endpoint, api_version=api_version, deployment_name=deployment_name ) self.settings = AzureRealtimeExecutionSettings( instructions=system_prompt, turn_detection={"type": "server_vad"}, voice="shimmer", input_audio_format="pcm16", output_audio_format="pcm16", input_audio_transcription={"model": "whisper-1"}, function_choice_behavior=FunctionChoiceBehavior.Auto(), ) async def create_session(self): await self.client.create_session( settings=self.settings ) async def close_session(self): await self.client.close_session() async def send(self, event: RealtimeEvents | np.ndarray) -> None: if isinstance(event, np.ndarray): event = self._cast_input_audio_to_event(event) print("Sending event to openAI") await self.client.send(event=event) async def receive( self, audio_output_callback: Callable[[ndarray], Coroutine[Any, Any, None]] | None = None, ) -> AsyncGenerator[RealtimeEvents, None]: async for event in self.client.receive( audio_output_callback=audio_output_callback ): yield event def _cast_input_audio_to_event(self, audio_frame) -> RealtimeAudioEvent: return RealtimeAudioEvent( audio=AudioContent( data=base64.b64encode(cast(Any, audio_frame)).decode("utf-8") ) )
内容的提问来源于stack exchange,提问作者Mattia Surricchio
相关产品推荐
相关产品推荐

