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

OpenAI实时语音Bot音频录制与播放并发问题求助

问题描述

我正使用OpenAI实时API结合Semantic Kernel开发异步实时语音Bot,初始运行正常,但遇到任务同步与事件循环相关问题:

  • 无法中断Bot或音频播放,必须等待播放完全结束才能进行下一步操作
  • 首次交互正常,但后续只能在Bot完成响应后才能再次提问
  • Bot响应时,麦克风录制的音频变得断断续续,仅在播放结束后恢复正常录制

排查后发现问题出在生成器同步上:当OpenAI发送事件时,接收生成器处理事件并播放响应,输入生成器会停止收集数据,仅偶尔短暂回到麦克风生成器,导致音频录制不完整。期望实现麦克风输入流持续录制与模型输出播放并行,考虑过多线程但受Python GIL限制,寻求解决方案或变通方法,以及适用于该场景的并发管理技术/库。

解决方案与分析

核心问题定位

当前代码存在两个关键阻塞点:

  1. VoiceBotService.run()中,麦克风输入的异步生成器循环是主阻塞循环,当receive_task中的音频播放操作占用事件循环时,麦克风读取协程无法被及时调度
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 13:45:53