Qwen3-Omni 实时全双工实现语音交互智能助手
#!/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;需要代理或切换网络。
更多推荐
所有评论(0)