Real-time speech recognition - Qwen

更新时间:
复制 MD 格式

The real-time speech recognition service receives an audio stream and transcribes it into punctuated text in real time. Use it for live captioning, online meetings, voice chat, smart assistants, and similar scenarios.

Overview

The service streams audio and returns transcribed text with low latency.

  • Recognizes Mandarin Chinese with high accuracy, plus Cantonese, Sichuanese, and other dialects.

  • Handles complex acoustic environments, with automatic language detection and intelligent filtering of non-speech audio.

  • Recognizes a range of emotional states, including surprise, calm, happiness, sadness, disgust, anger, and fear.

  • Supports custom hotwords to improve recognition accuracy for specific terms.

  • Supports context enhancement to improve recognition accuracy by passing in conversation history or domain terms.

  • Outputs timestamps to produce structured recognition results.

  • Accepts flexible sample rates and multiple audio formats to fit different recording environments.

For batch scenarios such as meeting transcription, call analysis, and subtitle generation, use Non-real-time speech recognition. For guidance on choosing a model, see Speech-to-text.

Prerequisites

Quick start

The following examples show how to call the real-time speech recognition service through the DashScope SDK.

Qwen-Audio-3.0-ASR-Flash-Streaming/Fun-ASR-Realtime

In addition to WebSocket, this model also supports the AOQ protocol. For client-side integration that prioritizes stable latency, resilience on weak networks, and built-in full-duplex noise suppression and echo cancellation, AOQ is recommended. For a protocol comparison, see Realtime API overview.

Recognize speech from a microphone

Recognize speech from a microphone and output text in real time, so words appear as the speaker talks.

Java

import com.alibaba.dashscope.audio.asr.recognition.Recognition;
import com.alibaba.dashscope.audio.asr.recognition.RecognitionParam;
import com.alibaba.dashscope.audio.asr.recognition.RecognitionResult;
import com.alibaba.dashscope.common.ResultCallback;
import com.alibaba.dashscope.utils.Constants;

import javax.sound.sampled.AudioFormat;
import javax.sound.sampled.AudioSystem;
import javax.sound.sampled.TargetDataLine;

import java.nio.ByteBuffer;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;

public class Main {
    public static void main(String[] args) throws InterruptedException {
        // The following is the configuration for the China (Beijing) region. When calling, replace "{WorkspaceId}" with your real workspace ID. Configurations differ by region.
        Constants.baseWebsocketApiUrl = "wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/inference";
        ExecutorService executorService = Executors.newSingleThreadExecutor();
        executorService.submit(new RealtimeRecognitionTask());
        executorService.shutdown();
        executorService.awaitTermination(1, TimeUnit.MINUTES);
        System.exit(0);
    }
}

class RealtimeRecognitionTask implements Runnable {
    @Override
    public void run() {
        RecognitionParam param = RecognitionParam.builder()
                .model("qwen-audio-3.0-asr-flash-streaming")
                // The API Key differs between the Singapore and Beijing regions. Get an API Key: https://help.aliyun.com/zh/model-studio/get-api-key
                // If you have not configured the environment variable, replace the following line with your Model Studio API Key: .apiKey("sk-xxx")
                .apiKey(System.getenv("DASHSCOPE_API_KEY"))
                .format("pcm")
                .sampleRate(16000)
                .build();
        Recognition recognizer = new Recognition();

        ResultCallback<RecognitionResult> callback = new ResultCallback<RecognitionResult>() {
            @Override
            public void onEvent(RecognitionResult result) {
                if (result.isSentenceEnd()) {
                    System.out.println("Final Result: " + result.getSentence().getText());
                } else {
                    System.out.println("Intermediate Result: " + result.getSentence().getText());
                }
            }

            @Override
            public void onComplete() {
                System.out.println("Recognition complete");
            }

            @Override
            public void onError(Exception e) {
                System.out.println("RecognitionCallback error: " + e.getMessage());
            }
        };
        try {
            recognizer.call(param, callback);
            // Create the audio format
            AudioFormat audioFormat = new AudioFormat(16000, 16, 1, true, false);
            // Match the default recording device based on the format
            TargetDataLine targetDataLine =
                    AudioSystem.getTargetDataLine(audioFormat);
            targetDataLine.open(audioFormat);
            // Start recording
            targetDataLine.start();
            ByteBuffer buffer = ByteBuffer.allocate(1024);
            long start = System.currentTimeMillis();
            // Record for 50s and perform real-time transcription
            while (System.currentTimeMillis() - start < 50000) {
                int read = targetDataLine.read(buffer.array(), 0, buffer.capacity());
                if (read > 0) {
                    buffer.limit(read);
                    // Send the recorded audio data to the streaming recognition service
                    recognizer.sendAudioFrame(buffer);
                    buffer = ByteBuffer.allocate(1024);
                    // The recording rate is limited; sleep for a short while to prevent excessive CPU usage
                    Thread.sleep(20);
                }
            }
            recognizer.stop();
        } catch (Exception e) {
            e.printStackTrace();
        } finally {
            // Close the WebSocket connection after the task is complete
            recognizer.getDuplexApi().close(1000, "bye");
        }

        System.out.println(
                "[Metric] requestId: "
                        + recognizer.getLastRequestId()
                        + ", first package delay ms: "
                        + recognizer.getFirstPackageDelay()
                        + ", last package delay ms: "
                        + recognizer.getLastPackageDelay());
    }
}

Python

Before you run the Python example, install the third-party audio playback and capture toolkit with pip install pyaudio.

import os
import signal  # for keyboard events handling (press "Ctrl+C" to terminate recording)
import sys

import dashscope
import pyaudio
from dashscope.audio.asr import *

mic = None
stream = None

# Set recording parameters
sample_rate = 16000  # sampling rate (Hz)
channels = 1  # mono channel
dtype = 'int16'  # data type
format_pcm = 'pcm'  # the format of the audio data
block_size = 3200  # number of frames per buffer

# Real-time speech recognition callback
class Callback(RecognitionCallback):
    def on_open(self) -> None:
        global mic
        global stream
        print('RecognitionCallback open.')
        mic = pyaudio.PyAudio()
        stream = mic.open(format=pyaudio.paInt16,
                          channels=1,
                          rate=16000,
                          input=True)

    def on_close(self) -> None:
        global mic
        global stream
        print('RecognitionCallback close.')
        stream.stop_stream()
        stream.close()
        mic.terminate()
        stream = None
        mic = None

    def on_complete(self) -> None:
        print('RecognitionCallback completed.')  # recognition completed

    def on_error(self, message) -> None:
        print('RecognitionCallback task_id: ', message.request_id)
        print('RecognitionCallback error: ', message.message)
        # Stop and close the audio stream if it is running
        if 'stream' in globals() and stream.active:
            stream.stop()
            stream.close()
        # Forcefully exit the program
        sys.exit(1)

    def on_event(self, result: RecognitionResult) -> None:
        sentence = result.get_sentence()
        if 'text' in sentence:
            print('RecognitionCallback text: ', sentence['text'])
            if RecognitionResult.is_sentence_end(sentence):
                print(
                    'RecognitionCallback sentence end, request_id:%s, usage:%s'
                    % (result.get_request_id(), result.get_usage(sentence)))

def signal_handler(sig, frame):
    print('Ctrl+C pressed, stop recognition ...')
    # Stop recognition
    recognition.stop()
    print('Recognition stopped.')
    print(
        '[Metric] requestId: {}, first package delay ms: {}, last package delay ms: {}'
        .format(
            recognition.get_last_request_id(),
            recognition.get_first_package_delay(),
            recognition.get_last_package_delay(),
        ))
    # Forcefully exit the program
    sys.exit(0)

# main function
if __name__ == '__main__':
    # The API Key differs between the Singapore and Beijing regions. Get an API Key: https://help.aliyun.com/zh/model-studio/get-api-key
    # If you have not configured the environment variable, replace the following line with your Model Studio API Key: dashscope.api_key = "sk-xxx"
    dashscope.api_key = os.environ.get('DASHSCOPE_API_KEY')

    # The following is the configuration for the China (Beijing) region. When calling, replace "{WorkspaceId}" with your actual workspace ID. The configuration differs by region.
    dashscope.base_websocket_api_url='wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/inference'

    # Create the recognition callback
    callback = Callback()

    # Call recognition service by async mode, you can customize the recognition parameters, like model, format,
    # sample_rate
    recognition = Recognition(
        model='qwen-audio-3.0-asr-flash-streaming',
        format=format_pcm,
        # 'pcm'、'wav'、'opus'、'speex'、'aac'、'amr', you can check the supported formats in the document
        sample_rate=sample_rate,
        # support 8000, 16000
        semantic_punctuation_enabled=False,
        callback=callback)

    # Start recognition
    recognition.start()

    signal.signal(signal.SIGINT, signal_handler)
    print("Press 'Ctrl+C' to stop recording and recognition...")
    # Create a keyboard listener until "Ctrl+C" is pressed

    while True:
        if stream:
            data = stream.read(3200, exception_on_overflow=False)
            recognition.send_audio_frame(data)
        else:
            break

    recognition.stop()

Recognize a local audio file

Recognize a local audio file and output the result. This suits shorter, near-real-time scenarios such as chat conversations, voice commands, voice input methods, and voice search.

Java

import com.alibaba.dashscope.api.GeneralApi;
import com.alibaba.dashscope.audio.asr.recognition.Recognition;
import com.alibaba.dashscope.audio.asr.recognition.RecognitionParam;
import com.alibaba.dashscope.audio.asr.recognition.RecognitionResult;
import com.alibaba.dashscope.base.HalfDuplexParamBase;
import com.alibaba.dashscope.common.GeneralListParam;
import com.alibaba.dashscope.common.ResultCallback;
import com.alibaba.dashscope.protocol.GeneralServiceOption;
import com.alibaba.dashscope.protocol.HttpMethod;
import com.alibaba.dashscope.protocol.Protocol;
import com.alibaba.dashscope.protocol.StreamingMode;
import com.alibaba.dashscope.utils.Constants;

import java.io.FileInputStream;
import java.nio.ByteBuffer;
import java.nio.file.Path;
import java.nio.file.Paths;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;

class TimeUtils {
    private static final DateTimeFormatter formatter =
            DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss.SSS");

    public static String getTimestamp() {
        return LocalDateTime.now().format(formatter);
    }
}

public class Main {
    public static void main(String[] args) throws InterruptedException {
        // The following is the configuration for the China (Beijing) region. When calling, replace "{WorkspaceId}" with your real workspace ID. Configurations differ by region.
        Constants.baseWebsocketApiUrl = "wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/inference";
        // In real applications, this method only needs to be executed once at the very beginning of the program; there is no need to execute it multiple times.
        warmUp();

        ExecutorService executorService = Executors.newSingleThreadExecutor();
        executorService.submit(new RealtimeRecognitionTask(Paths.get(System.getProperty("user.dir"), "{YOUR_AUDIO_FILE}")));
        executorService.shutdown();

        // wait for all tasks to complete
        executorService.awaitTermination(1, TimeUnit.MINUTES);
        System.exit(0);
    }

    public static void warmUp() {
        try {
            // Lightweight GET request to establish connection
            GeneralServiceOption warmupOption = GeneralServiceOption.builder()
                    .protocol(Protocol.HTTP)
                    .httpMethod(HttpMethod.GET)
                    .streamingMode(StreamingMode.OUT)
                    .path("assistants")
                    .build();

            warmupOption.setBaseHttpUrl(Constants.baseHttpApiUrl);
            GeneralApi<HalfDuplexParamBase> api = new GeneralApi<>();
            api.get(GeneralListParam.builder().limit(1L).build(), warmupOption);
        } catch (Exception e) {
            // Reset flag to allow retry if pre-warming failed
        }
    }
}

class RealtimeRecognitionTask implements Runnable {
    private Path filepath;

    public RealtimeRecognitionTask(Path filepath) {
        this.filepath = filepath;
    }

    @Override
    public void run() {
        RecognitionParam param = RecognitionParam.builder()
                .model("qwen-audio-3.0-asr-flash-streaming")
                // The API Key differs between the Singapore and Beijing regions. Get an API Key: https://help.aliyun.com/zh/model-studio/get-api-key
                // If you have not configured the environment variable, replace the following line with your Model Studio API Key: .apiKey("sk-xxx")
                .apiKey(System.getenv("DASHSCOPE_API_KEY"))
                .format("wav")
                .sampleRate(16000)
                .build();
        Recognition recognizer = new Recognition();

        String threadName = Thread.currentThread().getName();

        ResultCallback<RecognitionResult> callback = new ResultCallback<RecognitionResult>() {
            @Override
            public void onEvent(RecognitionResult message) {
                if (message.isSentenceEnd()) {

                    System.out.println(TimeUtils.getTimestamp()+" "+
                            "[process " + threadName + "] Final Result:" + message.getSentence().getText());
                } else {
                    System.out.println(TimeUtils.getTimestamp()+" "+
                            "[process " + threadName + "] Intermediate Result: " + message.getSentence().getText());
                }
            }

            @Override
            public void onComplete() {
                System.out.println(TimeUtils.getTimestamp()+" "+"[" + threadName + "] Recognition complete");
            }

            @Override
            public void onError(Exception e) {
                System.out.println(TimeUtils.getTimestamp()+" "+
                        "[" + threadName + "] RecognitionCallback error: " + e.getMessage());
            }
        };

        try {
            recognizer.call(param, callback);
            // Please replace the path with your audio file path
            System.out.println(TimeUtils.getTimestamp()+" "+"[" + threadName + "] Input file_path is: " + this.filepath);
            // Read file and send audio by chunks
            FileInputStream fis = new FileInputStream(this.filepath.toFile());
            byte[] allData = new byte[fis.available()];
            int ret = fis.read(allData);
            fis.close();

            int sendFrameLength = 3200;
            for (int i = 0; i * sendFrameLength < allData.length; i ++) {
                int start = i * sendFrameLength;
                int end = Math.min(start + sendFrameLength, allData.length);
                ByteBuffer byteBuffer = ByteBuffer.wrap(allData, start, end - start);
                recognizer.sendAudioFrame(byteBuffer);
                Thread.sleep(100);
            }

            System.out.println(TimeUtils.getTimestamp()+" "+LocalDateTime.now());
            recognizer.stop();
        } catch (Exception e) {
            e.printStackTrace();
        } finally {
            // Close the WebSocket connection after the task is complete
            recognizer.getDuplexApi().close(1000, "bye");
        }

        System.out.println(
                "["
                        + threadName
                        + "][Metric] requestId: "
                        + recognizer.getLastRequestId()
                        + ", first package delay ms: "
                        + recognizer.getFirstPackageDelay()
                        + ", last package delay ms: "
                        + recognizer.getLastPackageDelay());
    }
}

Python

import os
import time
import dashscope
from dashscope.audio.asr import *

# The API Key differs between the Singapore and Beijing regions. Get an API Key: https://help.aliyun.com/zh/model-studio/get-api-key
# If you have not configured the environment variable, replace the following line with your Model Studio API Key: dashscope.api_key = "sk-xxx"
dashscope.api_key = os.environ.get('DASHSCOPE_API_KEY')

# The following is the configuration for the China (Beijing) region. When calling, replace "{WorkspaceId}" with your actual workspace ID. The configuration differs by region.
dashscope.base_websocket_api_url = 'wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/inference'

from datetime import datetime

def get_timestamp():
    now = datetime.now()
    formatted_timestamp = now.strftime("[%Y-%m-%d %H:%M:%S.%f]")
    return formatted_timestamp

