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

Alibaba Cloud Model Studio:リアルタイム音声認識 - Qwen

最終更新日:Sep 09, 2026

リアルタイム音声認識サービスは、音声ストリームを受信し、リアルタイムで句読点付きのテキストに変換します。ライブキャプション、オンライン会議、音声チャット、スマートアシスタントなど、さまざまなシナリオにご利用いただけます。

概要

このサービスは音声をストリーミングで送信し、低遅延でテキストを返却します。

  • 標準中国語を高精度で認識し、広東語や四川語などの方言にも対応しています。
  • 複雑な音響環境にも対応し、自動言語検出と非音声音声のインテリジェントフィルタリングを実現します。
  • 驚き、平静、喜び、悲しみ、嫌悪、怒り、恐怖など、さまざまな感情状態を認識できます。
  • 特定の用語に対する認識精度を向上させるためのカスタムホットワードをサポートしています。
  • 会話履歴やドメイン用語を渡すことで認識精度を向上させるコンテキスト強化をサポートしています。
  • 構造化された認識結果を生成するためのタイムスタンプを出力します。
  • 録音環境に合わせて柔軟なサンプルレートと複数の音声フォーマットに対応しています。

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

前提条件

クイックスタート

以下の例では、DashScope SDK を使用してリアルタイム音声認識サービスを呼び出す方法を示します。

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

このモデルは WebSocket に加えて AOQ プロトコルもサポートしています。安定した遅延、弱いネットワークでの耐障害性、および内蔵の全二重ノイズ抑制とエコーキャンセレーションを優先するクライアント側統合には、AOQ を推奨します。プロトコル比較については、「リアルタイム API の概要」をご参照ください。

マイクからの音声認識

マイクから音声を認識し、リアルタイムでテキストを出力します。これにより、話者が話すと同時に単語が表示されます。

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("qwen-audio-3.0-asr-flash-streaming")
                // API キーはシンガポールリージョンと北京リージョンで異なります。API キーの取得: https://www.alibabacloud.com/help/zh/model-studio/get-api-key
                // 環境変数を設定していない場合は、次の行を Model Studio API キーに置き換えてください: .apiKey("sk-xxx")
                .apiKey(System.getenv("DASHSCOPE_API_KEY"))
                .format("pcm")
                .sampleRate(16000)
                .build();
        Recognition recognizer = new Recognition();

        ResultCallback<RecognitionResult> callback = new ResultCallback<RecognitionResult>() {
            @Override
            public void onEvent(RecognitionResult result) {
                if (result.isSentenceEnd()) {
                    System.out.println("最終結果: " + result.getSentence().getText());
                } else {
                    System.out.println("中間結果: " + result.getSentence().getText());
                }
            }

            @Override
            public void onComplete() {
                System.out.println("認識完了");
            }

            @Override
            public void onError(Exception e) {
                System.out.println("RecognitionCallback エラー: " + 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(
                "[メトリクス] requestId: "
                        + recognizer.getLastRequestId()
                        + ", 最初のパッケージ遅延 ms: "
                        + recognizer.getFirstPackageDelay()
                        + ", 最後のパッケージ遅延 ms: "
                        + recognizer.getLastPackageDelay());
    }
}

Python

Python の例を実行する前に、サードパーティの音声再生・キャプチャツールキットを 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.active:
            stream.stop()
            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 が押されました。認識を停止します...')
    # 認識を停止
    recognition.stop()
    print('認識が停止しました。')
    print(
        '[メトリクス] requestId: {}, 最初のパッケージ遅延 ms: {}, 最後のパッケージ遅延 ms: {}'
        .format(
            recognition.get_last_request_id(),
            recognition.get_first_package_delay(),
            recognition.get_last_package_delay(),
        ))
    # プログラムを強制終了
    sys.exit(0)

# main function
if __name__ == '__main__':
    # API キーはシンガポールリージョンと北京リージョンで異なります。API キーの取得: https://www.alibabacloud.com/help/zh/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='qwen-audio-3.0-asr-flash-streaming',
        format=format_pcm,
        # 'pcm'、'wav'、'opus'、'speex'、'aac'、'amr'。サポートされているフォーマットはドキュメントで確認してください
        sample_rate=sample_rate,
        # 8000、16000 をサポート
        semantic_punctuation_enabled=False,
        callback=callback)

    # 認識を開始
    recognition.start()

    signal.signal(signal.SIGINT, signal_handler)
    print("'Ctrl+C' を押して録音と認識を停止します...")
    # "Ctrl+C" が押されるまでキーボードリスナーを維持

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

    recognition.stop()

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

ローカル音声ファイルを認識し、結果を出力します。これはチャット会話、音声コマンド、音声入力メソッド、音声検索などの短時間のニアリアルタイムシナリオに適しています。

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"), "{YOUR_AUDIO_FILE}")));
        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("qwen-audio-3.0-asr-flash-streaming")
                // API キーはシンガポールリージョンと北京リージョンで異なります。API キーの取得: https://www.alibabacloud.com/help/zh/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 + "] 最終結果:" + message.getSentence().getText());
                } else {
                    System.out.println(TimeUtils.getTimestamp()+" "+
                            "[process " + threadName + "] 中間結果: " + message.getSentence().getText());
                }
            }

            @Override
            public void onComplete() {
                System.out.println(TimeUtils.getTimestamp()+" "+"[" + threadName + "] 認識完了");
            }

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

        try {
            recognizer.call(param, callback);
            // パスを音声ファイルパスに置き換えてください
            System.out.println(TimeUtils.getTimestamp()+" "+"[" + threadName + "] 入力 file_path: " + 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
                        + "][メトリクス] requestId: "
                        + recognizer.getLastRequestId()
                        + ", 最初のパッケージ遅延 ms: "
                        + recognizer.getFirstPackageDelay()
                        + ", 最後のパッケージ遅延 ms: "
                        + recognizer.getLastPackageDelay());
    }
}
import os
import time
import dashscope
from dashscope.audio.asr import *

# API キーはシンガポールリージョンと北京リージョンで異なります。API キーの取得: https://www.alibabacloud.com/help/zh/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() + ' 認識完了')  # 認識完了

    def on_error(self, result: RecognitionResult) -> None:
        print('認識 task_id: ', result.request_id)
        print('認識エラー: ', result.message)
        exit(0)

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

callback = Callback()

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

try:
    audio_data: bytes = None
    f = open("{YOUR_AUDIO_FILE}", 'rb')
    if os.path.getsize("{YOUR_AUDIO_FILE}"):
        # ファイルの全データを一度にバッファに読み込む
        file_buffer = f.read()
        f.close()
        print("認識を開始")
        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(
            '指定されたファイルは空です (0 バイト)')
except Exception as e:
    raise e

