#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
Qwen3-Omni 实时全双工示例(修正版)
- 修正回调缩进错误(on_open / on_close / on_event 同层级)
- 更安全地清空队列,优雅释放音频设备与线程
- 使用环境变量读取 DashScope API Key:DASHSCOPE_API_KEY
- 16kHz 输入(麦克风);24kHz 输出(TTS 播放);100ms 分块
"""

import os
import sys
import signal
import time
import base64
import threading
import queue
import contextlib

import pyaudio
import dashscope
from dashscope.audio.qwen_omni import *  # OmniRealtimeConversation, OmniRealtimeCallback, MultiModality, AudioFormat

# -----------------------------
# 配置
# -----------------------------
# 用环境变量传入 API Key:export DASHSCOPE_API_KEY="sk-xxx"
dashscope.api_key = "XXXXXX"
if not dashscope.api_key:
    print("[ERROR] 环境变量 DASHSCOPE_API_KEY 未设置。请先执行:export DASHSCOPE_API_KEY='sk-xxx'")
    sys.exit(1)

MODEL_NAME = "qwen3-omni-flash-realtime"
REALTIME_URL = "wss://dashscope.aliyuncs.com/api-ws/v1/realtime"  # 新加坡区域;北京区域请替换为北京的 base_url
VOICE_NAME = "Cherry"

# 采样参数:输入 16kHz、输出 24kHz;每 100ms 一块
MIC_SAMPLE_RATE = 16000
PLAY_SAMPLE_RATE = 24000
CHUNK_MS = 100
READ_FRAMES = MIC_SAMPLE_RATE * CHUNK_MS // 1000  # 1600 帧 ≈ 100ms
PLAY_CHUNK_BYTES = PLAY_SAMPLE_RATE * 2 * CHUNK_MS // 1000  # 4800 字节(16bit 单声道)


# -----------------------------
# 播放器:接收 base64 PCM → 解码 → 切块 → 播放
# -----------------------------
class B64PCMPlayer:
    def __init__(self, pya: pyaudio.PyAudio, sample_rate=PLAY_SAMPLE_RATE, chunk_size_bytes=PLAY_CHUNK_BYTES):
        self.pya = pya
        self.sample_rate = sample_rate
        self.chunk_size_bytes = chunk_size_bytes

        self.player_stream = pya.open(
            format=pyaudio.paInt16,
            channels=1,
            rate=self.sample_rate,
            output=True,
        )

        self.raw_audio_buffer: queue.Queue[bytes] = queue.Queue()
        self.b64_audio_buffer: queue.Queue[str] = queue.Queue()

        self._stop = threading.Event()
        self.decoder_thread = threading.Thread(target=self._decoder_loop, daemon=True)
        self.player_thread = threading.Thread(target=self._player_loop, daemon=True)
        self.decoder_thread.start()
        self.player_thread.start()

    def _decoder_loop(self):
        while not self._stop.is_set():
            b64 = self.b64_audio_buffer.get()
            if b64 is None:  # 作为关闭信号
                break
            try:
                raw = base64.b64decode(b64)
            except Exception:
                # 略过异常块,避免线程终止
                continue

            # 切成 100ms 的小块推给播放器
            for i in range(0, len(raw), self.chunk_size_bytes):
                self.raw_audio_buffer.put(raw[i:i + self.chunk_size_bytes])

    def _player_loop(self):
        while not self._stop.is_set():
            raw = self.raw_audio_buffer.get()
            if raw is None:  # 作为关闭信号
                break
            try:
                self.player_stream.write(raw)
            except Exception:
                # 播放异常不终止线程
                continue

    def add_data(self, b64_audio_chunk: str):
        """添加一段 base64 编码的 PCM 数据(模型的增量音频)"""
        self.b64_audio_buffer.put(b64_audio_chunk)

    def cancel_playing(self):
        """打断播放:安全清空两个队列"""
        with contextlib.suppress(queue.Empty):
            while True:
                self.b64_audio_buffer.get_nowait()
        with contextlib.suppress(queue.Empty):
            while True:
                self.raw_audio_buffer.get_nowait()

    def shutdown(self):
        """优雅关闭线程与设备"""
        self._stop.set()
        # 放入 None 作为关闭信号
        self.b64_audio_buffer.put(None)
        self.raw_audio_buffer.put(None)
        # 等待线程结束
        self.decoder_thread.join(timeout=1.0)
        self.player_thread.join(timeout=1.0)
        # 关闭播放流
        with contextlib.suppress(Exception):
            self.player_stream.close()


# -----------------------------
# 实时会话回调
# -----------------------------
class MyCallback(OmniRealtimeCallback):
    def __init__(self):
        super().__init__()
        self.pya_inst: pyaudio.PyAudio | None = None
        self.mic_stream: pyaudio.Stream | None = None
        self.b64_player: B64PCMPlayer | None = None

    def on_open(self) -> None:
        print("[INFO] connection opened, init microphone & speaker")
        # 初始化音频设备
        self.pya_inst = pyaudio.PyAudio()
        self.mic_stream = self.pya_inst.open(
            format=pyaudio.paInt16,
            channels=1,
            rate=MIC_SAMPLE_RATE,
            input=True,
        )
        self.b64_player = B64PCMPlayer(self.pya_inst)

    def on_close(self, close_status_code, close_msg) -> None:
        print(f"[INFO] connection closed: code={close_status_code}, msg={close_msg}")

    def on_event(self, response: dict) -> None:
        try:
            evt_type = response.get("type", "")
            if evt_type == "session.created":
                sid = response.get("session", {}).get("id", "")
                print(f"[INFO] start session: {sid}")

            elif evt_type == "conversation.item.input_audio_transcription.completed":
                print("[USER]", response.get("transcript", ""))

            elif evt_type == "response.audio_transcript.delta":
                # 模型文本增量
                print("[LLMΔ]", response.get("delta", ""))

            elif evt_type == "response.audio.delta":
                # 模型音频增量(base64 PCM 24kHz mono 16bit)
                if self.b64_player:
                    self.b64_player.add_data(response.get("delta", ""))

            elif evt_type == "input_audio_buffer.speech_started":
                # 你开始说话:打断当前播放,避免“你我同声”
                print("====== VAD Speech Start ======")
                if self.b64_player:
                    self.b64_player.cancel_playing()

            elif evt_type == "response.done":
                print("====== RESPONSE DONE ======")

        except Exception as e:
            print(f"[ERROR] on_event exception: {e}")


# -----------------------------
# 主程序
# -----------------------------
def main():
    print("[INFO] Initializing ...")
    callback = MyCallback()

    # 创建并连接实时会话
    conversation = OmniRealtimeConversation(
        model=MODEL_NAME,
        callback=callback,
        url=REALTIME_URL,
    )
    conversation.connect()

    # 会话参数:输入 16k PCM,输出 24k PCM + 文本
    conversation.update_session(
        output_modalities=[MultiModality.AUDIO, MultiModality.TEXT],
        voice=VOICE_NAME,
        input_audio_format=AudioFormat.PCM_16000HZ_MONO_16BIT,
        output_audio_format=AudioFormat.PCM_24000HZ_MONO_16BIT,
        enable_input_audio_transcription=True,
        input_audio_transcription_model="gummy-realtime-v1",
        enable_turn_detection=True,
        turn_detection_type="server_vad",
    )

    # Ctrl+C 优雅退出
    def handle_sigint(sig, frame):
        print("\n[INFO] Ctrl+C pressed, stopping ...")
        try:
            conversation.close()
        finally:
            # 关闭本地音频资源
            if callback.b64_player:
                callback.b64_player.shutdown()
            if callback.mic_stream:
                with contextlib.suppress(Exception):
                    callback.mic_stream.close()
            if callback.pya_inst:
                with contextlib.suppress(Exception):
                    callback.pya_inst.terminate()
        sys.exit(0)

    signal.signal(signal.SIGINT, handle_sigint)
    print("[INFO] Press 'Ctrl+C' to stop conversation...")

    # 主循环:读取 100ms 麦克风音频,编码 base64,发送给服务端
    while True:
        try:
            if callback.mic_stream is None:
                time.sleep(0.01)
                continue
            audio_data = callback.mic_stream.read(READ_FRAMES, exception_on_overflow=False)
            audio_b64 = base64.b64encode(audio_data).decode("ascii")
            conversation.append_audio(audio_b64)
        except KeyboardInterrupt:
            handle_sigint(None, None)
        except Exception as e:
            print(f"[WARN] audio loop exception: {e}")
            time.sleep(0.02)  # 简单退避避免热循环


if __name__ == "__main__":
    main()

下面把这份“Qwen3-Omni 实时全双工示例(修正版)”逐段讲透:它做的事、涉及到的线程/缓冲区、事件流,以及目前代码里仍然需要注意的坑与可选优化点。

1)它在做什么(一句话总览)
• 上行:从麦克风以 16 kHz、16-bit 单声道读取原始 PCM,每 ~100ms 打包一次 → base64 编码 → 通过 DashScope Realtime WebSocket 持续发送给模型。
• 下行:模型边生成文本、边生成 24 kHz 的 TTS 音频。文本以增量事件返回;音频以 base64 增量块返回 → 本地解码为 PCM → 按 100ms 小块实时播放。
• 并发结构:主线程负责采集和发送;播放器内部有两个后台线程(解码线程 + 播放线程)通过两个队列解耦。

2)顶部配置与常量

MODEL_NAME = “qwen3-omni-flash-realtime”
REALTIME_URL = “wss://dashscope.aliyuncs.com/api-ws/v1/realtime”
VOICE_NAME = “Cherry”
MIC_SAMPLE_RATE = 16000
PLAY_SAMPLE_RATE = 24000
CHUNK_MS = 100
READ_FRAMES = MIC_SAMPLE_RATE * CHUNK_MS // 1000 # 1600 帧 ≈ 100ms
PLAY_CHUNK_BYTES = PLAY_SAMPLE_RATE * 2 * CHUNK_MS // 1000 # 4800 字节

•	输入 16 kHz,输出 24 kHz;100ms 为一个“处理节拍”。
•	READ_FRAMES=1600 表示从 PyAudio 一次读取 1600 帧(不是字节)。16 kHz × 0.1s = 1600。
•	PLAY_CHUNK_BYTES=4800:播放时每块写入 4800 字节(24kHz × 0.1s × 16bit/采样 × 1声道)。

⚠️ 你在注释中写“使用环境变量”,但实际上代码里把 dashscope.api_key 硬编码成了一个值。这与注释不一致(见下文坑点)。

3)B64PCMPlayer:把下行音频从 base64 变成“能播的小块”

核心字段与线程:
• b64_audio_buffer(Queue[str]):接收模型返回的 base64 音频增量。
• raw_audio_buffer(Queue[bytes]):放已经解码好的原始 PCM 小块。
• decoder_thread:从 b64_audio_buffer 取出数据 → base64.b64decode → 按 PLAY_CHUNK_BYTES 切片 → 放进 raw_audio_buffer。
• player_thread:从 raw_audio_buffer 取出小块 → player_stream.write() 播放。

关键方法:
• add_data(b64):把一段模型的音频增量塞进 b64 队列。
• cancel_playing():线程安全地清空两个队列(使用 get_nowait() 循环并吞掉 queue.Empty),用于“你开始说话”时立刻打断播放,避免人机同声。
• shutdown():通过 Event + 向队列投递 None 作为哨兵来优雅停掉两条后台线程,并关闭 player_stream。

实现细节做得比较稳:
• 两个线程都用阻塞式 get()(没有忙等),CPU 友好。
• 出错用 try/except 兜底,不让线程直接崩。

4)MyCallback:处理实时会话事件(SDK 在内部线程里调用)
• on_open():WebSocket 建立 → 初始化音频设备
• pyaudio.PyAudio() 创建音频上下文。
• 打开麦克风输入流(16kHz、Int16、单声道)。
• 创建一个 B64PCMPlayer(内部会打开播放流 24kHz)。
• on_event(response):按 type 分发
• session.created:会话建立成功,打印 session.id。
• conversation.item.input_audio_transcription.completed:ASR 完成,服务端给出你刚才说的话文本。
• response.audio_transcript.delta:文本增量(模型回复的内容,每来一点就打印一点)。
• response.audio.delta:音频增量(base64 PCM 24kHz),推入 b64_player.add_data(…) 播放。
• input_audio_buffer.speech_started:VAD 检测到你开始说话,调用 cancel_playing() 立即打断当前 TTS 播放。
• response.done:本轮回答完成。
• on_close(code, msg):连接关闭时,仅打印信息(资源清理放在主程序的 Ctrl+C 处理里做)。

注意:这些回调运行在 SDK 的线程 中,不是主线程;所以回调里尽量避免重活,把主业务处理放在主线程更安全。此例只打印和把音频增量入队,比较轻。

5)主程序 main():连接 → 配置 → 循环采集并发送音频

流程:
1. conversation = OmniRealtimeConversation(…) 并 connect() 建立 WebSocket。
2. update_session(…) 配置:
• output_modalities=[AUDIO, TEXT] 让模型同时返回 TTS 音频与文本;
• 指定 voice;
• input_audio_format=PCM_16000…(与输入麦克风一致);
• output_audio_format=PCM_24000…(与播放器一致);
• enable_input_audio_transcription=True 打开 ASR;
• enable_turn_detection=True 且 turn_detection_type=‘server_vad’,让服务端做端点检测,自动切分说话轮次。
3. 注册 SIGINT(Ctrl+C)处理器:关闭会话 + 播放器 + 麦克风 + 终止 PyAudio。
4. 主循环:每次从麦克风读 1600 帧(≈100ms) → base64 编码 → conversation.append_audio(…) 发送给模型。
• 读不到麦克风就 sleep(10ms) 等待;
• 捕获异常后 sleep(20ms) 简单退避,避免热循环。

6)事件与数据流(从“说话”到“听到回复”)
1. 你对着麦克风说话 → 主线程每 100ms 读一块 16kHz PCM → base64 → 送到云端。
2. 云端做 VAD(判断你在说/不说)+ ASR(识别你说的文本)+ 生成回复(文本 + TTS)。
3. 回调收到事件:
• input_audio_buffer.speech_started:你开始说话 → 立刻 cancel_playing(),打断扬声器的播放。
• conversation.item.input_audio_transcription.completed:给出你这段话的识别文本。
• response.audio_transcript.delta:模型的文本回复“一个片段一个片段”地到达。
• response.audio.delta:模型的音频回复“一个片段一个片段”地到达 → 放入播放器队列。
4. 播放器后台线程边解码、边按 100ms 小块播放,形成流式 TTS效果。
5. response.done:这轮答复结束。

7)这份代码还有哪些“坑”与可选改进?

✅ 做得好的地方:
• 回调缩进正确、线程/队列关闭用哨兵(None)+ Event,比较稳;
• 播放端是阻塞式读取队列,避免忙等;
• VAD 触发时能安全清空两级队列,立即打断播音,体验好。

⚠️ 仍需注意/可改进:
1. API Key 硬编码
• 你注释里说“使用环境变量”,但代码实际写成了:

dashscope.api_key = “sk-xxxx”
if not dashscope.api_key: …

这样永远不会触发错误分支,也把 Key 暴露在源码里。

•	✅ 建议:

dashscope.api_key = os.getenv(“DASHSCOPE_API_KEY”)
if not dashscope.api_key:
print(“[ERROR] 请先: export DASHSCOPE_API_KEY=‘sk-xxx’”)
sys.exit(1)

并从环境注入 Key。切勿把真实 Key 放进仓库/日志。

2.	from ... import *
•	通配符导入不利于可读性与静态检查,且易污染命名空间。
•	✅ 建议显式导入:

from dashscope.audio.qwen_omni import (
OmniRealtimeConversation, OmniRealtimeCallback, MultiModality, AudioFormat
)

3.	播放队列无上限 → 极端情况下可能“越堆越多”
•	如果网络/本地播放瞬时跟不上模型推流,b64_audio_buffer/raw_audio_buffer 会无限长。
•	✅ 可选优化:给队列加 maxsize,或在入队前检测“滞后长度”,超限时丢弃旧块(或直接 cancel_playing() 重置),避免延迟滚雪球。
4.	on_close 时的资源回收
•	资源回收逻辑现在只在 Ctrl+C 的信号处理里做。如果是对端主动断开且你的进程不退出,就会遗留打开的音频设备。
•	✅ 建议在 on_close 里也触发一次播放器 shutdown()/麦克风 close()(注意线程安全)。
5.	设备选择与权限
•	现在用的是默认输入/输出设备。某些机器(多声卡/蓝牙设备)会选错设备;macOS 需要麦克风授权。
•	✅ 可选:通过 pyaudio.PyAudio().get_device_info_by_index() 枚举设备,并在 open(...) 指定 input_device_index/output_device_index。
6.	地区/可用性
•	你显式写死了新加坡的 REALTIME_URL;切换北京地域要记得替换 URL 与 Key。
•	✅ 可放到环境变量或配置文件中。
7.	健壮性
•	目前仅打印异常。生产里建议 logging,并在主循环中增加重连与退避策略(例如收到特定 close code 触发自动重连)。

8)关键参数怎么权衡?
• CHUNK_MS=100:100ms 是一个兼顾延迟与稳定的常见选择。想更快的响应,可试 40–60ms,但对网络/线程调度更敏感。
• MIC_SAMPLE_RATE=16000:ASR 常用 16 kHz,体量小、足够清晰。
• PLAY_SAMPLE_RATE=24000:TTS 常用 24 kHz,音质比 16 kHz 更好;与 update_session 的输出格式要一致。

9)常见“跑不起来”的原因速查
• 麦克风权限(尤其是 macOS):系统设置 → 隐私与安全性 → 麦克风(勾选你的终端/Python)。
• PortAudio 依赖:
• macOS:brew install portaudio && pip install pyaudio
• Linux:apt-get install -y portaudio19-dev && pip install pyaudio
• Key/区域不匹配:北京与新加坡Key 不是一套,URL 也不一样。
• 网络:办公网络可能拦截 WebSocket;需要代理或切换网络。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