class Callback(RecognitionCallback):
    def on_complete(self) -> None:
        print(get_timestamp() + ' Recognition completed')  # recognition complete

    def on_error(self, result: RecognitionResult) -> None:
        print('Recognition task_id: ', result.request_id)
        print('Recognition error: ', result.message)
        exit(0)

    def on_event(self, result: RecognitionResult) -> None:
        sentence = result.get_sentence()
        if 'text' in sentence:
            print(get_timestamp() + ' RecognitionCallback text: ', sentence['text'])
        if RecognitionResult.is_sentence_end(sentence):
            print(get_timestamp() +
                  'RecognitionCallback sentence end, request_id:%s, usage:%s'
                  % (result.get_request_id(), result.get_usage(sentence)))

callback = Callback()

recognition = Recognition(model='qwen-audio-3.0-asr-flash-streaming',
                          format='wav',
                          sample_rate=16000,
                          callback=callback)

try:
    audio_data: bytes = None
    f = open("{YOUR_AUDIO_FILE}", 'rb')
    if os.path.getsize("{YOUR_AUDIO_FILE}"):
        # Read all the file data into the buffer at once
        file_buffer = f.read()
        f.close()
        print("Start Recognition")
        recognition.start()

        # Send 3200 bytes from the buffer at a time
        buffer_size = len(file_buffer)
        offset = 0
        chunk_size = 3200

        while offset < buffer_size:
            # Calculate the size of the data chunk to send this time
            remaining_bytes = buffer_size - offset
            current_chunk_size = min(chunk_size, remaining_bytes)

            # Extract the current data chunk from the buffer
            audio_data = file_buffer[offset:offset + current_chunk_size]

            # Send the audio data frame
            recognition.send_audio_frame(audio_data)
            # Update the offset
            offset += current_chunk_size

            # Add a delay to simulate real-time transmission
            time.sleep(0.1)

        recognition.stop()
    else:
        raise Exception(
            'The supplied file was empty (zero bytes long)')
except Exception as e:
    raise e

print(
    '[Metric] requestId: {}, first package delay ms: {}, last package delay ms: {}'
    .format(
        recognition.get_last_request_id(),
        recognition.get_first_package_delay(),
        recognition.get_last_package_delay(),
    ))

Qwen3-ASR-Flash-Realtime

Note

The example code reads your_audio_file.pcm (PCM16, 16 kHz, mono). If you only have an MP3, WAV, or similar format, convert it with ffmpeg:

ffmpeg -i your_audio.mp3 -ar 16000 -ac 1 -f s16le your_audio_file.pcm

Java

import com.alibaba.dashscope.audio.omni.*;
import com.alibaba.dashscope.exception.NoApiKeyException;
import com.google.gson.JsonObject;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import javax.sound.sampled.LineUnavailableException;
import java.io.File;
import java.io.FileInputStream;
import java.util.Base64;
import java.util.Collections;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.atomic.AtomicReference;

public class Qwen3AsrRealtimeUsage {
    private static final Logger log = LoggerFactory.getLogger(Qwen3AsrRealtimeUsage.class);
    private static final int AUDIO_CHUNK_SIZE = 1024; // Audio chunk size in bytes
    private static final int SLEEP_INTERVAL_MS = 30;  // Sleep interval in milliseconds

    public static void main(String[] args) throws InterruptedException, LineUnavailableException {
        CountDownLatch finishLatch = new CountDownLatch(1);

        OmniRealtimeParam param = OmniRealtimeParam.builder()
                .model("qwen3-asr-flash-realtime")
                // The following is the configuration for the China (Beijing) region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
                .url("wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/realtime")
                // The API Key differs between the Singapore and Beijing regions. Get your API Key: https://help.aliyun.com/zh/model-studio/get-api-key
                // If you have not configured the environment variable, replace the following line with your Alibaba Cloud Model Studio API Key: .apikey("sk-xxx")
                .apikey(System.getenv("DASHSCOPE_API_KEY"))
                .build();

        OmniRealtimeConversation conversation = null;
        final AtomicReference<OmniRealtimeConversation> conversationRef = new AtomicReference<>(null);
        conversation = new OmniRealtimeConversation(param, new OmniRealtimeCallback() {
            @Override
            public void onOpen() {
                System.out.println("connection opened");
            }
            @Override
            public void onEvent(JsonObject message) {
                String type = message.get("type").getAsString();
                switch(type) {
                    case "session.created":
                        System.out.println("start session: " + message.get("session").getAsJsonObject().get("id").getAsString());
                        break;
                    case "conversation.item.input_audio_transcription.completed":
                        System.out.println("transcription: " + message.get("transcript").getAsString());
                        finishLatch.countDown();
                        break;
                    case "input_audio_buffer.speech_started":
                        System.out.println("======VAD Speech Start======");
                        break;
                    case "input_audio_buffer.speech_stopped":
                        System.out.println("======VAD Speech Stop======");
                        break;
                    case "conversation.item.input_audio_transcription.text":
                        System.out.println("transcription: " + message.get("text").getAsString() + message.get("stash").getAsString());
                        break;
                    default:
                        break;
                }
            }
            @Override
            public void onClose(int code, String reason) {
                System.out.println("connection closed code: " + code + ", reason: " + reason);
            }
        });
        conversationRef.set(conversation);
        try {
            conversation.connect();
        } catch (NoApiKeyException e) {
            throw new RuntimeException(e);
        }

        OmniRealtimeTranscriptionParam transcriptionParam = new OmniRealtimeTranscriptionParam();
        transcriptionParam.setLanguage("zh");
        transcriptionParam.setInputAudioFormat("pcm");
        transcriptionParam.setInputSampleRate(16000);

        OmniRealtimeConfig config = OmniRealtimeConfig.builder()
                .modalities(Collections.singletonList(OmniRealtimeModality.TEXT))
                .transcriptionConfig(transcriptionParam)
                .build();
        conversation.updateSession(config);

        String filePath = "your_audio_file.pcm";
        File audioFile = new File(filePath);
        if (!audioFile.exists()) {
            log.error("Audio file not found: {}", filePath);
            return;
        }

        try (FileInputStream audioInputStream = new FileInputStream(audioFile)) {
            byte[] audioBuffer = new byte[AUDIO_CHUNK_SIZE];
            int bytesRead;
            int totalBytesRead = 0;

            log.info("Starting to send audio data from: {}", filePath);

            // Read and send audio data in chunks
            while ((bytesRead = audioInputStream.read(audioBuffer)) != -1) {
                totalBytesRead += bytesRead;
                String audioB64 = Base64.getEncoder().encodeToString(audioBuffer);
                // Send audio chunk to conversation
                conversation.appendAudio(audioB64);

                // Add small delay to simulate real-time audio streaming
                Thread.sleep(SLEEP_INTERVAL_MS);
            }

            log.info("Finished sending audio data. Total bytes sent: {}", totalBytesRead);

        } catch (Exception e) {
            log.error("Error sending audio from file: {}", filePath, e);
        }

        //send session.finish and wait for finish and close
        conversation.endSession();
        log.info("task finished");

        System.exit(0);
    }
}

Python

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

def setup_logging():
    """Configure log output"""
    logger = logging.getLogger('dashscope')
    logger.setLevel(logging.DEBUG)
    handler = logging.StreamHandler(sys.stdout)
    handler.setLevel(logging.DEBUG)
    formatter = logging.Formatter('%(asctime)s - %(name)s - %(levelname)s - %(message)s')
    handler.setFormatter(formatter)
    logger.addHandler(handler)
    logger.propagate = False
    return logger

def init_api_key():
    """Initialize the API Key"""
    # The API Key differs between the Singapore and Beijing regions. Get your API Key: https://help.aliyun.com/zh/model-studio/get-api-key
    # If you have not configured the environment variable, replace the following line with your Alibaba Cloud Model Studio API Key: dashscope.api_key = "sk-xxx"
    dashscope.api_key = os.environ.get('DASHSCOPE_API_KEY', 'YOUR_API_KEY')
    if dashscope.api_key == 'YOUR_API_KEY':
        print('[Warning] Using placeholder API key, set DASHSCOPE_API_KEY environment variable.')

class MyCallback(OmniRealtimeCallback):
    """Real-time recognition callback handler"""
    def __init__(self, conversation):
        self.conversation = conversation
        self.handlers = {
            'session.created': self._handle_session_created,
            'conversation.item.input_audio_transcription.completed': self._handle_final_text,
            'conversation.item.input_audio_transcription.text': self._handle_transcription_text,
            'input_audio_buffer.speech_started': lambda r: print('======Speech Start======'),
            'input_audio_buffer.speech_stopped': lambda r: print('======Speech Stop======')
        }

    def on_open(self):
        print('Connection opened')

    def on_close(self, code, msg):
        print(f'Connection closed, code: {code}, msg: {msg}')

    def on_event(self, response):
        try:
            handler = self.handlers.get(response['type'])
            if handler:
                handler(response)
        except Exception as e:
            print(f'[Error] {e}')

    def _handle_session_created(self, response):
        print(f"Start session: {response['session']['id']}")

    def _handle_final_text(self, response):
        print(f"Final recognized text: {response['transcript']}")

    def _handle_transcription_text(self, response):
        print(f"Got transcription result: {response['text'] + response['stash']}")

def read_audio_chunks(file_path, chunk_size=3200):
    """Read the audio file in chunks"""
    with open(file_path, 'rb') as f:
        while chunk := f.read(chunk_size):
            yield chunk

def send_audio(conversation, file_path, delay=0.1):
    """Send audio data"""
    if not os.path.exists(file_path):
        raise FileNotFoundError(f"Audio file {file_path} does not exist.")

    print("Processing audio file... Press 'Ctrl+C' to stop.")
    for chunk in read_audio_chunks(file_path):
        audio_b64 = base64.b64encode(chunk).decode('ascii')
        conversation.append_audio(audio_b64)
        time.sleep(delay)

def main():
    setup_logging()
    init_api_key()

    audio_file_path = "./your_audio_file.pcm"
    callback = MyCallback(conversation=None)
    conversation = OmniRealtimeConversation(
        model='qwen3-asr-flash-realtime',
        # The following is the configuration for the China (Beijing) region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
        url='wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/realtime',
        callback=callback,
    )
    callback.conversation = conversation  # Inject conversation into the callback so its methods can be called within the callback

    def handle_exit(sig, frame):
        print('Ctrl+C pressed, exiting...')
        conversation.close()
        sys.exit(0)

    signal.signal(signal.SIGINT, handle_exit)

    conversation.connect()

    transcription_params = TranscriptionParams(
        language='zh',
        sample_rate=16000,
        input_audio_format="pcm"
    )

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

    try:
        send_audio(conversation, audio_file_path)
        # send session.finish and wait for finished and close
        conversation.end_session()
    except Exception as e:
        print(f"Error occurred: {e}")
    finally:
        conversation.close()
        print("Audio processing completed.")

if __name__ == '__main__':
    main()

Paraformer

The Paraformer example code is similar to that of Qwen-Audio-3.0-ASR-Flash-Streaming/Fun-ASR-Realtime. Replace the model name with a Paraformer model.

Recognition configuration

Qwen3-ASR-Flash-Realtime interaction modes

The Qwen3-ASR-Flash-Realtime Realtime API offers two interaction modes:

  • VAD mode (default): The server automatically detects the start and end of speech (segmentation). This mode suits real-time conversations, meeting notes, and similar scenarios. To enable it, configure the session.turn_detection parameter (enabled by default).

  • Manual mode: The client controls segmentation by sending input_audio_buffer.commit. This mode suits scenarios that require explicit control over when audio is sent, such as sending a voice message in a chat app. To enable it, set session.turn_detection to null.

Switch interaction modes:

  • WebSocket: Set the turn_detection field in a session.update event.

    {
        "type": "session.update",
        "session": {
            "turn_detection": null
        }
    }
  • Python SDK: Set the enable_turn_detection parameter in the update_session method.

    conversation.update_session(
        enable_turn_detection=False
    )
  • Java SDK: Set the enableTurnDetection parameter through OmniRealtimeConfig.builder().

    OmniRealtimeConfig config = OmniRealtimeConfig.builder()
            .enableTurnDetection(false)
            .build();
    conversation.updateSession(config);

For complete SDK code examples, see Qwen-ASR-Realtime Python SDK - API reference and Java SDK. For the WebSocket event lifecycle, see Event interaction flow.

VAD segmentation configuration

Voice Activity Detection (VAD) determines when a continuous segment of speech ends, which triggers the final recognition result event. All three model families enable server-side VAD by default, but their parameter names and tuning granularity differ:

  • Qwen-Audio-3.0-ASR-Flash-Streaming / Fun-ASR-Realtime / Paraformer: Configured through max_sentence_silence (the VAD silence threshold for segmentation, in milliseconds). When the silence after a segment of speech exceeds this threshold, the system treats the sentence as complete.

  • Qwen3-ASR-Flash-Realtime: Configured through session.turn_detection, which includes silence_duration_ms (the silence duration threshold that ends a turn when exceeded; server default 800, with 400 recommended for conversation and chat scenarios that need fast segmentation) and threshold (VAD detection sensitivity; server default 0.2). Qwen3-ASR-Flash-Realtime also supports Manual mode, which disables VAD and uses client-side commit for segmentation. For details, see Qwen3-ASR-Flash-Realtime interaction modes above.

Parameter names vary by protocol: the same concept is called max_sentence_silence in Qwen-Audio-3.0-ASR-Flash-Streaming / Fun-ASR-Realtime / Paraformer, and silence_duration_ms in Qwen3-ASR-Flash-Realtime. For the full field definitions, see API reference.

Advanced features

Improve accuracy with hotwords

Use hotwords to improve recognition accuracy for specific terms, such as brand names, personal names, and proper terminology.

For detailed hotword configuration and usage, see Improve recognition accuracy.

Improve accuracy with context enhancement

Context enhancement passes conversation history or domain terminology to the ASR model to significantly improve transcription accuracy for proper terms. For detailed usage and result examples, see Context enhancement.

Get timestamps

The Qwen-Audio-3.0-ASR-Flash-Streaming, Fun-ASR-Realtime, and Paraformer model families output timestamps at both the sentence level and the word level by default, which supports subtitle alignment, keyword highlighting, karaoke-style read-along, and similar scenarios. Qwen3-ASR-Flash-Realtime (qwen3-asr-flash-realtime) does not currently return timestamps. If you need timestamps, use Qwen-Audio-3.0-ASR-Flash-Streaming, Fun-ASR-Realtime, or Paraformer. For file transcription, the Qwen ASR recording-file transcription model qwen3-asr-flash-filetrans supports word-level timestamps. For details, see Non-real-time speech recognition.

Timestamps are returned in milliseconds at two levels:

  • Sentence level: payload.output.sentence.begin_time and payload.output.sentence.end_time mark the start and end of a full sentence in the audio. In an intermediate result, end_time may be null and is filled with the final value when the sentence ends (sentence_end = true).

  • Word level: The payload.output.sentence.words array, where each element contains begin_time, end_time, text (the word or character text), and punctuation (the punctuation that follows the word, or an empty string if none).

The following excerpt shows the response structure:

{
  "payload": {
    "output": {
      "sentence": {
        "begin_time": 170,
        "end_time": 920,
        "text": "OK, I got it",
        "sentence_end": true,
        "words": [
          { "begin_time": 170, "end_time": 295, "text": "OK", "punctuation": "," },
          { "begin_time": 295, "end_time": 503, "text": "I", "punctuation": "" },
          { "begin_time": 503, "end_time": 711, "text": "got", "punctuation": "" },
          { "begin_time": 711, "end_time": 920, "text": "it", "punctuation": "" }
        ]
      }
    }
  }
}