print(
    '[メトリクス] requestId: {}, 最初のパッケージ遅延 ms: {}, 最後のパッケージ遅延 ms: {}'
    .format(
        recognition.get_last_request_id(),
        recognition.get_first_package_delay(),
        recognition.get_last_package_delay(),
    ))

Qwen3-ASR-Flash-Realtime

注記サンプルコードは 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
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/zh/model-studio/get-api-key
                // 環境変数を設定していない場合は、次の行を Alibaba Cloud 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("接続が開かれました");
            }
            @Override
            public void onEvent(JsonObject message) {
                String type = message.get("type").getAsString();
                switch(type) {
                    case "session.created":
                        System.out.println("セッション開始: " + message.get("session").getAsJsonObject().get("id").getAsString());
                        break;
                    case "conversation.item.input_audio_transcription.completed":
                        System.out.println("文字起こし: " + message.get("transcript").getAsString());
                        finishLatch.countDown();
                        break;
                    case "input_audio_buffer.speech_started":
                        System.out.println("======VAD 音声開始======");
                        break;
                    case "input_audio_buffer.speech_stopped":
                        System.out.println("======VAD 音声停止======");
                        break;
                    case "conversation.item.input_audio_transcription.text":
                        System.out.println("文字起こし: " + message.get("text").getAsString() + message.get("stash").getAsString());
                        break;
                    default:
                        break;
                }
            }
            @Override
            public void onClose(int code, String reason) {
                System.out.println("接続が閉じられました コード: " + code + ", 理由: " + reason);
            }
        });
        conversationRef.set(conversation);
        try {
            conversation.connect();
        } catch (NoApiKeyException e) {
            throw new RuntimeException(e);
        }

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

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

        String filePath = "your_audio_file.pcm";
        File audioFile = new File(filePath);
        if (!audioFile.exists()) {
            log.error("音声ファイルが見つかりません: {}", filePath);
            return;
        }

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

            log.info("音声データの送信を開始: {}", 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("音声データの送信が完了しました。送信された総バイト数: {}", totalBytesRead);

        } catch (Exception e) {
            log.error("ファイルからの音声送信エラー: {}", filePath, e);
        }

        // session.finish を送信して終了を待ち、閉じる
        conversation.endSession();
        log.info("タスク完了");

        System.exit(0);
    }
}
        Constants.baseHttpApiUrl = "https://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api/v1";
import logging
import os
import base64
import signal
import sys
import time
import dashscope
from dashscope.audio.qwen_omni import *
from dashscope.audio.qwen_omni.omni_realtime import TranscriptionParams

def setup_logging():
    """ログ出力を設定"""
    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/zh/model-studio/get-api-key
    # 環境変数を設定していない場合は、次の行を Alibaba Cloud 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('[警告] プレースホルダー API キーを使用しています。DASHSCOPE_API_KEY 環境変数を設定してください。')

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('======音声開始======'),
            'input_audio_buffer.speech_stopped': lambda r: print('======音声停止======')
        }

    def on_open(self):
        print('接続が開かれました')

    def on_close(self, code, msg):
        print(f'接続が閉じられました、コード: {code}, メッセージ: {msg}')

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

    def _handle_session_created(self, response):
        print(f"セッション開始: {response['session']['id']}")

    def _handle_final_text(self, response):
        print(f"最終認識テキスト: {response['transcript']}")

    def _handle_transcription_text(self, response):
        print(f"文字起こし結果を取得: {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"音声ファイル {file_path} が存在しません。")

    print("音声ファイルを処理中... 'Ctrl+C' を押して停止します。")
    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 が押されました。終了します...')
        conversation.close()
        sys.exit(0)

    signal.signal(signal.SIGINT, handle_exit)

    conversation.connect()

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

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

    try:
        send_audio(conversation, audio_file_path)
        # session.finish を送信して終了を待ち、閉じる
        conversation.end_session()
    except Exception as e:
        print(f"エラーが発生しました: {e}")
    finally:
        conversation.close()
        print("音声処理が完了しました。")

if __name__ == '__main__':
    main()

Paraformer

Paraformer のサンプルコードは Qwen-Audio-3.0-ASR-Flash-Streaming/Fun-ASR-Realtime と同様です。モデル名を Paraformer モデルに置き換えてください。

認識設定

Qwen3-ASR-Flash-Realtime のインタラクションモード

Qwen3-ASR-Flash-Realtime リアルタイム API には 2 つのインタラクションモードがあります:

  • VAD モード (デフォルト): サーバーが自動的に音声の開始と終了 (セグメンテーション) を検出します。このモードはリアルタイム会話、会議メモなどのシナリオに適しています。session.turn_detection パラメーターを設定することで有効になります (デフォルトで有効)。
  • マニュアルモード: クライアントが input_audio_buffer.commit を送信することでセグメンテーションを制御します。このモードはチャットアプリで音声メッセージを送信するなど、音声送信タイミングを明示的に制御する必要があるシナリオに適しています。session.turn_detection を null に設定することで有効になります。

インタラクションモードの切り替え:

  • WebSocket: session.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: OmniRealtimeConfig.builder() を介して enableTurnDetection パラメーターを設定します。
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-Audio-3.0-ASR-Flash-Streaming / Fun-ASR-Realtime / Paraformer: max_sentence_silence (セグメンテーションの VAD サイレンスしきい値、ミリ秒単位) を介して設定します。音声セグメント後のサイレンスがこのしきい値を超えると、システムは文が完了したと判断します。
  • Qwen3-ASR-Flash-Realtime: session.turn_detection を介して設定します。これには silence_duration_ms (ターンを終了させるサイレンス持続時間しきい値。サーバーデフォルト 800、会話やチャットなど高速セグメンテーションが必要なシナリオでは 400 を推奨) と threshold (VAD 検出感度。サーバーデフォルト 0.2) が含まれます。Qwen3-ASR-Flash-Realtime はマニュアルモードもサポートしており、VAD を無効にしてクライアント側コミットによるセグメンテーションを使用します。詳細については、上記の「Qwen3-ASR-Flash-Realtime インタラクションモード」をご参照ください。

プロトコルによってパラメーター名が異なります。同じ概念は Qwen-Audio-3.0-ASR-Flash-Streaming / Fun-ASR-Realtime / Paraformer では max_sentence_silence、Qwen3-ASR-Flash-Realtime では silence_duration_ms と呼ばれます。フィールド定義の詳細については、「API リファレンス」をご参照ください。

高度な機能

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

ブランド名、個人名、専門用語などの特定の用語に対する認識精度を向上させるためにホットワードを使用します。

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

コンテキスト強化で精度を向上

コンテキスト強化は、会話履歴やドメイン用語を ASR モデルに渡すことで、固有名詞の文字起こし精度を大幅に向上させます。使用方法と結果例の詳細については、「コンテキスト強化」をご参照ください。

タイムスタンプの取得

Qwen-Audio-3.0-ASR-Flash-Streaming、Fun-ASR-Realtime、Paraformer モデルファミリーは、デフォルトで文レベル単語レベルの両方でタイムスタンプを出力します。これにより、字幕の配置、キーワードのハイライト、カラオケスタイルの読み上げなどのシナリオがサポートされます。Qwen3-ASR-Flash-Realtime (qwen3-asr-flash-realtime) は現在タイムスタンプを返しません。 タイムスタンプが必要な場合は、Qwen-Audio-3.0-ASR-Flash-Streaming、Fun-ASR-Realtime、または 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": "OK, I got it",
        "sentence_end": true,
        "words": [
          { "begin_time": 170, "end_time": 295, "text": "OK", "punctuation": "," },
          { "begin_time": 295, "end_time": 503, "text": "I", "punctuation": "" },
          { "begin_time": 503, "end_time": 711, "text": "got", "punctuation": "" },
          { "begin_time": 711, "end_time": 920, "text": "it", "punctuation": "" }
        ]
      }
    }
  }
}

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

