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
- An API key is Obtain an API key and set as an environment variable.
- To call the service through the DashScope SDK, install the latest SDK.
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 Singapore region. When calling, replace "{WorkspaceId}" with your real workspace ID. Configurations differ by region.
Constants.baseWebsocketApiUrl = "wss://{WorkspaceId}.ap-southeast-1.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://www.alibabacloud.com/help/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://www.alibabacloud.com/help/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 Singapore region. When calling, replace "{WorkspaceId}" with your actual workspace ID. The configuration differs by region.
dashscope.base_websocket_api_url='wss://{WorkspaceId}.ap-southeast-1.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.
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 Singapore region. When calling, replace "{WorkspaceId}" with your real workspace ID. Configurations differ by region.
Constants.baseWebsocketApiUrl = "wss://{WorkspaceId}.ap-southeast-1.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://www.alibabacloud.com/help/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());
}
}
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://www.alibabacloud.com/help/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 Singapore region. When calling, replace "{WorkspaceId}" with your actual workspace ID. The configuration differs by region.
dashscope.base_websocket_api_url='wss://{WorkspaceId}.ap-southeast-1.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
NoteThe 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
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 Singapore region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
.url("wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/realtime")
// The API Key differs between the Singapore and Beijing regions. Get your API Key: https://www.alibabacloud.com/help/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);
}
}
Constants.baseHttpApiUrl = "https://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api/v1";
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://www.alibabacloud.com/help/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 Singapore region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
url='wss://{WorkspaceId}.ap-southeast-1.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_detectionparameter (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, setsession.turn_detectionto null.
Switch interaction modes:
- WebSocket: Set the
turn_detectionfield in asession.updateevent.
{
"type": "session.update",
"session": {
"turn_detection": null
}
}
- Python SDK: Set the
enable_turn_detectionparameter in theupdate_sessionmethod.
conversation.update_session(
enable_turn_detection=False
)
- Java SDK: Set the
enableTurnDetectionparameter throughOmniRealtimeConfig.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 includessilence_duration_ms(the silence duration threshold that ends a turn when exceeded; server default800, with400recommended for conversation and chat scenarios that need fast segmentation) andthreshold(VAD detection sensitivity; server default0.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_timeandpayload.output.sentence.end_timemark the start and end of a full sentence in the audio. In an intermediate result,end_timemay benulland is filled with the final value when the sentence ends (sentence_end = true). - Word level: The
payload.output.sentence.wordsarray, where each element containsbegin_time,end_time,text(the word or character text), andpunctuation(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 tofalse. 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://www.alibabacloud.com/help/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 Singapore region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
url = 'wss://{WorkspaceId}.ap-southeast-1.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:
<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>
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://www.alibabacloud.com/help/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 Singapore region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
private static final String URL = "wss://{WorkspaceId}.ap-southeast-1.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://www.alibabacloud.com/help/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 Singapore region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
const url = 'wss://{WorkspaceId}.ap-southeast-1.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://www.alibabacloud.com/help/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 Singapore region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
private const string WebSocketUrl = "wss://{WorkspaceId}.ap-southeast-1.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://www.alibabacloud.com/help/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 Singapore region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
$websocket_url = 'wss://{WorkspaceId}.ap-southeast-1.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 Singapore region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
wsURL = "wss://{WorkspaceId}.ap-southeast-1.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://www.alibabacloud.com/help/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
NoteThe 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://www.alibabacloud.com/help/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 Singapore region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
baseUrl = "wss://{WorkspaceId}.ap-southeast-1.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:
<dependency>
<groupId>org.java-websocket</groupId>
<artifactId>Java-WebSocket</artifactId>
<version>1.5.6</version>
</dependency>
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://www.alibabacloud.com/help/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 Singapore region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
String baseUrl = "wss://{WorkspaceId}.ap-southeast-1.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://www.alibabacloud.com/help/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 Singapore region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
const baseUrl = 'wss://{WorkspaceId}.ap-southeast-1.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://www.alibabacloud.com/help/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 Singapore region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
private const string BaseUrl = "wss://{WorkspaceId}.ap-southeast-1.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://www.alibabacloud.com/help/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 Singapore region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
$base_url = 'wss://{WorkspaceId}.ap-southeast-1.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 Singapore region. When calling, replace "{WorkspaceId}" with your actual workspace ID. Configurations differ by region.
baseURL = "wss://{WorkspaceId}.ap-southeast-1.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://www.alibabacloud.com/help/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
- Wait for the server to return the
task-finishedevent before starting a new task. - Different tasks over a reused connection must use different
task_idvalues. - When a task fails, the server returns an error event and closes the connection. That connection cannot be reused.
- 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.
ImportantCurrently, only the Java SDK supports this feature.
Click to view high-concurrency best practices
Prerequisites
- Obtain an API key
- The DashScope SDK is installed and meets the version requirement. We recommend that you install the latest version: Java SDK version 2.16.9 or later.
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 ofRecognitionobjects whose connections are already established. Borrowing an object from the pool eliminates the connection setup latency and significantly reduces first-packet latency.
Implementation steps
-
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
- Open the
pom.xmlfile of your Maven project. - 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>- Save the
pom.xmlfile. - Run a Maven command (such as
mvn clean installormvn compile) to update the project dependencies.
Gradle
- Open the
build.gradlefile of your Gradle project. - Add the following dependencies inside the
dependenciesblock.
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' }- Save the
build.gradlefile. - On the command line, switch to the project root directory and run the following Gradle command to update the project dependencies.
./gradlew build --refresh-dependenciesOn Windows, use the following command:
gradlew build --refresh-dependencies - Open the
-
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.
-
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:
- The object pool size (
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;
}
}
-
Borrow a Recognition object from the object pool
When the number of unreturned objects exceeds the object pool limit, the system creates additional
Recognitionobjects. These new objects must reestablish a WebSocket connection and cannot be reused.
recognizer = RecognitionObjectPool.getInstance().borrowObject();
-
Perform speech recognition
Call the call or streamCall method of the
Recognitionobject to perform speech recognition. -
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.ImportantDo 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:
- Actively close the underlying WebSocket connection.
- 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:
- 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.
- 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:
- Catch exceptions: Implement the
on_errormethod in theCallbackclass. ThedashscopeSDK calls this method when it encounters a network error or another issue. - Signal the state: When
on_erroris triggered, set a reconnection signal. In Python, you can usethreading.Event, a thread-safe signal flag. - Reconnection loop: Wrap the main logic in a
forloop (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.
- Catch exceptions: Implement the
-
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
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)
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
API reference
- Real-time speech recognition - Qwen-Audio-3.0-ASR-Flash-Streaming/Fun-ASR-Realtime API reference
- Real-time speech recognition - Qwen3-ASR-Flash-Realtime API reference
- Real-time speech recognition - Paraformer API reference
- AOQ client SDK (for Qwen-Audio-3.0-ASR-Flash-Streaming/Fun-ASR-Realtime)
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.