The field names above follow the WebSocket JSON paths. Different SDKs expose these fields with their own naming conventions (dictionary keys, object properties, getter methods, and so on). For the complete field mapping, see the API reference for each SDK.

For the full field definitions, see API reference.

Emotion recognition

Qwen3-ASR-Flash-Realtime and some Paraformer models can include the speaker's emotional state in the transcription result, but the two differ in output granularity and in how the feature is enabled.

Qwen3-ASR-Flash-Realtime (qwen3-asr-flash-realtime): Always on, no configuration required. The emotion is returned through a top-level emotion field in both the conversation.item.input_audio_transcription.text and conversation.item.input_audio_transcription.completed events. The value is one of seven fine-grained emotions: surprised, neutral, happy, sad, disgusted, angry, and fearful.

{
  "type": "conversation.item.input_audio_transcription.text",
  "emotion": "neutral",
  "text": "The weather is nice today",
  "stash": ""
}

Paraformer (paraformer-realtime-8k-v2): This is the only Paraformer model that supports emotion recognition. The result is returned through payload.output.sentence.emo_tag and payload.output.sentence.emo_confidence. The value is one of three polarities: positive (such as happy or satisfied), negative (such as angry or subdued), and neutral (no clear emotion). The confidence ranges from 0.0 to 1.0.

Emotion recognition is returned only when all of the following conditions are met:

  • The model is paraformer-realtime-8k-v2.

  • Semantic segmentation is off: semantic_punctuation_enabled = false (false is the default, so no special setting is needed).

  • The result is returned only in the sentence-end event, where sentence_end = true.

To stop returning the emotion fields, set semantic_punctuation_enabled to true. This enables semantic segmentation and no longer returns the emo_tag and emo_confidence fields.

The field names above follow the WebSocket JSON paths. Different SDKs expose these fields with their own naming conventions (dictionary keys, object properties, getter methods, and so on). For the complete field mapping, see the API reference for each SDK.

For the full field definitions, value constraints, and examples, see API reference.

Sensitive word filtering

Sensitive word filtering replaces or removes sensitive words in the recognition result. Use it for call-center quality inspection, content compliance, subtitle review, and similar scenarios.

Supported models: Qwen-Audio-3.0-ASR-Flash-Streaming and Fun-ASR-Realtime only.

Limit: You can set up to 32 sensitive words.

Default behavior: When the special_word_filter parameter is not passed, no sensitive words are filtered.

How to configure: special_word_filter is a JSON object with three subfields:

  • filter_with_signed.word_list: A string array that lists the sensitive words to replace with an equal-length string of * characters. For example, with ["test"], "Help me test it" becomes "Help me **** it".

  • filter_with_empty.word_list: A string array that lists the sensitive words to remove entirely from the result. For example, with ["start"], "Is the game about to start" becomes "Is the game about to".

  • system_reserved_filter: A boolean that defaults to false. It determines whether sensitive word filtering is enabled.

Configuration example:

{
  "special_word_filter": {
    "filter_with_signed": {
      "word_list": ["test"]
    },
    "filter_with_empty": {
      "word_list": ["start", "occur"]
    },
    "system_reserved_filter": true
  }
}

Different SDKs expose these parameters with their own naming conventions (dictionary keys, object properties, methods, and so on). For the complete field mapping, see the API reference.

Call the raw WebSocket protocol

The following examples show how to connect directly to the server over the raw WebSocket protocol, for scenarios that do not use the DashScope SDK. Each example is a minimal, runnable implementation. For the WebSocket protocol, see the API reference of each model.

Click to view raw WebSocket protocol examples

Qwen-Audio-3.0-ASR-Flash-Streaming/Fun-ASR-Realtime

Python

Before you run the example, install the dependencies with the following commands:

pip uninstall websocket-client
pip uninstall websocket
pip install websocket-client

Do not name the example file websocket.py. This name conflicts with the websocket library and causes the following error: AttributeError: module 'websocket' has no attribute 'WebSocketApp'. Did you mean: 'WebSocket'?.

# pip install websocket-client
import os
import json
import time
import uuid
import threading
import websocket

# The API Key differs between the Singapore and Beijing regions. Get your API Key: https://help.aliyun.com/zh/model-studio/get-api-key
# If you have not configured the environment variable, replace the following line with your Alibaba Cloud Model Studio API Key: api_key = "sk-xxx"
api_key = os.environ.get('DASHSCOPE_API_KEY')
# The following is the configuration for the China (Beijing) region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
url = 'wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/inference/'  # WebSocket server address
audio_file = '{YOUR_AUDIO_FILE}'  # Replace with the path to your audio file

# Generate a 32-character random ID
TASK_ID = uuid.uuid4().hex[:32]

task_started = False  # Flag indicating whether the task has started

# Send the run-task instruction
def send_run_task(ws):
    run_task_message = {
        'header': {
            'action': 'run-task',
            'task_id': TASK_ID,
            'streaming': 'duplex'
        },
        'payload': {
            'task_group': 'audio',
            'task': 'asr',
            'function': 'recognition',
            'model': 'qwen-audio-3.0-asr-flash-streaming',
            'parameters': {
                'sample_rate': 16000,
                'format': 'wav'
            },
            'input': {}
        }
    }
    ws.send(json.dumps(run_task_message))

# Send the finish-task instruction
def send_finish_task(ws):
    finish_task_message = {
        'header': {
            'action': 'finish-task',
            'task_id': TASK_ID,
            'streaming': 'duplex'
        },
        'payload': {
            'input': {}
        }
    }
    ws.send(json.dumps(finish_task_message))

# Send the audio stream (send one binary chunk every 100ms)
def send_audio_stream(ws):
    chunk_size = 3200  # 100ms @ 16kHz 16bit mono
    try:
        with open(audio_file, 'rb') as f:
            while True:
                chunk = f.read(chunk_size)
                if not chunk:
                    break
                ws.send(chunk, opcode=websocket.ABNF.OPCODE_BINARY)
                time.sleep(0.1)
        print('Audio stream ended')
        send_finish_task(ws)
    except Exception as e:
        print('Error reading audio file:', e)
        ws.close()

# Send the run-task instruction when the connection opens
def on_open(ws):
    print('Connected to server')
    send_run_task(ws)

# Handle received messages
def on_message(ws, data):
    global task_started
    message = json.loads(data)
    event = message['header']['event']
    if event == 'task-started':
        print('Task started')
        task_started = True
        threading.Thread(target=send_audio_stream, args=(ws,), daemon=True).start()
    elif event == 'result-generated':
        print('Recognition result:', message['payload']['output']['sentence']['text'])
        if message['payload'].get('usage'):
            print('Task billing duration (seconds):', message['payload']['usage']['duration'])
    elif event == 'task-finished':
        print('Task finished')
        ws.close()
    elif event == 'task-failed':
        print('Task failed:', message['header'].get('error_message'))
        ws.close()
    else:
        print('Unknown event:', event)

# Close the connection if the task-started event is not received
def on_close(ws, close_status_code, close_msg):
    if not task_started:
        print('Task not started, closing connection')

# Error handling
def on_error(ws, error):
    print('WebSocket error:', error)

if __name__ == '__main__':
    ws = websocket.WebSocketApp(
        url,
        header={'Authorization': f'bearer {api_key}'},
        on_open=on_open,
        on_message=on_message,
        on_error=on_error,
        on_close=on_close
    )
    ws.run_forever()

Java

Before you run the example, install the Java-WebSocket dependency:

Maven

<dependency>
    <groupId>org.java-websocket</groupId>
    <artifactId>Java-WebSocket</artifactId>
    <version>1.5.6</version>
</dependency>
<dependency>
    <groupId>org.json</groupId>
    <artifactId>json</artifactId>
    <version>20240303</version>
</dependency>

Gradle

implementation 'org.java-websocket:Java-WebSocket:1.5.6'
implementation 'org.json:json:20240303'
import org.java_websocket.client.WebSocketClient;
import org.java_websocket.handshake.ServerHandshake;
import org.json.JSONObject;

import java.net.URI;
import java.nio.ByteBuffer;
import java.nio.file.Files;
import java.nio.file.Paths;
import java.util.UUID;
import java.util.concurrent.atomic.AtomicBoolean;

public class FunASRRealtimeClient {

    // The API Key differs between the Singapore and Beijing regions. Get your API Key: https://help.aliyun.com/zh/model-studio/get-api-key
    // If you have not configured the environment variable, replace the following line with your Alibaba Cloud Model Studio API Key: private static final String API_KEY = "sk-xxx";
    private static final String API_KEY = System.getenv().getOrDefault("DASHSCOPE_API_KEY", "sk-xxx");
    // The following is the configuration for the China (Beijing) region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
    private static final String URL = "wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/inference/";
    private static final String AUDIO_FILE = "{YOUR_AUDIO_FILE}"; // Replace with the path to your audio file
    private static final String MODEL = "qwen-audio-3.0-asr-flash-streaming";

    // Generate a 32-character random ID
    private static final String TASK_ID = UUID.randomUUID().toString().replace("-", "").substring(0, 32);

    private static final AtomicBoolean taskStarted = new AtomicBoolean(false);
    private static WebSocketClient client;

    public static void main(String[] args) throws Exception {
        client = new WebSocketClient(new URI(URL)) {
            @Override
            public void onOpen(ServerHandshake handshake) {
                System.out.println("Connected to server");
                sendRunTask();
            }

            @Override
            public void onMessage(String data) {
                JSONObject message = new JSONObject(data);
                String event = message.getJSONObject("header").getString("event");
                switch (event) {
                    case "task-started":
                        System.out.println("Task started");
                        taskStarted.set(true);
                        new Thread(FunASRRealtimeClient::sendAudioStream).start();
                        break;
                    case "result-generated":
                        JSONObject payload = message.getJSONObject("payload");
                        String text = payload.getJSONObject("output").getJSONObject("sentence").getString("text");
                        System.out.println("Recognition result: " + text);
                        if (payload.has("usage")) {
                            System.out.println("Task billing duration (seconds): " + payload.getJSONObject("usage").get("duration"));
                        }
                        break;
                    case "task-finished":
                        System.out.println("Task finished");
                        close();
                        break;
                    case "task-failed":
                        String errMsg = message.getJSONObject("header").optString("error_message");
                        System.err.println("Task failed: " + errMsg);
                        close();
                        break;
                    default:
                        System.out.println("Unknown event: " + event);
                }
            }

            @Override
            public void onClose(int code, String reason, boolean remote) {
                if (!taskStarted.get()) {
                    System.err.println("Task not started, closing connection");
                }
            }

            @Override
            public void onError(Exception ex) {
                System.err.println("WebSocket error: " + ex.getMessage());
            }
        };
        client.addHeader("Authorization", "bearer " + API_KEY);
        client.connectBlocking();
    }

    // Send the run-task instruction
    private static void sendRunTask() {
        JSONObject runTask = new JSONObject()
                .put("header", new JSONObject()
                        .put("action", "run-task")
                        .put("task_id", TASK_ID)
                        .put("streaming", "duplex"))
                .put("payload", new JSONObject()
                        .put("task_group", "audio")
                        .put("task", "asr")
                        .put("function", "recognition")
                        .put("model", MODEL)
                        .put("parameters", new JSONObject()
                                .put("sample_rate", 16000)
                                .put("format", "wav"))
                        .put("input", new JSONObject()));
        client.send(runTask.toString());
    }

    // Send the audio stream (send one binary chunk every 100ms)
    private static void sendAudioStream() {
        int chunkSize = 3200; // 100ms @ 16kHz 16bit mono
        try {
            byte[] audio = Files.readAllBytes(Paths.get(AUDIO_FILE));
            int offset = 0;
            while (offset < audio.length) {
                int end = Math.min(offset + chunkSize, audio.length);
                byte[] chunk = new byte[end - offset];
                System.arraycopy(audio, offset, chunk, 0, end - offset);
                client.send(ByteBuffer.wrap(chunk));
                offset = end;
                Thread.sleep(100);
            }
            System.out.println("Audio stream ended");
            sendFinishTask();
        } catch (Exception e) {
            System.err.println("Error reading audio file: " + e.getMessage());
            client.close();
        }
    }

    // Send the finish-task instruction
    private static void sendFinishTask() {
        JSONObject finishTask = new JSONObject()
                .put("header", new JSONObject()
                        .put("action", "finish-task")
                        .put("task_id", TASK_ID)
                        .put("streaming", "duplex"))
                .put("payload", new JSONObject()
                        .put("input", new JSONObject()));
        client.send(finishTask.toString());
    }
}

Node.js

Install the required dependencies:

npm install ws
npm install uuid

The example code is as follows:

const fs = require('fs');
const WebSocket = require('ws');
const { v4: uuidv4 } = require('uuid'); // Used to generate a UUID

// The API Key differs between the Singapore and Beijing regions. Get your API Key: https://help.aliyun.com/zh/model-studio/get-api-key
// If you have not configured the environment variable, replace the following line with your Alibaba Cloud Model Studio API Key: const apiKey = "sk-xxx"
const apiKey = process.env.DASHSCOPE_API_KEY;
// The following is the configuration for the China (Beijing) region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
const url = 'wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/inference/'; // WebSocket server address
const audioFile = '{YOUR_AUDIO_FILE}'; // Replace with the path to your audio file

// Generate a 32-character random ID
const TASK_ID = uuidv4().replace(/-/g, '').slice(0, 32);

// Create the WebSocket client
const ws = new WebSocket(url, {
  headers: {
    Authorization: `bearer ${apiKey}`
  }
});

let taskStarted = false; // Flag indicating whether the task has started

// Send the run-task instruction when the connection opens
ws.on('open', () => {
  console.log('Connected to server');
  sendRunTask();
});

// Handle received messages
ws.on('message', (data) => {
  const message = JSON.parse(data);
  switch (message.header.event) {
    case 'task-started':
      console.log('Task started');
      taskStarted = true;
      sendAudioStream();
      break;
    case 'result-generated':
      console.log('Recognition result:', message.payload.output.sentence.text);
      if (message.payload.usage) {
        console.log('Task billing duration (seconds):', message.payload.usage.duration);
      }
      break;
    case 'task-finished':
      console.log('Task finished');
      ws.close();
      break;
    case 'task-failed':
      console.error('Task failed:', message.header.error_message);
      ws.close();
      break;
    default:
      console.log('Unknown event:', message.header.event);
  }
});

// Close the connection if the task-started event is not received
ws.on('close', () => {
  if (!taskStarted) {
    console.error('Task not started, closing connection');
  }
});

// Send the run-task instruction
function sendRunTask() {
  const runTaskMessage = {
    header: {
      action: 'run-task',
      task_id: TASK_ID,
      streaming: 'duplex'
    },
    payload: {
      task_group: 'audio',
      task: 'asr',
      function: 'recognition',
      model: 'qwen-audio-3.0-asr-flash-streaming',
      parameters: {
        sample_rate: 16000,
        format: 'wav'
      },
      input: {}
    }
  };
  ws.send(JSON.stringify(runTaskMessage));
}

// Send the audio stream
function sendAudioStream() {
  const audioStream = fs.createReadStream(audioFile);
  let chunkCount = 0;

  function sendNextChunk() {
    const chunk = audioStream.read();
    if (chunk) {
      ws.send(chunk);
      chunkCount++;
      setTimeout(sendNextChunk, 100); // Send once every 100ms
    }
  }

  audioStream.on('readable', () => {
    sendNextChunk();
  });

  audioStream.on('end', () => {
    console.log('Audio stream ended');
    sendFinishTask();
  });

  audioStream.on('error', (err) => {
    console.error('Error reading audio file:', err);
    ws.close();
  });
}

