すべてのプロダクト
Search
ドキュメントセンター

Alibaba Cloud Model Studio:Real-time speech recognition - Qwen

最終更新日:Jul 16, 2026

リアルタイム音声認識は、WebSocket を介して音声ストリームを句読点付きのテキストに文字起こしする機能で、ライブキャプション、オンライン会議、音声チャット、スマートアシスタントに適した低レイテンシを実現します。

概要

WebSocket ストリーミングプロトコルは、音声をサービスに配信し、文字起こしされたテキストを低レイテンシで返します。

  • 標準中国語および広東語や四川語などの方言に対する高精度な認識。

  • 自動言語検出と非音声フィルタリングにより、複雑な音響環境でも安定したパフォーマンスを発揮。

  • 驚き、平静、喜び、悲しみ、嫌悪、怒り、恐怖などの状態にわたる感情認識。

  • 指定した用語の認識精度を向上させるカスタムホットワード。

  • 構造化された認識結果のためのタイムスタンプ出力。

  • さまざまな録音環境に合わせて設定可能なサンプルレートと複数の音声フォーマット。

会議の文字起こし、通話分析、字幕生成などのバッチシナリオでは、非リアルタイム音声認識をご利用ください。モデル選択のガイダンスについては、音声テキスト変換をご参照ください。

前提条件

クイックスタート

以下の例は、DashScope SDK を介してリアルタイム音声認識を呼び出す方法を示しています。

Fun-ASR

マイク音声の認識

マイクから音声をキャプチャし、リアルタイムで文字起こしされたテキストをストリーミングすることで、話者が話すと同時にテキストが表示されます。

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 {
        // 次の設定はシンガポールリージョン用です。「{WorkspaceId}」を実際のワークスペース ID に置き換えてください。設定はリージョンによって異なります。
        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("fun-asr-realtime")
                // シンガポールリージョンと北京リージョンでは API キーが異なります。API キーを取得するには、https://www.alibabacloud.com/help/model-studio/get-api-key をご参照ください。
                // 環境変数を設定していない場合は、次の行を Model Studio API キーに置き換えてください: .apiKey("sk-xxx")
                .apiKey(System.getenv("DASHSCOPE_API_KEY"))
                .format("wav")
                .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);
            // 音声フォーマットを作成します。
            AudioFormat audioFormat = new AudioFormat(16000, 16, 1, true, false);
            // フォーマットに基づいてデフォルトの録音デバイスを照合します。
            TargetDataLine targetDataLine =
                    AudioSystem.getTargetDataLine(audioFormat);
            targetDataLine.open(audioFormat);
            // 録音を開始します。
            targetDataLine.start();
            ByteBuffer buffer = ByteBuffer.allocate(1024);
            long start = System.currentTimeMillis();
            // 50 秒間録音し、リアルタイムで文字起こしを行います。
            while (System.currentTimeMillis() - start < 50000) {
                int read = targetDataLine.read(buffer.array(), 0, buffer.capacity());
                if (read > 0) {
                    buffer.limit(read);
                    // 録音した音声データをストリーミング認識サービスに送信します。
                    recognizer.sendAudioFrame(buffer);
                    buffer = ByteBuffer.allocate(1024);
                    // 録音レートは制限されています。CPU 使用率が高くなるのを防ぐため、短時間スリープします。
                    Thread.sleep(20);
                }
            }
            recognizer.stop();
        } catch (Exception e) {
            e.printStackTrace();
        } finally {
            // タスク完了後、WebSocket 接続を閉じます。
            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

Python の例では、音声キャプチャに pyaudio ライブラリが必要です。例を実行する前に、pip install pyaudio でインストールしてください。

import os
import signal  # キーボードイベント処理用 (「Ctrl+C」を押して録音を終了)
import sys

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

mic = None
stream = None

# 録音パラメータを設定
sample_rate = 16000  # サンプリングレート (Hz)
channels = 1  # モノラルチャンネル
dtype = 'int16'  # データ型
format_pcm = 'pcm'  # 音声データのフォーマット
block_size = 3200  # バッファごとのフレーム数

# リアルタイム音声認識コールバック
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.')  # 認識完了

    def on_error(self, message) -> None:
        print('RecognitionCallback task_id: ', message.request_id)
        print('RecognitionCallback error: ', message.message)
        # オーディオストリームが実行中の場合は停止して閉じる
        if 'stream' in globals() and stream.is_active():
            stream.stop_stream()
            stream.close()
        # プログラムを強制的に終了
        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 ...')
    # 認識を停止
    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(),
        ))
    # プログラムを強制的に終了
    sys.exit(0)

# main 関数
if __name__ == '__main__':
    # シンガポールリージョンと北京リージョンでは API キーが異なります。API キーを取得するには、https://www.alibabacloud.com/help/model-studio/get-api-key をご参照ください
    # 環境変数を設定していない場合は、次の行を Model Studio API キーに置き換えてください: dashscope.api_key = "sk-xxx"
    dashscope.api_key = os.environ.get('DASHSCOPE_API_KEY')

    # 次の設定はシンガポールリージョン用です。「{WorkspaceId}」を実際のワークスペース ID に置き換えてください。設定はリージョンによって異なります。
    dashscope.base_websocket_api_url='wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/inference'

    # 認識コールバックを作成
    callback = Callback()

    # 非同期モードで認識サービスを呼び出す。モデル、フォーマット、サンプルレートなどの認識パラメータをカスタマイズできます
    recognition = Recognition(
        model='fun-asr-realtime',
        format=format_pcm,
        # 'pcm'、'wav'、'opus'、'speex'、'aac'、'amr'。サポートされているフォーマットはドキュメントで確認できます。
        sample_rate=sample_rate,
        # 16000 をサポートします。
        semantic_punctuation_enabled=False,
        callback=callback)

    # 認識を開始
    recognition.start()

    signal.signal(signal.SIGINT, signal_handler)
    print("Press 'Ctrl+C' to stop recording and recognition...")
    # 「Ctrl+C」が押されるまでキーボードリスナーを作成

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

    recognition.stop()

ローカル音声ファイルの認識

ローカル音声ファイルを文字起こしします。これは、音声チャット、音声コマンド、音声入力、音声検索などの、短時間でニアリアルタイムのシナリオに適しています。

Java

この例で使用されている音声ファイルは asr_example.wav です。

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 {
        // 次の設定はシンガポールリージョン用です。「{WorkspaceId}」を実際のワークスペース ID に置き換えてください。設定はリージョンによって異なります。
        Constants.baseWebsocketApiUrl = "wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/inference";
        // 実際のアプリケーションでは、プログラム起動時にこのメソッドを一度だけ呼び出します。
        warmUp();

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

        // すべてのタスクが完了するのを待ちます。
        executorService.awaitTermination(1, TimeUnit.MINUTES);
        System.exit(0);
    }

    public static void warmUp() {
        try {
            // 接続を確立するための軽量な GET リクエスト。
            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) {
            // ウォームアップが失敗した場合でもリトライを許可します。
        }
    }
}

class RealtimeRecognitionTask implements Runnable {
    private Path filepath;

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

    @Override
    public void run() {
        RecognitionParam param = RecognitionParam.builder()
                .model("fun-asr-realtime")
                // シンガポールリージョンと北京リージョンでは API キーが異なります。API キーを取得するには、https://www.alibabacloud.com/help/model-studio/get-api-key をご参照ください。
                // 環境変数を設定していない場合は、次の行を Model Studio API キーに置き換えてください: .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);
            // パスを音声ファイルのパスに置き換えてください。
            System.out.println(TimeUtils.getTimestamp()+" "+"[" + threadName + "] Input file_path is: " + this.filepath);
            // ファイルを読み込み、音声をチャンクで送信します。
            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 {
            // タスク完了後、WebSocket 接続を閉じます。
            recognizer.getDuplexApi().close(1000, "bye");
        }

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

Python

この例で使用されている音声ファイルは asr_example.wav です。

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

# シンガポールリージョンと北京リージョンでは API キーが異なります。API キーを取得するには、https://www.alibabacloud.com/help/model-studio/get-api-key をご参照ください
# 環境変数を設定していない場合は、次の行を Model Studio API キーに置き換えてください: dashscope.api_key = "sk-xxx"
dashscope.api_key = os.environ.get('DASHSCOPE_API_KEY')

# 次の設定はシンガポールリージョン用です。「{WorkspaceId}」を実際のワークスペース ID に置き換えてください。設定はリージョンによって異なります。
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')  # 認識完了

    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='fun-asr-realtime',
                          format='wav',
                          sample_rate=16000,
                          callback=callback)

try:
    audio_data: bytes = None
    f = open("asr_example.wav", 'rb')
    if os.path.getsize("asr_example.wav"):
        # ファイル全体をバッファに読み込む
        file_buffer = f.read()
        f.close()
        print("Start Recognition")
        recognition.start()