フィールド定義の詳細については、「API リファレンス」をご参照ください。

感情認識

Qwen3-ASR-Flash-Realtime と一部の Paraformer モデルは、文字起こし結果に話者の感情状態を含めることができますが、出力粒度と機能の有効化方法が異なります。

Qwen3-ASR-Flash-Realtime (qwen3-asr-flash-realtime): 常にオンで、設定は不要です。感情は conversation.item.input_audio_transcription.textconversation.item.input_audio_transcription.completed イベントの両方でトップレベルの emotion フィールドを通じて返されます。値は 7 つの詳細な感情のいずれかです: surprisedneutralhappysaddisgustedangry、および fearful

{
  "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 を通じて返されます。値は 3 つの極性のいずれかです: positive (例: happy または satisfied)、negative (例: angry または subdued)、および neutral (明確な感情なし)。信頼度は 0.0 から 1.0 の範囲です。

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

  • モデルが paraformer-realtime-8k-v2 であること。
  • セマンティックセグメンテーションがオフであること: semantic_punctuation_enabled = false (false がデフォルトのため、特別な設定は不要)。
  • 結果は文終了イベント (sentence_end = true) のみで返されること。

感情フィールドの返却を停止するには、semantic_punctuation_enabledtrue に設定します。これによりセマンティックセグメンテーションが有効になり、emo_tagemo_confidence フィールドが返されなくなります。

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

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

禁止用語フィルタリング

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

サポートモデル: Qwen-Audio-3.0-ASR-Flash-Streaming および Fun-ASR-Realtime のみ。

制限: 最大 32 個の禁止用語を設定できます。

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

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

  • filter_with_signed.word_list: 等しい長さの * 文字で置き換える禁止用語の文字列配列。例: ["test"] の場合、"Help me test it" は "Help me **** it" になります。
  • filter_with_empty.word_list: 結果から完全に削除する禁止用語の文字列配列。例: ["start"] の場合、"Is the game about to start" は "Is the game about to" になります。
  • system_reserved_filter: ブール値で、デフォルトは false です。禁止用語フィルタリングを有効にするかどうかを決定します。

設定例:

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

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

生 WebSocket プロトコルの呼び出し

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

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

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

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/zh/model-studio/get-api-key
# 環境変数を設定していない場合は、次の行を Alibaba Cloud 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 = '{YOUR_AUDIO_FILE}'  # 音声ファイルのパスに置き換えてください

# 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': 'qwen-audio-3.0-asr-flash-streaming',
            '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))

# 音声ストリームを送信 (100ms ごとに 1 つのバイナリチャンクを送信)
def send_audio_stream(ws):
    chunk_size = 3200  # 16kHz 16bit モノラルで 100ms
    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('音声ストリーム終了')
        send_finish_task(ws)
    except Exception as e:
        print('音声ファイルの読み取りエラー:', e)
        ws.close()

# 接続が開いたときに run-task 命令を送信
def on_open(ws):
    print('サーバーに接続しました')
    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 = True
        threading.Thread(target=send_audio_stream, args=(ws,), daemon=True).start()
    elif event == 'result-generated':
        print('認識結果:', message['payload']['output']['sentence']['text'])
        if message['payload'].get('usage'):
            print('タスク課金時間 (秒):', message['payload']['usage']['duration'])
    elif event == 'task-finished':
        print('タスク完了')
        ws.close()
    elif event == 'task-failed':
        print('タスク失敗:', message['header'].get('error_message'))
        ws.close()
    else:
        print('不明なイベント:', event)

# task-started イベントが受信されない場合に接続を閉じる
def on_close(ws, close_status_code, close_msg):
    if not task_started:
        print('タスクが開始されていません。接続を閉じます')

# エラー処理
def on_error(ws, error):
    print('WebSocket エラー:', 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 依存関係をインストールしてください:

<dependency>
    <groupId>org.java-websocket</groupId>
    <artifactId>Java-WebSocket</artifactId>
    <version>1.5.6</version>
</dependency>
<dependency>
    <groupId>org.json</groupId>
    <artifactId>json</artifactId>
    <version>20240303</version>
</dependency>
implementation 'org.java-websocket:Java-WebSocket:1.5.6'
implementation 'org.json:json:20240303'
import org.java_websocket.client.WebSocketClient;
import org.java_websocket.handshake.ServerHandshake;
import org.json.JSONObject;

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

public class FunASRRealtimeClient {

    // API キーはシンガポールリージョンと北京リージョンで異なります。API キーの取得: https://www.alibabacloud.com/help/zh/model-studio/get-api-key
    // 環境変数を設定していない場合は、次の行を Alibaba Cloud 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 = "{YOUR_AUDIO_FILE}"; // 音声ファイルのパスに置き換えてください
    private static final String MODEL = "qwen-audio-3.0-asr-flash-streaming";

    // 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("サーバーに接続しました");
                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("タスク開始");
                        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("認識結果: " + text);
                        if (payload.has("usage")) {
                            System.out.println("タスク課金時間 (秒): " + payload.getJSONObject("usage").get("duration"));
                        }
                        break;
                    case "task-finished":
                        System.out.println("タスク完了");
                        close();
                        break;
                    case "task-failed":
                        String errMsg = message.getJSONObject("header").optString("error_message");
                        System.err.println("タスク失敗: " + errMsg);
                        close();
                        break;
                    default:
                        System.out.println("不明なイベント: " + event);
                }
            }

            @Override
            public void onClose(int code, String reason, boolean remote) {
                if (!taskStarted.get()) {
                    System.err.println("タスクが開始されていません。接続を閉じます");
                }
            }

            @Override
            public void onError(Exception ex) {
                System.err.println("WebSocket エラー: " + 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());
    }

    // 音声ストリームを送信 (100ms ごとに 1 つのバイナリチャンクを送信)
    private static void sendAudioStream() {
        int chunkSize = 3200; // 16kHz 16bit モノラルで 100ms
        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("音声ストリーム終了");
            sendFinishTask();
        } catch (Exception e) {
            System.err.println("音声ファイルの読み取りエラー: " + 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/zh/model-studio/get-api-key
// 環境変数を設定していない場合は、次の行を Alibaba Cloud 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 = '{YOUR_AUDIO_FILE}'; // 音声ファイルのパスに置き換えてください

// 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('サーバーに接続しました');
  sendRunTask();
});

// 受信メッセージを処理
ws.on('message', (data) => {
  const message = JSON.parse(data);
  switch (message.header.event) {
    case 'task-started':
      console.log('タスク開始');
      taskStarted = true;
      sendAudioStream();
      break;
    case 'result-generated':
      console.log('認識結果:', message.payload.output.sentence.text);
      if (message.payload.usage) {
        console.log('タスク課金時間 (秒):', message.payload.usage.duration);
      }
      break;
    case 'task-finished':
      console.log('タスク完了');
      ws.close();
      break;
    case 'task-failed':
      console.error('タスク失敗:', message.header.error_message);
      ws.close();
      break;
    default:
      console.log('不明なイベント:', message.header.event);
  }
});

// task-started イベントが受信されない場合に接続を閉じる
ws.on('close', () => {
  if (!taskStarted) {
    console.error('タスクが開始されていません。接続を閉じます');
  }
});

// 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: 'qwen-audio-3.0-asr-flash-streaming',
      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); // 100ms ごとに送信
    }
  }

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

  audioStream.on('end', () => {
    console.log('音声ストリーム終了');
    sendFinishTask();
  });

  audioStream.on('error', (err) => {
    console.error('音声ファイルの読み取りエラー:', 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);
});

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/zh/model-studio/get-api-key
    // 環境変数を設定していない場合は、次の行を Alibaba Cloud 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 環境変数が設定されていません。");

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

    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("タスクが正常に開始されました");
                            _taskStartedReceived = true;
                            break;
                        case "result-generated":
                            Console.WriteLine($"認識結果: {message["payload"]?["output"]?["sentence"]?["text"]?.GetValue<string>()}");
                            if (message["payload"]?["usage"] != null && message["payload"]?["usage"]?["duration"] != null) {
                                Console.WriteLine($"タスク課金時間 (秒): {message["payload"]?["usage"]?["duration"]?.GetValue<int>()}");
                            }
                            break;
                        case "task-finished":
                            Console.WriteLine("タスク完了");
                            _taskFinishedReceived = true;
                            _cancellationTokenSource.Cancel();
                            break;
                        case "task-failed":
                            Console.WriteLine($"タスク失敗: {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]; // 100ms 分の音声データを送信
            int bytesRead;

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

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

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

PHP

サンプルプロジェクトのディレクトリ構造は以下の通りです:

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/zh/model-studio/get-api-key
// 環境変数を設定していない場合は、次の行を Alibaba Cloud 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 = '{YOUR_AUDIO_FILE}'; // 音声ファイルのパスに置き換えてください

$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 "WebSocket サーバーに接続しました\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 "不明なメッセージ形式\n";
        }
    });

    // 接続クローズをリッスン
    $conn->on('close', function($code = null, $reason = null) {
        echo "接続が閉じられました\n";
        if ($code !== null) {
            echo "クローズコード: " . $code . "\n";
        }
        if ($reason !== null) {
            echo "クローズ理由: " . $reason . "\n";
        }
    });

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

    // run-task 命令を送信
    sendRunTaskMessage($conn, $taskId);

}, function ($e) {
    echo "接続できません: {$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" => "qwen-audio-3.0-asr-flash-streaming",
            "parameters" => [
                "format" => "wav",
                "sample_rate" => 16000
            ],
            "input" => []
        ]
    ]);
    echo "run-task 命令の送信準備: " . $runTaskMessage . "\n";
    $conn->send($runTaskMessage);
    echo "run-task 命令を送信しました\n";
}