// Send the finish-task instruction
function sendFinishTask() {
  const finishTaskMessage = {
    header: {
      action: 'finish-task',
      task_id: TASK_ID,
      streaming: 'duplex'
    },
    payload: {
      input: {}
    }
  };
  ws.send(JSON.stringify(finishTaskMessage));
}

// Error handling
ws.on('error', (error) => {
  console.error('WebSocket error:', error);
});

C#

The example code is as follows:

using System.Net.WebSockets;
using System.Text;
using System.Text.Json;
using System.Text.Json.Nodes;

class Program {
    private static ClientWebSocket _webSocket = new ClientWebSocket();
    private static CancellationTokenSource _cancellationTokenSource = new CancellationTokenSource();
    private static bool _taskStartedReceived = false;
    private static bool _taskFinishedReceived = false;
    // The API Key differs between the Singapore and Beijing regions. Get your API Key: https://help.aliyun.com/zh/model-studio/get-api-key
    // If you have not configured the environment variable, replace the following line with your Alibaba Cloud Model Studio API Key: private static readonly string ApiKey = "sk-xxx"
    private static readonly string ApiKey = Environment.GetEnvironmentVariable("DASHSCOPE_API_KEY") ?? throw new InvalidOperationException("DASHSCOPE_API_KEY environment variable is not set.");

    // The following is the configuration for the China (Beijing) region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
    private const string WebSocketUrl = "wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/inference/";
    // Replace with the path to your audio file
    private const string AudioFilePath = "{YOUR_AUDIO_FILE}";

    static async Task Main(string[] args) {
        // Establish the WebSocket connection and configure headers for authentication
        _webSocket.Options.SetRequestHeader("Authorization", $"bearer {ApiKey}");

        await _webSocket.ConnectAsync(new Uri(WebSocketUrl), _cancellationTokenSource.Token);

        // Start a thread to receive WebSocket messages asynchronously
        var receiveTask = ReceiveMessagesAsync();

        // Send the run-task instruction
        string _taskId = Guid.NewGuid().ToString("N"); // Generate a 32-character random ID
        var runTaskJson = GenerateRunTaskJson(_taskId);
        await SendAsync(runTaskJson);

        // Wait for the task-started event
        while (!_taskStartedReceived) {
            await Task.Delay(100, _cancellationTokenSource.Token);
        }

        // Read the local file and send the audio stream to be recognized to the server
        await SendAudioStreamAsync(AudioFilePath);

        // Send the finish-task instruction to end the task
        var finishTaskJson = GenerateFinishTaskJson(_taskId);
        await SendAsync(finishTaskJson);

        // Wait for the task-finished event
        while (!_taskFinishedReceived && !_cancellationTokenSource.IsCancellationRequested) {
            try {
                await Task.Delay(100, _cancellationTokenSource.Token);
            } catch (OperationCanceledException) {
                // The task has been canceled, exit the loop
                break;
            }
        }

        // Close the connection
        if (!_cancellationTokenSource.IsCancellationRequested) {
            await _webSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, "Closing", _cancellationTokenSource.Token);
        }

        _cancellationTokenSource.Cancel();
        try {
            await receiveTask;
        } catch (OperationCanceledException) {
            // Ignore the operation canceled exception
        }
    }

    private static async Task ReceiveMessagesAsync() {
        try {
            while (_webSocket.State == WebSocketState.Open && !_cancellationTokenSource.IsCancellationRequested) {
                var message = await ReceiveMessageAsync(_cancellationTokenSource.Token);
                if (message != null) {
                    var eventValue = message["header"]?["event"]?.GetValue<string>();
                    switch (eventValue) {
                        case "task-started":
                            Console.WriteLine("Task started successfully");
                            _taskStartedReceived = true;
                            break;
                        case "result-generated":
                            Console.WriteLine($"Recognition result: {message["payload"]?["output"]?["sentence"]?["text"]?.GetValue<string>()}");
                            if (message["payload"]?["usage"] != null && message["payload"]?["usage"]?["duration"] != null) {
                                Console.WriteLine($"Task billing duration (seconds): {message["payload"]?["usage"]?["duration"]?.GetValue<int>()}");
                            }
                            break;
                        case "task-finished":
                            Console.WriteLine("Task finished");
                            _taskFinishedReceived = true;
                            _cancellationTokenSource.Cancel();
                            break;
                        case "task-failed":
                            Console.WriteLine($"Task failed: {message["header"]?["error_message"]?.GetValue<string>()}");
                            _cancellationTokenSource.Cancel();
                            break;
                    }
                }
            }
        } catch (OperationCanceledException) {
            // Ignore the operation canceled exception
        }
    }

    private static async Task<JsonNode?> ReceiveMessageAsync(CancellationToken cancellationToken) {
        var buffer = new byte[1024 * 4];
        var segment = new ArraySegment<byte>(buffer);
        var result = await _webSocket.ReceiveAsync(segment, cancellationToken);

        if (result.MessageType == WebSocketMessageType.Close) {
            await _webSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, "Closing", cancellationToken);
            return null;
        }

        var message = Encoding.UTF8.GetString(buffer, 0, result.Count);
        return JsonNode.Parse(message);
    }

    private static async Task SendAsync(string message) {
        var buffer = Encoding.UTF8.GetBytes(message);
        var segment = new ArraySegment<byte>(buffer);
        await _webSocket.SendAsync(segment, WebSocketMessageType.Text, true, _cancellationTokenSource.Token);
    }

    private static async Task SendAudioStreamAsync(string filePath) {
        using (var audioStream = File.OpenRead(filePath)) {
            var buffer = new byte[1024]; // Send 100ms of audio data each time
            int bytesRead;

            while ((bytesRead = await audioStream.ReadAsync(buffer, 0, buffer.Length)) > 0) {
                var segment = new ArraySegment<byte>(buffer, 0, bytesRead);
                await _webSocket.SendAsync(segment, WebSocketMessageType.Binary, true, _cancellationTokenSource.Token);
                await Task.Delay(100); // 100ms interval
            }
        }
    }

    private static string GenerateRunTaskJson(string taskId) {
        var runTask = new JsonObject {
            ["header"] = new JsonObject {
                ["action"] = "run-task",
                ["task_id"] = taskId,
                ["streaming"] = "duplex"
            },
            ["payload"] = new JsonObject {
                ["task_group"] = "audio",
                ["task"] = "asr",
                ["function"] = "recognition",
                ["model"] = "qwen-audio-3.0-asr-flash-streaming",
                ["parameters"] = new JsonObject {
                    ["format"] = "wav",
                    ["sample_rate"] = 16000,
                },
                ["input"] = new JsonObject()
            }
        };
        return JsonSerializer.Serialize(runTask);
    }

    private static string GenerateFinishTaskJson(string taskId) {
        var finishTask = new JsonObject {
            ["header"] = new JsonObject {
                ["action"] = "finish-task",
                ["task_id"] = taskId,
                ["streaming"] = "duplex"
            },
            ["payload"] = new JsonObject {
                ["input"] = new JsonObject()
            }
        };
        return JsonSerializer.Serialize(finishTask);
    }
}

PHP

The example project has the following directory structure:

my-php-project/

├── composer.json

├── vendor/

└── index.php

The contents of composer.json are as follows. Adjust the dependency versions as needed:

{
    "require": {
        "react/event-loop": "^1.3",
        "react/socket": "^1.11",
        "react/stream": "^1.2",
        "react/http": "^1.1",
        "ratchet/pawl": "^0.4"
    },
    "autoload": {
        "psr-4": {
            "App\\": "src/"
        }
    }
}

The contents of index.php are as follows:

<?php

require __DIR__ . '/vendor/autoload.php';

use Ratchet\Client\Connector;
use React\EventLoop\Loop;
use React\Socket\Connector as SocketConnector;
use Ratchet\rfc6455\Messaging\Frame;

// The API Key differs between the Singapore and Beijing regions. Get your API Key: https://help.aliyun.com/zh/model-studio/get-api-key
// If you have not configured the environment variable, replace the following line with your Alibaba Cloud Model Studio API Key: $api_key = "sk-xxx"
$api_key = getenv("DASHSCOPE_API_KEY");
// The following is the configuration for the China (Beijing) region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
$websocket_url = 'wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/inference/';
$audio_file_path = '{YOUR_AUDIO_FILE}'; // Replace with the path to your audio file

$loop = Loop::get();

// Create a custom connector
$socketConnector = new SocketConnector($loop, [
    'tcp' => [
        'bindto' => '0.0.0.0:0',
    ],
    'tls' => [
        'verify_peer' => false,
        'verify_peer_name' => false,
    ],
]);

$connector = new Connector($loop, $socketConnector);

$headers = [
    'Authorization' => 'bearer ' . $api_key
];

$connector($websocket_url, [], $headers)->then(function ($conn) use ($loop, $audio_file_path) {
    echo "Connected to WebSocket server\n";

    // Start a thread to receive WebSocket messages asynchronously
    $conn->on('message', function($msg) use ($conn, $loop, $audio_file_path) {
        $response = json_decode($msg, true);

        if (isset($response['header']['event'])) {
            handleEvent($conn, $response, $loop, $audio_file_path);
        } else {
            echo "Unknown message format\n";
        }
    });

    // Listen for the connection close
    $conn->on('close', function($code = null, $reason = null) {
        echo "Connection closed\n";
        if ($code !== null) {
            echo "Close code: " . $code . "\n";
        }
        if ($reason !== null) {
            echo "Close reason: " . $reason . "\n";
        }
    });

    // Generate the task ID
    $taskId = generateTaskId();

    // Send the run-task instruction
    sendRunTaskMessage($conn, $taskId);

}, function ($e) {
    echo "Unable to connect: {$e->getMessage()}\n";
});

$loop->run();

/**
 * Generate the task ID
 * @return string
 */
function generateTaskId(): string {
    return bin2hex(random_bytes(16));
}

/**
 * Send the run-task instruction
 * @param $conn
 * @param $taskId
 */
function sendRunTaskMessage($conn, $taskId) {
    $runTaskMessage = json_encode([
        "header" => [
            "action" => "run-task",
            "task_id" => $taskId,
            "streaming" => "duplex"
        ],
        "payload" => [
            "task_group" => "audio",
            "task" => "asr",
            "function" => "recognition",
            "model" => "qwen-audio-3.0-asr-flash-streaming",
            "parameters" => [
                "format" => "wav",
                "sample_rate" => 16000
            ],
            "input" => []
        ]
    ]);
    echo "Preparing to send the run-task instruction: " . $runTaskMessage . "\n";
    $conn->send($runTaskMessage);
    echo "run-task instruction sent\n";
}

/**
 * Read the audio file
 * @param string $filePath
 * @return bool|string
 */
function readAudioFile(string $filePath) {
    $voiceData = file_get_contents($filePath);
    if ($voiceData === false) {
        echo "Unable to read the audio file\n";
    }
    return $voiceData;
}

/**
 * Split the audio data
 * @param string $data
 * @param int $chunkSize
 * @return array
 */
function splitAudioData(string $data, int $chunkSize): array {
    return str_split($data, $chunkSize);
}

/**
 * Send the finish-task instruction
 * @param $conn
 * @param $taskId
 */
function sendFinishTaskMessage($conn, $taskId) {
    $finishTaskMessage = json_encode([
        "header" => [
            "action" => "finish-task",
            "task_id" => $taskId,
            "streaming" => "duplex"
        ],
        "payload" => [
            "input" => []
        ]
    ]);
    echo "Preparing to send the finish-task instruction: " . $finishTaskMessage . "\n";
    $conn->send($finishTaskMessage);
    echo "finish-task instruction sent\n";
}

/**
 * Handle events
 * @param $conn
 * @param $response
 * @param $loop
 * @param $audio_file_path
 */
function handleEvent($conn, $response, $loop, $audio_file_path) {
    static $taskId;
    static $chunks;
    static $allChunksSent = false;

    if (is_null($taskId)) {
        $taskId = generateTaskId();
    }

    switch ($response['header']['event']) {
        case 'task-started':
            echo "Task started, sending audio data...\n";
            // Read the audio file
            $voiceData = readAudioFile($audio_file_path);
            if ($voiceData === false) {
                echo "Unable to read the audio file\n";
                $conn->close();
                return;
            }

            // Split the audio data
            $chunks = splitAudioData($voiceData, 1024);

            // Define the send function
            $sendChunk = function() use ($conn, &$chunks, $loop, &$sendChunk, &$allChunksSent, $taskId) {
                if (!empty($chunks)) {
                    $chunk = array_shift($chunks);
                    $binaryMsg = new Frame($chunk, true, Frame::OP_BINARY);
                    $conn->send($binaryMsg);
                    // Send the next chunk after 100ms
                    $loop->addTimer(0.1, $sendChunk);
                } else {
                    echo "All data chunks sent\n";
                    $allChunksSent = true;

                    // Send the finish-task instruction
                    sendFinishTaskMessage($conn, $taskId);
                }
            };

            // Start sending audio data
            $sendChunk();
            break;
        case 'result-generated':
            $result = $response['payload']['output']['sentence'];
            echo "Recognition result: " . $result['text'] . "\n";
            if (isset($response['payload']['usage']['duration'])) {
                echo "Task billing duration (seconds): " . $response['payload']['usage']['duration'] . "\n";
            }
            break;
        case 'task-finished':
            echo "Task finished\n";
            $conn->close();
            break;
        case 'task-failed':
            echo "Task failed\n";
            echo "Error code: " . $response['header']['error_code'] . "\n";
            echo "Error message: " . $response['header']['error_message'] . "\n";
            $conn->close();
            break;
        case 'error':
            echo "Error: " . $response['payload']['message'] . "\n";
            break;
        default:
            echo "Unknown event: " . $response['header']['event'] . "\n";
            break;
    }

    // If all data has been sent and the task has finished, close the connection
    if ($allChunksSent && $response['header']['event'] == 'task-finished') {
        // Wait 1 second to ensure all data has been transmitted
        $loop->addTimer(1, function() use ($conn) {
            $conn->close();
            echo "Client closed the connection\n";
        });
    }
}

Go

package main

import (
	"encoding/json"
	"fmt"
	"io"
	"log"
	"net/http"
	"os"
	"time"

	"github.com/google/uuid"
	"github.com/gorilla/websocket"
)

const (
	// The following is the configuration for the China (Beijing) region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
	wsURL     = "wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/inference/" // WebSocket server address
	audioFile = "{YOUR_AUDIO_FILE}"                                   // Replace with the path to your audio file
)

var dialer = websocket.DefaultDialer

func main() {
	// The API Key differs between the Singapore and Beijing regions. Get your API Key: https://help.aliyun.com/zh/model-studio/get-api-key
    // If you have not configured the environment variable, replace the following line with your Alibaba Cloud Model Studio API Key: apiKey := "sk-xxx"
	apiKey := os.Getenv("DASHSCOPE_API_KEY")

	// Connect to the WebSocket service
	conn, err := connectWebSocket(apiKey)
	if err != nil {
		log.Fatal("Failed to connect to WebSocket: ", err)
	}
	defer closeConnection(conn)

	// Start a goroutine to receive results
	taskStarted := make(chan bool)
	taskDone := make(chan bool)
	startResultReceiver(conn, taskStarted, taskDone)

	// Send the run-task instruction
	taskID, err := sendRunTaskCmd(conn)
	if err != nil {
		log.Fatal("Failed to send the run-task instruction: ", err)
	}

	// Wait for the task-started event
	waitForTaskStarted(taskStarted)

	// Send the audio file stream to be recognized
	if err := sendAudioData(conn); err != nil {
		log.Fatal("Failed to send audio: ", err)
	}

	// Send the finish-task instruction
	if err := sendFinishTaskCmd(conn, taskID); err != nil {
		log.Fatal("Failed to send the finish-task instruction: ", err)
	}

	// Wait for the task to finish or fail
	<-taskDone
}

