首先安装相关依赖:

pip install pyaudio

pip install dashscope

然后创建B64PCMPlayer.py

import contextlib
import time
import pyaudio
import threading
import queue
import base64

class B64PCMPlayer:
    def __init__(self, pya: pyaudio.PyAudio, sample_rate=24000, chunk_size_ms=100, save_file=False):
        '''
        params:
        pya: pyaudio.PyAudio
        sample_rate: int, sample rate of audio
        chunk_size_ms: int, chunk size of audio in milliseconds, this will effect cancel latency
        '''

        self.pya = pya
        self.sample_rate = sample_rate
        self.chunk_size_bytes = chunk_size_ms * sample_rate *2 // 1000
        self.player_stream = pya.open(format=pyaudio.paInt16,
                channels=1,
                rate=sample_rate,
                output=True)

        self.raw_audio_buffer: queue.Queue = queue.Queue()
        self.b64_audio_buffer: queue.Queue = queue.Queue()
        self.status_lock = threading.Lock()
        self.status = 'playing'
        self.decoder_thread = threading.Thread(target=self.decoder_loop)
        self.player_thread = threading.Thread(target=self.player_loop)
        self.decoder_thread.start()
        self.player_thread.start()
        self.complete_event: threading.Event = None
        self.save_file = save_file
        if self.save_file:
            self.out_file = open('result.pcm', 'wb')

    def decoder_loop(self):
        while self.status != 'stop':
            recv_audio_b64 = None
            with contextlib.suppress(queue.Empty):
                recv_audio_b64 = self.b64_audio_buffer.get(timeout=0.1)
            if recv_audio_b64 is None:
                continue
            recv_audio_raw = base64.b64decode(recv_audio_b64)
            # push raw audio data into queue by chunk
            for i in range(0, len(recv_audio_raw), self.chunk_size_bytes):
                chunk = recv_audio_raw[i:i + self.chunk_size_bytes]
                self.raw_audio_buffer.put(chunk)
                if self.save_file:
                    self.out_file.write(chunk)

    def player_loop(self):
        while self.status != 'stop':
            recv_audio_raw = None
            with contextlib.suppress(queue.Empty):
                recv_audio_raw = self.raw_audio_buffer.get(timeout=0.1)
            if recv_audio_raw is None:
                if self.complete_event:
                    self.complete_event.set()
                continue
            # write chunk to pyaudio audio player, wait until finish playing this chunk.
            self.player_stream.write(recv_audio_raw)

    def cancel_playing(self):
        self.b64_audio_buffer.queue.clear()
        self.raw_audio_buffer.queue.clear()

    def add_data(self, data):
        self.b64_audio_buffer.put(data)

    def wait_for_complete(self):
        self.complete_event = threading.Event()
        self.complete_event.wait()
        self.complete_event = None

    def shutdown(self):
        self.status = 'stop'
        self.decoder_thread.join()
        self.player_thread.join()
        self.player_stream.close()
        if self.save_file:
            self.out_file.close()

接着创建测试文件asr.py

import logging
import os
import base64
import signal
import sys
import time
import pyaudio
import dashscope
from dashscope.audio.qwen_omni import *

# 添加当前脚本所在目录到Python路径
sys.path.append(os.path.dirname(os.path.abspath(__file__)))
from B64PCMPlayer import B64PCMPlayer
from dashscope.audio.qwen_omni.omni_realtime import TranscriptionParams

# 配置日志 - 关键改进
logger = logging.getLogger('dashscope')
logger.setLevel(logging.DEBUG)

# 创建控制台处理器并设置级别为debug
console_handler = logging.StreamHandler(sys.stdout)  # 明确指定输出到stdout
console_handler.setLevel(logging.DEBUG)

# 创建格式化器
formatter = logging.Formatter(
    '%(asctime)s - %(name)s - %(levelname)s - %(message)s')
# 添加格式化器到处理器
console_handler.setFormatter(formatter)

# 添加处理器到logger
logger.addHandler(console_handler)

# 强制刷新日志输出
logger.propagate = False

pya = None
mic_stream = None
conversation = None

def init_dashscope_api_key():
    """
        Set your DashScope API-key. More information:
        https://github.com/aliyun/alibabacloud-bailian-speech-demo/blob/master/PREREQUISITES.md
    """

    if 'DASHSCOPE_API_KEY' in os.environ:
        dashscope.api_key = os.environ[
            'DASHSCOPE_API_KEY']  # load API-key from environment variable DASHSCOPE_API_KEY
    else:
        dashscope.api_key = 'sk-xxx'  # set API-key manually