        # 3200 バイトのチャンクでデータを送信
        buffer_size = len(file_buffer)
        offset = 0
        chunk_size = 3200

        while offset < buffer_size:
            # 現在のチャンクのサイズを計算
            remaining_bytes = buffer_size - offset
            current_chunk_size = min(chunk_size, remaining_bytes)

            # バッファから現在のチャンクを抽出
            audio_data = file_buffer[offset:offset + current_chunk_size]

            # 音声フレームを送信
            recognition.send_audio_frame(audio_data)
            # オフセットを更新
            offset += current_chunk_size

            # リアルタイム伝送をシミュレートするために遅延を追加
            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(),
    ))

Qwen-ASR

説明

この例では、your_audio_file.pcm (PCM16, 16 kHz, モノラル) を読み取ります。 MP3、WAV、またはその他のフォーマットから変換するには、ffmpeg を使用します:

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

Java

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

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

public class Qwen3AsrRealtimeUsage {
    private static final Logger log = LoggerFactory.getLogger(Qwen3AsrRealtimeUsage.class);
    private static final int AUDIO_CHUNK_SIZE = 1024; // 音声チャンクサイズ (バイト)
    private static final int SLEEP_INTERVAL_MS = 30;  // スリープ間隔 (ミリ秒)

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

        OmniRealtimeParam param = OmniRealtimeParam.builder()
                .model("qwen3-asr-flash-realtime")
                // 次の設定はシンガポールリージョン用です。「{WorkspaceId}」を実際のワークスペース ID に置き換えてください。設定はリージョンによって異なります。
                .url("wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/realtime")
                // シンガポールリージョンと北京リージョンでは API キーが異なります。API キーを取得するには、https://www.alibabacloud.com/help/model-studio/get-api-key をご参照ください
                // 環境変数を設定していない場合は、次の行を Model Studio API キーに置き換えてください: .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("en");
        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);

            // 音声データをチャンクで読み込んで送信
            while ((bytesRead = audioInputStream.read(audioBuffer)) != -1) {
                totalBytesRead += bytesRead;
                String audioB64 = Base64.getEncoder().encodeToString(audioBuffer);
                // 音声チャンクを会話に送信
                conversation.appendAudio(audioB64);

                // リアルタイムの音声ストリーミングをシミュレートするために短い遅延を追加
                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);
        }

        //session.finish を送信し、終了して閉じるのを待つ
        conversation.endSession();
        log.info("task finished");

        System.exit(0);
    }
}

Python

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

def setup_logging():
    """ロギング出力を設定"""
    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():
    """API キーを初期化"""
    # シンガポールリージョンと北京リージョンでは API キーが異なります。API キーを取得するには、https://www.alibabacloud.com/help/model-studio/get-api-key をご参照ください
    # 環境変数を設定していない場合は、次の行を Model Studio API キーに置き換えてください: 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):
    """リアルタイム認識コールバックハンドラ"""
    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):
    """音声ファイルをチャンクで読み込む"""
    with open(file_path, 'rb') as f:
        while chunk := f.read(chunk_size):
            yield chunk

def send_audio(conversation, file_path, delay=0.1):
    """音声データを送信"""
    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',
        # 次の設定はシンガポールリージョン用です。「{WorkspaceId}」を実際のワークスペース ID に置き換えてください。設定はリージョンによって異なります。
        url='wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/realtime',
        callback=callback,
    )
    callback.conversation = conversation  # コールバックからメソッドを呼び出せるように、コールバックに会話を注入

    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='en',
        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)
        # session.finish を送信し、終了して閉じるのを待つ
        conversation.end_session()
    except Exception as e:
        print(f"Error occurred: {e}")
    finally:
        conversation.close()
        print("Audio processing completed.")

if __name__ == '__main__':
    main()

Paraformer

Paraformer は Fun-ASR のサンプルコードを再利用します。model パラメーターを Paraformer のモデル名に置き換えます。

認識設定

Qwen-ASR の対話モード

Qwen-ASR Realtime API は、2つの対話モードをサポートしています:

  • VAD モード (デフォルト): サーバーが発話の開始と終了 (ターン検出) を自動的に検出します。リアルタイムの会話、会議の文字起こし、および同様のシナリオに適しています。有効にするには、session.turn_detection を設定します (デフォルトで有効になっています)。

  • 手動モード: クライアントは input_audio_buffer.commit を送信することで、発話区切り検出を制御します。チャットアプリの音声メッセージなど、送信タイミングを明示的に制御する必要があるシナリオに適しています。有効にするには、session.turn_detection を null に設定します。

モードの切り替え

  • WebSocketsession.update イベントの turn_detection フィールドを設定します。

    {
        "type": "session.update",
        "session": {
            "turn_detection": null
        }
    }
  • Python SDK: update_session メソッドで enable_turn_detection パラメーターを設定します。

    conversation.update_session(
        enable_turn_detection=False
    )
  • Java SDK: enableTurnDetection パラメーターを OmniRealtimeConfig.builder() を使用して設定します。

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

完全な SDK コード例については、Qwen-ASR-Realtime Python SDK - API リファレンス および Java SDK をご参照ください。WebSocket イベントのライフサイクルについては、イベントフローをご参照ください。

VAD 発話区間検出

音声アクティビティ検出 (VAD) は、連続した音声がいつ終了するかを判断し、最終結果イベントをトリガーします。3つのモデルファミリーすべてで、デフォルトでサーバー側の VAD が有効になっていますが、パラメータ名と調整可能性は異なります:

  • Qwen-ASR: session.turn_detection を使用して設定されます。これには、silence_duration_ms (サイレンス持続時間のしきい値。サイレンスがこの値を超えるとターンが終了します。サーバーのデフォルトは 800 です。高速なターン検出が必要な会話やチャットのシナリオでは 400 が推奨されます) と threshold (VAD 検出の秘密度。サーバーのデフォルトは 0.2) が含まれます。また、Qwen-ASR は手動モードもサポートしています。このモードでは VAD が無効になり、クライアントは commit を使用してターン検出をコントロールできます。上記の「Qwen-ASR のインタラクションモード」をご参照ください。

  • Fun-ASR と Paraformer: max_sentence_silence (VAD 発話区切り検出サイレンスしきい値、ミリ秒) で設定します。音声セグメントの後のサイレンスがこのしきい値を超えると、文が終了したと見なされます。

プロトコルによってパラメーター名が異なります。同じ概念は、Qwen-ASR では silence_duration_ms、Fun-ASR および Paraformer では max_sentence_silence となります。詳細については、「API リファレンス」をご参照ください。

高度な機能

ホットワードで精度を向上

Fun-ASR と Paraformer は、ブランド名、人名、固有名詞などの特定の用語の認識精度を向上させるためにホットワードをサポートしています。

ホットワードの設定と使用方法については、認識精度の向上をご参照ください。

タイムスタンプ

Fun-ASR と Paraformer は、デフォルトで文レベル単語レベルの 2 つの粒度でタイムスタンプを返します。これにより、字幕の配置、キーワードのハイライト、カラオケスタイルのシングアロングなどのユースケースが実現できます。Qwen-ASR Realtime (qwen3-asr-flash-realtime) は現在、タイムスタンプを返しません。タイムスタンプが必要な場合は、Fun-ASR または Paraformer を使用してください。Qwen-ASR のファイル文字起こしモデル qwen3-asr-flash-filetrans は、単語レベルのタイムスタンプをサポートしています。詳細については、「非リアルタイム音声認識」をご参照ください。