// Define structs to represent the JSON data
type Header struct {
	Action       string                 `json:"action"`
	TaskID       string                 `json:"task_id"`
	Streaming    string                 `json:"streaming"`
	Event        string                 `json:"event"`
	ErrorCode    string                 `json:"error_code,omitempty"`
	ErrorMessage string                 `json:"error_message,omitempty"`
	Attributes   map[string]interface{} `json:"attributes"`
}

type Output struct {
	Sentence struct {
		BeginTime int64  `json:"begin_time"`
		EndTime   *int64 `json:"end_time"`
		Text      string `json:"text"`
		Words     []struct {
			BeginTime   int64  `json:"begin_time"`
			EndTime     *int64 `json:"end_time"`
			Text        string `json:"text"`
			Punctuation string `json:"punctuation"`
		} `json:"words"`
	} `json:"sentence"`
}

type Payload struct {
	TaskGroup  string `json:"task_group"`
	Task       string `json:"task"`
	Function   string `json:"function"`
	Model      string `json:"model"`
	Parameters Params `json:"parameters"`
	Input      Input  `json:"input"`
	Output     Output `json:"output,omitempty"`
	Usage      *struct {
		Duration int `json:"duration"`
	} `json:"usage,omitempty"`
}

type Params struct {
	Format                   string `json:"format"`
	SampleRate               int    `json:"sample_rate"`
	DisfluencyRemovalEnabled bool   `json:"disfluency_removal_enabled"`
}

type Input struct {
}

type Event struct {
	Header  Header  `json:"header"`
	Payload Payload `json:"payload"`
}

// Connect to the WebSocket service
func connectWebSocket(apiKey string) (*websocket.Conn, error) {
	header := make(http.Header)
	header.Add("Authorization", fmt.Sprintf("bearer %s", apiKey))
	conn, _, err := dialer.Dial(wsURL, header)
	return conn, err
}

// Start a goroutine to receive WebSocket messages asynchronously
func startResultReceiver(conn *websocket.Conn, taskStarted chan<- bool, taskDone chan<- bool) {
	go func() {
		for {
			_, message, err := conn.ReadMessage()
			if err != nil {
				log.Println("Failed to parse server message: ", err)
				return
			}
			var event Event
			err = json.Unmarshal(message, &event)
			if err != nil {
				log.Println("Failed to parse event: ", err)
				continue
			}
			if handleEvent(conn, event, taskStarted, taskDone) {
				return
			}
		}
	}()
}

// Send the run-task instruction
func sendRunTaskCmd(conn *websocket.Conn) (string, error) {
	runTaskCmd, taskID, err := generateRunTaskCmd()
	if err != nil {
		return "", err
	}
	err = conn.WriteMessage(websocket.TextMessage, []byte(runTaskCmd))
	return taskID, err
}

// Generate the run-task instruction
func generateRunTaskCmd() (string, string, error) {
	taskID := uuid.New().String()
	runTaskCmd := Event{
		Header: Header{
			Action:    "run-task",
			TaskID:    taskID,
			Streaming: "duplex",
		},
		Payload: Payload{
			TaskGroup: "audio",
			Task:      "asr",
			Function:  "recognition",
			Model:     "qwen-audio-3.0-asr-flash-streaming",
			Parameters: Params{
				Format:     "wav",
				SampleRate: 16000,
			},
			Input: Input{},
		},
	}
	runTaskCmdJSON, err := json.Marshal(runTaskCmd)
	return string(runTaskCmdJSON), taskID, err
}

// Wait for the task-started event
func waitForTaskStarted(taskStarted chan bool) {
	select {
	case <-taskStarted:
		fmt.Println("Task started successfully")
	case <-time.After(10 * time.Second):
		log.Fatal("Timed out waiting for task-started; failed to start the task")
	}
}

// Send audio data
func sendAudioData(conn *websocket.Conn) error {
	file, err := os.Open(audioFile)
	if err != nil {
		return err
	}
	defer file.Close()

	buf := make([]byte, 1024)
	for {
		n, err := file.Read(buf)
		if n == 0 {
			break
		}
		if err != nil && err != io.EOF {
			return err
		}
		err = conn.WriteMessage(websocket.BinaryMessage, buf[:n])
		if err != nil {
			return err
		}
		time.Sleep(100 * time.Millisecond)
	}
	return nil
}

// Send the finish-task instruction
func sendFinishTaskCmd(conn *websocket.Conn, taskID string) error {
	finishTaskCmd, err := generateFinishTaskCmd(taskID)
	if err != nil {
		return err
	}
	err = conn.WriteMessage(websocket.TextMessage, []byte(finishTaskCmd))
	return err
}

// Generate the finish-task instruction
func generateFinishTaskCmd(taskID string) (string, error) {
	finishTaskCmd := Event{
		Header: Header{
			Action:    "finish-task",
			TaskID:    taskID,
			Streaming: "duplex",
		},
		Payload: Payload{
			Input: Input{},
		},
	}
	finishTaskCmdJSON, err := json.Marshal(finishTaskCmd)
	return string(finishTaskCmdJSON), err
}

// Handle events
func handleEvent(conn *websocket.Conn, event Event, taskStarted chan<- bool, taskDone chan<- bool) bool {
	switch event.Header.Event {
	case "task-started":
		fmt.Println("Received the task-started event")
		taskStarted <- true
	case "result-generated":
		if event.Payload.Output.Sentence.Text != "" {
			fmt.Println("Recognition result: ", event.Payload.Output.Sentence.Text)
		}
		if event.Payload.Usage != nil {
			fmt.Println("Task billing duration (seconds): ", event.Payload.Usage.Duration)
		}
	case "task-finished":
		fmt.Println("Task finished")
		taskDone <- true
		return true
	case "task-failed":
		handleTaskFailed(event, conn)
		taskDone <- true
		return true
	default:
		log.Printf("Unexpected event: %v", event)
	}
	return false
}

// Handle the task-failed event
func handleTaskFailed(event Event, conn *websocket.Conn) {
	if event.Header.ErrorMessage != "" {
		log.Fatalf("Task failed: %s", event.Header.ErrorMessage)
	} else {
		log.Fatal("The task failed for an unknown reason")
	}
}

// Close the connection
func closeConnection(conn *websocket.Conn) {
	if conn != nil {
		conn.Close()
	}
}

Qwen3-ASR-Flash-Realtime

Note

The example code reads your_audio_file.pcm (PCM16, 16 kHz, mono). If you only have an MP3, WAV, or similar format, convert it with ffmpeg:

ffmpeg -i your_audio.mp3 -ar 16000 -ac 1 -f s16le your_audio_file.pcm

Python

Before you run the example, install the dependencies with the following commands:

pip uninstall websocket-client
pip uninstall websocket
pip install websocket-client

Do not name the example file websocket.py. This name conflicts with the websocket library and causes the following error: AttributeError: module 'websocket' has no attribute 'WebSocketApp'. Did you mean: 'WebSocket'?.

# pip install websocket-client
import os
import time
import json
import threading
import base64
import websocket
import logging
import logging.handlers
from datetime import datetime

logger = logging.getLogger(__name__)
logger.setLevel(logging.DEBUG)

# The API Key differs between the Singapore and Beijing regions. Get your API Key: https://help.aliyun.com/zh/model-studio/get-api-key
# If you have not configured the environment variable, replace the following line with your Alibaba Cloud Model Studio API Key: API_KEY="sk-xxx"
API_KEY = os.environ.get("DASHSCOPE_API_KEY", "sk-xxx")
QWEN_MODEL = "qwen3-asr-flash-realtime"
# The following is the configuration for the China (Beijing) region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
baseUrl = "wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/realtime"
url = f"{baseUrl}?model={QWEN_MODEL}"
print(f"Connecting to server: {url}")

# Note: In non-VAD mode, it is recommended that the cumulative duration of continuously sent audio does not exceed 60s
enableServerVad = True
is_running = True  # Add a running flag

headers = [
    "Authorization: Bearer " + API_KEY,
    "OpenAI-Beta: realtime=v1"
]

def init_logger():
    formatter = logging.Formatter('%(asctime)s|%(levelname)s|%(message)s')
    f_handler = logging.handlers.RotatingFileHandler(
        "omni_tester.log", maxBytes=100 * 1024 * 1024, backupCount=3
    )
    f_handler.setLevel(logging.DEBUG)
    f_handler.setFormatter(formatter)

    console = logging.StreamHandler()
    console.setLevel(logging.DEBUG)
    console.setFormatter(formatter)

    logger.addHandler(f_handler)
    logger.addHandler(console)

def on_open(ws):
    logger.info("Connected to server.")

    # Session update event
    event_manual = {
        "event_id": "event_123",
        "type": "session.update",
        "session": {
            "modalities": ["text"],
            "input_audio_format": "pcm",
            "sample_rate": 16000,
            # "input_audio_transcription": {
            #     # Language identifier, optional; recommended to set it if the language is known
            #     "language": "zh"
            # },
            "turn_detection": None
        }
    }
    event_vad = {
        "event_id": "event_123",
        "type": "session.update",
        "session": {
            "modalities": ["text"],
            "input_audio_format": "pcm",
            "sample_rate": 16000,
            # "input_audio_transcription": {
            #     "language": "zh"
            # },
            "turn_detection": {
                "type": "server_vad",
                "threshold": 0.0,
                "silence_duration_ms": 400
            }
        }
    }
    if enableServerVad:
        logger.info(f"Sending event: {json.dumps(event_vad, indent=2)}")
        ws.send(json.dumps(event_vad))
    else:
        logger.info(f"Sending event: {json.dumps(event_manual, indent=2)}")
        ws.send(json.dumps(event_manual))

def on_message(ws, message):
    global is_running
    try:
        data = json.loads(message)
        logger.info(f"Received event: {json.dumps(data, ensure_ascii=False, indent=2)}")
        if data.get("type") == "conversation.item.input_audio_transcription.completed":
            logger.info(f"Final transcript: {data.get('transcript')}")
        elif data.get("type") == "session.finished":
            logger.info("Closing WebSocket connection after session finished...")
            is_running = False  # Stop the audio sending thread
            ws.close()
    except json.JSONDecodeError:
        logger.error(f"Failed to parse message: {message}")

def on_error(ws, error):
    logger.error(f"Error: {error}")

def on_close(ws, close_status_code, close_msg):
    logger.info(f"Connection closed: {close_status_code} - {close_msg}")

def send_audio(ws, local_audio_path):
    time.sleep(3)  # Wait for the session update to complete
    global is_running

    with open(local_audio_path, 'rb') as audio_file:
        logger.info(f"File reading started: {datetime.now().strftime('%Y-%m-%d %H:%M:%S.%f')[:-3]}")
        while is_running:
            audio_data = audio_file.read(3200)  # ~0.1s PCM16/16kHz
            if not audio_data:
                logger.info(f"File reading finished: {datetime.now().strftime('%Y-%m-%d %H:%M:%S.%f')[:-3]}")
                if ws.sock and ws.sock.connected:
                    if not enableServerVad:
                        commit_event = {
                            "event_id": "event_789",
                            "type": "input_audio_buffer.commit"
                        }
                        ws.send(json.dumps(commit_event))
                    finish_event = {
                        "event_id": "event_987",
                        "type": "session.finish"
                    }
                    ws.send(json.dumps(finish_event))
                break

            if not ws.sock or not ws.sock.connected:
                logger.info("WebSocket is closed, stopping audio sending.")
                break

            encoded_data = base64.b64encode(audio_data).decode('utf-8')
            eventd = {
                "event_id": f"event_{int(time.time() * 1000)}",
                "type": "input_audio_buffer.append",
                "audio": encoded_data
            }
            ws.send(json.dumps(eventd))
            logger.info(f"Sending audio event: {eventd['event_id']}")
            time.sleep(0.1)  # Simulate real-time capture

# Initialize logging
init_logger()
logger.info(f"Connecting to WebSocket server at {url}...")

local_audio_path = "your_audio_file.pcm"
ws = websocket.WebSocketApp(
    url,
    header=headers,
    on_open=on_open,
    on_message=on_message,
    on_error=on_error,
    on_close=on_close
)

thread = threading.Thread(target=send_audio, args=(ws, local_audio_path))
thread.start()
ws.run_forever()

Java

Before you run the example, install the Java-WebSocket dependency:

Maven

<dependency>
    <groupId>org.java-websocket</groupId>
    <artifactId>Java-WebSocket</artifactId>
    <version>1.5.6</version>
</dependency>

Gradle

implementation 'org.java-websocket:Java-WebSocket:1.5.6'
import org.java_websocket.client.WebSocketClient;
import org.java_websocket.handshake.ServerHandshake;
import org.json.JSONObject;

import java.net.URI;
import java.nio.file.Files;
import java.nio.file.Paths;
import java.util.Base64;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.logging.*;

public class QwenASRRealtimeClient {

    private static final Logger logger = Logger.getLogger(QwenASRRealtimeClient.class.getName());
    // The API Key differs between the Singapore and Beijing regions. Get your API Key: https://help.aliyun.com/zh/model-studio/get-api-key
    // If you have not configured the environment variable, replace the following line with your Alibaba Cloud Model Studio API Key: private static final String API_KEY = "sk-xxx"
    private static final String API_KEY = System.getenv().getOrDefault("DASHSCOPE_API_KEY", "sk-xxx");
    private static final String MODEL = "qwen3-asr-flash-realtime";

    // Controls whether to use VAD mode
    private static final boolean enableServerVad = true;

    private static final AtomicBoolean isRunning = new AtomicBoolean(true);
    private static WebSocketClient client;

    public static void main(String[] args) throws Exception {
        initLogger();

        // The following is the configuration for the China (Beijing) region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
        String baseUrl = "wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/realtime";
        String url = baseUrl + "?model=" + MODEL;
        logger.info("Connecting to server: " + url);

        client = new WebSocketClient(new URI(url)) {
            @Override
            public void onOpen(ServerHandshake handshake) {
                logger.info("Connected to server.");
                sendSessionUpdate();
            }

            @Override
            public void onMessage(String message) {
                try {
                    JSONObject data = new JSONObject(message);
                    String eventType = data.optString("type");

                    logger.info("Received event: " + data.toString(2));

                    // The final recognition result is in the transcription.completed event
                    if ("conversation.item.input_audio_transcription.completed".equals(eventType)) {
                        logger.info("Final transcript: " + data.optString("transcript"));
                    }

                    // On receiving the finished event -> stop the sending thread and close the connection
                    if ("session.finished".equals(eventType)) {
                        logger.info("Closing WebSocket connection after session finished...");

                        isRunning.set(false); // Stop the audio sending thread
                        if (this.isOpen()) {
                            this.close(1000, "ASR finished");
                        }
                    }
                } catch (Exception e) {
                    logger.severe("Failed to parse message: " + message);
                }
            }

            @Override
            public void onClose(int code, String reason, boolean remote) {
                logger.info("Connection closed: " + code + " - " + reason);
            }

            @Override
            public void onError(Exception ex) {
                logger.severe("Error: " + ex.getMessage());
            }
        };

        // Add request headers
        client.addHeader("Authorization", "Bearer " + API_KEY);
        client.addHeader("OpenAI-Beta", "realtime=v1");

        client.connectBlocking(); // Block until the connection is established

        // Replace with the path to the audio file to be recognized
        String localAudioPath = "your_audio_file.pcm";
        Thread audioThread = new Thread(() -> {
            try {
                sendAudio(localAudioPath);
            } catch (Exception e) {
                logger.severe("Audio sending thread error: " + e.getMessage());
            }
        });
        audioThread.start();
    }

