Qwen-Omni-Realtime は、ストリーミング音声および画像入力 (ビデオフレームを含む) を処理し、テキストと音声の応答をリアルタイムで生成します。
サポートされているリージョン: シンガポール、中国 (北京)。各リージョンには、独自の API キー が必要です。
使用方法
1. 接続の確立
Qwen-Omni-Realtime は WebSocket と WebRTC をサポートしています。WebSocket はサーバー側の統合に適しており、迅速にセットアップできます。WebRTC はブラウザベースの低遅延音声シナリオを対象としており、UDP 経由で音声を送信し、エコーキャンセレーションとノイズリダクション機能を内蔵しています。
WebSocket
ネイティブ WebSocket
接続パラメーター:
| パラメーター | 説明 |
|---|---|
エンドポイント | 中国 (北京) リージョン: wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/realtime シンガポールリージョン: wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/realtime
|
クエリパラメーター |
|
リクエストヘッダー | 認証には Bearer トークンを使用します:
|
# pip install websocket-client
import json
import websocket
import os
API_KEY=os.getenv("DASHSCOPE_API_KEY")
# シンガポールリージョン。{WorkspaceId} をお使いの Bailian ワークスペース ID に置き換えてください。URL はリージョンによって異なります。
API_URL = "wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/realtime?model=qwen3.5-omni-plus-realtime"
headers = [
"Authorization: Bearer " + API_KEY
]
def on_open(ws):
print(f"Connected to server: {API_URL}")
def on_message(ws, message):
data = json.loads(message)
print("Received event:", json.dumps(data, indent=2))
def on_error(ws, error):
print("Error:", error)
ws = websocket.WebSocketApp(
API_URL,
header=headers,
on_open=on_open,
on_message=on_message,
on_error=on_error
)
ws.run_forever()
DashScope Python SDK
# SDK バージョン 1.23.9 以降が必要です。
import os
import json
from dashscope.audio.qwen_omni import OmniRealtimeConversation,OmniRealtimeCallback
import dashscope
# シンガポールと中国 (北京) リージョンの API キーは異なります。API キーを取得するには、https://www.alibabacloud.com/help/model-studio/get-api-key をご参照ください。
# API キーを設定していない場合は、次の行を dashscope.api_key = "sk-xxx" に変更してください。
dashscope.api_key = os.getenv("DASHSCOPE_API_KEY")
class PrintCallback(OmniRealtimeCallback):
def on_open(self) -> None:
print("Connected Successfully")
def on_event(self, response: dict) -> None:
print("Received event:")
print(json.dumps(response, indent=2, ensure_ascii=False))
def on_close(self, close_status_code: int, close_msg: str) -> None:
print(f"Connection closed (code={close_status_code}, msg={close_msg}).")
callback = PrintCallback()
conversation = OmniRealtimeConversation(
model="qwen3.5-omni-plus-realtime",
callback=callback,
# 次の URL はシンガポールリージョン用です。呼び出し時に、{WorkspaceId} を実際のワークスペース ID に置き換えてください。URL はリージョンによって異なります。
url="wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/realtime"
)
try:
conversation.connect()
print("Conversation started. Press Ctrl+C to exit.")
conversation.thread.join()
except KeyboardInterrupt:
conversation.close()
DashScope Java SDK
// SDK バージョン 2.20.9 以降が必要です。
import com.alibaba.dashscope.audio.omni.*;
import com.alibaba.dashscope.exception.NoApiKeyException;
import com.google.gson.JsonObject;
import java.util.concurrent.CountDownLatch;
public class Main {
public static void main(String[] args) throws InterruptedException, NoApiKeyException {
CountDownLatch latch = new CountDownLatch(1);
OmniRealtimeParam param = OmniRealtimeParam.builder()
.model("qwen3.5-omni-plus-realtime")
.apikey(System.getenv("DASHSCOPE_API_KEY"))
// 次の URL はシンガポールリージョン用です。呼び出し時に、{WorkspaceId} を実際のワークスペース ID に置き換えてください。URL はリージョンによって異なります。
.url("wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/realtime")
.build();
OmniRealtimeConversation conversation = new OmniRealtimeConversation(param, new OmniRealtimeCallback() {
@Override
public void onOpen() {
System.out.println("Connected Successfully");
}
@Override
public void onEvent(JsonObject message) {
System.out.println(message);
}
@Override
public void onClose(int code, String reason) {
System.out.println("connection closed code: " + code + ", reason: " + reason);
latch.countDown();
}
});
conversation.connect();
latch.await();
conversation.close(1000, "bye");
System.exit(0);
}
}
WebRTC
WebRTC 接続の確立には 2 つの段階があります:
- SDP 交換 (HTTP):クライアントは、メディア機能とネットワークアドレス (Offer SDP) を HTTP POST 経由でサーバーに送信します。サーバーは、その情報 (Answer SDP) を返して、機能ネゴシエーションを完了します。
- 接続 (自動):ネゴシエーション後、WebRTC レイヤーは自動的に音声トランスポートチャネルを確立します。
SDP 交換の設定:
パラメーター | 説明 |
|---|---|
リクエスト URL | 中国 (北京) リージョン: https://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api/v1/webrtc/realtime シンガポールリージョン: https://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api/v1/webrtc/realtime
|
クエリパラメーター |
|
Content-Type | application/sdp |
リクエストヘッダー | Authorization: Bearer DASHSCOPE_API_KEY |
リクエストボディ | クライアントが生成した Offer SDP 文字列 |
レスポンス | 成功: HTTP 200 とサーバーの Answer SDP 文字列。失敗: HTTP 4xx と JSON エラーメッセージ。 |
接続コードの例:
# pip install aiortc aiohttp certifi
import asyncio, aiohttp, ssl, certifi
from aiortc import RTCPeerConnection, RTCConfiguration, RTCSessionDescription
from aiortc.mediastreams import AudioStreamTrack
API_KEY = "your-api-key"
MODEL = "qwen3.5-omni-plus-realtime"
# シンガポールリージョン。{WorkspaceId} をお使いの Bailian ワークスペース ID に置き換えてください。URL はリージョンによって異なります。
SIGNALING_URL = "https://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api/v1/webrtc/realtime?model=" + MODEL
async def connect():
pc = RTCPeerConnection(RTCConfiguration(iceServers=[]))
# Offer SDP に m=audio が含まれるようにオーディオトラックを追加 (サーバーで必須)
pc.addTrack(AudioStreamTrack())
# SDP ネゴシエーションをトリガーするために DataChannel を作成 (名前はカスタマイズ可能。サーバーは "txt" という名前のチャネルを通じてイベントをプッシュします)
pc.createDataChannel("oai-events")
# SDP 交換:Offer を作成し、サーバーに送信
offer = await pc.createOffer()
await pc.setLocalDescription(offer)
async with aiohttp.ClientSession() as session:
async with session.post(
SIGNALING_URL,
ssl=ssl.create_default_context(cafile=certifi.where()),
data=offer.sdp.encode("utf-8"),
headers={
"Content-Type": "application/sdp",
"Authorization": f"Bearer {API_KEY}",
},
) as resp:
if not resp.ok:
raise Exception(f"SDP exchange failed: {resp.status} {await resp.text()}")
answer_sdp = await resp.text()
print("=== Offer SDP ===")
print(offer.sdp)
print("=== Answer SDP ===")
print(answer_sdp)
# ICE 接続は自動的に確立されます
await pc.setRemoteDescription(RTCSessionDescription(sdp=answer_sdp, type="answer"))
print("WebRTC connection established")
return pc
const API_KEY = 'your-api-key';
// シンガポールリージョン。{WorkspaceId} をお使いの Bailian ワークスペース ID に置き換えてください。URL はリージョンによって異なります。
const API_URL = 'https://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api/v1/webrtc/realtime?model=qwen3.5-omni-plus-realtime';
async function connect() {
const pc = new RTCPeerConnection({ iceServers: [] });
// Offer SDP に m=audio が含まれるようにオーディオトラックを追加 (サーバーで必須)
const stream = await navigator.mediaDevices.getUserMedia({ audio: true });
stream.getAudioTracks().forEach(t => pc.addTrack(t, stream));
// SDP ネゴシエーションをトリガーするために DataChannel を作成 (名前はカスタマイズ可能。サーバーは "txt" という名前のチャネルを通じてイベントをプッシュします)
pc.createDataChannel('oai-events');
// Answer を取得するために Offer を送信する前に ICE 収集が完了するのを待つ
pc.onicegatheringstatechange = async () => {
if (pc.iceGatheringState !== 'complete') return;
const resp = await fetch(API_URL, {
method: 'POST',
headers: {
'Content-Type': 'application/sdp',
'Authorization': `Bearer ${API_KEY}`,
},
body: pc.localDescription.sdp,
});
if (!resp.ok) throw new Error('SDP exchange failed: ' + resp.status);
const answerSdp = await resp.text();
// ICE 接続は自動的に確立されます
await pc.setRemoteDescription({ type: 'answer', sdp: answerSdp });
console.log('WebRTC connection established');
};
// Offer を作成
const offer = await pc.createOffer();
await pc.setLocalDescription(offer);
return pc;
}
2. セッションの設定
session.update クライアントイベントを送信します:
{
// クライアントが生成したイベント ID。
"event_id": "event_ToPZqeobitzUJnt3QqtWg",
// イベントタイプ。「session.update」である必要があります。
"type": "session.update",
// セッション設定。
"session": {
// 出力モダリティ。テキストのみの出力の場合は ["text"] に、テキストと音声の両方を出力する場合は ["text", "audio"] に設定します。
"modalities": [
"text",
"audio"
],
// 音声出力の声。
"voice": "Ethan",
// 入力音声フォーマット。「pcm」のみがサポートされています。入力音声は、16 kHz のサンプルレートの PCM 音声ストリームである必要があります。
"input_audio_format": "pcm",
// 出力音声フォーマット。「pcm」のみがサポートされています。出力音声は、24 kHz のサンプルレートの PCM 音声ストリームです。
"output_audio_format": "pcm",
// モデルの目標や役割を定義するためのシステム命令。
"instructions": "あなたは五つ星ホテルの AI カスタマーサービスエージェントです。客室タイプ、施設、料金、予約ポリシーに関するお客様からのお問い合わせに、正確かつフレンドリーに回答してください。常にプロフェッショナルで親切な態度で対応してください。未確認の情報やホテルのサービスの範囲外の情報は提供しないでください。",
// サーバー側の音声区間検出 (VAD) を有効にします。有効にすると、サーバーは自動的に発話の開始と終了を検出します。
// null の場合、クライアントはモデルの応答をトリガーするタイミングを制御する必要があります。
"turn_detection": {
// VAD のタイプ。有効な値は "server_vad" と "semantic_vad" です。qwen3.5-omni-realtime シリーズのモデルには "semantic_vad" を推奨します。
"type": "semantic_vad",
// VAD 検出のしきい値。ノイズの多い環境ではこの値を大きくし、静かな環境では小さくすることを推奨します。
"threshold": 0.5,
// 発話の終了を示す無音の継続時間 (ミリ秒)。この時間を超えると、モデルは応答をトリガーします。
"silence_duration_ms": 800
}
}
}
3. 音声と画像の入力
音声入力は必須ですが、画像入力はオプションです。入力方法はプロトコルによって異なります。
WebSocket
input_audio_buffer.append および input_image_buffer.append イベントを使用して、Base64 エンコードされた音声および画像データをサーバーバッファーに送信します。
画像は、ローカルファイルまたはリアルタイムのビデオストリームキャプチャから取得できます。
サーバー側の VAD が有効になっている場合、サーバーは自動的にデータを送信し、発話の終了時に応答をトリガーします。VAD が無効 (手動モード) の場合、送信後に input_audio_buffer.commit イベントを呼び出してデータを送信します。
WebRTC
接続確立時に追加された音声およびビデオトラック (RTP メディアチャネル) は、データを自動的にサーバーに送信します。
- 音声:音声トラック (RTP) を介して直接送信されます。
input_audio_buffer.appendイベントは不要です。 - 画像:ビデオトラック (RTP) を介してビデオフレームとして送信されます。
input_image_buffer.appendはサポートされていません。
WebRTC はサーバー側の VAD モード (
server_vadまたはsemantic_vad) のみをサポートします。手動モードはサポートされていません。
4. モデル応答の受信
応答形式は、設定された出力モダリティによって異なります。
WebSocket
-
テキストのみ
response.text.delta イベントでストリーミングテキストを受信し、response.text.done イベントで完全なテキストを受信します。
-
テキストと音声
- テキスト:response.audio_transcript.delta イベントでストリーミングテキストを受信し、response.audio_transcript.done イベントで完全なテキストを受信します。
- 音声:response.audio.delta イベントで Base64 エンコードされたストリーミング音声を受信します。response.audio.done イベントは、音声生成が完了したことを示します。
WebRTC
-
テキストのみ
WebSocket と同じです。DataChannel を介してストリーミングテキストイベントを受信します。
-
テキストと音声
- テキスト:WebSocket と同じく、DataChannel を介してストリーミングテキストイベントとして受信されます。
- 音声:RTP トラックを介してリアルタイムで受信および再生されます。
response.audio.deltaイベントは不要です。
モデルの選択
Qwen3.5-Omni-Realtime は、Qwen3-Omni-Flash-Realtime に比べて以下の点で改善されています:
-
インテリジェンスレベル
Qwen3.5-Plus と同等です。
-
Web 検索
内蔵の Web 検索機能により、モデルは自律的に検索してリアルタイムの質問に答えます。詳細については、「Web 検索」をご参照ください。
-
ツール呼び出し
関数呼び出しにより、モデルは自律的に外部ツールを呼び出します。詳細については、「Qwen-Omni-Realtime シリーズ」をご参照ください。
-
セマンティック割り込み
会話の意図を識別し、相槌やバックグラウンドノイズによる中断を防ぎます。
-
音声コントロール
「もっと速く話して」「もっと大きな声で」「楽しいトーンで」などの音声コマンドで、音量、話す速度、感情をコントロールします。
-
サポートされている言語
113 の言語と方言の音声認識と、36 の言語と方言の音声生成をサポートします。
-
サポートされている音声
47 の多言語音声と 8 つの方言音声を含む 55 の音声をサポートします。完全なリストについては、「音声リスト」をご参照ください。
-
音声クローン
リアルタイムの会話にカスタムのクローン音声を使用します (Qwen3.5-omni-plus-realtime および Qwen3.5-omni-flash-realtime)。詳細については、「音声クローン」をご参照ください。
モデル名、コンテキスト、価格、スナップショットバージョンについては、Model Studio コンソールをご確認ください。同時実行数のレート制限については、「レート制限」をご参照ください。
制限事項
-
Web 検索とツール呼び出しは相互排他的です。
-
1 つの WebSocket セッションは最大 120 分まで継続できます。この制限に達すると、接続は自動的に閉じられます。
-
モデルは、以下のターン数と持続時間の制限まで会話履歴を保持します。超過した場合、最も古い履歴が破棄されます。最大持続時間は、コンテキストに保持される累積的な音声またはビデオ (画像フレーム) の持続時間です。
ビデオは抽出されたフレームとして入力されます (推奨:1 fps)。ビデオの最大持続時間は、保持される累積フレーム持続時間です。たとえば、240 秒は、最後の 240 秒間のフレームのみが保持されることを意味します。
qwen3-omni-flash-realtimeモデルには 8 対話ターン (通常、最初に到達) の制限があります。その持続時間の制限はモデルのコンテキスト長に依存し、別途記載されていません。モデル
音声の最大ターン数
ビデオの最大ターン数
音声の最大持続時間
ビデオの最大持続時間
qwen3.5-omni-plus-realtime
100 ターン
50 ターン
600 秒
240 秒
qwen3.5-omni-flash-realtime
80 ターン
50 ターン
480 秒
120 秒
qwen3-omni-flash-realtime
8 ターン
8 ターン
—
—
クイックスタート
プログラミング言語を選択し、手順に従ってリアルタイムチャットを開始します。
WebSocket
DashScope Python SDK
- 実行環境
Python 3.10 以降がインストールされていることを確認してください。
お使いのオペレーティングシステムに PyAudio をインストールします。
macOS
brew install portaudio && pip install pyaudio
Debian/Ubuntu
- 仮想環境を使用していない場合は、システムパッケージマネージャーを使用して直接インストールできます:
sudo apt-get install python3-pyaudio
- 仮想環境を使用している場合は、まずビルド依存関係をインストールします:
sudo apt update
sudo apt install -y python3-dev portaudio19-dev
次に、アクティブ化された仮想環境で pip を使用してインストールします:
pip install pyaudio
CentOS
sudo yum install -y portaudio portaudio-devel && pip install pyaudio
Windows
pip install pyaudio
その他の依存関係をインストールします:
pip install websocket-client dashscope
-
対話モード
-
VAD モード (音声区間検出、発話の開始と終了を自動的に検出)
サーバーは、ユーザーの発話の終了を検出した後に応答します。
-
手動モード (押して話す、離して送信)
クライアントが発話の開始と終了を制御します。発話後、アプリケーションはサーバーに通知する必要があります。
VAD モード
vad_dash.py という名前の Python ファイルを作成し、次のコードをファイルにコピーします:
vad_dash.py
# 依存関係:dashscope >= 1.23.9, pyaudio import os import base64 import time import pyaudio from dashscope.audio.qwen_omni import MultiModality, AudioFormat,OmniRealtimeCallback,OmniRealtimeConversation import dashscope # 設定:URL、API キー、音声、モデル、モデルの役割 # リージョンを指定します。「intl」はシンガポールリージョン、「cn」は中国 (北京) リージョンです。{WorkspaceId} をお使いの Bailian ワークスペース ID に置き換えてください。 region = 'intl' base_domain = '{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com' if region == 'intl' else '{WorkspaceId}.cn-beijing.maas.aliyuncs.com' url = f'wss://{base_domain}/api-ws/v1/realtime' # API キーを設定します。環境変数が設定されていない場合は、次の行を API キーに置き換えてください:dashscope.api_key = "sk-xxx" dashscope.api_key = os.getenv('DASHSCOPE_API_KEY') # 音声を指定します。 voice = 'Ethan' # モデルを指定します。 model = 'qwen3.5-omni-plus-realtime' # モデルの役割を指定します。 instructions = "あなたはシャオユンという名のパーソナルアシスタントです。ユーザーの質問に、ユーモアとウィットに富んだ方法で答えてください。" class SimpleCallback(OmniRealtimeCallback): def __init__(self, pya): self.pya = pya self.out = None def on_open(self): # 音声出力ストリームを初期化します。 self.out = self.pya.open( format=pyaudio.paInt16, channels=1, rate=24000, output=True ) def on_event(self, response): if response['type'] == 'response.audio.delta': # 音声を再生します。 self.out.write(base64.b64decode(response['delta'])) elif response['type'] == 'conversation.item.input_audio_transcription.delta': # ストリーミングプレビュー:text は確定したプレフィックス、stash は未確定のサフィックスです。 preview = response.get('text', '') + response.get('stash', '') print(f"\r[User] {preview}", end='', flush=True) elif response['type'] == 'conversation.item.input_audio_transcription.completed': # 文字起こしが完了しました。最終的なテキストを印刷して改行します。 print(f"\r[User] {response['transcript']}") elif response['type'] == 'response.audio_transcript.done': # アシスタントの応答テキストを印刷します。 print(f"[LLM] {response['transcript']}") # 1. オーディオデバイスを初期化します。 pya = pyaudio.PyAudio() # 2. コールバック関数と会話を作成します。 callback = SimpleCallback(pya) conv = OmniRealtimeConversation(model=model, callback=callback, url=url) # 3. 接続してセッションを設定します。 conv.connect() conv.update_session(output_modalities=[MultiModality.AUDIO, MultiModality.TEXT], voice=voice, instructions=instructions) # 4. 音声入力ストリームを初期化します。 mic = pya.open(format=pyaudio.paInt16, channels=1, rate=16000, input=True) # 5. 音声入力を処理するメインループ。 print("Conversation started. Speak into the microphone (Ctrl+C to exit)...") try: while True: audio_data = mic.read(3200, exception_on_overflow=False) conv.append_audio(base64.b64encode(audio_data).decode()) time.sleep(0.01) except KeyboardInterrupt: # リソースをクリーンアップします。 conv.close() mic.close() callback.out.close() pya.terminate() print("\nConversation ended")vad_dash.pyを実行して、マイクを通じてリアルタイムの会話を開始します。システムは発話を検出し、音声をサーバーにストリーミングします。手動モード
manual_dash.pyという名前の Python ファイルを作成し、次のコードをファイルにコピーします:manual_dash.py
# 依存関係:dashscope >= 1.23.9, pyaudio import os import base64 import sys import threading import pyaudio from dashscope.audio.qwen_omni import * import dashscope # 環境変数が設定されていない場合は、次の行を API キーに置き換えてください:dashscope.api_key = "sk-xxx" dashscope.api_key = os.getenv('DASHSCOPE_API_KEY') voice = 'Ethan' class MyCallback(OmniRealtimeCallback): """最小限のコールバック:接続時にスピーカーを初期化し、返された音声をイベントハンドラで直接再生します。""" def __init__(self, ctx): super().__init__() self.ctx = ctx def on_open(self) -> None: # 接続が確立された後、PyAudio とスピーカー (24 kHz, モノラル, 16 ビット) を初期化します。 print('connection opened') try: self.ctx['pya'] = pyaudio.PyAudio() self.ctx['out'] = self.ctx['pya'].open( format=pyaudio.paInt16, channels=1, rate=24000, output=True ) print('audio output initialized') except Exception as e: print('[Error] audio init failed: {}'.format(e)) def on_close(self, close_status_code, close_msg) -> None: print('connection closed with code: {}, msg: {}'.format(close_status_code, close_msg)) sys.exit(0) def on_event(self, response: str) -> None: try: t = response['type'] handlers = { 'session.created': lambda r: print('start session: {}'.format(r['session']['id'])), 'conversation.item.input_audio_transcription.delta': lambda r: print('\rquestion: {}'.format(r.get('text', '') + r.get('stash', '')), end='', flush=True), 'conversation.item.input_audio_transcription.completed': self._transcription_completed, 'response.audio_transcript.delta': lambda r: print('llm text: {}'.format(r['delta'])), 'response.audio.delta': self._play_audio, 'response.done': self._response_done, } h = handlers.get(t) if h: h(response) except Exception as e: print('[Error] {}'.format(e)) def _transcription_completed(self, response): print() self.ctx['transcription_done'].set() def _play_audio(self, response): # base64 データをデコードし、再生のために出力ストリームに直接書き込みます。 if self.ctx['out'] is None: return try: data = base64.b64decode(response['delta']) self.ctx['out'].write(data) except Exception as e: print('[Error] audio playback failed: {}'.format(e)) def _response_done(self, response): # 現在のターンを完了としてマークし、メインループが続行できるようにします。 if self.ctx['conv'] is not None: print('[Metric] response: {}, first text delay: {}, first audio delay: {}'.format( self.ctx['conv'].get_last_response_id(), self.ctx['conv'].get_last_first_text_delay(), self.ctx['conv'].get_last_first_audio_delay(), )) if self.ctx['resp_done'] is not None: self.ctx['resp_done'].set() def shutdown_ctx(ctx): """オーディオと PyAudio リソースを安全に解放します。""" try: if ctx['out'] is not None: ctx['out'].close() ctx['out'] = None except Exception: pass try: if ctx['pya'] is not None: ctx['pya'].terminate() ctx['pya'] = None except Exception: pass def stream_record_and_send(pya_inst, conversation, sample_rate=16000, chunk_size=3200): stop_evt = threading.Event() stream = pya_inst.open( format=pyaudio.paInt16, channels=1, rate=sample_rate, input=True, frames_per_buffer=chunk_size ) def _reader(): while not stop_evt.is_set(): try: data = stream.read(chunk_size, exception_on_overflow=False) conversation.append_audio(base64.b64encode(data).decode()) except Exception: break t = threading.Thread(target=_reader, daemon=True) t.start() input() stop_evt.set() t.join(timeout=1.0) stream.close() if __name__ == '__main__': print('Initializing ...') # 実行時コンテキスト:オーディオと会話のハンドルを格納します。 ctx = {'pya': None, 'out': None, 'conv': None, 'resp_done': threading.Event(), 'transcription_done': threading.Event()} callback = MyCallback(ctx) conversation = OmniRealtimeConversation( model='qwen3.5-omni-plus-realtime', callback=callback, # 次の URL はシンガポールリージョン用です。呼び出し時に、{WorkspaceId} を実際のワークスペース ID に置き換えてください。URL はリージョンによって異なります。 url="wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/realtime", ) try: conversation.connect() except Exception as e: print('[Error] connect failed: {}'.format(e)) sys.exit(1) ctx['conv'] = conversation # セッション設定:テキストと音声出力を有効にし、手動録音のためにサーバー側の VAD を無効にします。 conversation.update_session( output_modalities=[MultiModality.AUDIO, MultiModality.TEXT], voice=voice, enable_input_audio_transcription=True, input_audio_transcription_model='qwen3-asr-flash-realtime', enable_turn_detection=False, instructions="あなたはシャオユンという名のパーソナルアシスタントです。ユーザーの質問に正確かつフレンドリーに回答し、常に親切な態度で対応してください。" ) try: turn = 1 while True: print(f"\n--- turn {turn} ---") print("Press Enter to start recording (enter q and press Enter to exit)...") user_input = input() if user_input.strip().lower() in ['q', 'quit']: print("User requested to exit...") break print("Recording... Press Enter again to stop recording.") if ctx['pya'] is None: ctx['pya'] = pyaudio.PyAudio() stream_record_and_send(ctx['pya'], conversation) ctx['transcription_done'].clear() ctx['resp_done'].clear() conversation.commit() ctx['transcription_done'].wait(timeout=10) print("Waiting for model response...") conversation.create_response() ctx['resp_done'].wait() turn += 1 except KeyboardInterrupt: print("\nProgram interrupted by user.") finally: shutdown_ctx(ctx) print("Program exited.")manual_dash.pyを実行します。Enter キーを押して録音を開始し、もう一度 Enter キーを押して停止して送信します。モデルの音声応答は自動的に再生されます。 -
DashScope Java SDK
対話モードの選択-
VAD モード (音声区間検出、発話の開始と終了を自動的に検出)
Realtime API は、発話の開始と停止を検出し、応答します。
-
手動モード (押して話す、離して送信)
クライアントが発話の開始と終了を制御します。発話後、クライアントはサーバーにメッセージを送信する必要があります。
VAD モード
OmniServerVad.java
import com.alibaba.dashscope.audio.omni.*;
import com.alibaba.dashscope.exception.NoApiKeyException;
import com.google.gson.JsonObject;
import javax.sound.sampled.*;
import java.nio.ByteBuffer;
import java.util.Arrays;
import java.util.Base64;
import java.util.Map;
import java.util.Queue;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.atomic.AtomicBoolean;
public class OmniServerVad {
static class SequentialAudioPlayer {
private final SourceDataLine line;
private final Queue<byte[]> audioQueue = new ConcurrentLinkedQueue<>();
private final Thread playerThread;
private final AtomicBoolean shouldStop = new AtomicBoolean(false);
public SequentialAudioPlayer() throws LineUnavailableException {
AudioFormat format = new AudioFormat(24000, 16, 1, true, false);
line = AudioSystem.getSourceDataLine(format);
line.open(format);
line.start();
playerThread = new Thread(() -> {
while (!shouldStop.get()) {
byte[] audio = audioQueue.poll();
if (audio != null) {
line.write(audio, 0, audio.length);
} else {
try { Thread.sleep(10); } catch (InterruptedException ignored) {}
}
}
}, "AudioPlayer");
playerThread.start();
}
public void play(String base64Audio) {
try {
byte[] audio = Base64.getDecoder().decode(base64Audio);
audioQueue.add(audio);
} catch (Exception e) {
System.err.println("Failed to decode audio: " + e.getMessage());
}
}
public void cancel() {
audioQueue.clear();
line.flush();
}
public void close() {
shouldStop.set(true);
try { playerThread.join(1000); } catch (InterruptedException ignored) {}
line.drain();
line.close();
}
}
public static void main(String[] args) {
try {
SequentialAudioPlayer player = new SequentialAudioPlayer();
AtomicBoolean userIsSpeaking = new AtomicBoolean(false);
AtomicBoolean shouldStop = new AtomicBoolean(false);
OmniRealtimeParam param = OmniRealtimeParam.builder()
.model("qwen3.5-omni-plus-realtime")
.apikey(System.getenv("DASHSCOPE_API_KEY"))
// 次の URL はシンガポールリージョン用です。呼び出し時に、{WorkspaceId} を実際のワークスペース ID に置き換えてください。URL はリージョンによって異なります。
.url("wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/realtime")
.build();
OmniRealtimeConversation conversation = new OmniRealtimeConversation(param, new OmniRealtimeCallback() {
@Override public void onOpen() {
System.out.println("Connection established.");
}
@Override public void onClose(int code, String reason) {
System.out.println("Connection closed (" + code + "): " + reason);
shouldStop.set(true);
}
@Override public void onEvent(JsonObject event) {
handleEvent(event, player, userIsSpeaking);
}
});
conversation.connect();
conversation.updateSession(OmniRealtimeConfig.builder()
.modalities(Arrays.asList(OmniRealtimeModality.AUDIO, OmniRealtimeModality.TEXT))
.voice("Ethan")
.enableTurnDetection(true)
.enableInputAudioTranscription(true)
.parameters(Map.of("instructions",
"あなたは五つ星ホテルの AI カスタマーサービスエージェントです。客室タイプ、施設、料金、予約ポリシーに関するお客様からのお問い合わせに、正確かつフレンドリーに回答してください。常にプロフェッショナルで親切な態度で対応してください。未確認の情報やホテルのサービスの範囲外の情報は提供しないでください。"))
.build()
);
System.out.println("Start speaking (speech start/end is automatically detected, press Ctrl+C to exit)...");
AudioFormat format = new AudioFormat(16000, 16, 1, true, false);
TargetDataLine mic = AudioSystem.getTargetDataLine(format);
mic.open(format);
mic.start();
ByteBuffer buffer = ByteBuffer.allocate(3200);
while (!shouldStop.get()) {
int bytesRead = mic.read(buffer.array(), 0, buffer.capacity());
if (bytesRead > 0) {
try {
conversation.appendAudio(Base64.getEncoder().encodeToString(buffer.array()));
} catch (Exception e) {
if (e.getMessage() != null && e.getMessage().contains("closed")) {
System.out.println("Conversation closed. Stopping recording.");
break;
}
}
}
Thread.sleep(20);
}
conversation.close(1000, "Normal termination");
player.close();
mic.close();
System.out.println("\nProgram exited.");
} catch (NoApiKeyException e) {
System.err.println("API key not found: Set the DASHSCOPE_API_KEY environment variable.");
System.exit(1);
} catch (Exception e) {
e.printStackTrace();
}
}
private static void handleEvent(JsonObject event, SequentialAudioPlayer player, AtomicBoolean userIsSpeaking) {
String type = event.get("type").getAsString();
switch (type) {
case "input_audio_buffer.speech_started":
System.out.println("\n[User started speaking]");
player.cancel();
userIsSpeaking.set(true);
break;
case "input_audio_buffer.speech_stopped":
System.out.println("[User stopped speaking]");
userIsSpeaking.set(false);
break;
case "response.audio.delta":
if (!userIsSpeaking.get()) {
player.play(event.get("delta").getAsString());
}
break;
case "conversation.item.input_audio_transcription.delta":
// ストリーミングプレビュー:text は確定したプレフィックス、stash は未確定のサフィックス
String preview2 = event.get("text").getAsString() + event.get("stash").getAsString();
System.out.print("\rUser: " + preview2);
break;
case "conversation.item.input_audio_transcription.completed":
System.out.println();
break;
case "response.audio_transcript.done":
System.out.println("Assistant: " + event.get("transcript").getAsString());
break;
case "response.done":
System.out.println("Response completed.");
break;
}
}
}
OmniServerVad.main() を実行して、マイクを通じてリアルタイムの会話を開始します。システムは発話を検出し、音声をサーバーに送信します。
手動モード
OmniWithoutServerVad.java
// DashScope Java SDK 2.20.9 以降が必要です。
import com.alibaba.dashscope.audio.omni.*;
import com.alibaba.dashscope.exception.NoApiKeyException;
import com.google.gson.JsonObject;
import javax.sound.sampled.*;
import java.io.IOException;
import java.util.Arrays;
import java.util.Base64;
import java.util.HashMap;
import java.util.Queue;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicReference;
public class Main {
public static class RealtimePcmPlayer {
private int sampleRate;
private SourceDataLine line;
private AudioFormat audioFormat;
private Thread decoderThread;
private Thread playerThread;
private AtomicBoolean stopped = new AtomicBoolean(false);
private Queue<String> b64AudioBuffer = new ConcurrentLinkedQueue<>();
private Queue<byte[]> RawAudioBuffer = new ConcurrentLinkedQueue<>();
public RealtimePcmPlayer(int sampleRate) throws LineUnavailableException {
this.sampleRate = sampleRate;
this.audioFormat = new AudioFormat(this.sampleRate, 16, 1, true, false);
DataLine.Info info = new DataLine.Info(SourceDataLine.class, audioFormat);
line = (SourceDataLine) AudioSystem.getLine(info);
line.open(audioFormat);
line.start();
decoderThread = new Thread(new Runnable() {
@Override
public void run() {
while (!stopped.get()) {
String b64Audio = b64AudioBuffer.poll();
if (b64Audio != null) {
byte[] rawAudio = Base64.getDecoder().decode(b64Audio);
RawAudioBuffer.add(rawAudio);
} else {
try {
Thread.sleep(100);
} catch (InterruptedException e) {
throw new RuntimeException(e);
}
}
}
}
});
playerThread = new Thread(new Runnable() {
@Override
public void run() {
while (!stopped.get()) {
byte[] rawAudio = RawAudioBuffer.poll();
if (rawAudio != null) {
try {
playChunk(rawAudio);
} catch (IOException e) {
throw new RuntimeException(e);
} catch (InterruptedException e) {
throw new RuntimeException(e);
}
} else {
try {
Thread.sleep(100);
} catch (InterruptedException e) {
throw new RuntimeException(e);
}
}
}
}
});
decoderThread.start();
playerThread.start();
}
// オーディオチャンクを再生し、再生が完了するまでブロックします。
private void playChunk(byte[] chunk) throws IOException, InterruptedException {
if (chunk == null || chunk.length == 0) return;
int bytesWritten = 0;
while (bytesWritten < chunk.length) {
bytesWritten += line.write(chunk, bytesWritten, chunk.length - bytesWritten);
}
int audioLength = chunk.length / (this.sampleRate*2/1000);
// バッファ内のオーディオが再生し終わるのを待ちます。
Thread.sleep(audioLength - 10);
}
public void write(String b64Audio) {
b64AudioBuffer.add(b64Audio);
}
public void cancel() {
b64AudioBuffer.clear();
RawAudioBuffer.clear();
}
public void waitForComplete() throws InterruptedException {
while (!b64AudioBuffer.isEmpty() || !RawAudioBuffer.isEmpty()) {
Thread.sleep(100);
}
line.drain();
}
public void shutdown() throws InterruptedException {
stopped.set(true);
decoderThread.join();
playerThread.join();
if (line != null && line.isRunning()) {
line.drain();
line.close();
}
}
}
// オーディオを録音し、リアルタイムで会話にストリーミングします。
private static void recordAndSend(TargetDataLine line, OmniRealtimeConversation conversation) {
byte[] buffer = new byte[3200];
AtomicBoolean stopRecording = new AtomicBoolean(false);
// Enter キーをリッスンするスレッドを開始します。
Thread enterKeyListener = new Thread(() -> {
try {
System.in.read();
stopRecording.set(true);
} catch (IOException e) {
e.printStackTrace();
}
});
enterKeyListener.start();
// リアルタイムでオーディオを録音して送信します。
while (!stopRecording.get()) {
int count = line.read(buffer, 0, buffer.length);
if (count > 0) {
byte[] chunk = new byte[count];
System.arraycopy(buffer, 0, chunk, 0, count);
conversation.appendAudio(Base64.getEncoder().encodeToString(chunk));
}
}
}
public static void main(String[] args) throws InterruptedException, LineUnavailableException {
OmniRealtimeParam param = OmniRealtimeParam.builder()
.model("qwen3.5-omni-plus-realtime")
// シンガポールと中国 (北京) の API キーは異なります。API キーは https://www.alibabacloud.com/help/model-studio/get-api-key から取得してください
// DASHSCOPE_API_KEY 環境変数が設定されていない場合は、以下の apikey 呼び出しを Model Studio API キーに置き換えてください。例:.apikey("sk-xxx")
.apikey(System.getenv("DASHSCOPE_API_KEY"))
// 次の URL はシンガポールリージョン用です。呼び出し時に、{WorkspaceId} を実際のワークスペース ID に置き換えてください。URL はリージョンによって異なります。
.url("wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/realtime")
.build();
AtomicReference<CountDownLatch> responseDoneLatch = new AtomicReference<>(null);
responseDoneLatch.set(new CountDownLatch(1));
AtomicReference<CountDownLatch> transcriptionDoneLatch = new AtomicReference<>(null);
transcriptionDoneLatch.set(new CountDownLatch(1));
RealtimePcmPlayer audioPlayer = new RealtimePcmPlayer(24000);
final AtomicReference<OmniRealtimeConversation> conversationRef = new AtomicReference<>(null);
OmniRealtimeConversation conversation = new OmniRealtimeConversation(param, new OmniRealtimeCallback() {
@Override
public void onOpen() {
System.out.println("connection opened");
}
@Override
public void onEvent(JsonObject message) {
String type = message.get("type").getAsString();
switch(type) {
case "session.created":
System.out.println("start session: " + message.get("session").getAsJsonObject().get("id").getAsString());
break;
case "conversation.item.input_audio_transcription.delta":
// ストリーミングプレビュー:text は確定したプレフィックス、stash は未確定のサフィックス
String transcriptPreview = message.get("text").getAsString() + message.get("stash").getAsString();
System.out.print("\rquestion: " + transcriptPreview);
break;
case "conversation.item.input_audio_transcription.completed":
System.out.println();
transcriptionDoneLatch.get().countDown();
break;
case "response.audio_transcript.delta":
System.out.println("got llm response delta: " + message.get("delta").getAsString());
break;
case "response.audio.delta":
String recvAudioB64 = message.get("delta").getAsString();
audioPlayer.write(recvAudioB64);
break;
case "response.done":
System.out.println("======RESPONSE DONE======");
if (conversationRef.get() != null) {
System.out.println("[Metric] response: " + conversationRef.get().getResponseId() +
", first text delay: " + conversationRef.get().getFirstTextDelay() +
" ms, first audio delay: " + conversationRef.get().getFirstAudioDelay() + " ms");
}
responseDoneLatch.get().countDown();
break;
default:
break;
}
}
@Override
public void onClose(int code, String reason) {
System.out.println("connection closed code: " + code + ", reason: " + reason);
}
});
conversationRef.set(conversation);
try {
conversation.connect();
} catch (NoApiKeyException e) {
throw new RuntimeException(e);
}
OmniRealtimeConfig config = OmniRealtimeConfig.builder()
.modalities(Arrays.asList(OmniRealtimeModality.AUDIO, OmniRealtimeModality.TEXT))
.voice("Ethan")
.enableTurnDetection(false)
// モデルの役割を設定します。
.parameters(new HashMap<String, Object>() {{
put("instructions","あなたはシャオユンという名のパーソナルアシスタントです。ユーザーの質問に正確かつフレンドリーに回答し、常に親切な態度で対応してください。");
}})
.build();
conversation.updateSession(config);
// 録音用にマイクを設定します。
AudioFormat format = new AudioFormat(16000, 16, 1, true, false);
DataLine.Info info = new DataLine.Info(TargetDataLine.class, format);
if (!AudioSystem.isLineSupported(info)) {
System.out.println("Line not supported");
return;
}
TargetDataLine line = null;
try {
line = (TargetDataLine) AudioSystem.getLine(info);
line.open(format);
line.start();
while (true) {
System.out.println("Press Enter to start recording...");
try {
System.in.read();
} catch (IOException e) {
System.err.println("Error reading input: " + e.getMessage());
break; // エラーが発生した場合はループを終了します。
}
System.out.println("Recording started. Speak now... Press Enter again to stop recording and send.");
recordAndSend(line, conversation);
conversation.commit();
// 出力が交錯するのを避けるため、モデルの応答をトリガーする前に文字起こしが完了するのを待ちます
transcriptionDoneLatch.get().await(10, TimeUnit.SECONDS);
System.out.println("Waiting for model response...");
conversation.createResponse(null, null);
responseDoneLatch.get().await();
// 次のターンのためにラッチをリセットします。
responseDoneLatch.set(new CountDownLatch(1));
transcriptionDoneLatch.set(new CountDownLatch(1));
}
} catch (LineUnavailableException e) {
e.printStackTrace();
} finally {
if (line != null) {
line.stop();
line.close();
}
}
}}
OmniWithoutServerVad.main() を実行します。Enter キーを押して録音を開始し、もう一度押して停止して送信します。モデルの応答は自動的に再生されます。
WebSocket (Python)
-
実行環境の準備
Python 3.10 以降がインストールされていることを確認してください。
お使いのオペレーティングシステムに pyaudio をインストールします。
macOS
brew install portaudio && pip install pyaudioDebian/Ubuntu
sudo apt-get install python3-pyaudio or pip install pyaudiopip install pyaudioの使用を推奨します。インストールに失敗した場合は、まずお使いのオペレーティングシステムのportaudio依存関係をインストールしてください。CentOS
sudo yum install -y portaudio portaudio-devel && pip install pyaudioWindows
pip install pyaudioWebSocket の依存関係をインストールします:
pip install websockets==15.0.1
-
クライアントの作成
omni_realtime_client.pyという名前のファイルを作成し、次のコードをコピーします:omni_realtime_client.py
import asyncio import websockets import json import base64 import time from typing import Optional, Callable, List, Dict, Any from enum import Enum class TurnDetectionMode(Enum): SERVER_VAD = "server_vad" SEMANTIC_VAD = "semantic_vad" # qwen3.5-omni-realtime シリーズのモデルに推奨 MANUAL = "manual" class OmniRealtimeClient: def __init__( self, base_url, api_key: str, model: str = "", voice: str = "Ethan", instructions: str = "You are a helpful assistant.", turn_detection_mode: TurnDetectionMode = TurnDetectionMode.SERVER_VAD, on_text_delta: Optional[Callable[[str], None]] = None, on_audio_delta: Optional[Callable[[bytes], None]] = None, on_input_transcript: Optional[Callable[[str], None]] = None, on_output_transcript: Optional[Callable[[str], None]] = None, extra_event_handlers: Optional[Dict[str, Callable[[Dict[str, Any]], None]]] = None ): self.base_url = base_url self.api_key = api_key self.model = model self.voice = voice self.instructions = instructions self.ws = None self.on_text_delta = on_text_delta self.on_audio_delta = on_audio_delta self.on_input_transcript = on_input_transcript self.on_output_transcript = on_output_transcript self.turn_detection_mode = turn_detection_mode self.extra_event_handlers = extra_event_handlers or {} # 現在の応答ステータス self._current_response_id = None self._current_item_id = None self._is_responding = False # 入力/出力トランスクリプトの印刷ステータス self._print_input_transcript = True self._output_transcript_buffer = "" async def connect(self) -> None: """Realtime API との WebSocket 接続を確立します。""" url = f"{self.base_url}?model={self.model}" headers = { "Authorization": f"Bearer {self.api_key}" } self.ws = await websockets.connect(url, additional_headers=headers) # セッション設定 session_config = { "modalities": ["text", "audio"], "voice": self.voice, "instructions": self.instructions, "input_audio_format": "pcm", "output_audio_format": "pcm", "input_audio_transcription": { "model": "qwen3-asr-flash-realtime" } } if self.turn_detection_mode == TurnDetectionMode.MANUAL: session_config['turn_detection'] = None await self.update_session(session_config) elif self.turn_detection_mode == TurnDetectionMode.SERVER_VAD: session_config['turn_detection'] = { "type": "server_vad", "threshold": 0.1, "prefix_padding_ms": 500, "silence_duration_ms": 900 } await self.update_session(session_config) elif self.turn_detection_mode == TurnDetectionMode.SEMANTIC_VAD: session_config['turn_detection'] = { "type": "semantic_vad", "threshold": 0.1, "prefix_padding_ms": 500, "silence_duration_ms": 900 } await self.update_session(session_config) else: raise ValueError(f"Invalid turn detection mode: {self.turn_detection_mode}") async def send_event(self, event) -> None: event['event_id'] = "event_" + str(int(time.time() * 1000)) await self.ws.send(json.dumps(event)) async def update_session(self, config: Dict[str, Any]) -> None: """セッション設定を更新します。""" event = { "type": "session.update", "session": config } await self.send_event(event) async def stream_audio(self, audio_chunk: bytes) -> None: """RAW オーディオデータを API にストリーミングします。""" # 16 ビット、16 kHz、モノラル PCM のみがサポートされています。 audio_b64 = base64.b64encode(audio_chunk).decode() append_event = { "type": "input_audio_buffer.append", "audio": audio_b64 } await self.send_event(append_event) async def commit_audio_buffer(self) -> None: """オーディオバッファをコミットして処理をトリガーします。""" event = { "type": "input_audio_buffer.commit" } await self.send_event(event) async def append_image(self, image_chunk: bytes) -> None: """画像データを画像バッファに追加します。 画像データは、ローカルファイルまたはリアルタイムのビデオストリームから取得できます。 注: - 画像フォーマットは JPG または JPEG である必要があります。解像度は 480p または 720p を推奨します。最大サポート解像度は 1080p です。 - Base64 エンコード後の単一画像のサイズは 256 KB を超えてはなりません。エンコード前の RAW 画像サイズは 190 KB 未満に保つことを推奨します。 - 送信する前に画像データを Base64 にエンコードします。 - 1 秒あたり 1 フレームのレートで画像を送信することを推奨します。 - 画像データを送信する前に、少なくとも 1 回は音声データを送信する必要があります。 """ image_b64 = base64.b64encode(image_chunk).decode() event = { "type": "input_image_buffer.append", "image": image_b64 } await self.send_event(event) async def create_response(self) -> None: """API に応答の生成をリクエストします。これは手動モードでのみ必要です。""" event = { "type": "response.create" } await self.send_event(event) async def cancel_response(self) -> None: """現在の応答をキャンセルします。""" event = { "type": "response.cancel" } await self.send_event(event) async def handle_interruption(self): """現在の応答に対するユーザーの割り込みを処理します。""" if not self._is_responding: return # 1. 現在の応答をキャンセルします。 if self._current_response_id: await self.cancel_response() self._is_responding = False self._current_response_id = None self._current_item_id = None async def handle_messages(self) -> None: try: async for message in self.ws: event = json.loads(message) event_type = event.get("type") if event_type == "error": print(" Error: ", event['error']) continue elif event_type == "response.created": self._current_response_id = event.get("response", {}).get("id") self._is_responding = True elif event_type == "response.output_item.added": self._current_item_id = event.get("item", {}).get("id") elif event_type == "response.done": self._is_responding = False self._current_response_id = None self._current_item_id = None elif event_type == "input_audio_buffer.speech_started": print("Speech start detected") if self._is_responding: print("Handling interruption") await self.handle_interruption() elif event_type == "input_audio_buffer.speech_stopped": print("Speech end detected") elif event_type == "response.text.delta": if self.on_text_delta: self.on_text_delta(event["delta"]) elif event_type == "response.audio.delta": if self.on_audio_delta: audio_bytes = base64.b64decode(event["delta"]) self.on_audio_delta(audio_bytes) elif event_type == "conversation.item.input_audio_transcription.delta": preview = event.get("text", "") + event.get("stash", "") print(f"\rUser: {preview}", end='', flush=True) elif event_type == "conversation.item.input_audio_transcription.completed": transcript = event.get("transcript", "") print() if self.on_input_transcript: await asyncio.to_thread(self.on_input_transcript, transcript) self._print_input_transcript = True elif event_type == "response.audio_transcript.delta": if self.on_output_transcript: delta = event.get("delta", "") if not self._print_input_transcript: self._output_transcript_buffer += delta else: if self._output_transcript_buffer: await asyncio.to_thread(self.on_output_transcript, self._output_transcript_buffer) self._output_transcript_buffer = "" await asyncio.to_thread(self.on_output_transcript, delta) elif event_type == "response.audio_transcript.done": print(f"assistant: {event.get('transcript', '')}") self._print_input_transcript = False elif event_type in self.extra_event_handlers: self.extra_event_handlers[event_type](event) except websockets.exceptions.ConnectionClosed: print(" Connection closed") except Exception as e: print(" Error in message handling: ", str(e)) async def close(self) -> None: """WebSocket 接続を閉じます。""" if self.ws: await self.ws.close() -
対話モードの選択
-
VAD モード (音声区間検出、発話の開始と終了を自動的に検出)
Realtime API は、発話の開始と停止を検出し、応答を生成します。
-
手動モード (押して話す、離して送信)
音声の送信を開始および停止するタイミングを制御します。発話後、クライアントはサーバーにメッセージを送信して応答を生成する必要があります。
VAD モード
omni_realtime_client.pyと同じディレクトリに、vad_mode.pyという名前のファイルを作成し、次のコードをコピーします:vad_mode.py
# -- coding: utf-8 -- import os, asyncio, pyaudio, queue, threading from omni_realtime_client import OmniRealtimeClient, TurnDetectionMode # 割り込みを処理するオーディオプレーヤークラス class AudioPlayer: def __init__(self, pyaudio_instance, rate=24000): self.stream = pyaudio_instance.open(format=pyaudio.paInt16, channels=1, rate=rate, output=True) self.queue = queue.Queue() self.stop_evt = threading.Event() self.interrupt_evt = threading.Event() threading.Thread(target=self._run, daemon=True).start() def _run(self): while not self.stop_evt.is_set(): try: data = self.queue.get(timeout=0.5) if data is None: break if not self.interrupt_evt.is_set(): self.stream.write(data) self.queue.task_done() except queue.Empty: continue def add_audio(self, data): self.queue.put(data) def handle_interrupt(self): self.interrupt_evt.set(); self.queue.queue.clear() def stop(self): self.stop_evt.set(); self.queue.put(None); self.stream.stop_stream(); self.stream.close() # マイクから録音して音声を送信 async def record_and_send(client): p = pyaudio.PyAudio() stream = p.open(format=pyaudio.paInt16, channels=1, rate=16000, input=True, frames_per_buffer=3200) print("Recording started. Please speak...") try: while True: audio_data = stream.read(3200) await client.stream_audio(audio_data) await asyncio.sleep(0.02) finally: stream.stop_stream(); stream.close(); p.terminate() async def main(): p = pyaudio.PyAudio() player = AudioPlayer(pyaudio_instance=p) client = OmniRealtimeClient( # これはシンガポールリージョンの base_url です。中国 (北京) リージョンの base_url は wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/realtime です。 base_url="wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/realtime", api_key=os.environ.get("DASHSCOPE_API_KEY"), model="qwen3.5-omni-plus-realtime", voice="Ethan", instructions="あなたはシャオユンという名の、ウィットに富んだユーモラスなアシスタントです。", # qwen3.5-omni-realtime のようなモデルには SEMANTIC_VAD を推奨します。 turn_detection_mode=TurnDetectionMode.SEMANTIC_VAD, on_text_delta=lambda t: print(f"\nassistant: {t}", end="", flush=True), on_audio_delta=player.add_audio, ) await client.connect() print("Connection successful. Starting the real-time conversation...") # 同時実行 await asyncio.gather(client.handle_messages(), record_and_send(client)) if __name__ == "__main__": try: asyncio.run(main()) except KeyboardInterrupt: print("\nProgram exited.")vad_mode.pyを実行して、マイクを通じてリアルタイムの会話を開始します。システムは発話を検出し、音声をサーバーにストリーミングします。手動モード
omni_realtime_client.pyと同じディレクトリに、manual_mode.pyという名前のファイルを作成し、次のコードをコピーします:manual_mode.py
# -- coding: utf-8 -- import os import asyncio import time import threading import queue import pyaudio from omni_realtime_client import OmniRealtimeClient, TurnDetectionMode class AudioPlayer: """リアルタイムオーディオプレーヤークラス。""" def __init__(self, sample_rate=24000, channels=1, sample_width=2): self.sample_rate = sample_rate self.channels = channels self.sample_width = sample_width # 16 ビットの場合は 2 バイト self.audio_queue = queue.Queue() self.is_playing = False self.play_thread = None self.pyaudio_instance = None self.stream = None self._lock = threading.Lock() # 同期アクセスのためのロックを追加します。 self._last_data_time = time.time() # 最後のデータを受信した時刻を記録します。 self._response_done = False # 応答が完了したことを示すフラグを追加します。 self._waiting_for_response = False # クライアントがサーバーの応答を待っているかどうかを示すフラグ。 # オーディオストリームへの最後の書き込み時刻と、より正確な再生終了検出のための最後のオーディオチャンクの持続時間を記録します。 self._last_play_time = time.time() self._last_chunk_duration = 0.0 def start(self): """オーディオプレーヤーを開始します。""" with self._lock: if self.is_playing: return self.is_playing = True try: self.pyaudio_instance = pyaudio.PyAudio() # オーディオ出力ストリームを作成します。 self.stream = self.pyaudio_instance.open( format=pyaudio.paInt16, # 16 ビット channels=self.channels, rate=self.sample_rate, output=True, frames_per_buffer=1024 ) # 再生スレッドを開始します。 self.play_thread = threading.Thread(target=self._play_audio) self.play_thread.daemon = True self.play_thread.start() print("Audio player started") except Exception as e: print(f"Failed to start audio player: {e}") self._cleanup_resources() raise def stop(self): """オーディオプレーヤーを停止します。""" with self._lock: if not self.is_playing: return self.is_playing = False # キューをクリアします。 while not self.audio_queue.empty(): try: self.audio_queue.get_nowait() except queue.Empty: break # 再生スレッドが終了するのを待ちます。デッドロックを避けるためにロックの外で待ちます。 if self.play_thread and self.play_thread.is_alive(): self.play_thread.join(timeout=2.0) # リソースをクリーンアップするために再度ロックを取得します。 with self._lock: self._cleanup_resources() print("Audio player stopped") def _cleanup_resources(self): """オーディオリソースをクリーンアップします。これはロック内で呼び出す必要があります。""" try: # オーディオストリームを閉じます。 if self.stream: if not self.stream.is_stopped(): self.stream.stop_stream() self.stream.close() self.stream = None except Exception as e: print(f"Error closing audio stream: {e}") try: if self.pyaudio_instance: self.pyaudio_instance.terminate() self.pyaudio_instance = None except Exception as e: print(f"Error terminating PyAudio: {e}") def add_audio_data(self, audio_data): """再生キューにオーディオデータを追加します。""" if self.is_playing and audio_data: self.audio_queue.put(audio_data) with self._lock: self._last_data_time = time.time() # 最後のデータを受信した時刻を更新します。 self._waiting_for_response = False # データを受信したので、待機を解除します。 def stop_receiving_data(self): """これ以上新しいオーディオデータを受信しないことをマークします。""" with self._lock: self._response_done = True self._waiting_for_response = False # 応答が終了したので、待機を解除します。 def prepare_for_next_turn(self): """次の会話ターンのためにプレーヤーの状態をリセットします。""" with self._lock: self._response_done = False self._last_data_time = time.time() self._last_play_time = time.time() self._last_chunk_duration = 0.0 self._waiting_for_response = True # 次の応答を待機開始します。 # 前のターンから残っているオーディオデータをクリアします。 while not self.audio_queue.empty(): try: self.audio_queue.get_nowait() except queue.Empty: break def is_finished_playing(self): """すべてのオーディオデータが再生されたかどうかを確認します。""" with self._lock: queue_size = self.audio_queue.qsize() time_since_last_data = time.time() - self._last_data_time time_since_last_play = time.time() - self._last_play_time # ---------------------- スマート終了検出 ---------------------- # 1. 推奨方法:サーバーが完了をマークし、再生キューが空の場合、 # 最新のオーディオチャンクが再生し終わるのを待ちます (チャンク持続時間 + 0.1 秒の許容範囲)。 if self._response_done and queue_size == 0: min_wait = max(self._last_chunk_duration + 0.1, 0.5) # 少なくとも 0.5 秒待ちます。 if time_since_last_play >= min_wait: return True # 2. フォールバック方法:1 秒以上新しいデータが受信されず、再生キューが空の場合。 # このロジックは、サーバーが明示的に `response.done` を送信しない場合のセーフガードとして機能します。 if not self._waiting_for_response and queue_size == 0 and time_since_last_data > 1.0: print("\n(No new audio received for a while, assuming playback is finished)") return True return False def _play_audio(self): """オーディオデータを再生するためのワーカースレッド。""" while True: # 停止すべきかどうかを確認します。 with self._lock: if not self.is_playing: break stream_ref = self.stream # ストリームへの参照を取得します。 try: # キューからオーディオデータを取得します。タイムアウトは 0.1 秒です。 audio_data = self.audio_queue.get(timeout=0.1) # ステータスとストリームの有効性を再度確認します。 with self._lock: if self.is_playing and stream_ref and not stream_ref.is_stopped(): try: # オーディオデータを再生します。 stream_ref.write(audio_data) # 最新の再生情報を更新します。 self._last_play_time = time.time() self._last_chunk_duration = len(audio_data) / ( self.channels * self.sample_width) / self.sample_rate except Exception as e: print(f"Error writing to audio stream: {e}") break # このデータブロックを処理済みとしてマークします。 self.audio_queue.task_done() except queue.Empty: # キューが空の場合は待機を続けます。 continue except Exception as e: print(f"Error playing audio: {e}") break class MicrophoneRecorder: """リアルタイムマイクレコーダー。""" def __init__(self, sample_rate=16000, channels=1, chunk_size=3200): self.sample_rate = sample_rate self.channels = channels self.chunk_size = chunk_size self.pyaudio_instance = None self.stream = None self.frames = [] self._is_recording = False self._record_thread = None def _recording_thread(self): """録音ワーカースレッド。""" # _is_recording が True の間、オーディオストリームからデータを継続的に読み取ります。 while self._is_recording: try: # バッファオーバーフローによるクラッシュを避けるために exception_on_overflow=False を使用します。 data = self.stream.read(self.chunk_size, exception_on_overflow=False) self.frames.append(data) except (IOError, OSError) as e: # ストリームが閉じられると、読み取り操作でエラーが発生する可能性があります。 print(f"Error reading from recording stream, it might be closed: {e}") break def start(self): """録音を開始します。""" if self._is_recording: print("Recording is already in progress.") return self.frames = [] self._is_recording = True try: self.pyaudio_instance = pyaudio.PyAudio() self.stream = self.pyaudio_instance.open( format=pyaudio.paInt16, channels=self.channels, rate=self.sample_rate, input=True, frames_per_buffer=self.chunk_size ) self._record_thread = threading.Thread(target=self._recording_thread) self._record_thread.daemon = True self._record_thread.start() print("Microphone recording started...") except Exception as e: print(f"Failed to start microphone: {e}") self._is_recording = False self._cleanup() raise def stop(self): """録音を停止し、オーディオデータを返します。""" if not self._is_recording: return None self._is_recording = False # 録音スレッドが安全に終了するのを待ちます。 if self._record_thread: self._record_thread.join(timeout=1.0) self._cleanup() print("Microphone recording stopped.") return b''.join(self.frames) def _cleanup(self): """PyAudio リソースを安全にクリーンアップします。""" if self.stream: try: if self.stream.is_active(): self.stream.stop_stream() self.stream.close() except Exception as e: print(f"Error closing audio stream: {e}") if self.pyaudio_instance: try: self.pyaudio_instance.terminate() except Exception as e: print(f"Error terminating PyAudio instance: {e}") self.stream = None self.pyaudio_instance = None async def interactive_test(): """ オーディオと画像のサポートを備えたマルチターン会話のインタラクティブテスト。 """ # ------------------- 1. 初期化と接続 (1 回限り) ------------------- # シンガポールと中国 (北京) リージョンの API キーは異なります。API キーを取得するには、https://www.alibabacloud.com/help/model-studio/get-api-key をご参照ください。 api_key = os.environ.get("DASHSCOPE_API_KEY") if not api_key: print("Please set the DASHSCOPE_API_KEY environment variable.") return print("--- Real-time Multimodal Audio/Video Chat Client ---") print("Initializing audio player and client...") audio_player = AudioPlayer() audio_player.start() def on_audio_received(audio_data): audio_player.add_audio_data(audio_data) def on_response_done(event): print("\n(Received response end marker)") audio_player.stop_receiving_data() realtime_client = OmniRealtimeClient( # これはシンガポールリージョンの base_url です。中国 (北京) リージョンのモデルを使用する場合は、base_url を wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/realtime に置き換えてください。 base_url="wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/realtime", api_key=api_key, model="qwen3.5-omni-plus-realtime", voice="Ethan", instructions="あなたはシャオユンという名のパーソナルアシスタントです。ユーザーの質問に正確かつフレンドリーに回答し、常に親切な態度で対応してください。", # モデルの役割を設定します。 on_text_delta=lambda text: print(f"assistant: {text}", end="", flush=True), on_audio_delta=on_audio_received, turn_detection_mode=TurnDetectionMode.MANUAL, extra_event_handlers={"response.done": on_response_done} ) message_handler_task = None try: await realtime_client.connect() print("Connected to the server. Enter 'q' or 'quit' to exit at any time.") message_handler_task = asyncio.create_task(realtime_client.handle_messages()) await asyncio.sleep(0.5) turn_counter = 1 # ------------------- 2. マルチターン会話ループ ------------------- while True: print(f"\n--- Turn {turn_counter} ---") audio_player.prepare_for_next_turn() recorded_audio = None image_paths = [] # --- ユーザー入力を取得:マイクから録音 --- loop = asyncio.get_event_loop() recorder = MicrophoneRecorder(sample_rate=16000) # 音声認識には 16k のサンプルレートを推奨します。 print("Ready to record. Press Enter to start recording (or enter 'q' to exit)...") user_input = await loop.run_in_executor(None, input) if user_input.strip().lower() in ['q', 'quit']: print("User requested to exit...") return try: recorder.start() except Exception: print("Could not start recording. Please check your microphone permissions and device. Skipping this turn.") continue print("Recording... Press Enter again to stop recording.") await loop.run_in_executor(None, input) recorded_audio = recorder.stop() if not recorded_audio or len(recorded_audio) == 0: print("No valid audio was recorded. Please start this turn again.") continue # --- 画像入力を取得 (オプション) --- # 画像入力機能はデフォルトで無効になっています。有効にするには、以下のコードのコメントを解除してください。 # print("\nEnter the absolute path of an [image file] on each line (optional). When finished, enter 's' or press Enter to send the request.") # while True: # path = input("Image path: ").strip() # if path.lower() == 's' or path == '': # break # if path.lower() in ['q', 'quit']: # print("User requested to exit...") # return # # if not os.path.isabs(path): # print("Error: Please enter an absolute path.") # continue # if not os.path.exists(path): # print(f"Error: File not found -> {path}") # continue # image_paths.append(path) # print(f"Image added: {os.path.basename(path)}") # --- 3. データを送信して応答を取得 --- print("\n--- Input Confirmation ---") print(f"Audio to process: 1 (from microphone), Images: {len(image_paths)}") print("------------------") # 3.1 録音した音声を送信します。 try: print(f"Sending microphone recording ({len(recorded_audio)} bytes)") await realtime_client.stream_audio(recorded_audio) await asyncio.sleep(0.1) except Exception as e: print(f"Failed to send microphone recording: {e}") continue # 3.2 すべての画像ファイルを送信します。 # 画像入力機能はデフォルトで無効になっています。有効にするには、以下のコードのコメントを解除してください。 # for i, path in enumerate(image_paths): # try: # with open(path, "rb") as f: # data = f.read() # print(f"Sending image {i+1}: {os.path.basename(path)} ({len(data)} bytes)") # await realtime_client.append_image(data) # await asyncio.sleep(0.1) # except Exception as e: # print(f"Failed to send image {os.path.basename(path)}: {e}") # 3.3 送信して応答を待ちます。 print("Submitting all inputs, requesting server response...") await realtime_client.commit_audio_buffer() await realtime_client.create_response() print("Waiting for and playing server response audio...") start_time = time.time() max_wait_time = 60 while not audio_player.is_finished_playing(): if time.time() - start_time > max_wait_time: print(f"\nWait timed out ({max_wait_time} seconds). Moving to the next turn.") break await asyncio.sleep(0.2) print("\nAudio playback for this turn is complete!") turn_counter += 1 except (asyncio.CancelledError, KeyboardInterrupt): print("\nProgram was interrupted.") except Exception as e: print(f"An unhandled error occurred: {e}") finally: # ------------------- 4. リソースのクリーンアップ ------------------- print("\nClosing connection and cleaning up resources...") if message_handler_task and not message_handler_task.done(): message_handler_task.cancel() if 'realtime_client' in locals() and realtime_client.ws and not realtime_client.ws.close: await realtime_client.close() print("Connection closed.") audio_player.stop() print("Program exited.") if __name__ == "__main__": try: asyncio.run(interactive_test()) except KeyboardInterrupt: print("\nProgram was forcibly exited by the user.")manual_mode.pyを実行します。Enter キーを押して録音を開始し、もう一度 Enter キーを押して停止して送信します。 -
WebRTC
Python
-
実行環境
Python 3.10 以降が必要です。次の依存関係をインストールします:
pip install aiortc aiohttp sounddevice numpy certifi av
-
デモの実行
webrtc_demo.pyという名前の Python ファイルを作成し、次のコードを貼り付けます:webrtc_demo.py
# 依存関係:pip install aiortc aiohttp sounddevice numpy certifi av import asyncio import json import os import queue import ssl import threading import aiohttp import certifi import numpy as np import sounddevice as sd from aiortc import RTCPeerConnection, RTCConfiguration, RTCSessionDescription from aiortc.contrib.media import MediaPlayer from av import AudioFrame # API キーに置き換えるか、DASHSCOPE_API_KEY 環境変数を設定してください API_KEY = os.getenv("DASHSCOPE_API_KEY", "your-api-key") MODEL = "qwen3.5-omni-plus-realtime" # シンガポールリージョン。{WorkspaceId} をお使いの Bailian ワークスペース ID に置き換えてください。URL はリージョンによって異なります。 SIGNALING_URL = "https://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api/v1/webrtc/realtime?model=" + MODEL # --------------- オーディオフレームの解析 --------------- def _nb_channels(frame: AudioFrame) -> int: """オーディオフレームのチャネル数を取得します。さまざまな PyAV バージョンと互換性があります""" if hasattr(frame.layout, "nb_channels"): return int(frame.layout.nb_channels) ch = getattr(frame.layout, "channels", 1) if isinstance(ch, (tuple, list)): return len(ch) return int(ch) def audioframe_to_s16_samples(frame: AudioFrame) -> np.ndarray: """ サーバーのオーディオフレームはステレオインターリーブです。直接 reshape するとチャネルのずれが生じます。 実際のチャネル数に基づいて (サンプル数, チャネル数) に再配置します。 異なる aiortc デコーダーバージョンは、同じオーディオに対して異なる配列形状を返すため、 ここで統一された処理が必要です。 """ arr = np.asarray(frame.to_ndarray()) ch = _nb_channels(frame) samples = int(frame.samples) if arr.ndim == 2 and arr.shape[0] == ch and arr.shape[1] == samples: return arr.T.copy() if arr.ndim == 2 and arr.shape[0] == 1 and arr.shape[1] == samples * ch: return arr.reshape(-1).reshape(samples, ch).copy() if arr.ndim == 1 and arr.shape[0] == samples * ch: return arr.reshape(samples, ch).copy() flat = arr.reshape(-1) if ch > 0 and flat.size % ch == 0: return flat.reshape(flat.size // ch, ch).copy() raise ValueError(f"unexpected shape={arr.shape}, ch={ch}, samples={samples}") # --------------- 低遅延オーディオプレーヤー --------------- class RemoteAudioPlayer: """ 遅延を最小限に抑えるために 5ms のオーディオブロックを再生する低遅延オーディオプレーヤー。 音声割り込みをサポート:ユーザーが話し始めるとバッファをクリアし、 古いモデル応答の再生を停止します。 サーバーのステレオオーディオをモノラル (左右チャネルの平均) にマージして再生します。 """ def __init__(self, samplerate=48000, out_channels=1, blocksize=240, max_seconds=0.2): self.samplerate = samplerate self.out_channels = out_channels self.blocksize = blocksize self._q = queue.Queue(maxsize=max(5, int(max_seconds * samplerate / blocksize) + 5)) self._lock = threading.Lock() self._rb_size = max(1, int(max_seconds * samplerate)) self._rb = np.zeros((self._rb_size, out_channels), dtype=np.int16) self._rb_w = 0 self._rb_r = 0 self._rb_len = 0 self._stream = None self._closed = False def start(self): if self._stream: return def callback(outdata, frames, _time, status): if self._closed: outdata[:] = np.zeros((frames, self.out_channels), dtype=np.int16) return while True: try: chunk = self._q.get_nowait() except queue.Empty: break with self._lock: self._write_rb(chunk) with self._lock: out = self._read_rb(frames) outdata[:] = out self._stream = sd.OutputStream( samplerate=self.samplerate, channels=self.out_channels, dtype="int16", blocksize=self.blocksize, callback=callback, ) self._stream.start() def clear(self): """音声割り込みのために再生バッファをクリアします""" try: while True: self._q.get_nowait() except queue.Empty: pass with self._lock: self._rb_w = 0 self._rb_r = 0 self._rb_len = 0 self._rb[:] = 0 def _write_rb(self, chunk: np.ndarray): n = int(chunk.shape[0]) if n <= 0: return overflow = max(0, self._rb_len + n - self._rb_size) if overflow > 0: self._rb_r = (self._rb_r + overflow) % self._rb_size self._rb_len -= overflow end = self._rb_size - self._rb_w if n <= end: self._rb[self._rb_w:self._rb_w + n] = chunk else: self._rb[self._rb_w:] = chunk[:end] self._rb[:n - end] = chunk[end:] self._rb_w = (self._rb_w + n) % self._rb_size self._rb_len += n def _read_rb(self, frames: int) -> np.ndarray: if self._rb_len <= 0: return np.zeros((frames, self.out_channels), dtype=np.int16) n = min(frames, self._rb_len) out = np.zeros((frames, self.out_channels), dtype=np.int16) end = self._rb_size - self._rb_r if n <= end: out[:n] = self._rb[self._rb_r:self._rb_r + n] else: out[:end] = self._rb[self._rb_r:] out[end:n] = self._rb[:n - end] self._rb_r = (self._rb_r + n) % self._rb_size self._rb_len -= n return out async def push_frame(self, frame: AudioFrame): """オーディオフレームを受信し、自動的にチャネルをマージしてキューに入れます""" if self._closed: return pcm = audioframe_to_s16_samples(frame) in_ch = pcm.shape[1] if self.out_channels == 1: if in_ch == 1: out = pcm else: out = np.mean(pcm.astype(np.int32), axis=1).astype(np.int16).reshape(-1, 1) else: if in_ch == self.out_channels: out = pcm elif in_ch == 1 and self.out_channels == 2: out = np.repeat(pcm, 2, axis=1) else: out = pcm[:, :self.out_channels] try: self._q.put_nowait(out) except queue.Full: try: self._q.get_nowait() except queue.Empty: pass try: self._q.put_nowait(out) except queue.Full: pass async def close(self): self._closed = True if self._stream: self._stream.stop() self._stream.close() self._stream = None # --------------- main --------------- async def main(): pc = RTCPeerConnection(RTCConfiguration(iceServers=[])) # オーディオプレーヤーを初期化 (モノラル出力、低遅延のための 5ms ブロックサイズ) speaker = RemoteAudioPlayer(samplerate=48000, out_channels=1, blocksize=240, max_seconds=0.2) speaker.start() # マイクを初期化 (macOS avfoundation; Linux では pulse または alsa を使用) mic = MediaPlayer("none:0", format="avfoundation", options={"sample_rate": "48000", "channels": "1"}) if not mic.audio: raise RuntimeError("No microphone detected. Check the avfoundation audio device index.") pc.addTrack(mic.audio) # クライアントが DataChannel を作成 (名前はカスタマイズ可能); サーバーは "txt" という名前のチャネルを通じてイベントをプッシュします pc.createDataChannel("oai-events") remote_dc = None got_first_txt_msg = False def make_session_update() -> dict: """session.update 設定を構築:音声、オーディオフォーマット、VAD 戦略、推論パラメーター""" return { "type": "session.update", "session": { "modalities": ["text", "audio"], "voice": "Tina", "input_audio_format": "pcm", "output_audio_format": "pcm", "instructions": "You are a friendly AI assistant.", "turn_detection": {"type": "server_vad", "threshold": 0.5, "silence_duration_ms": 800}, "max_tokens": 16384, "temperature": 0.9, }, } # サーバーがプッシュした DataChannel イベントを処理 @pc.on("datachannel") def on_datachannel(ch): nonlocal remote_dc, got_first_txt_msg print(f"[DC] Received server DataChannel: {ch.label}") if ch.label == "txt": remote_dc = ch @ch.on("message") def on_msg(msg): nonlocal got_first_txt_msg try: evt = json.loads(msg) except Exception: return print(f"[{ch.label}] {evt.get('type')}") # ユーザーが話し始めたときに再生バッファをクリア (音声割り込み) if isinstance(evt, dict) and evt.get("type") == "input_audio_buffer.speech_started": speaker.clear() print("[Playback] User speech detected, clearing buffer (interruption)") # txt チャネルで最初のメッセージを受信した後に session.update を送信 if ch.label == "txt" and not got_first_txt_msg: got_first_txt_msg = True if remote_dc and remote_dc.readyState == "open": remote_dc.send(json.dumps(make_session_update(), ensure_ascii=False)) print("[DC] session.update sent") # サーバーのオーディオを受信し、低遅延で再生 @pc.on("track") async def on_track(track): if track.kind == "audio": async def _play(): try: while True: frame = await track.recv() await speaker.push_frame(frame) except Exception: pass asyncio.create_task(_play()) @pc.on("iceconnectionstatechange") def on_ice(): print(f"[ICE] {pc.iceConnectionState}") @pc.on("connectionstatechange") async def on_conn(): print(f"[Connection] {pc.connectionState}") if pc.connectionState in ("failed", "closed", "disconnected"): await pc.close() # SDP 交換:Offer を作成し、シグナリングサーバーに POST し、Answer を取得 offer = await pc.createOffer() await pc.setLocalDescription(offer) async with aiohttp.ClientSession() as session: async with session.post( SIGNALING_URL, ssl=ssl.create_default_context(cafile=certifi.where()), data=offer.sdp.encode("utf-8"), headers={ "Content-Type": "application/sdp", "Authorization": f"Bearer {API_KEY}", }, timeout=aiohttp.ClientTimeout(total=10), ) as resp: if not resp.ok: raise Exception(f"SDP exchange failed: {resp.status} {await resp.text()}") answer_sdp = await resp.text() await pc.setRemoteDescription(RTCSessionDescription(sdp=answer_sdp, type="answer")) print("SDP exchange complete, waiting for connection...") try: await asyncio.Event().wait() except (KeyboardInterrupt, asyncio.CancelledError): pass finally: print(f"\nExiting. Final state: connection={pc.connectionState}, ICE={pc.iceConnectionState}") await speaker.close() try: if mic and mic.audio: mic.audio.stop() except Exception: pass await pc.close() asyncio.run(main())webrtc_demo.pyを実行して、マイクを通じて Qwen-Omni-Realtime モデルとのリアルタイム会話を開始します。システムは発話の開始を検出し、自動的に音声をサーバーに送信します。
JavaScript
-
前提条件
- WebRTC をサポートする最新のブラウザ (Chrome、Edge、Firefox、Safari など) を使用します。
- ブラウザにはマイクの権限が必要です。
- ブラウザのクロスオリジンセキュリティポリシーにより、ブラウザは接続リクエストを直接サーバーに送信できません。ターミナルで curl コマンドを実行して接続設定を完了する必要があります。
-
デモの実行
webrtc_demo.htmlという名前の HTML ファイルを作成し、次のコードを貼り付けます:webrtc_demo.html
<!DOCTYPE html> <html lang="ja"> <head> <meta charset="UTF-8" /> <title>WebRTC リアルタイム音声チャット</title> <style> * { box-sizing: border-box; margin: 0; padding: 0; } body { font-family: -apple-system, BlinkMacSystemFont, "Segoe UI", Roboto, "Helvetica Neue", Arial, sans-serif; background: #f5f7fa; color: #1d2129; padding: 24px; line-height: 1.6; } .container { max-width: 800px; margin: 0 auto; } h1 { font-size: 22px; font-weight: 600; margin-bottom: 20px; color: #1d2129; } /* スティッキートップバー */ .sticky-top { position: sticky; top: 0; z-index: 100; background: #f5f7fa; margin: 0 -24px 16px; padding: 12px 24px; border-bottom: 1px solid transparent; transition: border-color .2s; } .sticky-top.scrolled { border-bottom-color: #e5e6eb; } /* ツールバー */ .toolbar { display: flex; align-items: center; gap: 10px; flex-wrap: wrap; margin-bottom: 12px; } .toolbar label { display: flex; align-items: center; gap: 6px; font-size: 13px; color: #4e5969; cursor: pointer; } /* ボタン */ button { padding: 8px 18px; font-size: 13px; font-weight: 500; border: 1px solid #c9cdd4; border-radius: 6px; background: #fff; color: #1d2129; cursor: pointer; transition: all .15s; } button:hover:not(:disabled) { border-color: #165dff; color: #165dff; } button:disabled { opacity: .4; cursor: not-allowed; } .btn-primary { background: #165dff; border-color: #165dff; color: #fff; } .btn-primary:hover:not(:disabled) { background: #4080ff; border-color: #4080ff; color: #fff; } .btn-danger { border-color: #f53f3f; color: #f53f3f; } .btn-danger:hover:not(:disabled) { background: #f53f3f; color: #fff; } /* ステータスインジケーター */ .status-bar { display: flex; align-items: center; gap: 8px; padding: 10px 14px; border-radius: 8px; background: #fff; border: 1px solid #e5e6eb; font-size: 13px; } .status-dot { width: 8px; height: 8px; border-radius: 50%; background: #c9cdd4; flex-shrink: 0; } .status-dot.connected { background: #00b42a; } .status-dot.connecting { background: #ff7d00; animation: pulse 1s infinite; } .status-dot.error { background: #f53f3f; } @keyframes pulse { 0%,100% { opacity: 1; } 50% { opacity: .4; } } /* SDP カード */ .card { background: #fff; border: 1px solid #e5e6eb; border-radius: 10px; padding: 16px; margin-bottom: 16px; } .card-title { font-size: 13px; font-weight: 600; color: #4e5969; margin-bottom: 8px; } .step-num { display: inline-flex; align-items: center; justify-content: center; width: 20px; height: 20px; border-radius: 50%; background: #165dff; color: #fff; font-size: 11px; font-weight: 600; margin-right: 6px; } .card-hint { font-size: 12px; color: #86909c; margin-top: 6px; } textarea { width: 100%; font-family: "SF Mono", "Fira Code", "Fira Mono", Menlo, Consolas, monospace; font-size: 12px; padding: 10px; border: 1px solid #e5e6eb; border-radius: 6px; resize: vertical; background: #f7f8fa; color: #1d2129; transition: border-color .15s; } textarea:focus { outline: none; border-color: #165dff; background: #fff; } /* ビデオ */ .video-section { margin-bottom: 16px; } .video-label { font-size: 13px; color: #86909c; margin-bottom: 6px; } video { width: 320px; max-width: 100%; background: #000; border-radius: 8px; display: block; } /* イベントパネル */ .events-title { font-size: 14px; font-weight: 600; color: #1d2129; margin-bottom: 10px; } .events-container { display: flex; flex-direction: column; gap: 6px; } .event-item { background: #fff; border: 1px solid #e5e6eb; border-radius: 8px; overflow: hidden; } .event-header { display: flex; align-items: center; gap: 8px; padding: 8px 12px; cursor: pointer; user-select: none; font-size: 12px; } .event-header:hover { background: #f7f8fa; } .event-arrow { font-size: 14px; font-weight: 700; width: 18px; text-align: center; } .event-arrow.server { color: #00b42a; } .event-arrow.client { color: #165dff; } .event-label { color: #4e5969; } .event-time { color: #c9cdd4; margin-left: auto; font-size: 11px; } .event-body { display: none; padding: 10px 12px; background: #f7f8fa; border-top: 1px solid #e5e6eb; } .event-body pre { margin: 0; font-size: 11px; font-family: "SF Mono", Menlo, Consolas, monospace; color: #4e5969; white-space: pre-wrap; word-break: break-all; } .events-empty { font-size: 13px; color: #c9cdd4; padding: 16px 0; text-align: center; } </style> </head> <body> <div class="container"> <h1>WebRTC リアルタイム音声チャット</h1> <div class="sticky-top"> <div class="toolbar"> <button id="startBtn" class="btn-primary">セッション開始</button> <button id="setAnswerBtn" disabled>Answer 設定</button> <button id="endBtn" class="btn-danger" disabled>セッション終了</button> <button id="downloadBtn" disabled>リモート音声ダウンロード</button> <label> <input id="sendVideoCheckbox" type="checkbox" /> ビデオ有効化 </label> </div> <div class="status-bar"> </div> </div> <div class="card"> <div class="card-title">Offer SDP</div> <div style="margin-bottom: 8px;"> <button id="copyOfferBtn" disabled>Offer SDP をコピー</button> </div> <textarea id="offerBox" rows="6" readonly placeholder="「セッション開始」をクリックすると自動生成されます"></textarea> <div class="card-hint">ICE 収集完了後に自動生成されます。コピーして curl でサーバーに送信し、Answer を取得します。</div> </div> <div class="card"> <div class="card-title">curl コマンド</div> <div style="margin-bottom: 8px;"> <button id="copyCurlBtn" disabled>curl コマンドをコピー</button> </div> <div class="card-hint" style="margin-bottom: 4px;">シンガポールリージョン。{WorkspaceId} をお使いの Bailian ワークスペース ID に置き換えてください。URL はリージョンによって異なります。</div> <textarea id="curlBox" rows="6" placeholder="Offer SDP が生成されると自動入力されます"></textarea> <div class="card-hint">{WorkspaceId} をワークスペース ID に置き換え、このコマンドをターミナルにコピーします。返された Answer SDP を以下に貼り付けます。</div> </div> <div class="card"> <div class="card-title">Answer SDP</div> <textarea id="answerBox" rows="6" placeholder="curl が返した Answer SDP をここに貼り付けます"></textarea> <div class="card-hint">貼り付け後、上の「Answer 設定」をクリックして接続を確立します。</div> </div> <div class="video-section" id="videoSection" style="display:none;"> <div class="video-label">ローカルビデオプレビュー</div> <video id="localVideo" autoplay playsinline muted></video> </div> <div class="events-title">イベント (DataChannel)</div> <div id="events" class="events-container"></div> </div> <script> const eventsDiv = document.getElementById('events'); const startBtn = document.getElementById('startBtn'); const setAnswerBtn = document.getElementById('setAnswerBtn'); const endBtn = document.getElementById('endBtn'); const downloadBtn = document.getElementById('downloadBtn'); const copyOfferBtn = document.getElementById('copyOfferBtn'); const statusDot = document.getElementById('statusDot'); const statusText = document.getElementById('statusText'); const copyCurlBtn = document.getElementById('copyCurlBtn'); const curlBox = document.getElementById('curlBox'); const sendVideoCheckbox = document.getElementById('sendVideoCheckbox'); const localVideo = document.getElementById('localVideo'); const offerBox = document.getElementById('offerBox'); const answerBox = document.getElementById('answerBox'); let pc = null; let hiddenRemoteAudioEl = null; let mediaRecorder = null; let recordedChunks = []; let audioBlob = null; let localStream = null; let sendCanvas = null; let sendCanvasCtx = null; let sendCanvasStream = null; let sendRafId = 0; let gatedAudioTracks = []; let gatedVideoTracks = []; let audioSender = null; let videoSender = null; let audioTrack = null; let videoTrack = null; function setStatus(text, state) { statusText.textContent = text; statusDot.className = 'status-dot' + (state ? ' ' + state : ''); } function gateMedia(on) { for (const t of gatedAudioTracks) t.enabled = !!on; for (const t of gatedVideoTracks) t.enabled = !!on; } function sendUpdate(channel) { const update = { event_id: `event_${Date.now()}`, type: "session.update", session: { input_audio_format: "pcm", input_audio_transcription: { model: "qwen3-asr-flash-realtime" }, instructions: "You are a helpful assistant.", modalities: ["text", "audio"], output_audio_format: "pcm", smooth_output: false, turn_detection: { prefix_padding_ms: 500, silence_duration_ms: 800, threshold: 0.5, type: "server_vad", }, }, }; if (channel && channel.readyState === "open") channel.send(JSON.stringify(update)); } // ===== イベントパネル ===== const events = []; function nowTs() { return new Date().toLocaleTimeString(); } function renderEvents() { eventsDiv.innerHTML = ""; if (events.length === 0) { const empty = document.createElement("div"); empty.className = "events-empty"; empty.textContent = "イベントを待機中..."; eventsDiv.appendChild(empty); return; } for (const item of events) { const { event, timestamp } = item; const isClient = event?.type?.includes("update") || event?.type?.includes("create"); const wrap = document.createElement("div"); wrap.className = "event-item"; const header = document.createElement("div"); header.className = "event-header"; const arrow = document.createElement("span"); arrow.className = "event-arrow " + (isClient ? "client" : "server"); arrow.textContent = isClient ? "↓" : "↑"; const label = document.createElement("span"); label.className = "event-label"; const who = isClient ? "client" : "server"; const type = event?.type ?? "message"; label.textContent = `${who}: ${type}`; const time = document.createElement("span"); time.className = "event-time"; time.textContent = timestamp; const body = document.createElement("div"); body.className = "event-body"; const pre = document.createElement("pre"); pre.textContent = JSON.stringify(event, null, 2); body.appendChild(pre); header.onclick = () => { body.style.display = body.style.display === "block" ? "none" : "block"; }; header.appendChild(arrow); header.appendChild(label); header.appendChild(time); wrap.appendChild(header); wrap.appendChild(body); eventsDiv.appendChild(wrap); } } function clearUIEvents() { events.length = 0; renderEvents(); } function pushEventFromDataChannel(eventObj) { const ts = eventObj.timestamp || nowTs(); if (!eventObj.timestamp) eventObj.timestamp = ts; events.unshift({ event: eventObj, timestamp: ts }); renderEvents(); } function normalizeSdpForSetRemote(sdp) { sdp = String(sdp).trim().replace(/\r?\n/g, "\r\n"); if (!sdp.endsWith("\r\n")) sdp += "\r\n"; return sdp; } // ===== WebRTC ===== startBtn.onclick = () => startSession().catch(err => console.log("startSession error:", err)); endBtn.onclick = () => endSession(); setAnswerBtn.onclick = () => setRemoteAnswerFromUI().catch(err => console.log("setRemoteAnswer error:", err)); copyOfferBtn.onclick = async () => { const txt = offerBox.value; if (!txt) return; await navigator.clipboard.writeText(txt); alert("Offer SDP をコピーしました"); }; copyCurlBtn.onclick = async () => { const txt = curlBox.value; if (!txt) return; await navigator.clipboard.writeText(txt); alert("curl コマンドをコピーしました。ターミナルで実行してください。"); }; downloadBtn.onclick = () => { if (audioBlob) downloadBlob(audioBlob, 'remote-audio.webm'); else alert('利用可能な音声録音はありません'); }; async function startSession() { if (pc) return; pc = new RTCPeerConnection({ iceServers: [] }); clearUIEvents(); setStatus('マイクへのアクセスを要求中...', 'connecting'); offerBox.value = ""; answerBox.value = ""; curlBox.value = ""; setAnswerBtn.disabled = true; copyOfferBtn.disabled = true; copyCurlBtn.disabled = true; endBtn.disabled = false; downloadBtn.disabled = true; pc.onconnectionstatechange = () => { if (!pc) return; if (pc.connectionState === 'connected') { setStatus('接続済み。話し始めてください。', 'connected'); } else if (["failed", "closed", "disconnected"].includes(pc.connectionState)) { console.log("onconnectionstatechange:", pc.connectionState); endSession(true); } }; pc.ontrack = async (e) => { const stream = e.streams[0]; ensureHiddenAudioEl(); hiddenRemoteAudioEl.srcObject = stream; try { await hiddenRemoteAudioEl.play(); } catch {} startRecordingRemoteStream(stream); }; const wantVideo = !!sendVideoCheckbox.checked; const localPreviewFps = 30; const sendFps = 2; const constraints = wantVideo ? { audio: true, video: { facingMode: { ideal: "user" }, frameRate: { ideal: localPreviewFps, max: localPreviewFps }, width: { ideal: 640 }, height: { ideal: 480 }, } } : { audio: true }; localStream = await navigator.mediaDevices.getUserMedia(constraints); const videoSection = document.getElementById('videoSection'); if (wantVideo) { localVideo.srcObject = localStream; localVideo.style.display = "block"; videoSection.style.display = ""; try { await localVideo.play(); } catch {} } else { localVideo.srcObject = null; localVideo.style.display = "none"; videoSection.style.display = "none"; } gatedAudioTracks = []; gatedVideoTracks = []; localStream.getAudioTracks().forEach(t => { pc.addTrack(t, localStream); gatedAudioTracks.push(t); }); if (wantVideo) { if (sendRafId) cancelAnimationFrame(sendRafId); sendRafId = 0; if (sendCanvasStream) sendCanvasStream.getTracks().forEach(t => t.stop()); sendCanvasStream = null; sendCanvasCtx = null; sendCanvas = null; const settings = localStream.getVideoTracks()[0].getSettings(); sendCanvas = document.createElement("canvas"); sendCanvas.width = settings.width || 640; sendCanvas.height = settings.height || 480; sendCanvasCtx = sendCanvas.getContext("2d", { alpha: false }); sendCanvasStream = sendCanvas.captureStream(sendFps); const lowFpsTrack = sendCanvasStream.getVideoTracks()[0]; pc.addTrack(lowFpsTrack, sendCanvasStream); gatedVideoTracks.push(lowFpsTrack); const pump = () => { if (!sendCanvasCtx || !sendCanvas) return; try { sendCanvasCtx.drawImage(localVideo, 0, 0, sendCanvas.width, sendCanvas.height); } catch {} sendRafId = requestAnimationFrame(pump); }; sendRafId = requestAnimationFrame(pump); } gateMedia(false); audioSender = pc.getSenders().find(s => s.track?.kind === 'audio'); videoSender = pc.getSenders().find(s => s.track?.kind === 'video'); audioTrack = audioSender?.track; videoTrack = videoSender?.track; await audioSender?.replaceTrack(null); await videoSender?.replaceTrack(videoTrack ? null : undefined); const dc = pc.createDataChannel('oai-events'); dc.onopen = () => console.log("DC open"); dc.onmessage = (e) => { handleDcMessage(e.data, dc); }; pc.ondatachannel = (event) => { const ch = event.channel; ch.onmessage = (e) => { handleDcMessage(e.data, ch); }; }; function handleDcMessage(data, channel) { let obj; try { obj = JSON.parse(data); } catch (err) { pushEventFromDataChannel({ type: "raw", data: String(data), parseError: String(err) }); return; } pushEventFromDataChannel(obj); if (obj?.type === "session.created") { console.log("Session created, opening media gate."); gateMedia(true); if(audioSender) audioSender.replaceTrack(audioTrack); if(videoSender && videoTrack) videoSender.replaceTrack(videoTrack); sendUpdate(channel); } } pc.onicegatheringstatechange = () => { if (!pc) return; if (pc.iceGatheringState === "complete" && pc.localDescription?.sdp) { const sdp = pc.localDescription.sdp; offerBox.value = sdp; copyOfferBtn.disabled = false; setAnswerBtn.disabled = false; const escapedSdp = sdp.replace(/'/g, "'\\''"); // シンガポールリージョン。{WorkspaceId} をお使いの Bailian ワークスペース ID に置き換えてください。URL はリージョンによって異なります。 curlBox.value = `curl -X POST 'https://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api/v1/webrtc/realtime?model=qwen3.5-omni-plus-realtime' \\\n -H 'Content-Type: application/sdp' \\\n -H 'Authorization: Bearer $DASHSCOPE_API_KEY' \\\n --data-binary '${escapedSdp}'`; copyCurlBtn.disabled = false; setStatus('Offer SDP が生成されました。curl コマンドをターミナルにコピーして Answer SDP を取得してください。', 'connecting'); console.log("ICE Gathering Complete. Ready to set remote description."); } }; const offer = await pc.createOffer(); await pc.setLocalDescription(offer); } async function setRemoteAnswerFromUI() { if (!pc) return alert('最初に「セッション開始」をクリックして Offer を生成してください。'); const txt = answerBox.value.trim(); if (!txt) return alert("Answer SDP を貼り付けてください"); const answerSdp = normalizeSdpForSetRemote(txt); try { await pc.setRemoteDescription({ type: 'answer', sdp: answerSdp }); setStatus('接続を確立中...', 'connecting'); } catch (e) { alert("Answer の設定に失敗しました: " + e.message); console.error(e); } } function endSession(silent = false) { if (sendRafId) cancelAnimationFrame(sendRafId); sendRafId = 0; if (sendCanvasStream) { sendCanvasStream.getTracks().forEach(t => t.stop()); } sendCanvasStream = null; sendCanvasCtx = null; sendCanvas = null; try { if (mediaRecorder && mediaRecorder.state !== "inactive") mediaRecorder.stop(); } catch {} mediaRecorder = null; if (localStream) { localStream.getTracks().forEach(t => t.stop()); localStream = null; } localVideo.srcObject = null; localVideo.style.display = "none"; document.getElementById('videoSection').style.display = "none"; if (pc) { try { pc.close(); } catch {} pc = null; } gatedAudioTracks = []; gatedVideoTracks = []; if (hiddenRemoteAudioEl) { try { hiddenRemoteAudioEl.pause(); } catch {} hiddenRemoteAudioEl.srcObject = null; hiddenRemoteAudioEl.remove(); hiddenRemoteAudioEl = null; } endBtn.disabled = true; setAnswerBtn.disabled = true; copyOfferBtn.disabled = true; copyCurlBtn.disabled = true; downloadBtn.disabled = !audioBlob; setStatus('切断済み', ''); if (!silent) console.log("session ended"); } function ensureHiddenAudioEl() { if (hiddenRemoteAudioEl) return; hiddenRemoteAudioEl = document.createElement("audio"); hiddenRemoteAudioEl.autoplay = true; hiddenRemoteAudioEl.playsInline = true; hiddenRemoteAudioEl.muted = false; hiddenRemoteAudioEl.style.display = "none"; document.body.appendChild(hiddenRemoteAudioEl); } function startRecordingRemoteStream(remoteStream) { const audioTracks = remoteStream.getAudioTracks(); if (!audioTracks.length) return; const audioStream = new MediaStream(audioTracks); recordedChunks = []; audioBlob = null; downloadBtn.disabled = true; try { mediaRecorder = new MediaRecorder(audioStream, { mimeType: 'audio/webm' }); } catch (err) { console.log("MediaRecorder の作成に失敗しました:", err); return; } mediaRecorder.ondataavailable = (e) => { if (e.data && e.data.size > 0) recordedChunks.push(e.data); }; mediaRecorder.onstop = () => { audioBlob = new Blob(recordedChunks, { type: 'audio/webm' }); downloadBtn.disabled = !audioBlob || audioBlob.size === 0; }; mediaRecorder.start(); } function downloadBlob(blob, filename) { const url = URL.createObjectURL(blob); const a = document.createElement('a'); a.style.display = 'none'; a.href = url; a.download = filename; document.body.appendChild(a); a.click(); URL.revokeObjectURL(url); a.remove(); } renderEvents(); const stickyTop = document.querySelector('.sticky-top'); window.addEventListener('scroll', () => { stickyTop.classList.toggle('scrolled', window.scrollY > 10); }, { passive: true }); </script> </body> </html>このファイルをブラウザで開き、次の手順に従います:
- [セッション開始] をクリックします。ページは自動的に Offer SDP と対応する curl コマンドを生成します。
- [curl コマンドをコピー] をクリックし、ターミナルで実行します。出力は Answer SDP です。
- Answer SDP を Answer SDP テキストボックスに貼り付け、[Answer 設定] をクリックして接続を確立し、音声チャットを開始します。
対話フロー
VAD モード
session.update の session.turn_detection.type を "server_vad" または "semantic_vad" に設定して、VAD モードを有効化します。音声通話シナリオに適しています。WebSocket と WebRTC はどちらも同じサーバーイベントで VAD モードをサポートしますが、音声と画像の送信方法のみが異なります。
WebRTC は VAD モードのみをサポートし、手動モードはサポートしていません。WebRTC では、音声は
input_audio_buffer.appendイベントを送信せずに RTP 経由で直接送信され、画像はinput_image_buffer.appendイベントをサポートせずにビデオトラック経由で送信されます。制御コマンドとサーバーイベントは、WebSocket と同じイベントタイプで DataChannel 経由で送信されます。
対話フローは次のとおりです:
- クライアントは音声データを送信します。WebSocket は input_audio_buffer.append イベント経由で送信し、WebRTC は手動でイベントを送信せずに音声トラック (RTP) 経由で自動的に送信します。
- サーバーは発話の開始を検出し、DataChannel (WebRTC) または WebSocket 経由で input_audio_buffer.speech_started イベントを送信します。
- サーバーは発話の終了を検出し、input_audio_buffer.speech_stopped イベントを送信します。
- サーバーは音声バッファをコミットし、input_audio_buffer.committed イベントを送信します。
- サーバーは応答の生成を開始し、conversation.item.created などのイベントを送信します。音声応答は、WebSocket の
response.audio.deltaイベントを介して増分的に返されるか、WebRTC の音声トラック (RTP) を介して直接送信されます。 - 応答中、サーバーは
response.audio_transcript.deltaイベントを介して増分的なテキスト文字起こしを返し、応答が完了するとresponse.doneイベントを送信します。
| ライフサイクル | クライアントイベント | サーバーイベント |
|---|---|---|
セッションの初期化 |
|
|
ユーザーの音声入力 |
| input_audio_buffer.speech_started
input_audio_buffer.speech_stopped
|
サーバーの音声出力 | なし |
response.audio_transcript.delta
response.audio_transcript.done
|
手動モード
session.update の session.turn_detection を null に設定して手動モードにします。クライアントは input_audio_buffer.commit と response.create を送信して応答を要求します。チャットアプリの音声メッセージなどのプッシュツートークシナリオに適しています。
対話フローは次のとおりです:
-
クライアントはいつでも input_audio_buffer.append および input_image_buffer.append イベントを送信して、音声と画像をバッファに追加できます。
input_image_buffer.appendイベントを送信する前に、少なくとも 1 つのinput_audio_buffer.appendイベントを送信する必要があります。 -
クライアントは input_audio_buffer.commit イベントを送信して音声および画像バッファをコミットし、現在のターンのすべてのユーザー入力 (音声および画像) が送信されたことをサーバーに通知します。
-
サーバーは input_audio_buffer.committed イベントで応答します。
-
クライアントは response.create イベントを送信し、サーバーがモデルの出力を返すのを待ちます。
-
サーバーは conversation.item.created イベントで応答します。
| ライフサイクル | クライアントイベント | サーバーイベント |
|---|---|---|
セッションの初期化 |
|
|
ユーザーの音声入力 |
|
|
サーバーの音声出力 |
|
response.audio_transcript.delta
response.audio_transcript.done
|
Web 検索
Web 検索により、モデルはリアルタイムデータを使用して、株価や天気などのタイムリーな情報に関する質問に答えることができます。モデルは検索が必要かどうかを自動的に判断します。
qwen3.5-omni-plus-realtimeのみが Web 検索をサポートしています。デフォルトでは無効になっており、session.updateで有効にします。
課金については、課金ルール内の
agent
Web 検索の有効化
session.update イベントに次のパラメーターを追加します:
enable_search:Web 検索機能を有効にするにはtrueに設定します。search_options.enable_source:応答に検索結果のソースを含めるにはtrueに設定します。
その他のパラメーターについては、「session.update」をご参照ください。
応答形式
Web 検索が有効な場合、response.done の usage オブジェクトには、検索の計測情報を含む plugins フィールドが含まれます:
{
"usage": {
"total_tokens": 2937,
"input_tokens": 2554,
"output_tokens": 383,
"input_tokens_details": {
"text_tokens": 2512,
"audio_tokens": 42
},
"output_tokens_details": {
"text_tokens": 90,
"audio_tokens": 293
},
"plugins": {
"search": {
"count": 1,
"strategy": "agent"
}
}
}
}
コード例
リアルタイム会話で Web 検索を有効にします:
DashScope Python SDK
update_session 呼び出しで enable_search および search_options パラメーターを渡します:
import os
import base64
import time
import json
import pyaudio
from dashscope.audio.qwen_omni import MultiModality, AudioFormat, OmniRealtimeCallback, OmniRealtimeConversation
import dashscope
dashscope.api_key = os.getenv('DASHSCOPE_API_KEY')
# シンガポールリージョン。{WorkspaceId} をお使いの Bailian ワークスペース ID に置き換えてください。URL はリージョンによって異なります。
url = 'wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/realtime'
model = 'qwen3.5-omni-plus-realtime'
voice = 'Tina'
class SearchCallback(OmniRealtimeCallback):
def __init__(self, pya):
self.pya = pya
self.out = None
def on_open(self):
self.out = self.pya.open(format=pyaudio.paInt16, channels=1, rate=24000, output=True)
def on_event(self, response):
if response['type'] == 'response.audio.delta':
self.out.write(base64.b64decode(response['delta']))
elif response['type'] == 'conversation.item.input_audio_transcription.delta':
preview = response.get('text', '') + response.get('stash', '')
print(f"\r[User] {preview}", end='', flush=True)
elif response['type'] == 'conversation.item.input_audio_transcription.completed':
print(f"\r[User] {response['transcript']}")
elif response['type'] == 'response.audio_transcript.done':
print(f"[LLM] {response['transcript']}")
elif response['type'] == 'response.done':
usage = response.get('response', {}).get('usage', {})
plugins = usage.get('plugins', {})
if plugins.get('search'):
print(f"[Search] count={plugins['search']['count']}, strategy={plugins['search']['strategy']}")
pya = pyaudio.PyAudio()
callback = SearchCallback(pya)
conv = OmniRealtimeConversation(model=model, callback=callback, url=url)
conv.connect()
conv.update_session(
output_modalities=[MultiModality.AUDIO, MultiModality.TEXT],
voice=voice,
instructions="あなたはシャオユンという名のパーソナルアシスタントです。",
enable_search=True,
search_options={'enable_source': True}
)
mic = pya.open(format=pyaudio.paInt16, channels=1, rate=16000, input=True)
print("Web search is enabled. Speak into the microphone (Ctrl+C to exit)...")
try:
while True:
audio_data = mic.read(3200, exception_on_overflow=False)
conv.append_audio(base64.b64encode(audio_data).decode())
time.sleep(0.01)
except KeyboardInterrupt:
conv.close()
mic.close()
callback.out.close()
pya.terminate()
print("\nConversation ended.")
DashScope Java SDK
updateSession メソッドで、Web 検索設定を parameters 引数に渡します:
import com.alibaba.dashscope.audio.omni.*;
import com.alibaba.dashscope.exception.NoApiKeyException;
import com.google.gson.JsonObject;
import javax.sound.sampled.*;
import java.nio.ByteBuffer;
import java.util.*;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.atomic.AtomicBoolean;
public class OmniSearch {
static class SequentialAudioPlayer {
private final SourceDataLine line;
private final Queue<byte[]> audioQueue = new ConcurrentLinkedQueue<>();
private final Thread playerThread;
private final AtomicBoolean shouldStop = new AtomicBoolean(false);
public SequentialAudioPlayer() throws LineUnavailableException {
AudioFormat format = new AudioFormat(24000, 16, 1, true, false);
line = AudioSystem.getSourceDataLine(format);
line.open(format);
line.start();
playerThread = new Thread(() -> {
while (!shouldStop.get()) {
byte[] audio = audioQueue.poll();
if (audio != null) {
line.write(audio, 0, audio.length);
} else {
try { Thread.sleep(10); } catch (InterruptedException ignored) {}
}
}
}, "AudioPlayer");
playerThread.start();
}
public void play(String base64Audio) {
audioQueue.add(Base64.getDecoder().decode(base64Audio));
}
public void close() {
shouldStop.set(true);
try { playerThread.join(1000); } catch (InterruptedException ignored) {}
line.drain();
line.close();
}
}
public static void main(String[] args) {
try {
SequentialAudioPlayer player = new SequentialAudioPlayer();
AtomicBoolean shouldStop = new AtomicBoolean(false);
OmniRealtimeParam param = OmniRealtimeParam.builder()
.model("qwen3.5-omni-plus-realtime")
.apikey(System.getenv("DASHSCOPE_API_KEY"))
// シンガポールリージョン。{WorkspaceId} をお使いの Bailian ワークスペース ID に置き換えてください。URL はリージョンによって異なります。
.url("wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/realtime")
.build();
OmniRealtimeConversation conversation = new OmniRealtimeConversation(param, new OmniRealtimeCallback() {
@Override public void onOpen() {
System.out.println("Connection established.");
}
@Override public void onClose(int code, String reason) {
System.out.println("Connection closed.");
shouldStop.set(true);
}
@Override public void onEvent(JsonObject event) {
String type = event.get("type").getAsString();
if ("response.audio.delta".equals(type)) {
player.play(event.get("delta").getAsString());
} else if ("response.audio_transcript.done".equals(type)) {
System.out.println("[LLM] " + event.get("transcript").getAsString());
} else if ("response.done".equals(type)) {
JsonObject response = event.getAsJsonObject("response");
if (response != null && response.has("usage")) {
JsonObject usage = response.getAsJsonObject("usage");
if (usage.has("plugins")) {
JsonObject plugins = usage.getAsJsonObject("plugins");
if (plugins.has("search")) {
JsonObject search = plugins.getAsJsonObject("search");
System.out.println("[Search] count=" + search.get("count").getAsInt()
+ ", strategy=" + search.get("strategy").getAsString());
}
}
}
}
}
});
conversation.connect();
conversation.updateSession(OmniRealtimeConfig.builder()
.modalities(Arrays.asList(OmniRealtimeModality.AUDIO, OmniRealtimeModality.TEXT))
.voice("Tina")
.enableTurnDetection(true)
.enableInputAudioTranscription(true)
.parameters(Map.of(
"instructions", "あなたはシャオユンという名のパーソナルアシスタントです。",
"enable_search", true,
"search_options", Map.of("enable_source", true)
))
.build()
);
System.out.println("Web search is enabled. Start speaking (press Ctrl+C to exit)...");
AudioFormat format = new AudioFormat(16000, 16, 1, true, false);
TargetDataLine mic = AudioSystem.getTargetDataLine(format);
mic.open(format);
mic.start();
ByteBuffer buffer = ByteBuffer.allocate(3200);
while (!shouldStop.get()) {
int bytesRead = mic.read(buffer.array(), 0, buffer.capacity());
if (bytesRead > 0) {
conversation.appendAudio(Base64.getEncoder().encodeToString(buffer.array()));
}
Thread.sleep(20);
}
conversation.close(1000, "Normal termination");
player.close();
mic.close();
} catch (NoApiKeyException e) {
System.err.println("API key not found: Set the DASHSCOPE_API_KEY environment variable.");
} catch (Exception e) {
e.printStackTrace();
}
}
}
WebSocket (Python)
session.update イベントの JSON ペイロードに enable_search および search_options フィールドを追加します:
import json
import os
import websocket
import base64
import pyaudio
import threading
API_KEY = os.getenv("DASHSCOPE_API_KEY")
# シンガポールリージョン。{WorkspaceId} をお使いの Bailian ワークスペース ID に置き換えてください。URL はリージョンによって異なります。
API_URL = "wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/realtime?model=qwen3.5-omni-plus-realtime"
pya = pyaudio.PyAudio()
out_stream = pya.open(format=pyaudio.paInt16, channels=1, rate=24000, output=True)
def on_open(ws):
ws.send(json.dumps({
"type": "session.update",
"session": {
"modalities": ["text", "audio"],
"voice": "Tina",
"instructions": "あなたはシャオユンという名のパーソナルアシスタントです。",
"input_audio_format": "pcm",
"output_audio_format": "pcm",
"enable_search": True,
"search_options": {
"enable_source": True
}
}
}))
print("Web search is enabled. Speak into the microphone...")
def send_audio():
mic = pya.open(format=pyaudio.paInt16, channels=1, rate=16000, input=True)
try:
while True:
audio = mic.read(3200, exception_on_overflow=False)
ws.send(json.dumps({
"type": "input_audio_buffer.append",
"audio": base64.b64encode(audio).decode()
}))
except Exception:
mic.close()
threading.Thread(target=send_audio, daemon=True).start()
def on_message(ws, message):
event = json.loads(message)
if event["type"] == "response.audio.delta":
out_stream.write(base64.b64decode(event["delta"]))
elif event["type"] == "response.audio_transcript.done":
print(f"[LLM] {event['transcript']}")
elif event["type"] == "response.done":
usage = event.get("response", {}).get("usage", {})
plugins = usage.get("plugins", {})
if plugins.get("search"):
print(f"[Search] count={plugins['search']['count']}, strategy={plugins['search']['strategy']}")
def on_error(ws, error):
print(f"Error: {error}")
headers = ["Authorization: Bearer " + API_KEY]
ws = websocket.WebSocketApp(API_URL, header=headers, on_open=on_open, on_message=on_message, on_error=on_error)
ws.run_forever()
API リファレンス
課金とレート制限
課金
課金はトークンベースで、モダリティ (音声、画像、テキスト) ごとに計測されます。価格については Model Studio コンソールをご確認ください。
注記マルチターンリアルタイム会話では、モデルが応答を生成するたびに、コンテキストウィンドウ内のすべての過去の会話コンテンツ (前のターンの音声、画像、テキストを含む) を、現在のターンの新しい入力とともに、入力トークンとして処理します。その結果、入力トークンは現在のターンの新しい入力のみでカウントされるのではなく、ターンごとに蓄積されます。
たとえば、10 秒の音声入力が 70 トークン (Qwen3.5-Omni-Realtime) に変換され、その音声がターン 3 でまだコンテキストウィンドウ内にある場合、それはターン 3 の入力トークンとしてカウントされます。実際に請求される入力トークン = コンテキストウィンドウ内のすべての過去のターンからのトークン + 現在のターンの新しい入力からのトークン。
音声と画像をトークンに変換するルール
音声
-
Qwen3.5-Omni-Realtime:- 入力音声:
合計トークン = 音声の持続時間 (秒) * 7 - 出力音声:
合計トークン = 音声の持続時間 (秒) * 12.5
- 入力音声:
-
Qwen3-Omni-Flash-Realtime:入力および出力音声は同じ数式を使用します:合計トークン = 音声の持続時間 (秒) * 12.5 -
Qwen-Omni-Turbo-Realtime:入力および出力音声は同じ数式を使用します:合計トークン = 音声の持続時間 (秒) * 251 秒未満の音声持続時間は 1 秒として課金されます。
画像
Qwen3.5-Omni-Realtimeシリーズモデルは、32x32ピクセルあたり 1 トークンを使用しますQwen3-Omni-Flash-Realtimeモデルは、32x32ピクセルあたり 1 トークンを使用しますQwen-Omni-Turbo-Realtimeモデルは、28x28ピクセルあたり 1 トークンを使用します
1 つの画像は 4~1,280 トークンを消費します。次のコードを使用して、画像の寸法とセッション期間に基づいてトークン消費量を見積もります:
# Pillow ライブラリをインストールするには、pip install Pillow を実行します
from PIL import Image
import math
# Qwen-Omni-Turbo-Realtime モデルの場合、スケーリング係数は 28 です。
# factor = 28
# Qwen3-Omni-Flash-Realtime および Qwen3.5-Omni-Realtime モデルの場合、スケーリング係数は 32 です。
factor = 32
def token_calculate(image_path='', duration=10):
"""
:param image_path: 画像パス
:param duration: セッション期間
:return: セッション期間に基づいた画像の合計トークン
"""
if len(image_path) > 0:
# 画像ファイルを開きます。
image = Image.open(image_path)
# 画像の元の寸法を取得します。
height = image.height
width = image.width
print(f"スケーリング前の画像の寸法: height={height}, width={width}")
# 高さを factor の倍数に調整します。
h_bar = round(height / factor) * factor
# 幅を factor の倍数に調整します。
w_bar = round(width / factor) * factor
# 画像トークンの下限:4 トークン。
min_pixels = factor * factor * 4
# 画像トークンの上限:1,280 トークン。
max_pixels = 1280 * factor * factor
# ピクセル数の制限に収まるように画像をスケーリングします。
if h_bar * w_bar > max_pixels:
# max_pixels を超えないようにスケーリング係数 beta を計算します。
beta = math.sqrt((height * width) / max_pixels)
# 調整後の高さが factor の整数倍になるように再計算します。
h_bar = math.floor(height / beta / factor) * factor
# 調整後の幅が factor の整数倍になるように再計算します。
w_bar = math.floor(width / beta / factor) * factor
elif h_bar * w_bar < min_pixels:
# スケーリングされた画像のピクセル数が min_pixels 未満にならないようにスケーリング係数 beta を計算します。
beta = math.sqrt(min_pixels / (height * width))
# 調整後の高さが factor の整数倍になるように再計算します。
h_bar = math.ceil(height * beta / factor) * factor
# 調整後の幅が factor の整数倍になるように再計算します。
w_bar = math.ceil(width * beta / factor) * factor
print(f"スケーリング後の画像の寸法: height={h_bar}, width={w_bar}")
# 画像トークンを計算します。
token = int((h_bar * w_bar) / (factor * factor))
print(f"スケーリング後のトークン数: {token}")
total_token = token * math.ceil(duration / 2)
print(f"合計トークン: {total_token}")
return total_token
else:
print("Error: image_path is empty. Cannot calculate tokens.")
return 0
if __name__ == "__main__":
total_token = token_calculate(image_path="xxx/test.jpg", duration=10)
レート制限
モデルのレート制限については、「レート制限」をご参照ください。
エラーコード
モデルの呼び出しに失敗し、エラーメッセージが返された場合は、「エラーコード」を参照して解決してください。
音声リスト
Qwen-Omni-Realtime モデルで利用可能な音声のリストについては、「音声」をご参照ください。