タイムスタンプはミリ秒単位で2つのレベルで報告されます:

  • 文レベル: payload.output.sentence.begin_timepayload.output.sentence.end_time は、音声内の文の開始と終了を示します。中間結果では、end_timenull になることがあります。最終値は、文が終了したとき (sentence_end = true) に設定されます。

  • 単語レベル: payload.output.sentence.words 配列です。各要素には begin_timeend_timetext (文字または単語のテキスト)、および punctuation (文字の後に続く句読点。ない場合は空の文字列) が含まれます。

応答例 (抜粋):

{
  "payload": {
    "output": {
      "sentence": {
        "begin_time": 170,
        "end_time": 920,
        "text": "Okay, I got it.",
        "sentence_end": true,
        "words": [
          { "begin_time": 170, "end_time": 295, "text": "Okay", "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": "" }
        ]
      }
    }
  }
}

上記のフィールド名は WebSocket JSON パスを使用しています。各 SDK は、独自の命名規則 (辞書キー、オブジェクトプロパティ、ゲッターメソッドなど) でこれらのフィールドを公開します。完全なフィールドマッピングについては、各 SDK の API リファレンスをご参照ください。

完全なフィールド定義については、API リファレンスをご参照ください。

感情認識

一部の Qwen-ASR および Paraformer モデルには、文字起こし結果に話者の感情状態が含まれます。出力の粒度と有効化する方法は、両者で異なります。

Qwen-ASR (qwen3-asr-flash-realtime): 常時オンで、構成は不要です。トップレベルの emotion フィールドは、conversation.item.input_audio_transcription.text イベントと conversation.item.input_audio_transcription.completed イベントの両方で返されます。値は、surprisedneutralhappysaddisgustedangry、および fearful の 7 つの詳細な感情のうちの 1 つです。

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

Paraformer (paraformer-realtime-8k-v2): この Paraformer モデルのみが感情認識をサポートしています。結果は payload.output.sentence.emo_tagpayload.output.sentence.emo_confidence を介して返されます。値は、positive (幸福や満足などの肯定的な感情)、negative (怒りや落ち込みなどの否定的な感情)、および neutral (明確な感情がない) の3つの感情のいずれかです。信頼度の範囲は [0.0, 1.0] です。

感情認識は、以下のすべての条件が満たされた場合にのみ返されます:

  • モデルは paraformer-realtime-8k-v2 です。

  • 意味的なターンの検出は無効になっています: semantic_punctuation_enabled = false (デフォルト。特別な設定は不要です)。

  • 結果は、sentence_end = true のセンテンス終了イベントでのみ返されます。

感情認識フィールドを抑制するには、semantic_punctuation_enabledtrue に設定します。これによりセマンティックな発話ターン検出が有効になり、emo_tag フィールドおよび emo_confidence フィールドは返されなくなります。

上記のフィールド名は WebSocket JSON パスを使用しています。各 SDK は、独自の命名規則 (辞書キー、オブジェクトプロパティ、ゲッターメソッドなど) でこれらのフィールドを公開します。完全なフィールドマッピングについては、各 SDK の API リファレンスをご参照ください。

完全なフィールド定義、値の制約、および例については、API リファレンスをご参照ください。

禁止用語フィルタリング

禁止用語フィルタリングは、認識結果から禁止用語を置き換えたり削除したりします。これは、顧客サービスの品質検査、コンテンツのコンプライアンス、字幕レビューなどのシナリオに適しています。

サポートされているモデル: Fun-ASR のみ。

制限: 最大 32 個の禁止用語。

デフォルトの動作: special_word_filter パラメーターが指定されていない場合、禁止用語フィルタリングは適用されません。

設定: special_word_filter は、3 つのサブフィールドを持つ JSON オブジェクトです:

  • filter_with_signed.word_list: 同じ長さのアスタリスク (*) で置き換える禁止用語をリストする文字列配列です。たとえば、["test"] を指定した場合、「please help me test it」 は 「please help me **** it」 になります。

  • filter_with_empty.word_list: 結果から完全に削除する禁止用語を指定する文字列配列です。 たとえば、["start"] が指定されている場合、"is the game about to start now" は "is the game about to now" になります。

  • system_reserved_filter: ブール値。デフォルトは false。禁止用語のフィルタリングを有効にするかどうかを指定します。

設定例:

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

異なる SDK は、これらのパラメータを異なる命名規則 (辞書キー、オブジェクトプロパティ、メソッドなど) で公開します。完全なフィールドマッピングについては、API リファレンスをご参照ください。

生の WebSocket プロトコル

以下の例は、DashScope SDK を使用しないシナリオに適した、生の WebSocket プロトコルを介してサーバーに接続する方法を示しています。これらは最小限の実行可能な実装です。各モデルの WebSocket プロトコルについては、API リファレンスをご参照ください。

生の WebSocket プロトコルの例を表示

Fun-ASR

この例では、音声ファイル asr_example.wav を使用します。

Python

例を実行する前に、必要な依存関係をインストールしてください:

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

サンプルファイルに websocket.py という名前を付けないでください。この名前は websocket ライブラリと競合し、以下のエラーを引き起こします: 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

# シンガポールリージョンと北京リージョンでは API キーが異なります。API キーを取得するには、https://www.alibabacloud.com/help/model-studio/get-api-key をご参照ください
# 環境変数を設定していない場合は、次の行を Model Studio API キーに置き換えてください: api_key = "sk-xxx"
api_key = os.environ.get('DASHSCOPE_API_KEY')
# 次の設定はシンガポールリージョン用です。「{WorkspaceId}」を実際のワークスペース ID に置き換えてください。設定はリージョンによって異なります。
url = 'wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/inference/'  # WebSocket サーバーアドレス
audio_file = 'asr_example.wav'  # 音声ファイルのパスに置き換えてください

# 32 文字のランダムな ID を生成
TASK_ID = uuid.uuid4().hex[:32]

task_started = False  # タスクが開始されたかどうかを追跡

# run-task コマンドを送信
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': 'fun-asr-realtime',
            'parameters': {
                'sample_rate': 16000,
                'format': 'wav'
            },
            'input': {}
        }
    }
    ws.send(json.dumps(run_task_message))

# finish-task コマンドを送信
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))

# 音声ストリームを送信 (100 ミリ秒ごとにバイナリチャンク)
def send_audio_stream(ws):
    chunk_size = 3200  # 100ms @ 16kHz 16bit モノラル
    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('Failed to read audio file: ', e)
        ws.close()

# 接続が開いたら、run-task コマンドを送信
def on_open(ws):
    print('Connected to server')
    send_run_task(ws)

# 受信したメッセージを処理
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)

# task-started が受信されない場合は接続を閉じる
def on_close(ws, close_status_code, close_msg):
    if not task_started:
        print('Task not started, closing connection')

# エラー処理
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

例を実行する前に、Java-WebSocket の依存関係をインストールしてください:

Maven

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

Gradle

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

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

public class FunASRRealtimeClient {

    // シンガポールリージョンと北京リージョンでは API キーが異なります。API キーを取得するには、https://www.alibabacloud.com/help/model-studio/get-api-key をご参照ください
    // 環境変数を設定していない場合は、次の行を Model Studio API キーに置き換えてください: private static final String API_KEY = "sk-xxx";
    private static final String API_KEY = System.getenv().getOrDefault("DASHSCOPE_API_KEY", "sk-xxx");
    // 次の設定はシンガポールリージョン用です。「{WorkspaceId}」を実際のワークスペース ID に置き換えてください。設定はリージョンによって異なります。
    private static final String URL = "wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/inference/";
    private static final String AUDIO_FILE = "asr_example.wav"; // 音声ファイルのパスに置き換えてください
    private static final String MODEL = "fun-asr-realtime";