    /** Session update event (enable/disable VAD) */
    private static void sendSessionUpdate() {
        JSONObject eventNoVad = new JSONObject()
                .put("event_id", "event_123")
                .put("type", "session.update")
                .put("session", new JSONObject()
                        .put("modalities", new String[]{"text"})
                        .put("input_audio_format", "pcm")
                        .put("sample_rate", 16000)
                        // .put("input_audio_transcription", new JSONObject()
                        //         .put("language", "zh"))
                        .put("turn_detection", JSONObject.NULL) // Manual mode
                );

        JSONObject eventVad = new JSONObject()
                .put("event_id", "event_123")
                .put("type", "session.update")
                .put("session", new JSONObject()
                        .put("modalities", new String[]{"text"})
                        .put("input_audio_format", "pcm")
                        .put("sample_rate", 16000)
                        // .put("input_audio_transcription", new JSONObject()
                        //         .put("language", "zh"))
                        .put("turn_detection", new JSONObject()
                                .put("type", "server_vad")
                                .put("threshold", 0.0)
                                .put("silence_duration_ms", 400))
                );

        if (enableServerVad) {
            logger.info("Sending event (VAD):\n" + eventVad.toString(2));
            client.send(eventVad.toString());
        } else {
            logger.info("Sending event (Manual):\n" + eventNoVad.toString(2));
            client.send(eventNoVad.toString());
        }
    }

    /** Send the audio file stream */
    private static void sendAudio(String localAudioPath) throws Exception {
        Thread.sleep(3000); // Wait for the session to be ready
        byte[] allBytes = Files.readAllBytes(Paths.get(localAudioPath));
        logger.info("File reading started");

        int offset = 0;
        while (isRunning.get() && offset < allBytes.length) {
            int chunkSize = Math.min(3200, allBytes.length - offset);
            byte[] chunk = new byte[chunkSize];
            System.arraycopy(allBytes, offset, chunk, 0, chunkSize);
            offset += chunkSize;

            if (client != null && client.isOpen()) {
                String encoded = Base64.getEncoder().encodeToString(chunk);
                JSONObject eventd = new JSONObject()
                        .put("event_id", "event_" + System.currentTimeMillis())
                        .put("type", "input_audio_buffer.append")
                        .put("audio", encoded);

                client.send(eventd.toString());
                logger.info("Sending audio event: " + eventd.getString("event_id"));
            } else {
                break; // Avoid continuing to send after disconnection
            }

            Thread.sleep(100); // Simulate real-time sending
        }

        logger.info("File reading finished");

        if (client != null && client.isOpen()) {
            // Commit is required in non-VAD mode
            if (!enableServerVad) {
                JSONObject commitEvent = new JSONObject()
                        .put("event_id", "event_789")
                        .put("type", "input_audio_buffer.commit");
                client.send(commitEvent.toString());
                logger.info("Sent commit event for manual mode.");
            }

            JSONObject finishEvent = new JSONObject()
                    .put("event_id", "event_987")
                    .put("type", "session.finish");
            client.send(finishEvent.toString());
            logger.info("Sent finish event.");
        }
    }

    /** Initialize logging */
    private static void initLogger() {
        logger.setLevel(Level.ALL);
        Logger rootLogger = Logger.getLogger("");
        for (Handler h : rootLogger.getHandlers()) {
            rootLogger.removeHandler(h);
        }

        Handler consoleHandler = new ConsoleHandler();
        consoleHandler.setLevel(Level.ALL);
        consoleHandler.setFormatter(new SimpleFormatter());
        logger.addHandler(consoleHandler);
    }
}

Node.js

Before you run the example, install the dependencies with the following command:

npm install ws
/**
 * Qwen-ASR Realtime WebSocket client (Node.js version)
 * Features:
 * - Supports VAD mode and Manual mode
 * - Sends session.update to start the session
 * - Continuously sends audio chunks via input_audio_buffer.append
 * - In Manual mode, sends input_audio_buffer.commit
 * - Sends the session.finish event
 * - Closes the connection after receiving the session.finished event
 */

import WebSocket from 'ws';
import fs from 'fs';

// ===== Configuration =====
// The API Key differs between the Singapore and Beijing regions. Get your API Key: https://help.aliyun.com/zh/model-studio/get-api-key
// If you have not configured the environment variable, replace the following line with your Alibaba Cloud Model Studio API Key: const API_KEY = "sk-xxx"
const API_KEY = process.env.DASHSCOPE_API_KEY || 'sk-xxx';
const MODEL = 'qwen3-asr-flash-realtime';
const enableServerVad = true; // true for VAD mode, false for Manual mode
const localAudioPath = 'your_audio_file.pcm'; // Path to the PCM16, 16kHz audio file

// The following is the configuration for the China (Beijing) region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
const baseUrl = 'wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/realtime';
const url = `${baseUrl}?model=${MODEL}`;

console.log(`Connecting to server: ${url}`);

// ===== State control =====
let isRunning = true;

// ===== Establish the connection =====
const ws = new WebSocket(url, {
    headers: {
        'Authorization': `Bearer ${API_KEY}`,
        'OpenAI-Beta': 'realtime=v1'
    }
});

// ===== Event bindings =====
ws.on('open', () => {
    console.log('[WebSocket] Connected to server.');
    sendSessionUpdate();
    // Start the audio sending thread
    sendAudio(localAudioPath);
});

ws.on('message', (message) => {
    try {
        const data = JSON.parse(message);
        console.log('[Received Event]:', JSON.stringify(data, null, 2));

        // The final recognition result is in the transcription.completed event
        if (data.type === 'conversation.item.input_audio_transcription.completed') {
            console.log(`[Final Transcript] ${data.transcript}`);
        }

        // Received the finished event
        if (data.type === 'session.finished') {
            console.log('[Action] Closing WebSocket connection after session finished...');

            if (ws.readyState === WebSocket.OPEN) {
                ws.close(1000, 'ASR finished');
            }
        }
    } catch (e) {
        console.error('[Error] Failed to parse message:', message);
    }
});

ws.on('close', (code, reason) => {
    console.log(`[WebSocket] Connection closed: ${code} - ${reason}`);
});

ws.on('error', (err) => {
    console.error('[WebSocket Error]', err);
});

// ===== Session update =====
function sendSessionUpdate() {
    const eventNoVad = {
        event_id: 'event_123',
        type: 'session.update',
        session: {
            modalities: ['text'],
            input_audio_format: 'pcm',
            sample_rate: 16000,
            // input_audio_transcription: {
            //     language: 'zh'
            // },
            turn_detection: null
        }
    };

    const eventVad = {
        event_id: 'event_123',
        type: 'session.update',
        session: {
            modalities: ['text'],
            input_audio_format: 'pcm',
            sample_rate: 16000,
            // input_audio_transcription: {
            //     language: 'zh'
            // },
            turn_detection: {
                type: 'server_vad',
                threshold: 0.0,
                silence_duration_ms: 400
            }
        }
    };

    if (enableServerVad) {
        console.log('[Send Event] VAD Mode:\n', JSON.stringify(eventVad, null, 2));
        ws.send(JSON.stringify(eventVad));
    } else {
        console.log('[Send Event] Manual Mode:\n', JSON.stringify(eventNoVad, null, 2));
        ws.send(JSON.stringify(eventNoVad));
    }
}

// ===== Send the audio file stream =====
function sendAudio(audioPath) {
    setTimeout(() => {
        console.log(`[File Read Start] ${audioPath}`);
        const buffer = fs.readFileSync(audioPath);

        let offset = 0;
        const chunkSize = 3200; // About 0.1s of PCM16 audio

        function sendChunk() {
            if (!isRunning) return;
            if (offset >= buffer.length) {
                isRunning = false; // Stop sending audio
                console.log('[File Read End]');
                if (ws.readyState === WebSocket.OPEN) {
                    if (!enableServerVad) {
                        const commitEvent = {
                            event_id: 'event_789',
                            type: 'input_audio_buffer.commit'
                        };
                        ws.send(JSON.stringify(commitEvent));
                        console.log('[Send Commit Event]');
                    }

                    const finishEvent = {
                        event_id: 'event_987',
                        type: 'session.finish'
                    };
                    ws.send(JSON.stringify(finishEvent));
                    console.log('[Send Finish Event]');
                }
                
                return;
            }

            if (ws.readyState !== WebSocket.OPEN) {
                console.log('[Stop] WebSocket is not open.');
                return;
            }

            const chunk = buffer.slice(offset, offset + chunkSize);
            offset += chunkSize;

            const encoded = chunk.toString('base64');
            const appendEvent = {
                event_id: `event_${Date.now()}`,
                type: 'input_audio_buffer.append',
                audio: encoded
            };

            ws.send(JSON.stringify(appendEvent));
            console.log(`[Send Audio Event] ${appendEvent.event_id}`);

            setTimeout(sendChunk, 100); // Simulate real-time sending
        }

        sendChunk();
    }, 3000); // Wait for the session configuration to complete
}

C#

The example code is as follows:

using System.Net.WebSockets;
using System.Text;
using System.Text.Json.Nodes;

class Program {
    private static ClientWebSocket _webSocket = new ClientWebSocket();
    private static CancellationTokenSource _cts = new CancellationTokenSource();
    private static bool _sessionFinished = false;
    private static bool _isRunning = true;

    // Controls whether to use VAD mode
    private const bool EnableServerVad = true;

    // The API Key differs between the Singapore and Beijing regions. Get your API Key: https://help.aliyun.com/zh/model-studio/get-api-key
    // If you have not configured the environment variable, replace the following line with your Alibaba Cloud Model Studio API Key: private static readonly string ApiKey = "sk-xxx"
    private static readonly string ApiKey = Environment.GetEnvironmentVariable("DASHSCOPE_API_KEY") ?? throw new InvalidOperationException("DASHSCOPE_API_KEY environment variable is not set.");
    private const string Model = "qwen3-asr-flash-realtime";
    // The following is the configuration for the China (Beijing) region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
    private const string BaseUrl = "wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/realtime";
    private const string AudioFilePath = "your_audio_file.pcm"; // Replace with the path to your PCM audio file

    static async Task Main(string[] args) {
        var url = $"{BaseUrl}?model={Model}";
        Console.WriteLine($"Connecting to server: {url}");

        // Set authentication headers
        _webSocket.Options.SetRequestHeader("Authorization", $"Bearer {ApiKey}");
        _webSocket.Options.SetRequestHeader("OpenAI-Beta", "realtime=v1");

        await _webSocket.ConnectAsync(new Uri(url), _cts.Token);
        Console.WriteLine("Connected to server.");

        // Start the message receiving task
        var receiveTask = ReceiveMessagesAsync();

        // Send the session.update configuration
        await SendSessionUpdateAsync();

        // Send the audio stream
        await SendAudioStreamAsync();

        // Wait for the session.finished event
        while (!_sessionFinished && !_cts.IsCancellationRequested) {
            await Task.Delay(100);
        }

        if (_webSocket.State == WebSocketState.Open) {
            await _webSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, "ASR finished", _cts.Token);
        }
    }

    private static async Task SendAsync(string text) {
        var bytes = Encoding.UTF8.GetBytes(text);
        await _webSocket.SendAsync(new ArraySegment<byte>(bytes), WebSocketMessageType.Text, true, _cts.Token);
    }

    // Send the session.update event
    private static async Task SendSessionUpdateAsync() {
        var session = new JsonObject {
            ["modalities"] = new JsonArray { "text" },
            ["input_audio_format"] = "pcm",
            ["sample_rate"] = 16000,
            // ["input_audio_transcription"] = new JsonObject { ["language"] = "zh" }
        };
        if (EnableServerVad) {
            session["turn_detection"] = new JsonObject {
                ["type"] = "server_vad",
                ["threshold"] = 0.0,
                ["silence_duration_ms"] = 400
            };
        } else {
            session["turn_detection"] = null;
        }
        var payload = new JsonObject {
            ["event_id"] = "event_123",
            ["type"] = "session.update",
            ["session"] = session
        };
        Console.WriteLine($"Sending session.update: {payload.ToJsonString()}");
        await SendAsync(payload.ToJsonString());
    }

    // Send the audio stream (send one PCM chunk every 100ms)
    private static async Task SendAudioStreamAsync() {
        await Task.Delay(3000); // Wait for the session configuration to complete
        const int chunkSize = 3200; // 100ms @ 16kHz 16bit mono
        using var fs = new FileStream(AudioFilePath, FileMode.Open, FileAccess.Read);
        var buffer = new byte[chunkSize];
        int read;
        while (_isRunning && (read = await fs.ReadAsync(buffer, 0, chunkSize)) > 0) {
            if (_webSocket.State != WebSocketState.Open) break;
            string b64 = Convert.ToBase64String(buffer, 0, read);
            var append = new JsonObject {
                ["event_id"] = $"event_{DateTimeOffset.Now.ToUnixTimeMilliseconds()}",
                ["type"] = "input_audio_buffer.append",
                ["audio"] = b64
            };
            await SendAsync(append.ToJsonString());
            await Task.Delay(100);
        }
        Console.WriteLine("File read end.");
        if (_webSocket.State == WebSocketState.Open) {
            if (!EnableServerVad) {
                var commit = new JsonObject {
                    ["event_id"] = "event_789",
                    ["type"] = "input_audio_buffer.commit"
                };
                await SendAsync(commit.ToJsonString());
            }
            var finish = new JsonObject {
                ["event_id"] = "event_987",
                ["type"] = "session.finish"
            };
            await SendAsync(finish.ToJsonString());
        }
    }

    // Receive and handle server-side events
    private static async Task ReceiveMessagesAsync() {
        var buffer = new byte[16384];
        var sb = new StringBuilder();
        while (_webSocket.State == WebSocketState.Open && !_cts.IsCancellationRequested) {
            try {
                var result = await _webSocket.ReceiveAsync(new ArraySegment<byte>(buffer), _cts.Token);
                if (result.MessageType == WebSocketMessageType.Close) {
                    await _webSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, "Closing", _cts.Token);
                    break;
                }
                sb.Append(Encoding.UTF8.GetString(buffer, 0, result.Count));
                if (!result.EndOfMessage) continue;
                string text = sb.ToString();
                sb.Clear();
                var data = JsonNode.Parse(text);
                string? type = data?["type"]?.GetValue<string>();
                Console.WriteLine($"Received event: {type}");
                if (type == "conversation.item.input_audio_transcription.completed") {
                    Console.WriteLine($"Final transcript: {data!["transcript"]}");
                } else if (type == "session.finished") {
                    Console.WriteLine("Session finished, closing...");
                    _sessionFinished = true;
                    _isRunning = false;
                    break;
                }
            } catch (Exception ex) {
                Console.WriteLine($"Receive error: {ex.Message}");
                break;
            }
        }
    }
}

PHP

The example project has the following directory structure:

my-php-project/

├── composer.json

├── vendor/

└── index.php

The contents of composer.json are as follows. Adjust the dependency versions as needed:

{
    "require": {
        "react/event-loop": "^1.3",
        "react/socket": "^1.11",
        "ratchet/pawl": "^0.4"
    }
}

The contents of index.php are as follows:

<?php

require __DIR__ . '/vendor/autoload.php';

use Ratchet\Client\Connector;
use React\EventLoop\Loop;
use React\Socket\Connector as SocketConnector;