class MyCallback(OmniRealtimeCallback):
    def on_open(self) -> None:
        global pya
        global mic_stream
        print('connection opened, init microphone')
        pya = pyaudio.PyAudio()
        mic_stream = pya.open(format=pyaudio.paInt16,
                              channels=1,
                              rate=16000,
                              input=True)

    def on_close(self, close_status_code, close_msg) -> None:
        print('connection closed with code: {}, msg: {}, destroy microphone'.format(close_status_code, close_msg))
        sys.exit(0)

    def on_event(self, response: str) -> None:
        try:
            global conversation
            type = response['type']
            if 'session.created' == type:
                print('start session: {}'.format(response['session']['id']))
            if 'conversation.item.input_audio_transcription.completed' == type:
                print('final recognized text: {}'.format(response['transcript']))
            if 'conversation.item.input_audio_transcription.text' == type:
                text = response['stash']
                print("got stash result: {}".format(text))
            if 'input_audio_buffer.speech_started' == type:
                print('======Speech Start======')
            if 'input_audio_buffer.speech_stopped' == type:
                print('======Speech Stop======')
            if 'response.done' == type:
                print('======RESPONSE DONE======')
                print('[Metric] response: {}, first text delay: {}, first audio delay: {}'.format(
                    conversation.get_last_response_id(),
                    conversation.get_last_first_text_delay(),
                    conversation.get_last_first_audio_delay(),
                ))
        except Exception as e:
            print('[Error] {}'.format(e))
            return


if __name__ == '__main__':
    init_dashscope_api_key()

    print('Initializing ...')

    record_pcm_file = open('./record_16khz.pcm', 'wb')

    callback = MyCallback()

    conversation = OmniRealtimeConversation(
        model='qwen3-asr-flash-realtime',
        url='wss://dashscope.aliyuncs.com/api-ws/v1/realtime',
        callback=callback,
    )

    conversation.connect()

    transcription_params = TranscriptionParams(
        language='zh',
        sample_rate=16000,
        input_audio_format="pcm",
        corpus_text="这是一段中文对话"
    )

    conversation.update_session(
        output_modalities=[MultiModality.TEXT],
        enable_input_audio_transcription=True,
        transcription_params=transcription_params,
    )

    def signal_handler(sig, frame):
        print('Ctrl+C pressed, stop recognition ...')
        # Stop recognition
        conversation.close()
        print('omni realtime stopped.')
        # Forcefully exit the program
        sys.exit(0)


    signal.signal(signal.SIGINT, signal_handler)
    print("Press 'Ctrl+C' to stop conversation...")

    while True:
        if mic_stream:
            audio_data = mic_stream.read(3200, exception_on_overflow=False)
            record_pcm_file.write(audio_data)
            audio_b64 = base64.b64encode(audio_data).decode('ascii')
            conversation.append_audio(audio_b64)

        else:
            break

运行测试文件python asr.py

用户可以对着电脑麦克风说话,此时控制台打印:

2025-12-06 14:36:56,446 - dashscope - DEBUG - [omni realtime] append audio: 8536
2025-12-06 14:36:56,466 - dashscope - DEBUG - [omni realtime] receive string {"event_id":"event_O6PrUcbFba9VgAXDhhQyx","type":"conversation.item.input_audio_transcription.text","item_id":"item_OQ66CPHYSfTiQt9ayfs9e","content_index":0,"text":"","stash":"你好,小志","language":"zh","emotion":"neutral"}
got stash result: 你好,小志
2025-12-06 14:36:56,468 - dashscope - DEBUG - [omni realtime] receive string {"event_id":"event_WukDe9RPmo7Yhj7NKBNN4","type":"conversation.item.input_audio_transcription.text","item_id":"item_OQ66CPHYSfTiQt9ayfs9e","content_index":0,"text":"","stash":"你好,小志。","language":"zh","emotion":"neutral"}
got stash result: 你好,小志。
2025-12-06 14:36:56,472 - dashscope - DEBUG - [omni realtime] receive string {"event_id":"event_CiQnBIxhnOh9yrobPa64n","type":"conversation.item.input_audio_transcription.completed","item_id":"item_OQ66CPHYSfTiQt9ayfs9e","content_index":0,"transcript":"你好,小志。","language":"zh","emotion":"neutral","usage":{"duration":11}}
final recognized text: 你好,小志。
2025-12-06 14:36:56,655 - dashscope - DEBUG - [omni realtime] append audio: 8536
2025-12-06 14:36:56,846 - dashscope - DEBUG - [omni realtime] append audio: 8536
Ctrl+C pressed, stop recognition ...

Logo

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

更多推荐