    // 32 文字のランダムな 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();
    }

    // run-task コマンドを送信
    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());
    }

    // 音声ストリームを送信 (100 ミリ秒ごとにバイナリチャンク)
    private static void sendAudioStream() {
        int chunkSize = 3200; // 100ms @ 16kHz 16bit モノラル
        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("Failed to read audio file: " + e.getMessage());
            client.close();
        }
    }

    // finish-task コマンドを送信
    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

必要な依存関係をインストールします:

npm install ws
npm install uuid

コード例:

const fs = require('fs');
const WebSocket = require('ws');
const { v4: uuidv4 } = require('uuid'); // UUID の生成に使用

// シンガポールリージョンと北京リージョンでは API キーが異なります。API キーを取得するには、https://www.alibabacloud.com/help/model-studio/get-api-key をご参照ください
// 環境変数を設定していない場合は、次の行を Model Studio API キーに置き換えてください: const apiKey = "sk-xxx"
const apiKey = process.env.DASHSCOPE_API_KEY;
// 次の設定はシンガポールリージョン用です。「{WorkspaceId}」を実際のワークスペース ID に置き換えてください。設定はリージョンによって異なります。
const url = 'wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/inference/'; // WebSocket サーバーアドレス
const audioFile = 'asr_example.wav'; // 音声ファイルのパスに置き換えてください

// 32 文字のランダムな ID を生成
const TASK_ID = uuidv4().replace(/-/g, '').slice(0, 32);

// WebSocket クライアントを作成
const ws = new WebSocket(url, {
  headers: {
    Authorization: `bearer ${apiKey}`
  }
});

let taskStarted = false; // タスクが開始されたかどうかを示すフラグ

// 接続が開かれたときに run-task コマンドを送信
ws.on('open', () => {
  console.log('Connected to server');
  sendRunTask();
});

// 受信したメッセージを処理
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);
  }
});

// task-started イベントが受信されない場合は接続を閉じる
ws.on('close', () => {
  if (!taskStarted) {
    console.error('Task not started, closing connection');
  }
});

// run-task コマンドを送信
function sendRunTask() {
  const runTaskMessage = {
    header: {
      action: 'run-task',
      task_id: TASK_ID,
      streaming: 'duplex'
    },
    payload: {
      task_group: 'audio',
      task: 'asr',
      function: 'recognition',
      model: 'fun-asr-realtime',
      parameters: {
        sample_rate: 16000,
        format: 'wav'
      },
      input: {}
    }
  };
  ws.send(JSON.stringify(runTaskMessage));
}

// 音声ストリームを送信
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); // 100 ミリ秒ごとに送信
    }
  }

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

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

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

// finish-task コマンドを送信
function sendFinishTask() {
  const finishTaskMessage = {
    header: {
      action: 'finish-task',
      task_id: TASK_ID,
      streaming: 'duplex'
    },
    payload: {
      input: {}
    }
  };
  ws.send(JSON.stringify(finishTaskMessage));
}

// エラー処理
ws.on('error', (error) => {
  console.error('WebSocket error: ', error);
});

C#

コード例:

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;
    // シンガポールリージョンと北京リージョンでは API キーが異なります。API キーを取得するには、https://www.alibabacloud.com/help/model-studio/get-api-key をご参照ください
    // 環境変数を設定していない場合は、次の行を Model Studio API キーに置き換えてください: 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.");

    // 次の設定はシンガポールリージョン用です。「{WorkspaceId}」を実際のワークスペース ID に置き換えてください。設定はリージョンによって異なります。
    private const string WebSocketUrl = "wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/inference/";
    // 音声ファイルのパスに置き換えてください
    private const string AudioFilePath = "asr_example.wav";

    static async Task Main(string[] args) {
        // WebSocket 接続を確立し、認証用のヘッダーを設定
        _webSocket.Options.SetRequestHeader("Authorization", $"bearer {ApiKey}");

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

        // WebSocket メッセージを非同期で受信するスレッドを開始
        var receiveTask = ReceiveMessagesAsync();

        // run-task コマンドを送信
        string _taskId = Guid.NewGuid().ToString("N"); // 32 文字のランダムな ID を生成
        var runTaskJson = GenerateRunTaskJson(_taskId);
        await SendAsync(runTaskJson);

        // task-started イベントを待機
        while (!_taskStartedReceived) {
            await Task.Delay(100, _cancellationTokenSource.Token);
        }

        // ローカルファイルを読み取り、認識対象の音声ストリームをサーバーに送信
        await SendAudioStreamAsync(AudioFilePath);

        // finish-task コマンドを送信してタスクを終了
        var finishTaskJson = GenerateFinishTaskJson(_taskId);
        await SendAsync(finishTaskJson);

        // task-finished イベントを待機
        while (!_taskFinishedReceived && !_cancellationTokenSource.IsCancellationRequested) {
            try {
                await Task.Delay(100, _cancellationTokenSource.Token);
            } catch (OperationCanceledException) {
                // タスクがキャンセルされたため、ループを終了
                break;
            }
        }

        // 接続を閉じる
        if (!_cancellationTokenSource.IsCancellationRequested) {
            await _webSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, "Closing", _cancellationTokenSource.Token);
        }

        _cancellationTokenSource.Cancel();
        try {
            await receiveTask;
        } catch (OperationCanceledException) {
            // 操作キャンセル例外を無視
        }
    }

    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) {
            // 操作キャンセル例外を無視
        }
    }

    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]; // 毎回 100 ミリ秒の音声データを送信
            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); // 100 ミリ秒の間隔
            }
        }
    }

    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"] = "fun-asr-realtime",
                ["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

このサンプルコードは以下のディレクトリ構造を使用します:

my-php-project/

├── composer.json

├── vendor/

└── index.php

composer.json ファイルは以下のとおりです。ご使用のプロジェクトに合わせて、依存関係のバージョンを設定してください。

{
    "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/"
        }
    }
}

index.php ファイル:

<?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;

// シンガポールリージョンと北京リージョンでは API キーが異なります。API キーを取得するには、https://www.alibabacloud.com/help/model-studio/get-api-key をご参照ください
// 環境変数を設定していない場合は、次の行を Model Studio API キーに置き換えてください: $api_key = "sk-xxx"
$api_key = getenv("DASHSCOPE_API_KEY");
// 次の設定はシンガポールリージョン用です。「{WorkspaceId}」を実際のワークスペース ID に置き換えてください。設定はリージョンによって異なります。
$websocket_url = 'wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/inference/';
$audio_file_path = 'asr_example.wav'; // 音声ファイルのパスに置き換えてください

$loop = Loop::get();

// カスタムコネクタを作成
$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";

    // WebSocket メッセージを受信するための非同期スレッドを開始
    $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";
        }
    });

    // 接続クローズをリッスン
    $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";
        }
    });

    // タスク ID を生成
    $taskId = generateTaskId();

    // run-task コマンドを送信
    sendRunTaskMessage($conn, $taskId);

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

$loop->run();

/**
 * タスク ID を生成
 * @return string
 */
function generateTaskId(): string {
    return bin2hex(random_bytes(16));
}

/**
 * run-task コマンドを送信
 * @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" => "fun-asr-realtime",
            "parameters" => [
                "format" => "wav",
                "sample_rate" => 16000
            ],
            "input" => []
        ]
    ]);
    echo "Preparing to send run-task command: " . $runTaskMessage . "\n";
    $conn->send($runTaskMessage);
    echo "run-task command sent\n";
}

/**
 * 音声ファイルを読み込む
 * @param string $filePath
 * @return bool|string
 */
function readAudioFile(string $filePath) {
    $voiceData = file_get_contents($filePath);
    if ($voiceData === false) {
        echo "Failed to read audio file\n";
    }
    return $voiceData;
}

/**
 * 音声データを分割
 * @param string $data
 * @param int $chunkSize
 * @return array
 */
function splitAudioData(string $data, int $chunkSize): array {
    return str_split($data, $chunkSize);
}