// The API Key differs between the Singapore and Beijing regions. Get your API Key: https://help.aliyun.com/zh/model-studio/get-api-key
// If you have not configured the environment variable, replace the following line with your Alibaba Cloud Model Studio API Key: $api_key = "sk-xxx"
$api_key = getenv("DASHSCOPE_API_KEY");
$model = 'qwen3-asr-flash-realtime';
// The following is the configuration for the China (Beijing) region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
$base_url = 'wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/realtime';
$websocket_url = $base_url . '?model=' . $model;
$audio_file_path = 'your_audio_file.pcm'; // Replace with the path to your PCM audio file

// Controls whether to use VAD mode
$enable_server_vad = true;

$loop = Loop::get();
$socketConnector = new SocketConnector($loop, [
    'tls' => ['verify_peer' => false, 'verify_peer_name' => false],
]);
$connector = new Connector($loop, $socketConnector);

$headers = [
    'Authorization' => 'Bearer ' . $api_key,
    'OpenAI-Beta' => 'realtime=v1',
];

$is_running = true;

$connector($websocket_url, [], $headers)->then(function ($conn) use ($loop, $audio_file_path, $enable_server_vad, &$is_running) {
    echo "Connected to WebSocket server\n";

    // Listen for server-side events
    $conn->on('message', function($msg) use ($conn, &$is_running) {
        $event = json_decode($msg, true);
        if (!isset($event['type'])) {
            return;
        }
        echo "Received event: {$event['type']}\n";
        if ($event['type'] === 'conversation.item.input_audio_transcription.completed') {
            echo "Final transcript: {$event['transcript']}\n";
        } elseif ($event['type'] === 'session.finished') {
            echo "Session finished, closing...\n";
            $is_running = false;
            $conn->close();
        }
    });
    $conn->on('close', function() {
        echo "Connection closed\n";
    });

    // Send the session.update event
    sendSessionUpdate($conn, $enable_server_vad);

    // Start sending audio after the session configuration completes
    $loop->addTimer(3, function () use ($conn, $audio_file_path, $enable_server_vad, $loop, &$is_running) {
        sendAudioStream($conn, $audio_file_path, $enable_server_vad, $loop, $is_running);
    });
}, function ($e) {
    echo "Unable to connect: {$e->getMessage()}\n";
});

$loop->run();

// Send the session.update event
function sendSessionUpdate($conn, $enable_server_vad) {
    $session = [
        'modalities' => ['text'],
        'input_audio_format' => 'pcm',
        'sample_rate' => 16000,
        // 'input_audio_transcription' => ['language' => 'zh'],
        'turn_detection' => $enable_server_vad ? [
            'type' => 'server_vad',
            'threshold' => 0.0,
            'silence_duration_ms' => 400,
        ] : null,
    ];
    $event = [
        'event_id' => 'event_123',
        'type' => 'session.update',
        'session' => $session,
    ];
    $conn->send(json_encode($event));
    echo "Sent session.update\n";
}

// Send the audio stream (send one PCM chunk every 100ms)
function sendAudioStream($conn, $audio_file_path, $enable_server_vad, $loop, &$is_running) {
    $fp = fopen($audio_file_path, 'rb');
    if (!$fp) {
        echo "Unable to open the audio file\n";
        return;
    }
    $send_chunk = function() use ($conn, $fp, $enable_server_vad, $loop, &$send_chunk, &$is_running) {
        if (!$is_running) {
            fclose($fp);
            return;
        }
        $chunk = fread($fp, 3200); // 100ms @ 16kHz 16bit mono
        if ($chunk === false || strlen($chunk) === 0) {
            fclose($fp);
            echo "Audio stream ended\n";
            if (!$enable_server_vad) {
                $conn->send(json_encode([
                    'event_id' => 'event_789',
                    'type' => 'input_audio_buffer.commit',
                ]));
            }
            $conn->send(json_encode([
                'event_id' => 'event_987',
                'type' => 'session.finish',
            ]));
            return;
        }
        $append = [
            'event_id' => 'event_' . round(microtime(true) * 1000),
            'type' => 'input_audio_buffer.append',
            'audio' => base64_encode($chunk),
        ];
        $conn->send(json_encode($append));
        $loop->addTimer(0.1, $send_chunk);
    };
    $send_chunk();
}

Go

Before you run the example, install the required dependency:

go get github.com/gorilla/websocket
package main

import (
	"encoding/base64"
	"encoding/json"
	"fmt"
	"io"
	"log"
	"net/http"
	"os"
	"time"

	"github.com/gorilla/websocket"
)

const (
	// The following is the configuration for the China (Beijing) region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
	baseURL         = "wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/realtime"
	model           = "qwen3-asr-flash-realtime"
	audioFile       = "your_audio_file.pcm" // Replace with the path to your PCM audio file
	enableServerVad = true                  // Controls whether to use VAD mode
)

// Server-side event structure
type ServerEvent struct {
	Type       string `json:"type"`
	Transcript string `json:"transcript,omitempty"`
}

func main() {
	// The API Key differs between the Singapore and Beijing regions. Get your API Key: https://help.aliyun.com/zh/model-studio/get-api-key
	// If you have not configured the environment variable, replace the following line with your Alibaba Cloud Model Studio API Key: apiKey := "sk-xxx"
	apiKey := os.Getenv("DASHSCOPE_API_KEY")

	url := baseURL + "?model=" + model
	log.Printf("Connecting to server: %s", url)

	conn, err := connect(url, apiKey)
	if err != nil {
		log.Fatal("Failed to connect to WebSocket: ", err)
	}
	defer conn.Close()

	// Start a goroutine to receive messages
	sessionFinished := make(chan bool, 1)
	go receiveMessages(conn, sessionFinished)

	// Send session.update
	if err := sendSessionUpdate(conn); err != nil {
		log.Fatal("Failed to send session.update: ", err)
	}

	// Wait for the session configuration to complete
	time.Sleep(3 * time.Second)

	// Send the audio stream
	if err := sendAudioStream(conn); err != nil {
		log.Fatal("Failed to send audio: ", err)
	}

	// Wait for session.finished
	<-sessionFinished
}

// Establish the WebSocket connection
func connect(url, apiKey string) (*websocket.Conn, error) {
	headers := http.Header{}
	headers.Set("Authorization", "Bearer "+apiKey)
	headers.Set("OpenAI-Beta", "realtime=v1")
	conn, _, err := websocket.DefaultDialer.Dial(url, headers)
	return conn, err
}

// Send the session.update event
func sendSessionUpdate(conn *websocket.Conn) error {
	session := map[string]interface{}{
		"modalities":         []string{"text"},
		"input_audio_format": "pcm",
		"sample_rate":        16000,
		// "input_audio_transcription": map[string]interface{}{
		// 	"language": "zh",
		// },
	}
	if enableServerVad {
		session["turn_detection"] = map[string]interface{}{
			"type":                "server_vad",
			"threshold":           0.0,
			"silence_duration_ms": 400,
		}
	} else {
		session["turn_detection"] = nil
	}
	event := map[string]interface{}{
		"event_id": "event_123",
		"type":     "session.update",
		"session":  session,
	}
	payload, _ := json.Marshal(event)
	log.Printf("Sending session.update: %s", string(payload))
	return conn.WriteMessage(websocket.TextMessage, payload)
}

// Send the audio stream (send one PCM chunk every 100ms)
func sendAudioStream(conn *websocket.Conn) error {
	f, err := os.Open(audioFile)
	if err != nil {
		return err
	}
	defer f.Close()

	chunk := make([]byte, 3200) // 100ms @ 16kHz 16bit mono
	for {
		n, err := f.Read(chunk)
		if n > 0 {
			event := map[string]interface{}{
				"event_id": fmt.Sprintf("event_%d", time.Now().UnixMilli()),
				"type":     "input_audio_buffer.append",
				"audio":    base64.StdEncoding.EncodeToString(chunk[:n]),
			}
			payload, _ := json.Marshal(event)
			if err := conn.WriteMessage(websocket.TextMessage, payload); err != nil {
				return err
			}
			time.Sleep(100 * time.Millisecond)
		}
		if err == io.EOF {
			break
		}
		if err != nil {
			return err
		}
	}
	log.Println("Audio stream ended")
	if !enableServerVad {
		commitEvt := map[string]interface{}{
			"event_id": "event_789",
			"type":     "input_audio_buffer.commit",
		}
		payload, _ := json.Marshal(commitEvt)
		if err := conn.WriteMessage(websocket.TextMessage, payload); err != nil {
			return err
		}
	}
	finishEvt := map[string]interface{}{
		"event_id": "event_987",
		"type":     "session.finish",
	}
	payload, _ := json.Marshal(finishEvt)
	return conn.WriteMessage(websocket.TextMessage, payload)
}

// Receive and handle server-side events
func receiveMessages(conn *websocket.Conn, sessionFinished chan<- bool) {
	for {
		_, msg, err := conn.ReadMessage()
		if err != nil {
			log.Println("Error reading message: ", err)
			sessionFinished <- true
			return
		}
		var evt ServerEvent
		if err := json.Unmarshal(msg, &evt); err != nil {
			log.Println("Error parsing message: ", err)
			continue
		}
		log.Printf("Received event: %s", evt.Type)
		switch evt.Type {
		case "conversation.item.input_audio_transcription.completed":
			log.Printf("Final transcript: %s", evt.Transcript)
		case "session.finished":
			log.Println("Session finished")
			sessionFinished <- true
			return
		}
	}
}

Paraformer

The Paraformer example code is similar to that of Qwen-Audio-3.0-ASR-Flash-Streaming/Fun-ASR-Realtime. Replace the model name with a Paraformer model.

Apply in production

Reuse connections (WebSocket)

The WebSocket connections for Qwen-Audio-3.0-ASR-Flash-Streaming/Fun-ASR-Realtime and Paraformer support reuse: after one recognition task finishes, you can start the next task without reestablishing the connection.

Reuse flow: The client sends finish-task. After the server returns task-finished, the client can send run-task again to start a new task.

Important
  1. Wait for the server to return the task-finished event before starting a new task.

  2. Different tasks over a reused connection must use different task_id values.

  3. When a task fails, the server returns an error event and closes the connection. That connection cannot be reused.

  4. If no new task starts within 60 seconds after a task ends, the connection closes automatically.

Qwen3-ASR-Flash-Realtime uses a session model and does not support connection reuse. Close the connection after each session ends.

For the events of each model, see the corresponding API reference.

High-concurrency best practices

The DashScope SDK includes a built-in pooling mechanism that reuses WebSocket connections and recognition objects, which avoids the overhead of frequent creation and destruction.

Important

Currently, only the Paraformer Java SDK supports this feature.

Click to view high-concurrency best practices

Prerequisites

The Java SDK combines a built-in connection pool with a custom object pool to achieve optimal performance:

  • Connection pool: The OkHttp3 connection pool integrated in the SDK manages and reuses the underlying WebSocket connections, which reduces network handshake overhead. This feature is enabled by default.

  • Object pool: Built on commons-pool2, the object pool maintains a set of Recognition objects whose connections are already established. Borrowing an object from the pool eliminates the connection setup latency and significantly reduces first-packet latency.

Implementation steps

  1. Add dependencies

    Add dashscope-sdk-java and commons-pool2 to your dependency configuration file, based on your project's build tool.

    The following examples show the configuration for Maven and Gradle:

    Maven

    1. Open the pom.xml file of your Maven project.

    2. Add the following dependencies inside the <dependencies> tag.

    <dependency>
        <groupId>com.alibaba</groupId>
        <artifactId>dashscope-sdk-java</artifactId>
        <!-- Replace 'the-latest-version' with version 2.16.9 or later. You can look up version numbers at: https://mvnrepository.com/artifact/com.alibaba/dashscope-sdk-java -->
        <version>the-latest-version</version>
    </dependency>
    
    <dependency>
        <groupId>org.apache.commons</groupId>
        <artifactId>commons-pool2</artifactId>
        <!-- Replace 'the-latest-version' with the latest version. You can look up version numbers at: https://mvnrepository.com/artifact/org.apache.commons/commons-pool2 -->
        <version>the-latest-version</version>
    </dependency>
    1. Save the pom.xml file.

    2. Run a Maven command (such as mvn clean install or mvn compile) to update the project dependencies.

    Gradle

    1. Open the build.gradle file of your Gradle project.

    2. Add the following dependencies inside the dependencies block.

      dependencies {
          // Replace 'the-latest-version' with version 2.16.9 or later. You can look up version numbers at: https://mvnrepository.com/artifact/com.alibaba/dashscope-sdk-java
          implementation group: 'com.alibaba', name: 'dashscope-sdk-java', version: 'the-latest-version'
      
          // Replace 'the-latest-version' with the latest version. You can look up version numbers at: https://mvnrepository.com/artifact/org.apache.commons/commons-pool2
          implementation group: 'org.apache.commons', name: 'commons-pool2', version: 'the-latest-version'
      }
    3. Save the build.gradle file.

    4. On the command line, switch to the project root directory and run the following Gradle command to update the project dependencies.

      ./gradlew build --refresh-dependencies

      On Windows, use the following command:

      gradlew build --refresh-dependencies
  2. Configure the connection pool

    Configure the key connection pool parameters through environment variables:

    Environment variable

    Description

    DASHSCOPE_CONNECTION_POOL_SIZE

    The connection pool size.

    Recommended value: at least twice the peak concurrency.

    Default value: 32.

    DASHSCOPE_MAXIMUM_ASYNC_REQUESTS

    The maximum number of asynchronous requests.

    Recommended value: the same as DASHSCOPE_CONNECTION_POOL_SIZE.

    Default value: 32.

    DASHSCOPE_MAXIMUM_ASYNC_REQUESTS_PER_HOST

    The maximum number of asynchronous requests per host.

    Recommended value: the same as DASHSCOPE_CONNECTION_POOL_SIZE.

    Default value: 32.

  3. Configure the object pool

    Configure the object pool size through an environment variable:

    Environment variable

    Description

    RECOGNITION_OBJECTPOOL_SIZE

    The object pool size.

    Recommended value: 1.5 to 2 times the peak concurrency.

    Default value: 500.

    Important
    • The object pool size (RECOGNITION_OBJECTPOOL_SIZE) must be less than or equal to the connection pool size (DASHSCOPE_CONNECTION_POOL_SIZE). Otherwise, when the object pool requests an object and the connection pool is full, the calling thread blocks.

    • The object pool size must not exceed your account's queries per second (QPS) limit.

    Create the object pool with the following code:

    class RecognitionObjectPool {
        // ... For the full example, see the complete code.
        public static GenericObjectPool<Recognition> getInstance() {
            lock.lock();
            if (recognitionGenericObjectPool == null) {
                int objectPoolSize = getObjectivePoolSize();
                RecognitionObjectFactory recognitionObjectFactory =
                        new RecognitionObjectFactory();
                GenericObjectPoolConfig<Recognition> config =
                        new GenericObjectPoolConfig<>();
                config.setMaxTotal(objectPoolSize);
                config.setMaxIdle(objectPoolSize);
                config.setMinIdle(objectPoolSize);
                recognitionGenericObjectPool =
                        new GenericObjectPool<>(recognitionObjectFactory, config);
            }
            lock.unlock();
            return recognitionGenericObjectPool;
        }
    }
  4. Borrow a Recognition object from the object pool

    When the number of unreturned objects exceeds the object pool limit, the system creates additional Recognition objects. These new objects must reestablish a WebSocket connection and cannot be reused.

    recognizer = RecognitionObjectPool.getInstance().borrowObject();
  5. Perform speech recognition

    Call the call or streamCall method of the Recognition object to perform speech recognition.

  6. Return the Recognition object

    After the speech recognition task finishes, return the Recognition object so that it can be reused. Do not return objects with unfinished or failed tasks.

    RecognitionObjectPool.getInstance().returnObject(recognizer);