/**
 * 音声ファイルを読み込む
 * @param string $filePath
 * @return bool|string
 */
function readAudioFile(string $filePath) {
    $voiceData = file_get_contents($filePath);
    if ($voiceData === false) {
        echo "音声ファイルを読み込めません\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 "finish-task 命令の送信準備: " . $finishTaskMessage . "\n";
    $conn->send($finishTaskMessage);
    echo "finish-task 命令を送信しました\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 "タスク開始、音声データを送信中...\n";
            // 音声ファイルを読み込む
            $voiceData = readAudioFile($audio_file_path);
            if ($voiceData === false) {
                echo "音声ファイルを読み込めません\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);
                    // 次のチャンクを 100ms 後に送信
                    $loop->addTimer(0.1, $sendChunk);
                } else {
                    echo "すべてのデータチャンクを送信しました\n";
                    $allChunksSent = true;

                    // finish-task 命令を送信
                    sendFinishTaskMessage($conn, $taskId);
                }
            };

            // 音声データの送信を開始
            $sendChunk();
            break;
        case 'result-generated':
            $result = $response['payload']['output']['sentence'];
            echo "認識結果: " . $result['text'] . "\n";
            if (isset($response['payload']['usage']['duration'])) {
                echo "タスク課金時間 (秒): " . $response['payload']['usage']['duration'] . "\n";
            }
            break;
        case 'task-finished':
            echo "タスク完了\n";
            $conn->close();
            break;
        case 'task-failed':
            echo "タスク失敗\n";
            echo "エラーコード: " . $response['header']['error_code'] . "\n";
            echo "エラーメッセージ: " . $response['header']['error_message'] . "\n";
            $conn->close();
            break;
        case 'error':
            echo "エラー: " . $response['payload']['message'] . "\n";
            break;
        default:
            echo "不明なイベント: " . $response['header']['event'] . "\n";
            break;
    }

    // すべてのデータが送信され、タスクが完了した場合、接続を閉じる
    if ($allChunksSent && $response['header']['event'] == 'task-finished') {
        // すべてのデータが送信されたことを確認するために 1 秒待機
        $loop->addTimer(1, function() use ($conn) {
            $conn->close();
            echo "クライアントが接続を閉じました\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 = "{YOUR_AUDIO_FILE}"                                   // 音声ファイルのパスに置き換えてください
)

var dialer = websocket.DefaultDialer

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

	// WebSocket サービスに接続
	conn, err := connectWebSocket(apiKey)
	if err != nil {
		log.Fatal("WebSocket への接続に失敗しました: ", err)
	}
	defer closeConnection(conn)

	// 結果を受信するゴルーチンを開始
	taskStarted := make(chan bool)
	taskDone := make(chan bool)
	startResultReceiver(conn, taskStarted, taskDone)

	// run-task 命令を送信
	taskID, err := sendRunTaskCmd(conn)
	if err != nil {
		log.Fatal("run-task 命令の送信に失敗しました: ", err)
	}

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

	// 認識対象の音声ファイルストリームを送信
	if err := sendAudioData(conn); err != nil {
		log.Fatal("音声の送信に失敗しました: ", err)
	}

	// finish-task 命令を送信
	if err := sendFinishTaskCmd(conn, taskID); err != nil {
		log.Fatal("finish-task 命令の送信に失敗しました: ", 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 メッセージを非同期で受信するゴルーチンを開始
func startResultReceiver(conn *websocket.Conn, taskStarted chan<- bool, taskDone chan<- bool) {
	go func() {
		for {
			_, message, err := conn.ReadMessage()
			if err != nil {
				log.Println("サーバーメッセージの解析に失敗しました: ", err)
				return
			}
			var event Event
			err = json.Unmarshal(message, &event)
			if err != nil {
				log.Println("イベントの解析に失敗しました: ", 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:     "qwen-audio-3.0-asr-flash-streaming",
			Parameters: Params{
				Format:     "wav",
				SampleRate: 16000,
			},
			Input: Input{},
		},
	}
	runTaskCmdJSON, err := json.Marshal(runTaskCmd)
	return string(runTaskCmdJSON), taskID, err
}

// task-started イベントを待機
func waitForTaskStarted(taskStarted chan bool) {
	select {
	case <-taskStarted:
		fmt.Println("タスクが正常に開始されました")
	case <-time.After(10 * time.Second):
		log.Fatal("task-started の待機がタイムアウトしました。タスクの開始に失敗しました")
	}
}

// 音声データを送信
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("task-started イベントを受信しました")
		taskStarted <- true
	case "result-generated":
		if event.Payload.Output.Sentence.Text != "" {
			fmt.Println("認識結果: ", event.Payload.Output.Sentence.Text)
		}
		if event.Payload.Usage != nil {
			fmt.Println("タスク課金時間 (秒): ", event.Payload.Usage.Duration)
		}
	case "task-finished":
		fmt.Println("タスク完了")
		taskDone <- true
		return true
	case "task-failed":
		handleTaskFailed(event, conn)
		taskDone <- true
		return true
	default:
		log.Printf("予期しないイベント: %v", event)
	}
	return false
}

// task-failed イベントを処理
func handleTaskFailed(event Event, conn *websocket.Conn) {
	if event.Header.ErrorMessage != "" {
		log.Fatalf("タスク失敗: %s", event.Header.ErrorMessage)
	} else {
		log.Fatal("タスクが不明な理由で失敗しました")
	}
}

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

Qwen3-ASR-Flash-Realtime

注記サンプルコードは 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/zh/model-studio/get-api-key
# 環境変数を設定していない場合は、次の行を Alibaba Cloud 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"サーバーに接続中: {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("サーバーに接続しました。")

    # セッション更新イベント
    event_manual = {
        "event_id": "event_123",
        "type": "session.update",
        "session": {
            "modalities": ["text"],
            "input_audio_format": "pcm",
            "sample_rate": 16000,
            # "input_audio_transcription": {
            #     # 言語識別子、オプション。言語がわかっている場合は設定することを推奨
            #     "language": "zh"
            # },
            "turn_detection": None
        }
    }
    event_vad = {
        "event_id": "event_123",
        "type": "session.update",
        "session": {
            "modalities": ["text"],
            "input_audio_format": "pcm",
            "sample_rate": 16000,
            # "input_audio_transcription": {
            #     "language": "zh"
            # },
            "turn_detection": {
                "type": "server_vad",
                "threshold": 0.0,
                "silence_duration_ms": 400
            }
        }
    }
    if enableServerVad:
        logger.info(f"イベントを送信: {json.dumps(event_vad, indent=2)}")
        ws.send(json.dumps(event_vad))
    else:
        logger.info(f"イベントを送信: {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"イベントを受信: {json.dumps(data, ensure_ascii=False, indent=2)}")
        if data.get("type") == "conversation.item.input_audio_transcription.completed":
            logger.info(f"最終文字起こし: {data.get('transcript')}")
        elif data.get("type") == "session.finished":
            logger.info("セッション完了後に WebSocket 接続を閉じています...")
            is_running = False  # 音声送信スレッドを停止
            ws.close()
    except json.JSONDecodeError:
        logger.error(f"メッセージの解析に失敗しました: {message}")

def on_error(ws, error):
    logger.error(f"エラー: {error}")

def on_close(ws, close_status_code, close_msg):
    logger.info(f"接続が閉じられました: {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"ファイル読み取り開始: {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"ファイル読み取り完了: {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 が閉じられています。音声送信を停止します。")
                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"音声イベントを送信: {eventd['event_id']}")
            time.sleep(0.1)  # リアルタイムキャプチャをシミュレート

# ロギングを初期化
init_logger()
logger.info(f"WebSocket サーバー {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 依存関係をインストールしてください:

<dependency>
    <groupId>org.java-websocket</groupId>
    <artifactId>Java-WebSocket</artifactId>
    <version>1.5.6</version>
</dependency>
implementation 'org.java-websocket:Java-WebSocket:1.5.6'
import org.java_websocket.client.WebSocketClient;
import org.java_websocket.handshake.ServerHandshake;
import org.json.JSONObject;

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

public class QwenASRRealtimeClient {

    private static final Logger logger = Logger.getLogger(QwenASRRealtimeClient.class.getName());
    // API キーはシンガポールリージョンと北京リージョンで異なります。API キーの取得: https://www.alibabacloud.com/help/zh/model-studio/get-api-key
    // 環境変数を設定していない場合は、次の行を Alibaba Cloud 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("サーバーに接続中: " + url);

        client = new WebSocketClient(new URI(url)) {
            @Override
            public void onOpen(ServerHandshake handshake) {
                logger.info("サーバーに接続しました。");
                sendSessionUpdate();
            }

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

                    logger.info("イベントを受信: " + data.toString(2));

                    // 最終認識結果は transcription.completed イベントにあります
                    if ("conversation.item.input_audio_transcription.completed".equals(eventType)) {
                        logger.info("最終文字起こし: " + data.optString("transcript"));
                    }

                    // finished イベントを受信したら → 送信スレッドを停止し、接続を閉じる
                    if ("session.finished".equals(eventType)) {
                        logger.info("セッション完了後に WebSocket 接続を閉じています...");

                        isRunning.set(false); // 音声送信スレッドを停止
                        if (this.isOpen()) {
                            this.close(1000, "ASR finished");
                        }
                    }
                } catch (Exception e) {
                    logger.severe("メッセージの解析に失敗しました: " + message);
                }
            }

            @Override
            public void onClose(int code, String reason, boolean remote) {
                logger.info("接続が閉じられました: " + code + " - " + reason);
            }

            @Override
            public void onError(Exception ex) {
                logger.severe("エラー: " + 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("音声送信スレッドエラー: " + e.getMessage());
            }
        });
        audioThread.start();
    }

    /** セッション更新イベント (VAD の有効化/無効化) */
    private static void sendSessionUpdate() {
        JSONObject eventNoVad = new JSONObject()
                .put("event_id", "event_123")
                .put("type", "session.update")
                .put("session", new JSONObject()
                        .put("modalities", new String[]{"text"})
                        .put("input_audio_format", "pcm")
                        .put("sample_rate", 16000)
                        // .put("input_audio_transcription", new JSONObject()
                        //         .put("language", "zh"))
                        .put("turn_detection", JSONObject.NULL) // マニュアルモード
                );

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

        if (enableServerVad) {
            logger.info("イベントを送信 (VAD):\n" + eventVad.toString(2));
            client.send(eventVad.toString());
        } else {
            logger.info("イベントを送信 (マニュアル):\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("ファイル読み取り開始");

        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("音声イベントを送信: " + eventd.getString("event_id"));
            } else {
                break; // 切断後は送信を継続しない
            }

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

        logger.info("ファイル読み取り完了");

        if (client != null && client.isOpen()) {
            // 非 VAD モードではコミットが必要
            if (!enableServerVad) {
                JSONObject commitEvent = new JSONObject()
                        .put("event_id", "event_789")
                        .put("type", "input_audio_buffer.commit");
                client.send(commitEvent.toString());
                logger.info("マニュアルモードのコミットイベントを送信しました。");
            }

            JSONObject finishEvent = new JSONObject()
                    .put("event_id", "event_987")
                    .put("type", "session.finish");
            client.send(finishEvent.toString());
            logger.info("終了イベントを送信しました。");
        }
    }

    /** ロギングを初期化 */
    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 リアルタイム WebSocket クライアント (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/zh/model-studio/get-api-key
// 環境変数を設定していない場合は、次の行を Alibaba Cloud 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; // true の場合 VAD モード、false の場合マニュアルモード
const localAudioPath = 'your_audio_file.pcm'; // PCM16、16kHz 音声ファイルのパス

// 以下はシンガポールリージョンの設定です。呼び出し時に "{WorkspaceId}" を実際のワークスペース ID に置き換えてください。リージョンによって設定が異なります。
const baseUrl = 'wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/realtime';
const url = `${baseUrl}?model=${MODEL}`;

console.log(`サーバーに接続中: ${url}`);

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

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

// ===== イベントバインディング =====
ws.on('open', () => {
    console.log('[WebSocket] サーバーに接続しました。');
    sendSessionUpdate();
    // 音声送信スレッドを開始
    sendAudio(localAudioPath);
});

ws.on('message', (message) => {
    try {
        const data = JSON.parse(message);
        console.log('[受信イベント]:', JSON.stringify(data, null, 2));

        // 最終認識結果は transcription.completed イベントにあります
        if (data.type === 'conversation.item.input_audio_transcription.completed') {
            console.log(`[最終文字起こし] ${data.transcript}`);
        }

        // finished イベントを受信
        if (data.type === 'session.finished') {
            console.log('[アクション] セッション完了後に WebSocket 接続を閉じています...');

            if (ws.readyState === WebSocket.OPEN) {
                ws.close(1000, 'ASR finished');
            }
        }
    } catch (e) {
        console.error('[エラー] メッセージの解析に失敗しました:', message);
    }
});

ws.on('close', (code, reason) => {
    console.log(`[WebSocket] 接続が閉じられました: ${code} - ${reason}`);
});

ws.on('error', (err) => {
    console.error('[WebSocket エラー]', 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: 'zh'
            // },
            turn_detection: null
        }
    };

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

    if (enableServerVad) {
        console.log('[イベント送信] VAD モード:\n', JSON.stringify(eventVad, null, 2));
        ws.send(JSON.stringify(eventVad));
    } else {
        console.log('[イベント送信] マニュアルモード:\n', JSON.stringify(eventNoVad, null, 2));
        ws.send(JSON.stringify(eventNoVad));
    }
}

// ===== 音声ファイルストリームを送信 =====
function sendAudio(audioPath) {
    setTimeout(() => {
        console.log(`[ファイル読み取り開始] ${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('[ファイル読み取り完了]');
                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('[コミットイベント送信]');
                    }

                    const finishEvent = {
                        event_id: 'event_987',
                        type: 'session.finish'
                    };
                    ws.send(JSON.stringify(finishEvent));
                    console.log('[終了イベント送信]');
                }

                return;
            }

            if (ws.readyState !== WebSocket.OPEN) {
                console.log('[停止] WebSocket が開かれていません。');
                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(`[音声イベント送信] ${appendEvent.event_id}`);

            setTimeout(sendChunk, 100); // リアルタイム送信をシミュレート
        }

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

C#

サンプルコードは以下の通りです:

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/zh/model-studio/get-api-key
    // 環境変数を設定していない場合は、次の行を Alibaba Cloud 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 環境変数が設定されていません。");
    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($"サーバーに接続中: {url}");

        // 認証ヘッダーを設定
        _webSocket.Options.SetRequestHeader("Authorization", $"Bearer {ApiKey}");
        _webSocket.Options.SetRequestHeader("OpenAI-Beta", "realtime=v1");

        await _webSocket.ConnectAsync(new Uri(url), _cts.Token);
        Console.WriteLine("サーバーに接続しました。");

        // メッセージ受信タスクを開始
        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($"session.update を送信: {payload.ToJsonString()}");
        await SendAsync(payload.ToJsonString());
    }

    // 音声ストリームを送信 (100ms ごとに 1 つの PCM チャンクを送信)
    private static async Task SendAudioStreamAsync() {
        await Task.Delay(3000); // セッション設定の完了を待機
        const int chunkSize = 3200; // 16kHz 16bit モノラルで 100ms
        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("ファイル読み取り完了。");
        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($"イベントを受信: {type}");
                if (type == "conversation.item.input_audio_transcription.completed") {
                    Console.WriteLine($"最終文字起こし: {data!["transcript"]}");
                } else if (type == "session.finished") {
                    Console.WriteLine("セッション完了、終了中...");
                    _sessionFinished = true;
                    _isRunning = false;
                    break;
                }
            } catch (Exception ex) {
                Console.WriteLine($"受信エラー: {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/zh/model-studio/get-api-key
// 環境変数を設定していない場合は、次の行を Alibaba Cloud 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 "WebSocket サーバーに接続しました\n";

    // サーバー側イベントをリッスン
    $conn->on('message', function($msg) use ($conn, &$is_running) {
        $event = json_decode($msg, true);
        if (!isset($event['type'])) {
            return;
        }
        echo "イベントを受信: {$event['type']}\n";
        if ($event['type'] === 'conversation.item.input_audio_transcription.completed') {
            echo "最終文字起こし: {$event['transcript']}\n";
        } elseif ($event['type'] === 'session.finished') {
            echo "セッション完了、終了中...\n";
            $is_running = false;
            $conn->close();
        }
    });
    $conn->on('close', function() {
        echo "接続が閉じられました\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 "接続できません: {$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 "session.update を送信しました\n";
}

// 音声ストリームを送信 (100ms ごとに 1 つの PCM チャンクを送信)
function sendAudioStream($conn, $audio_file_path, $enable_server_vad, $loop, &$is_running) {
    $fp = fopen($audio_file_path, 'rb');
    if (!$fp) {
        echo "音声ファイルを開けません\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); // 16kHz 16bit モノラルで 100ms
        if ($chunk === false || strlen($chunk) === 0) {
            fclose($fp);
            echo "音声ストリーム終了\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/zh/model-studio/get-api-key
	// 環境変数を設定していない場合は、次の行を Alibaba Cloud Model Studio API キーに置き換えてください: apiKey := "sk-xxx"
	apiKey := os.Getenv("DASHSCOPE_API_KEY")

	url := baseURL + "?model=" + model
	log.Printf("サーバーに接続中: %s", url)

	conn, err := connect(url, apiKey)
	if err != nil {
		log.Fatal("WebSocket への接続に失敗しました: ", err)
	}
	defer conn.Close()

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

	// session.update を送信
	if err := sendSessionUpdate(conn); err != nil {
		log.Fatal("session.update の送信に失敗しました: ", err)
	}

	// セッション設定の完了を待機
	time.Sleep(3 * time.Second)

	// 音声ストリームを送信
	if err := sendAudioStream(conn); err != nil {
		log.Fatal("音声の送信に失敗しました: ", 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("session.update を送信: %s", string(payload))
	return conn.WriteMessage(websocket.TextMessage, payload)
}

// 音声ストリームを送信 (100ms ごとに 1 つの PCM チャンクを送信)
func sendAudioStream(conn *websocket.Conn) error {
	f, err := os.Open(audioFile)
	if err != nil {
		return err
	}
	defer f.Close()

	chunk := make([]byte, 3200) // 16kHz 16bit モノラルで 100ms
	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("音声ストリーム終了")
	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("メッセージの読み取りエラー: ", err)
			sessionFinished <- true
			return
		}
		var evt ServerEvent
		if err := json.Unmarshal(msg, &evt); err != nil {
			log.Println("メッセージの解析エラー: ", err)
			continue
		}
		log.Printf("イベントを受信: %s", evt.Type)
		switch evt.Type {
		case "conversation.item.input_audio_transcription.completed":
			log.Printf("最終文字起こし: %s", evt.Transcript)
		case "session.finished":
			log.Println("セッション完了")
			sessionFinished <- true
			return
		}
	}
}

Paraformer

Paraformer のサンプルコードは、Qwen-Audio-3.0-ASR-Flash-Streaming/Fun-ASR-Realtime のものと似ています。モデル名を Paraformer モデルに置き換えてください。

本番環境に適用

接続の再利用 (WebSocket)

Qwen-Audio-3.0-ASR-Flash-Streaming/Fun-ASR-Realtime および Paraformer の WebSocket 接続は再利用をサポートしています。1 つの認識タスクが完了した後、接続を再確立せずに次のタスクを開始できます。

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

重要

  1. 新しいタスクを開始する前に、サーバーが task-finished イベントを返すのを待機してください。
  2. 再利用された接続上の異なるタスクは、異なる task_id 値を使用する必要があります。
  3. タスクが失敗すると、サーバーはエラーイベントを返して接続を閉じます。その接続は再利用できません。
  4. タスク終了後 60 秒以内に新しいタスクが開始されない場合、接続は自動的に閉じられます。

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

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

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

DashScope SDK には組み込みのプーリング機構があり、WebSocket 接続と認識オブジェクトを再利用することで、頻繁な作成と破棄のオーバーヘッドを回避します。

重要現在、Paraformer Java SDK のみがこの機能をサポートしています。

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

前提条件

  • API キーを取得してください。
  • DashScope SDK がインストールされており、バージョン要件を満たしていること。最新バージョンをインストールすることを推奨します: Java SDK バージョン 2.16.9 以降。

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

  • 接続プール: SDK に統合された OkHttp3 接続プールが基盤となる WebSocket 接続を管理・再利用し、ネットワークハンドシェイクのオーバーヘッドを削減します。この機能はデフォルトで有効です。
  • オブジェクトプール: commons-pool2 上に構築されたオブジェクトプールが、すでに接続が確立された Recognition オブジェクトのセットを維持します。プールからオブジェクトを借用することで、接続設定の遅延を排除し、最初のパケット遅延を大幅に削減できます。

実装手順

  1. 依存関係を追加

    プロジェクトのビルドツールに基づいて、依存関係設定ファイルに dashscope-sdk-java と commons-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. Maven コマンド (例: mvn clean install または mvn compile) を実行して、プロジェクトの依存関係を更新します。

    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'
    }
    
    1. build.gradle ファイルを保存します。
    2. コマンドラインでプロジェクトのルートディレクトリに移動し、以下の 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;
    }
}
  1. オブジェクトプールから Recognition オブジェクトを借用

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

recognizer = RecognitionObjectPool.getInstance().borrowObject();
  1. 音声認識を実行

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

  2. 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;
import com.alibaba.dashscope.utils.Constants;

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

    public static void main(String[] args)
            throws NoApiKeyException, InterruptedException {
        // 以下は中国 (北京) リージョンの設定です。呼び出し時に "{WorkspaceId}" を実際のワークスペース ID に置き換えてください。リージョンによって設定が異なります。
        Constants.baseHttpApiUrl = "https://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api/v1";
        checkoutEnv("DASHSCOPE_CONNECTION_POOL_SIZE", 32);
        checkoutEnv("DASHSCOPE_MAXIMUM_ASYNC_REQUESTS", 32);
        checkoutEnv("DASHSCOPE_MAXIMUM_ASYNC_REQUESTS_PER_HOST", 32);
        checkoutEnv(RecognitionObjectPool.RECOGNITION_OBJECTPOOL_SIZE_ENV,
                RecognitionObjectPool.DEFAULT_OBJECT_POOL_SIZE);

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

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

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

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

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

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

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

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

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

    private static String getDashScopeApiKey() throws NoApiKeyException {
        String dashScopeApiKey = null;
        try {
            ApiKey apiKey = new ApiKey();
            dashScopeApiKey = ApiKey.getApiKey(null);
        } catch (NoApiKeyException e) {
            System.out.println("環境変数に API キーが見つかりません。");
        }
        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
                                        + "] 認識完了");
                            }

                            @Override
                            public void onError(Exception e) {
                                System.out.println("[" + threadName
                                        + "] RecognitionCallback error: " + e.getMessage());
                                hasError[0] = true;
                            }
                        };
                System.out.println("[" + threadName
                        + "] 入力 file_path: " + filePath);
                FileInputStream fis = null;
                try {
                    fis = new FileInputStream(filePath.toFile());
                } catch (Exception e) {
                    System.out.println("ファイルの読み込みエラー: " + filePath);
                    e.printStackTrace();
                }
                recognizer.call(param, callback);

                // 16KHz サンプルレートでチャンクサイズを 100 ms に設定
                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 + "] 音声送信完了");
                recognizer.stop();
                System.out.println("[" + threadName + "] asr タスク完了");
            } 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 vCPU、8 GiB

100

500

2000

8 vCPU、16 GiB

200

500

2000

16 vCPU、32 GiB

400

500

2000

リソース管理とエラー処理

  • タスク成功時: GenericObjectPool.returnObject() を呼び出して Recognition オブジェクトをプールに戻し、再利用可能にします。

    重要未完了または失敗したタスクの Recognition オブジェクトは返却しないでください。

  • タスク失敗時: SDK またはビジネスロジックによって例外がスローされタスクが中断された場合、以下の 2 つのアクションを実行します:

    1. 基盤となる WebSocket 接続を積極的に閉じます。
    2. オブジェクトプール内のオブジェクトを無効化し、再利用されないようにします。
// 接続を閉じる。
recognizer.getDuplexApi().close(1000, "bye");
// オブジェクトプール内の失敗した recognizer を無効化。
RecognitionObjectPool.getInstance().invalidateObject(recognizer);
  • サービスが TaskFailed エラーを返す場合、追加の処理は不要です。

ウォームアップと遅延測定

DashScope Java SDK の同時実行呼び出し遅延などのパフォーマンスを評価する際は、正式なテストの前に十分なウォームアップを実行することを推奨します。

接続再利用メカニズム

DashScope Java SDK は、グローバルシングルトン接続プールを通じて WebSocket 接続を管理・再利用します。このメカニズムは以下の通りです:

  • オンデマンド作成: SDK はサービス起動時に WebSocket 接続を事前作成しません。代わりに、最初の呼び出し時に必要に応じて接続を確立します。

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

    • 60 秒以内に新しいリクエストが到着した場合、SDK は既存の接続を再利用し、繰り返しハンドシェイクのオーバーヘッドを回避します。
    • 接続が 60 秒以上アイドル状態になると、SDK は自動的に接続を閉じてリソースを解放します。
ウォームアップの重要性

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

  • アプリケーションは起動したばかりで、まだ呼び出しを行っていません。
  • サービスが 60 秒以上アイドル状態だったため、タイムアウトにより接続プールが閉じられました。

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

推奨アプローチ

正式な負荷テストまたは遅延測定を実行する前に、以下のウォームアップ手順に従ってください:

  1. 正式なテストの同時実行レベルをシミュレートするために、事前に十分な数の呼び出しを送信します (例: 1~2 分間)。これにより接続プールが完全に充填されます。
  2. 接続プールが十分なアクティブ接続を確立・維持していることを確認した後、正式なパフォーマンスデータの収集を開始します。

認識精度の向上

  • サンプルレートに一致するモデルを選択: 8 kHz の電話音声の場合は、直接 8 kHz モデルを使用してください。これにより、16 kHz にアップサンプリングする際に発生する情報損失を回避できます。
  • 入力音声品質の向上: 高品質なマイクを使用し、信号対雑音比が高くエコーのない環境で録音してください。アプリケーション層では、ノイズリダクション (例: RNNoise) や音響エコーキャンセレーション (AEC) などのアルゴリズムを統合して前処理を行えます。

フォールトトレランス戦略の設定

  • クライアント側再接続: クライアントはネットワークジッターを処理するために自動再接続を実装する必要があります。Python SDK の参考実装は以下の通りです:

    1. 例外をキャッチ: on_error メソッドを Callback クラスに実装します。dashscope SDK はネットワークエラーやその他の問題に遭遇した際にこのメソッドを呼び出します。
    2. 状態を通知: on_error がトリガーされたら、再接続シグナルを設定します。Python では、threading.Event (スレッドセーフなシグナルフラグ) を使用できます。
    3. 再接続ループ: メインロジックを for ループ (例: 3 回リトライ) でラップします。再接続シグナルが検出されたら、現在の認識ラウンドを中断し、リソースをクリーンアップした後、数秒待ってループを再度実行してまったく新しい接続を作成します。
  • 接続維持のためのハートビート設定: サーバーとの長時間接続を維持するには、heartbeat パラメーターを true に設定します。これにより、音声に長時間無音が含まれていてもサーバーへの接続が維持されます。

  • モデルのレート制限: モデル API を呼び出す際は、モデルのレート制限ルールに注意してください。

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

シンガポール

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

  • Qwen-Audio-3.0-ASR-Flash-Streaming: qwen-audio-3.0-asr-flash-streaming
  • Fun-ASR-Realtime: 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 キーを使用してください:

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

  • Fun-ASR-Realtime: 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 リファレンス

よくある質問

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

Qwen-Audio-3.0-ASR-Flash-Streaming、Fun-ASR-Realtime、Paraformer モデルは pcm、wav、mp3、opus、speex、aac、amr フォーマットをサポートしています。Qwen3-ASR-Flash-Realtime モデルでは、pcm または opus フォーマットを推奨します。 他のフォーマット (wav、aac、amr など) は session.update 検証レイヤーで受け入れられますが、サーバー側のデコードが失敗する可能性があります。送信前に音声ストリームが推奨フォーマットを使用していることを確認してください。

SDK と WebSocket API の違いは何ですか? どのように選べばよいですか?

DashScope SDK は WebSocket 接続管理、認証、再接続などの詳細をカプセル化しているため、迅速な統合に適しています。WebSocket API に直接接続すると、より細かい制御が可能で、SDK がカバーしていないプログラミング言語やカスタム接続管理が必要なシナリオに適しています。まずは SDK の使用を推奨します。

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

ホットワードまたはコンテキスト強化を使用してください。詳細な設定方法と使用上の注意点については、「認識精度の向上」をご参照ください。

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

クライアント側の再接続を実装し、長時間音声がない場合でも接続が切断されないようにハートビートパラメーター (heartbeat=true) を有効にしてください。詳細なフォールトトレランス戦略については、「本番環境に適用」をご参照ください。