/**
 * finish-task コマンドを送信
 * @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 finish-task command: " . $finishTaskMessage . "\n";
    $conn->send($finishTaskMessage);
    echo "finish-task command sent\n";
}

/**
 * イベントを処理
 * @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";
            // 音声ファイルを読み込む
            $voiceData = readAudioFile($audio_file_path);
            if ($voiceData === false) {
                echo "Failed to read audio file\n";
                $conn->close();
                return;
            }

            // 音声データを分割
            $chunks = splitAudioData($voiceData, 1024);

            // 送信関数を定義
            $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);
                    // 100 ミリ秒後に次のチャンクを送信
                    $loop->addTimer(0.1, $sendChunk);
                } else {
                    echo "All data chunks sent\n";
                    $allChunksSent = true;

                    // finish-task コマンドを送信
                    sendFinishTaskMessage($conn, $taskId);
                }
            };

            // 音声データの送信を開始
            $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 ($allChunksSent && $response['header']['event'] == 'task-finished') {
        // すべてのデータが送信されたことを確認するために 1 秒待機
        $loop->addTimer(1, function() use ($conn) {
            $conn->close();
            echo "Client closes 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 (
	// 次の設定はシンガポールリージョン用です。「{WorkspaceId}」を実際のワークスペース ID に置き換えてください。設定はリージョンによって異なります。
	wsURL     = "wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/inference/" // WebSocket サーバーアドレス
	audioFile = "asr_example.wav"                                   // 音声ファイルのパスに置き換えてください
)

var dialer = websocket.DefaultDialer

func main() {
	// シンガポールリージョンと北京リージョンでは API キーが異なります。API キーを取得するには、https://www.alibabacloud.com/help/model-studio/get-api-key をご参照ください
	// 環境変数を設定していない場合は、次の行を Model Studio API キーに置き換えてください: apiKey := "sk-xxx"
	apiKey := os.Getenv("DASHSCOPE_API_KEY")

	// WebSocket サービスに接続
	conn, err := connectWebSocket(apiKey)
	if err != nil {
		log.Fatal("Failed to connect WebSocket: ", err)
	}
	defer closeConnection(conn)

	// 結果を受信するための goroutine を開始
	taskStarted := make(chan bool)
	taskDone := make(chan bool)
	startResultReceiver(conn, taskStarted, taskDone)

	// run-task コマンドを送信
	taskID, err := sendRunTaskCmd(conn)
	if err != nil {
		log.Fatal("Failed to send run-task command: ", err)
	}

	// task-started イベントを待機
	waitForTaskStarted(taskStarted)

	// 認識対象の音声ファイルストリームを送信
	if err := sendAudioData(conn); err != nil {
		log.Fatal("Failed to send audio: ", err)
	}

	// finish-task コマンドを送信
	if err := sendFinishTaskCmd(conn, taskID); err != nil {
		log.Fatal("Failed to send finish-task command: ", err)
	}

	// タスクが終了または失敗するのを待機
	<-taskDone
}

// JSON データを表す構造体を定義
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"`
}

// WebSocket サービスに接続
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
}

// WebSocket メッセージを非同期で受信する goroutine を開始
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
			}
		}
	}()
}

// run-task コマンドを送信
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
}

// run-task コマンドを生成
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:     "fun-asr-realtime",
			Parameters: Params{
				Format:     "wav",
				SampleRate: 16000,
			},
			Input: Input{},
		},
	}
	runTaskCmdJSON, err := json.Marshal(runTaskCmd)
	return string(runTaskCmdJSON), taskID, err
}

// task-started イベントを待機
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 task")
	}
}

// 音声データを送信
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
}

// finish-task コマンドを送信
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
}

// finish-task コマンドを生成
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
}

// イベントを処理
func handleEvent(conn *websocket.Conn, event Event, taskStarted chan<- bool, taskDone chan<- bool) bool {
	switch event.Header.Event {
	case "task-started":
		fmt.Println("Received 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
}

// task-failed イベントを処理
func handleTaskFailed(event Event, conn *websocket.Conn) {
	if event.Header.ErrorMessage != "" {
		log.Fatalf("Task failed: %s", event.Header.ErrorMessage)
	} else {
		log.Fatal("Task failed for unknown reason")
	}
}

// 接続を閉じる
func closeConnection(conn *websocket.Conn) {
	if conn != nil {
		conn.Close()
	}
}

Qwen-ASR

説明

この例では your_audio_file.pcm (PCM16、16 kHz、モノラル) を読み取ります。 MP3、WAV、またはその他のフォーマットから変換するには、ffmpeg を使用します:

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

Python

例を実行する前に、依存関係をインストールしてください:

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

サンプルファイルに websocket.py という名前を付けないでください。この名前は websocket ライブラリと競合し、以下のエラーの原因となります: 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)

# シンガポールリージョンと北京リージョンでは API キーが異なります。API キーを取得するには、https://www.alibabacloud.com/help/model-studio/get-api-key をご参照ください
# 環境変数を設定していない場合は、次の行を Model Studio API キーに置き換えてください: API_KEY="sk-xxx"
API_KEY = os.environ.get("DASHSCOPE_API_KEY", "sk-xxx")
QWEN_MODEL = "qwen3-asr-flash-realtime"
# 次の設定はシンガポールリージョン用です。「{WorkspaceId}」を実際のワークスペース ID に置き換えてください。設定はリージョンによって異なります。
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}")

# 注:非 VAD モードでは、送信される累積音声時間は 60 秒を超えてはなりません
enableServerVad = True
is_running = True  # 実行フラグ

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_manual = {
        "event_id": "event_123",
        "type": "session.update",
        "session": {
            "modalities": ["text"],
            "input_audio_format": "pcm",
            "sample_rate": 16000,
            "input_audio_transcription": {
                # 言語識別子 (オプション)。言語がわかっている場合は設定を推奨
                "language": "en"
            },
            "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": "en"
            },
            "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  # 音声送信スレッドを停止
            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)  # セッション更新が完了するのを待つ
    global is_running

    with open(local_audio_path, 'rb') as audio_file:
        logger.info(f"File read start: {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 read complete: {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)  # リアルタイムキャプチャをシミュレート

# ロガーを初期化
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

例を実行する前に、Java-WebSocket の依存関係をインストールしてください:

Maven

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

Gradle

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

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

public class QwenASRRealtimeClient {

    private static final Logger logger = Logger.getLogger(QwenASRRealtimeClient.class.getName());
    // シンガポールリージョンと北京リージョンでは API キーが異なります。API キーを取得するには、https://www.alibabacloud.com/help/model-studio/get-api-key をご参照ください
    // 環境変数を設定していない場合は、次の行を Model Studio API キーに置き換えてください: 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";

    // VAD モードを有効にするかどうかを制御
    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();

        // 次の設定はシンガポールリージョン用です。「{WorkspaceId}」を実際のワークスペース ID に置き換えてください。設定はリージョンによって異なります。
        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));

                    // 最終的な文字起こし結果は transcription.completed イベントに含まれる
                    if ("conversation.item.input_audio_transcription.completed".equals(eventType)) {
                        logger.info("Final transcript: " + data.optString("transcript"));
                    }

                    // 終了イベントで、送信スレッドを停止し、接続を閉じる
                    if ("session.finished".equals(eventType)) {
                        logger.info("Closing WebSocket connection after session finished...");

                        isRunning.set(false); // 音声送信スレッドを停止
                        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());
            }
        };

        // リクエストヘッダーを追加
        client.addHeader("Authorization", "Bearer " + API_KEY);
        client.addHeader("OpenAI-Beta", "realtime=v1");

        client.connectBlocking(); // 接続が確立されるまでブロック

        // 認識対象の音声ファイルのパスに置き換えてください
        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 イベント (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", "en"))
                        .put("turn_detection", JSONObject.NULL) // 手動モード
                );

        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", "en"))
                        .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());
        }
    }

    /** 音声ファイルストリームを送信 */
    private static void sendAudio(String localAudioPath) throws Exception {
        Thread.sleep(3000); // セッションの準備が整うのを待つ
        byte[] allBytes = Files.readAllBytes(Paths.get(localAudioPath));
        logger.info("File read start");

        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; // 切断後の送信を避ける
            }

            Thread.sleep(100); // リアルタイム送信をシミュレート
        }

        logger.info("File read complete");

        if (client != null && client.isOpen()) {
            // 非 VAD モードでは commit が必要
            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.");
        }
    }

    /** ロガーを初期化 */
    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