Complete code

package org.alibaba.bailian.example.examples;

import com.alibaba.dashscope.audio.asr.recognition.Recognition;
import com.alibaba.dashscope.audio.asr.recognition.RecognitionParam;
import com.alibaba.dashscope.audio.asr.recognition.RecognitionResult;
import com.alibaba.dashscope.common.ResultCallback;
import com.alibaba.dashscope.exception.NoApiKeyException;
import com.alibaba.dashscope.utils.ApiKey;
import org.apache.commons.pool2.BasePooledObjectFactory;
import org.apache.commons.pool2.PooledObject;
import org.apache.commons.pool2.impl.DefaultPooledObject;
import org.apache.commons.pool2.impl.GenericObjectPool;
import org.apache.commons.pool2.impl.GenericObjectPoolConfig;

import java.io.FileInputStream;
import java.nio.ByteBuffer;
import java.nio.file.Path;
import java.nio.file.Paths;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.Lock;
import com.alibaba.dashscope.utils.Constants;

public class Main {
    public static void checkoutEnv(String envName, int defaultSize) {
        if (System.getenv(envName) != null) {
            System.out.println("[ENV CHECK]: " + envName + " "
                    + System.getenv(envName));
        } else {
            System.out.println("[ENV CHECK]: " + envName
                    + " Using Default which is " + defaultSize);
        }
    }

    public static void main(String[] args)
            throws NoApiKeyException, InterruptedException {
        // The following is the configuration for the China (Beijing) region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
        Constants.baseHttpApiUrl = "https://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api/v1";
        checkoutEnv("DASHSCOPE_CONNECTION_POOL_SIZE", 32);
        checkoutEnv("DASHSCOPE_MAXIMUM_ASYNC_REQUESTS", 32);
        checkoutEnv("DASHSCOPE_MAXIMUM_ASYNC_REQUESTS_PER_HOST", 32);
        checkoutEnv(RecognitionObjectPool.RECOGNITION_OBJECTPOOL_SIZE_ENV,
                RecognitionObjectPool.DEFAULT_OBJECT_POOL_SIZE);

        int threadNums = 3;
        String currentDir = System.getProperty("user.dir");
        Path[] filePaths = {
                Paths.get(currentDir, "{YOUR_AUDIO_FILE}"),
                Paths.get(currentDir, "{YOUR_AUDIO_FILE}"),
                Paths.get(currentDir, "{YOUR_AUDIO_FILE}"),
        };
        ExecutorService executorService = Executors.newFixedThreadPool(threadNums);
        for (int i = 0; i < threadNums; i++) {
            executorService.submit(new RealtimeRecognizeTask(filePaths));
        }
        executorService.shutdown();
        executorService.awaitTermination(10, TimeUnit.MINUTES);
        System.exit(0);
    }
}

class RecognitionObjectFactory extends BasePooledObjectFactory<Recognition> {
    public RecognitionObjectFactory() {
        super();
    }

    @Override
    public Recognition create() throws Exception {
        return new Recognition();
    }

    @Override
    public PooledObject<Recognition> wrap(Recognition obj) {
        return new DefaultPooledObject<>(obj);
    }
}

class RecognitionObjectPool {
    public static GenericObjectPool<Recognition> recognitionGenericObjectPool;
    public static String RECOGNITION_OBJECTPOOL_SIZE_ENV =
            "RECOGNITION_OBJECTPOOL_SIZE";
    public static int DEFAULT_OBJECT_POOL_SIZE = 500;
    private static Lock lock = new java.util.concurrent.locks.ReentrantLock();

    public static int getObjectivePoolSize() {
        try {
            Integer n = Integer.parseInt(
                    System.getenv(RECOGNITION_OBJECTPOOL_SIZE_ENV));
            return n;
        } catch (NumberFormatException e) {
            return DEFAULT_OBJECT_POOL_SIZE;
        }
    }

    public static GenericObjectPool<Recognition> getInstance() {
        lock.lock();
        if (recognitionGenericObjectPool == null) {
            int objectPoolSize = getObjectivePoolSize();
            System.out.println("RECOGNITION_OBJECTPOOL_SIZE: "
                    + objectPoolSize);
            RecognitionObjectFactory recognitionObjectFactory =
                    new RecognitionObjectFactory();
            GenericObjectPoolConfig<Recognition> config =
                    new GenericObjectPoolConfig<>();
            config.setMaxTotal(objectPoolSize);
            config.setMaxIdle(objectPoolSize);
            config.setMinIdle(objectPoolSize);
            recognitionGenericObjectPool =
                    new GenericObjectPool<>(recognitionObjectFactory, config);
        }
        lock.unlock();
        return recognitionGenericObjectPool;
    }
}

class RealtimeRecognizeTask implements Runnable {
    private static final Object lock = new Object();
    private Path[] filePaths;

    public RealtimeRecognizeTask(Path[] filePaths) {
        this.filePaths = filePaths;
    }

    private static String getDashScopeApiKey() throws NoApiKeyException {
        String dashScopeApiKey = null;
        try {
            ApiKey apiKey = new ApiKey();
            dashScopeApiKey = ApiKey.getApiKey(null);
        } catch (NoApiKeyException e) {
            System.out.println("No API key found in environment.");
        }
        if (dashScopeApiKey == null) {
            dashScopeApiKey = "your-dashscope-apikey";
        }
        return dashScopeApiKey;
    }

    public void runCallback() {
        for (Path filePath : filePaths) {
            RecognitionParam param = null;
            try {
                param = RecognitionParam.builder()
                        .model("paraformer-realtime-v2")
                        .format("pcm")
                        .sampleRate(16000)
                        .apiKey(getDashScopeApiKey())
                        .build();
            } catch (Exception e) {
                throw new RuntimeException(e);
            }

            Recognition recognizer = null;
            final boolean[] hasError = {false};
            try {
                recognizer = RecognitionObjectPool.getInstance().borrowObject();
                String threadName = Thread.currentThread().getName();

                ResultCallback<RecognitionResult> callback =
                        new ResultCallback<RecognitionResult>() {
                            @Override
                            public void onEvent(RecognitionResult message) {
                                synchronized (lock) {
                                    if (message.isSentenceEnd()) {
                                        System.out.println("[process " + threadName
                                                + "] Fix:" + message.getSentence().getText());
                                    } else {
                                        System.out.println("[process " + threadName
                                                + "] Result: " + message.getSentence().getText());
                                    }
                                }
                            }

                            @Override
                            public void onComplete() {
                                System.out.println("[" + threadName
                                        + "] Recognition complete");
                            }

                            @Override
                            public void onError(Exception e) {
                                System.out.println("[" + threadName
                                        + "] RecognitionCallback error: " + e.getMessage());
                                hasError[0] = true;
                            }
                        };
                System.out.println("[" + threadName
                        + "] Input file_path is: " + filePath);
                FileInputStream fis = null;
                try {
                    fis = new FileInputStream(filePath.toFile());
                } catch (Exception e) {
                    System.out.println("Error when loading file: " + filePath);
                    e.printStackTrace();
                }
                recognizer.call(param, callback);

                // chunk size set to 100 ms for 16KHz sample rate
                byte[] buffer = new byte[3200];
                int bytesRead;
                while ((bytesRead = fis.read(buffer)) != -1) {
                    ByteBuffer byteBuffer;
                    if (bytesRead < buffer.length) {
                        byteBuffer = ByteBuffer.wrap(buffer, 0, bytesRead);
                    } else {
                        byteBuffer = ByteBuffer.wrap(buffer);
                    }
                    recognizer.sendAudioFrame(byteBuffer);
                    Thread.sleep(100);
                    buffer = new byte[3200];
                }
                System.out.println("[" + threadName + "] send audio done");
                recognizer.stop();
                System.out.println("[" + threadName + "] asr task finished");
            } catch (Exception e) {
                e.printStackTrace();
                hasError[0] = true;
            }
            if (recognizer != null) {
                try {
                    if (hasError[0] == true) {
                        recognizer.getDuplexApi().close(1000, "bye");
                        RecognitionObjectPool.getInstance()
                                .invalidateObject(recognizer);
                    } else {
                        RecognitionObjectPool.getInstance()
                                .returnObject(recognizer);
                    }
                } catch (Exception e) {
                    e.printStackTrace();
                }
            }
        }
    }

    @Override
    public void run() {
        runCallback();
    }
}

Recommended configuration

The following configurations are based on test results from running only the Paraformer real-time speech recognition service on Alibaba Cloud servers of the specified specifications. Single-machine concurrency is the number of Paraformer real-time speech recognition tasks running at the same time (that is, the number of worker threads).

Machine specification (Alibaba Cloud)

Max single-machine concurrency

Object pool size

Connection pool size

4 vCPUs, 8 GiB

100

500

2000

8 vCPUs, 16 GiB

200

500

2000

16 vCPUs, 32 GiB

400

500

2000

Resource management and error handling

  • Task succeeds: Call GenericObjectPool.returnObject() to return the Recognition object to the pool for reuse.

    Important

    Do not return Recognition objects with unfinished or failed tasks.

  • Task fails: When an exception thrown by the SDK or your business logic interrupts a task, perform the following two actions:

    1. Actively close the underlying WebSocket connection.

    2. Invalidate the object in the object pool to prevent it from being reused.

    // Close the connection.
    recognizer.getDuplexApi().close(1000, "bye");
    // Invalidate the failed recognizer in the object pool.
    RecognitionObjectPool.getInstance().invalidateObject(recognizer);
  • When the service returns a TaskFailed error, no extra handling is required.

Warm-up and latency measurement

When you evaluate performance such as concurrent call latency for the DashScope Java SDK, we recommend that you run a sufficient warm-up before the formal test.

Connection reuse mechanism

The DashScope Java SDK manages and reuses WebSocket connections through a global singleton connection pool. This mechanism works as follows:

  • On-demand creation: The SDK does not pre-create WebSocket connections at service startup. Instead, it establishes connections on demand at the first call.

  • Time-limited reuse: After a request completes, the connection stays in the pool for up to 60 seconds for reuse.

    • If a new request arrives within 60 seconds, the SDK reuses the existing connection and avoids the overhead of a repeated handshake.

    • If a connection stays idle for more than 60 seconds, the SDK closes it automatically to release resources.

Why warm-up matters

In the following scenarios, the connection pool might not have an active connection to reuse, so a request has to create a new connection:

  • The application has just started and has not made any calls yet.

  • The service has been idle for more than 60 seconds, so pooled connections have closed due to timeout.

In these scenarios, the first or early requests trigger the full WebSocket connection process (including the TCP handshake, TLS negotiation, and protocol upgrade). Their end-to-end latency is significantly higher than that of later requests that reuse connections.

Recommended approach

Before you run a formal load test or measure latency, follow these warm-up steps:

  1. Simulate the concurrency level of the formal test by sending a number of calls in advance (for example, for 1 to 2 minutes) to fully populate the connection pool.

  2. After you confirm that the connection pool has established and maintained enough active connections, start collecting the formal performance data.

Improve recognition accuracy

  • Choose a model that matches the sample rate: For 8 kHz telephone audio, use an 8 kHz model directly. This avoids the information loss caused by upsampling to 16 kHz.

  • Improve the input audio quality: Use a high-quality microphone and record in an environment with a high signal-to-noise ratio and no echo. At the application layer, you can integrate algorithms such as noise reduction (for example, RNNoise) and acoustic echo cancellation (AEC) for preprocessing.

Set up a fault-tolerance strategy

  • Client-side reconnection: The client should implement automatic reconnection to handle network jitter. The following is a reference implementation for the Python SDK:

    1. Catch exceptions: Implement the on_error method in the Callback class. The dashscope SDK calls this method when it encounters a network error or another issue.

    2. Signal the state: When on_error is triggered, set a reconnection signal. In Python, you can use threading.Event, a thread-safe signal flag.

    3. Reconnection loop: Wrap the main logic in a for loop (for example, retry 3 times). When the reconnection signal is detected, the current recognition round is interrupted, resources are cleaned up, and after a few seconds the loop runs again to create a brand-new connection.

  • Set a heartbeat to keep the connection alive: To maintain a long-lived connection with the server, set the heartbeat parameter to true. The connection to the server then stays open even when the audio contains no sound for a long time.

  • Model rate limits: When you call the model API, note the model's Rate limiting rules.

Supported models and regions

China (Beijing)

To call the following models, use an API Key for the China (Beijing) region:

  • Qwen-Audio-3.0-ASR-Flash-Streaming: qwen-audio-3.0-asr-flash-streaming

  • Fun-ASR-Realtime:

    • fun-asr-realtime (stable version, currently equivalent to fun-asr-realtime-2025-11-07), fun-asr-realtime-2026-02-28 (latest snapshot version), fun-asr-realtime-2025-11-07 (snapshot version), fun-asr-realtime-2025-09-15 (snapshot version)

    • fun-asr-flash-8k-realtime (stable version, currently equivalent to fun-asr-flash-8k-realtime-2026-01-28), fun-asr-flash-8k-realtime-2026-01-28

  • Qwen3-ASR-Flash-Realtime: qwen3-asr-flash-realtime (stable version, currently equivalent to qwen3-asr-flash-realtime-2025-10-27), qwen3-asr-flash-realtime-2026-02-10 (latest snapshot version), qwen3-asr-flash-realtime-2025-10-27 (snapshot version)

  • Paraformer: paraformer-realtime-v2, paraformer-realtime-v1, paraformer-realtime-8k-v2, paraformer-realtime-8k-v1

Singapore

To call the following models, use an API Key for the Singapore region:

  • Qwen-Audio-3.0-ASR-Flash-Streaming: qwen-audio-3.0-asr-flash-streaming

  • Fun-ASR-Realtime: fun-asr-realtime (stable version, currently equivalent to fun-asr-realtime-2025-11-07), fun-asr-realtime-2025-11-07 (snapshot version)

  • Qwen3-ASR-Flash-Realtime: qwen3-asr-flash-realtime (stable version, currently equivalent to qwen3-asr-flash-realtime-2025-10-27), qwen3-asr-flash-realtime-2026-02-10 (latest snapshot version), qwen3-asr-flash-realtime-2025-10-27 (snapshot version)

API reference

FAQ

Which audio formats does real-time speech recognition support?

The Qwen-Audio-3.0-ASR-Flash-Streaming, Fun-ASR-Realtime, and Paraformer models support the pcm, wav, mp3, opus, speex, aac, and amr formats. For the Qwen3-ASR-Flash-Realtime model, we recommend the pcm or opus format. Other formats (such as wav, aac, and amr) are accepted by the session.update validation layer, but the server-side decoding might fail. Confirm that the audio stream uses a recommended format before you send it.

What's the difference between the SDK and the WebSocket API, and how do I choose?

The DashScope SDK encapsulates details such as WebSocket connection management, authentication, and reconnection, which makes it a good fit for quick integration. Connecting directly to the WebSocket API provides finer-grained control and suits programming languages that the SDK does not cover or scenarios that require custom connection management. We recommend that you use the SDK first.

How do I improve recognition accuracy for proper nouns?

Use hotwords or context enhancement. For detailed configuration methods and usage notes, see Improve recognition accuracy.

What should I do when the connection drops frequently?

Implement client-side reconnection and enable the heartbeat parameter (heartbeat=true) to prevent the connection from dropping when there is no audio for a long time. For detailed fault-tolerance strategies, see Apply in production.

Model application listing and filing

See Application compliance filing.