例を実行する前に、依存関係をインストールしてください:

npm install ws
/**
 * Qwen-ASR Realtime WebSocket Client (Node.js 版)
 * 特徴:
 * - VAD モードと手動モードをサポート
 * - session.update を送信してセッションを開始
 * - input_audio_buffer.append を介して音声チャンクを継続的に送信
 * - 手動モードでは、input_audio_buffer.commit を送信
 * - session.finish イベントを送信
 * - session.finished を受信後に接続を閉じる
 */

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

// ===== 設定 =====
// シンガポールリージョンと北京リージョンでは API キーが異なります。API キーを取得するには、https://www.alibabacloud.com/help/model-studio/get-api-key をご参照ください
// 環境変数を設定していない場合は、次の行を Model Studio API キーに置き換えてください: const API_KEY = "sk-xxx"
const API_KEY = process.env.DASHSCOPE_API_KEY || 'sk-xxx';
const MODEL = 'qwen3-asr-flash-realtime';
const enableServerVad = true; // VAD モードの場合は true、手動モードの場合は false
const localAudioPath = 'your_audio_file.pcm'; // PCM16、16 kHz の音声ファイルパス

// 次の設定はシンガポールリージョン用です。「{WorkspaceId}」を実際のワークスペース ID に置き換えてください。設定はリージョンによって異なります。
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}`);

// ===== 状態制御 =====
let isRunning = true;

// ===== 接続確立 =====
const ws = new WebSocket(url, {
    headers: {
        'Authorization': `Bearer ${API_KEY}`,
        'OpenAI-Beta': 'realtime=v1'
    }
});

// ===== イベントバインディング =====
ws.on('open', () => {
    console.log('[WebSocket] Connected to server.');
    sendSessionUpdate();
    // 音声送信スレッドを開始
    sendAudio(localAudioPath);
});

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

        // 最終的な文字起こし結果は transcription.completed イベントに含まれる
        if (data.type === 'conversation.item.input_audio_transcription.completed') {
            console.log(`[Final Transcript] ${data.transcript}`);
        }

        // 終了イベント時
        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);
});

// ===== セッション更新 =====
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: 'en'
            },
            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: 'en'
            },
            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));
    }
}

// ===== 音声ファイルストリームを送信 =====
function sendAudio(audioPath) {
    setTimeout(() => {
        console.log(`[File Read Start] ${audioPath}`);
        const buffer = fs.readFileSync(audioPath);

        let offset = 0;
        const chunkSize = 3200; // ~0.1s の PCM16 音声

        function sendChunk() {
            if (!isRunning) return;
            if (offset >= buffer.length) {
                isRunning = false; // 音声送信を停止
                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); // リアルタイム送信をシミュレート
        }

        sendChunk();
    }, 3000); // セッション設定が完了するのを待つ
}

C#

次の例は、サーバーに接続し、PCM 音声ファイルを文字起こしします:

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;

    // VAD モードを使用するかどうかを制御
    private const bool EnableServerVad = true;

    // シンガポールリージョンと北京リージョンでは API キーが異なります。API キーを取得するには、https://www.alibabacloud.com/help/model-studio/get-api-key をご参照ください
    // 環境変数を設定していない場合は、次の行を Model Studio API キーに置き換えてください: 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";
    // 次の設定はシンガポールリージョン用です。「{WorkspaceId}」を実際のワークスペース ID に置き換えてください。設定はリージョンによって異なります。
    private const string BaseUrl = "wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/realtime";
    private const string AudioFilePath = "your_audio_file.pcm"; // PCM 音声ファイルのパスに置き換えてください

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

        // 認証ヘッダーを設定
        _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.");

        // メッセージ受信タスクを開始
        var receiveTask = ReceiveMessagesAsync();

        // session.update 設定を送信
        await SendSessionUpdateAsync();

        // 音声ストリームを送信
        await SendAudioStreamAsync();

        // session.finished イベントを待機
        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);
    }

    // session.update イベントを送信
    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());
    }

    // 音声ストリームを送信 (100 ミリ秒ごとに PCM チャンク)
    private static async Task SendAudioStreamAsync() {
        await Task.Delay(3000); // セッションが設定されるのを待つ
        const int chunkSize = 3200; // 100ms @ 16kHz 16bit モノラル
        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("Audio file read.");
        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());
        }
    }

    // サーバーイベントを受信して処理
    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

この例では、次のディレクトリ構造を使用します:

my-php-project/
├── composer.json
├── vendor/
└── index.php

環境に合わせて、composer.json 内の依存関係のバージョンを調整してください:

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

index.php ファイルは以下のとおりです:

<?php

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

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

// シンガポールリージョンと北京リージョンでは API キーが異なります。API キーを取得するには、https://www.alibabacloud.com/help/model-studio/get-api-key をご参照ください
// 環境変数を設定していない場合は、次の行を Model Studio API キーに置き換えてください: $api_key = "sk-xxx"
$api_key = getenv("DASHSCOPE_API_KEY");
$model = 'qwen3-asr-flash-realtime';
// 次の設定はシンガポールリージョン用です。「{WorkspaceId}」を実際のワークスペース ID に置き換えてください。設定はリージョンによって異なります。
$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'; // PCM 音声ファイルのパスに置き換えてください

// VAD モードを使用するかどうかを制御
$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 the WebSocket server\n";

    // サーバーイベントをリッスン
    $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";
    });

    // session.update イベントを送信
    sendSessionUpdate($conn, $enable_server_vad);

    // セッションが設定されるのを待ってから、音声を送信
    $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 "Failed to connect: {$e->getMessage()}\n";
});

$loop->run();

// session.update イベントを送信
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";
}

// 音声ストリームを送信 (100 ミリ秒ごとに PCM チャンク)
function sendAudioStream($conn, $audio_file_path, $enable_server_vad, $loop, &$is_running) {
    $fp = fopen($audio_file_path, 'rb');
    if (!$fp) {
        echo "Failed 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 モノラル
        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

例を実行する前に、依存関係をインストールしてください:

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 (
	// 次の設定はシンガポールリージョン用です。「{WorkspaceId}」を実際のワークスペース ID に置き換えてください。設定はリージョンによって異なります。
	baseURL         = "wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/realtime"
	model           = "qwen3-asr-flash-realtime"
	audioFile       = "your_audio_file.pcm" // PCM 音声ファイルのパスに置き換えてください
	enableServerVad = true                  // VAD モードを使用するかどうかを制御
)

// サーバーイベント構造体
type ServerEvent struct {
	Type       string `json:"type"`
	Transcript string `json:"transcript,omitempty"`
}

func main() {
	// シンガポールリージョンと北京リージョンでは API キーが異なります。API キーを取得するには、https://www.alibabacloud.com/help/model-studio/get-api-key をご参照ください
	// 環境変数を設定していない場合は、次の行を Model Studio API キーに置き換えてください: 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 the WebSocket server: ", err)
	}
	defer conn.Close()

	// メッセージを受信する goroutine を開始
	sessionFinished := make(chan bool, 1)
	go receiveMessages(conn, sessionFinished)

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

	// セッションが設定されるのを待つ
	time.Sleep(3 * time.Second)

	// 音声ストリームを送信
	if err := sendAudioStream(conn); err != nil {
		log.Fatal("Failed to send audio: ", err)
	}

	// session.finished を待つ
	<-sessionFinished
}

// WebSocket 接続を確立
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
}

// session.update イベントを送信
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)
}

// 音声ストリームを送信 (100 ミリ秒ごとに PCM チャンク)
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 モノラル
	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)
}

// サーバーイベントを受信して処理
func receiveMessages(conn *websocket.Conn, sessionFinished chan<- bool) {
	for {
		_, msg, err := conn.ReadMessage()
		if err != nil {
			log.Println("Failed to read message: ", err)
			sessionFinished <- true
			return
		}
		var evt ServerEvent
		if err := json.Unmarshal(msg, &evt); err != nil {
			log.Println("Failed to parse 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

Paraformer は Fun-ASR のサンプルコードを再利用します。model パラメーターを Paraformer モデル名に置き換えます。

本番環境に適用

接続の再利用 (WebSocket)

Fun-ASR と Paraformer は WebSocket 接続の再利用をサポートしています。1つの認識タスクが終了した後、再接続せずに同じ接続で次のタスクを開始できます。

再利用フロー: クライアントが finish-task を送信し、サーバーが task-finished を返した後、クライアントは再び run-task を送信して新しいタスクを開始します。

重要
  1. 新しいタスクを開始する前に、サーバーが task-finished イベントを返すのをお待ちください。

  2. 再利用する接続では、タスクごとに異なる task_id を使用する必要があります。

  3. タスクが失敗した場合、サーバーはエラーイベントを返し、接続を閉じます。その接続は再利用できません。

  4. タスク終了後 60 秒以内に新しいタスクが開始されない場合、接続は自動的に閉じられます。

Qwen-ASR Realtime はセッションモデルを使用します。各セッションが終了した後に接続を閉じてください。接続の再利用はサポートされていません。

各モデルのイベントの説明については、対応するAPI リファレンスをご参照ください。

高同時実行数のベストプラクティス

DashScope SDK には、WebSocket 接続と認識オブジェクトを再利用する組み込みのプーリングが含まれており、それらを繰り返し作成および破棄するオーバーヘッドを回避します。現在、この機能をサポートしているのは Paraformer Java SDK のみです。

高同時実行数のベストプラクティスを表示

前提条件

Java SDK は、組み込みの接続プールとカスタムオブジェクトプールを連携させることで、最適なパフォーマンスを実現します:

  • 接続プール:SDK に組み込まれた OkHttp3 接続プールが、基盤となる WebSocket 接続を管理および再利用し、ネットワークハンドシェイクのオーバーヘッドを削減します。これはデフォルトで有効になっています。

  • オブジェクトプール: commons-pool2 上に構築され、事前接続済みの Recognition オブジェクトのセットを維持します。プールからオブジェクトを借用すると、接続確立レイテンシーが解消され、最初のパケットのレイテンシーが大幅に削減されます。

実装手順

  1. 依存関係の追加

    お使いのビルドツールに応じて、プロジェクトの依存関係設定ファイルに dashscope-sdk-javacommons-pool2 を追加します。

    Maven と Gradle の設定を以下に示します:

    Maven

    1. Maven プロジェクトの pom.xml ファイルを開きます。

    2. <dependencies> タグ内に以下の依存関係を追加します。

    <dependency>
        <groupId>com.alibaba</groupId>
        <artifactId>dashscope-sdk-java</artifactId>
        <!-- 'the-latest-version' を 2.16.9 以降に置き換えてください。利用可能なバージョンについては 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>
        <!-- 'the-latest-version' を最新バージョンに置き換えてください。利用可能なバージョンについては https://mvnrepository.com/artifact/org.apache.commons/commons-pool2 をご参照ください。 -->
        <version>the-latest-version</version>
    </dependency>
    1. pom.xml ファイルを保存します。

    2. mvn clean installmvn compile などの Maven コマンドを使用して、プロジェクトの依存関係を更新します。

    Gradle

    1. Gradle プロジェクトの build.gradle ファイルを開きます。

    2. 以下の依存関係を dependencies ブロック内に追加します。

      dependencies {
          // 'the-latest-version' を 2.16.9 以降に置き換えてください。利用可能なバージョンについては https://mvnrepository.com/artifact/com.alibaba/dashscope-sdk-java をご参照ください。
          implementation group: 'com.alibaba', name: 'dashscope-sdk-java', version: 'the-latest-version'
      
          // 'the-latest-version' を最新バージョンに置き換えてください。利用可能なバージョンについては https://mvnrepository.com/artifact/org.apache.commons/commons-pool2 をご参照ください。
          implementation group: 'org.apache.commons', name: 'commons-pool2', version: 'the-latest-version'
      }
    3. build.gradleファイルを保存します。

    4. コマンドラインを開き、プロジェクトのルートディレクトリに移動して、以下の Gradle コマンドを実行して依存関係を更新します。

      ./gradlew build --refresh-dependencies

      Windows の場合:

      gradlew build --refresh-dependencies
  2. 接続プールの設定

    環境変数を通じて接続プールを設定します:

    環境変数

    説明

    DASHSCOPE_CONNECTION_POOL_SIZE

    接続プールサイズ。

    推奨値:ピーク時の同時実行数の少なくとも 2 倍。

    デフォルト値:32。

    DASHSCOPE_MAXIMUM_ASYNC_REQUESTS

    同時非同期リクエストの最大数。

    推奨値: DASHSCOPE_CONNECTION_POOL_SIZE と同じです。

    デフォルト値:32。

    DASHSCOPE_MAXIMUM_ASYNC_REQUESTS_PER_HOST

    ホストごとの非同期リクエストの最大数。

    推奨値: DASHSCOPE_CONNECTION_POOL_SIZE と同じ値。

    デフォルト値:32。

  3. オブジェクトプールの設定

    環境変数を通じてオブジェクトプールのサイズを設定します:

    環境変数

    説明

    RECOGNITION_OBJECTPOOL_SIZE

    オブジェクトプールサイズ。

    推奨値:ピーク時の同時実行数の 1.5~2 倍。

    デフォルト値:500。

    重要
    • オブジェクトプールサイズ (RECOGNITION_OBJECTPOOL_SIZE) は、接続プールサイズ (DASHSCOPE_CONNECTION_POOL_SIZE) 以下にする必要があります。そうでない場合、オブジェクトプールがオブジェクトをリクエストする際に接続プールが満杯だと、呼び出しスレッドがブロックされます。

    • オブジェクトプールのサイズは、アカウントの QPS (クエリ/秒) 制限を超えないようにしてください。

    以下のコードを使用してオブジェクトプールを作成します:

    class RecognitionObjectPool {
        // ... 完全な例については、完全なコードセクションをご参照ください
        public static GenericObjectPool<Recognition> getInstance() {
            lock.lock();
            if (recognitionGenericObjectPool == null) {
                int objectPoolSize = getObjectivePoolSize();
                RecognitionObjectFactory recognitionObjectFactory =
                        new RecognitionObjectFactory();
                GenericObjectPoolConfig<Recognition> config =
                        new GenericObjectPoolConfig<>();
                config.setMaxTotal(objectPoolSize);
                config.setMaxIdle(objectPoolSize);
                config.setMinIdle(objectPoolSize);
                recognitionGenericObjectPool =
                        new GenericObjectPool<>(recognitionObjectFactory, config);
            }
            lock.unlock();
            return recognitionGenericObjectPool;
        }
    }
  4. プールから Recognition オブジェクトを借用

    返却されていないオブジェクトの数がプールサイズを超えると、システムは追加の Recognition オブジェクトを作成します。これらの追加オブジェクトは、新しい WebSocket 接続を確立する必要があり、再利用できません。

    recognizer = RecognitionObjectPool.getInstance().borrowObject();
  5. 音声認識の実行

    Recognition オブジェクトの call メソッドまたは streamCall メソッドを呼び出して、音声認識を実行します。

  6. Recognition オブジェクトの返却

    音声認識タスクが終了したら、Recognition オブジェクトを返却して再利用できるようにします。タスクが完了しなかったり失敗したりしたオブジェクトは返却しないでください。

    RecognitionObjectPool.getInstance().returnObject(recognizer);

完全なコード

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;

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 {
        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, "asr_example.wav"),
                Paths.get(currentDir, "asr_example.wav"),
                Paths.get(currentDir, "asr_example.wav"),
        };
        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);

                // チャンクサイズを 16KHz サンプルレートで 100 ミリ秒に設定
                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();
    }
}

推奨構成

以下の構成は、指定されたスペックの Alibaba Cloud サーバー上で Paraformer リアルタイム音声認識サービスのみを実行した場合のテスト結果です。単一ホストの同時実行数は、単一ホスト上で同時に実行される Paraformer リアルタイム音声認識タスクの数 (ワーカースレッドの数) です。

サーバー構成 (Alibaba Cloud)

ホストあたりの最大同時実行数

オブジェクトプールサイズ

接続プールサイズ

4 コア、8 GiB

100

500

2000

8 コア、16 GiB

200

500

2000

16 コア、32 GiB

400

500

2000

リソース管理と例外処理

  • タスク成功: GenericObjectPool.returnObject() を呼び出し、Recognition オブジェクトを再利用のためにプールに返却します。

    重要

    タスクが完了しなかったり失敗したりした Recognition オブジェクトは返却しないでください。

  • タスク失敗:SDK またはビジネスロジックがタスクを中断する例外をスローした場合、以下の両方を行います:

    1. 基盤となる WebSocket 接続を閉じます。

    2. 再利用を防ぐために、プール内のオブジェクトを無効にします。

    // 接続を閉じる
    recognizer.getDuplexApi().close(1000, "bye");
    // オブジェクトプール内の失敗したレコグナイザを無効にする
    RecognitionObjectPool.getInstance().invalidateObject(recognizer);
  • サービス側の TaskFailed エラーについては、追加の処理は不要です。

ウォームアップとレイテンシ測定

DashScope Java SDK の同時呼び出しレイテンシやその他のパフォーマンスメトリクスを測定する際は、正式なテストの前に十分なウォームアップを行ってください。

接続再利用メカニズム

DashScope Java SDK は、グローバルなシングルトン接続プールを使用して WebSocket 接続を管理および再利用します。このメカニズムは次のように動作します:

  • オンデマンド作成:SDK はサービス起動時に WebSocket 接続を事前に作成しません。接続は最初の呼び出し時にオンデマンドで確立されます。

  • 時間制限付き再利用:リクエストが完了した後、接続は最大 60 秒間プールに保持され、再利用されます。

    • 60 秒以内に新しいリクエストが到着した場合、既存の接続が再利用され、ハンドシェイクのオーバーヘッドが回避されます。

    • 接続が 60 秒以上アイドル状態の場合、リソースを解放するために自動的に閉じられます。

ウォームアップが重要な理由

以下のシナリオでは、接続プールに再利用可能なアクティブな接続がない可能性があり、リクエストは新しい接続を作成する必要があります:

  • アプリケーションが起動したばかりで、まだ呼び出しを行っていない。

  • サービスが 60 秒以上アイドル状態だったため、プールされた接続がタイムアウトにより閉じられた。

これらのシナリオでは、最初または初期のリクエストが完全な WebSocket 接続プロセス (TCP ハンドシェイク、TLS ネゴシエーション、プロトコルアップグレード) をトリガーします。エンドツーエンドのレイテンシは、接続を再利用する後続のリクエストよりも大幅に高くなります。

推奨アプローチ

正式なパフォーマンステストやレイテンシ測定の前に、以下のウォームアップ手順に従ってください:

  1. 実際のテストの同時実行レベルに合わせて、事前に呼び出しを実行し (たとえば、1~2 分間)、接続プールを完全に満たします。

  2. 接続プールが確立され、十分なアクティブな接続が維持されていることを確認してから、正式なパフォーマンスデータ収集を開始します。

認識精度の向上

  • モデルとサンプルレートを一致させる:8 kHz の電話音声には、8 kHz モデルを直接使用してください。16 kHz にアップサンプリングすると、信号が歪むため行わないでください。

  • 入力音声の品質を向上させる:信号対雑音比が高く、エコーのない録音環境で高品質のマイクを使用してください。アプリケーション層では、ノイズリダクション (例:RNNoise) や音響エコーキャンセル (AEC) などの前処理を統合してください。

回復性の設定

  • クライアント側の再接続:クライアントは、ネットワークのジッターに対応するために自動再接続を実装する必要があります。Python SDK の参照実装:

    1. 例外のキャッチ: Callback クラスに on_error メソッドを実装します。dashscope SDK は、ネットワークエラーまたはその他の問題が発生した場合にこのメソッドを呼び出します。

    2. 状態の通知: on_error がトリガーされると、再接続信号をセットします。Python では、スレッドセーフな信号フラグである threading.Event を使用します。

    3. 再接続ループ: メインロジックを for ループでラップします (例: 3 回リトライ)。再接続信号が検出されると、現在の認識が中断されてリソースがクリーンアップされ、数秒後にループが新しい接続を作成します。

  • ハートビートを使用して接続を維持する:サーバーへの長時間持続する接続が必要な場合は、ハートビート パラメーターを true に設定します。音声が長時間無音の場合でも、接続は維持されます。

  • レート制限:モデルインターフェイスを呼び出す際は、モデルのレート制限ルールを遵守してください。

サポートされているモデルとリージョン

シンガポール

以下のいずれかのモデルを呼び出すには、シンガポールリージョンのAPI キーを使用してください:

  • Fun-ASR:fun-asr-realtime (安定版、現在 fun-asr-realtime-2025-11-07 と同等)、fun-asr-realtime-2025-11-07 (スナップショット)

  • Qwen3-ASR-Flash-Realtime:qwen3-asr-flash-realtime (安定版、現在 qwen3-asr-flash-realtime-2025-10-27 と同等)、qwen3-asr-flash-realtime-2026-02-10 (最新スナップショット)、qwen3-asr-flash-realtime-2025-10-27 (スナップショット)

中国 (北京)

以下のいずれかのモデルを呼び出すには、北京リージョンの API キーを使用します:

  • Fun-ASR:fun-asr-realtime (安定版、現在 fun-asr-realtime-2025-11-07 と同等)、fun-asr-realtime-2026-02-28 (最新スナップショット)、fun-asr-realtime-2025-11-07 (スナップショット)、fun-asr-realtime-2025-09-15 (スナップショット)

    • fun-asr-flash-8k-realtime (安定版、現在 fun-asr-flash-8k-realtime-2026-01-28 と同等)、fun-asr-flash-8k-realtime-2026-01-28

  • Qwen3-ASR-Flash-Realtime:qwen3-asr-flash-realtime (安定版、現在 qwen3-asr-flash-realtime-2025-10-27 と同等)、qwen3-asr-flash-realtime-2026-02-10 (最新スナップショット)、qwen3-asr-flash-realtime-2025-10-27 (スナップショット)

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

API リファレンス

よくある質問

リアルタイム音声認識はどの音声フォーマットをサポートしていますか?

Fun-ASR と Paraformer は pcm、wav、mp3、opus、speex、aac、amr をサポートしています。Qwen-ASR では、pcm または opus を使用してください。その他のフォーマット (wav、aac、amr など) は session.update の検証に合格しますが、サーバー側のデコーディングに失敗する可能性があります。送信する前に、音声ストリームが推奨フォーマットを使用していることを確認してください。

SDK と WebSocket API はいつ使用すべきですか?

DashScope SDK は WebSocket 接続管理、認証、再接続をラップしており、統合への最速のパスとなります。WebSocket API は、直接的で詳細な制御を提供します。SDK がご使用の言語をカバーしていない場合、またはカスタムの接続処理が必要な場合に使用してください。ほとんどのユースケースでは SDK が推奨される選択肢です。

固有名詞の認識精度を向上させるにはどうすればよいですか?

ホットワードを使用します (Fun-ASR と Paraformer でサポート)。ホットワードの設定と使用方法については、認識精度の向上をご参照ください。

接続が頻繁に切断される場合はどうすればよいですか?

クライアント側の再接続ロジックを実装し、ハートビートパラメーター (heartbeat=true) を有効にして、長い無通信期間中に接続が切断されるのを防ぎます。詳細な耐障害性戦略については、「本番環境に適用する」をご参照ください。