Qwen-Omni-Realtime memproses input audio dan gambar streaming (termasuk frame video) serta menghasilkan respons teks dan audio secara real-time.
Wilayah yang didukung: Singapura, China (Beijing). Setiap wilayah memerlukan API key sendiri.
Cara menggunakan
1. Membuat koneksi
Qwen-Omni-Realtime mendukung WebSocket dan WebRTC. WebSocket cocok untuk integrasi sisi server dengan pengaturan cepat. WebRTC ditujukan untuk skenario suara berlatensi rendah berbasis browser, mentransmisikan audio melalui UDP dengan pembatalan gema dan reduksi kebisingan bawaan.
WebSocket
WebSocket Native
Parameter koneksi:
| Parameter | Deskripsi |
|---|---|
Endpoint | Wilayah China (Beijing): wss://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api-ws/v1/realtime Wilayah Singapura: wss://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api-ws/v1/realtime Ganti |
Parameter kueri | Gunakan parameter kueri |
Header permintaan | Gunakan token Bearer untuk autentikasi:
|
# pip install websocket-client
import json
import websocket
import os
API_KEY=os.getenv("DASHSCOPE_API_KEY")
# Wilayah Singapura. Ganti {WorkspaceId} dengan workspace ID Bailian Anda. URL berbeda-beda tergantung wilayah.
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
# Diperlukan SDK versi 1.23.9 atau lebih baru.
import os
import json
from dashscope.audio.qwen_omni import OmniRealtimeConversation,OmniRealtimeCallback
import dashscope
# API key untuk wilayah Singapura dan China (Beijing) berbeda. Untuk mendapatkan API key, lihat https://www.alibabacloud.com/help/en/model-studio/get-api-key.
# Jika Anda belum mengonfigurasi API key, ubah baris berikut menjadi 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 berikut untuk wilayah Singapura. Saat memanggil, ganti {WorkspaceId} dengan workspace ID Anda yang sebenarnya. URL berbeda-beda tergantung wilayah.
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
// Diperlukan SDK versi 2.20.9 atau lebih baru.
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 berikut untuk wilayah Singapura. Saat memanggil, ganti {WorkspaceId} dengan workspace ID Anda yang sebenarnya. URL berbeda-beda tergantung wilayah.
.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
Membuat koneksi WebRTC melibatkan dua tahap:
- Pertukaran SDP (HTTP): Klien mengirim kemampuan media dan alamat jaringannya (Offer SDP) ke server melalui HTTP POST. Server mengembalikan informasinya (Answer SDP) untuk menyelesaikan negosiasi kemampuan.
- Koneksi (otomatis): Setelah negosiasi, lapisan WebRTC secara otomatis membuat saluran transportasi audio.
Konfigurasi pertukaran SDP:
Parameter | Deskripsi |
|---|---|
URL permintaan | Wilayah China (Beijing): https://{WorkspaceId}.cn-beijing.maas.aliyuncs.com/api/v1/webrtc/realtime Wilayah Singapura: https://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api/v1/webrtc/realtime Ganti |
Parameter kueri | Gunakan parameter kueri |
Content-Type | application/sdp |
Header permintaan | Authorization: Bearer DASHSCOPE_API_KEY |
Body permintaan | String Offer SDP yang dihasilkan klien |
Respons | Sukses: HTTP 200 dengan string Answer SDP server. Gagal: HTTP 4xx dengan pesan kesalahan JSON. |
Contoh kode koneksi:
# 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"
# Wilayah Singapura. Ganti {WorkspaceId} dengan workspace ID Bailian Anda. URL berbeda-beda tergantung wilayah.
SIGNALING_URL = "https://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api/v1/webrtc/realtime?model=" + MODEL
async def connect():
pc = RTCPeerConnection(RTCConfiguration(iceServers=[]))
# Tambahkan track audio untuk memastikan Offer SDP berisi m=audio (diperlukan oleh server)
pc.addTrack(AudioStreamTrack())
# Buat DataChannel untuk memicu negosiasi SDP (nama dapat disesuaikan; server mendorong event melalui channel bernama "txt")
pc.createDataChannel("oai-events")
# Pertukaran SDP: buat Offer dan kirim ke server
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)
# Koneksi ICE dibuat secara otomatis
await pc.setRemoteDescription(RTCSessionDescription(sdp=answer_sdp, type="answer"))
print("WebRTC connection established")
return pc
const API_KEY = 'your-api-key';
// Wilayah Singapura. Ganti {WorkspaceId} dengan workspace ID Bailian Anda. URL berbeda-beda tergantung wilayah.
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: [] });
// Tambahkan track audio untuk memastikan Offer SDP berisi m=audio (diperlukan oleh server)
const stream = await navigator.mediaDevices.getUserMedia({ audio: true });
stream.getAudioTracks().forEach(t => pc.addTrack(t, stream));
// Buat DataChannel untuk memicu negosiasi SDP (nama dapat disesuaikan; server mendorong event melalui channel bernama "txt")
pc.createDataChannel('oai-events');
// Tunggu hingga pengumpulan ICE selesai sebelum mengirim Offer untuk mendapatkan Answer
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();
// Koneksi ICE dibuat secara otomatis
await pc.setRemoteDescription({ type: 'answer', sdp: answerSdp });
console.log('WebRTC connection established');
};
// Buat Offer
const offer = await pc.createOffer();
await pc.setLocalDescription(offer);
return pc;
}
2. Mengonfigurasi sesi
Kirim event klien session.update:
{
// ID event yang dihasilkan klien.
"event_id": "event_ToPZqeobitzUJnt3QqtWg",
// Jenis event. Harus "session.update".
"type": "session.update",
// Konfigurasi sesi.
"session": {
// Modalitas output. Atur ke ["text"] untuk output teks saja, atau ["text", "audio"] untuk output teks dan audio.
"modalities": [
"text",
"audio"
],
// Suara untuk output audio.
"voice": "Ethan",
// Format audio input. Hanya "pcm" yang didukung. Input audio harus berupa aliran audio PCM dengan laju sampel 16 kHz.
"input_audio_format": "pcm",
// Format audio output. Hanya "pcm" yang didukung. Output audio adalah aliran audio PCM dengan laju sampel 24 kHz.
"output_audio_format": "pcm",
// Instruksi sistem untuk menentukan tujuan atau peran model.
"instructions": "Anda adalah agen layanan pelanggan AI untuk hotel bintang lima. Jawab pertanyaan pelanggan tentang jenis kamar, fasilitas, harga, dan kebijakan pemesanan secara akurat dan ramah. Selalu tanggapi dengan sikap profesional dan membantu. Jangan memberikan informasi yang belum dikonfirmasi atau informasi di luar cakupan layanan hotel.",
// Mengaktifkan deteksi aktivitas suara (VAD) di sisi server. Jika diaktifkan, server secara otomatis mendeteksi awal dan akhir ucapan.
// Jika null, klien mengontrol kapan memicu respons model.
"turn_detection": {
// Jenis VAD. Nilai valid: "server_vad" dan "semantic_vad". Kami merekomendasikan "semantic_vad" untuk model seri qwen3.5-omni-realtime.
"type": "semantic_vad",
// Ambang batas deteksi VAD. Kami merekomendasikan menaikkan nilai ini di lingkungan berisik dan menurunkannya di lingkungan tenang.
"threshold": 0.5,
// Durasi diam dalam milidetik (ms) yang menandakan akhir ucapan. Model memicu respons jika durasi ini terlampaui.
"silence_duration_ms": 800
}
}
}
3. Input audio dan gambar
Input audio wajib; input gambar opsional. Metode input tergantung pada protokol.
WebSocket
Kirim data audio dan gambar yang diencode Base64 ke buffer server menggunakan event input_audio_buffer.append dan input_image_buffer.append.
Gambar dapat berasal dari file lokal atau tangkapan aliran video real-time.
Dengan VAD sisi server diaktifkan, server secara otomatis mengirimkan data dan memicu respons saat akhir ucapan. Dengan VAD dinonaktifkan (mode manual), panggil event input_audio_buffer.commit untuk mengirimkan data setelah pengiriman.
WebRTC
Track audio dan video (saluran media RTP) yang ditambahkan saat pembuatan koneksi mentransmisikan data ke server secara otomatis.
- Audio: Ditransmisikan langsung melalui track audio (RTP). Tidak perlu event
input_audio_buffer.append. - Gambar: Dikirim sebagai frame video melalui track video (RTP).
input_image_buffer.appendtidak didukung.
WebRTC hanya mendukung mode VAD sisi server (
server_vadatausemantic_vad). Mode manual tidak didukung.
4. Menerima respons model
Format respons tergantung pada modalitas output yang dikonfigurasi.
WebSocket
-
Hanya teks
Terima teks streaming dengan event response.text.delta, dan teks lengkap dengan event response.text.done.
-
Teks dan audio
- Teks: Terima teks streaming dengan event response.audio_transcript.delta, dan teks lengkap dengan event response.audio_transcript.done.
- Audio: Terima audio streaming yang diencode Base64 dengan event response.audio.delta. Event response.audio.done menunjukkan bahwa generasi audio telah selesai.
WebRTC
-
Hanya teks
Sama seperti WebSocket. Terima event teks streaming melalui DataChannel.
-
Teks dan audio
- Teks: Diterima melalui DataChannel sebagai event teks streaming, sama seperti WebSocket.
- Audio: Diterima dan diputar secara real-time melalui track RTP. Tidak perlu event
response.audio.delta.
Pemilihan model
Qwen3.5-Omni-Realtime meningkatkan Qwen3-Omni-Flash-Realtime di bidang-bidang berikut:
-
Tingkat kecerdasan
Setara dengan Qwen3.5-Plus.
-
Pencarian web
Pencarian web bawaan — model secara otonom mencari untuk menjawab pertanyaan real-time. Untuk detailnya, lihat Pencarian web.
-
Pemanggilan alat
Pemanggilan fungsi — model secara otonom memanggil alat eksternal. Untuk detailnya, lihat Seri Qwen-Omni-Realtime.
-
Interrupsi semantik
Mengidentifikasi maksud percakapan untuk mencegah interupsi dari backchanneling dan kebisingan latar belakang.
-
Kontrol suara
Kontrol volume, laju berbicara, dan emosi melalui perintah suara (misalnya, "berbicara lebih cepat", "lebih keras", "dengan nada bahagia").
-
Bahasa yang didukung
Mendukung pengenalan suara untuk 113 bahasa dan dialek serta generasi suara untuk 36 bahasa dan dialek.
-
Suara yang didukung
Mendukung 55 suara, termasuk 47 suara multibahasa dan 8 suara dialek. Untuk daftar lengkap, lihat Daftar suara.
-
Kloning suara
Gunakan suara kloning kustom untuk percakapan real-time (Qwen3.5-omni-plus-realtime dan Qwen3.5-omni-flash-realtime). Untuk detailnya, lihat Kloning suara.
Periksa konsol Model Studio untuk nama model, konteks, harga, dan versi snapshot. Untuk batas laju konkurensi, lihat Pembatasan laju.
Batasan
-
Pencarian web dan pemanggilan alat saling eksklusif.
-
Sesi WebSocket tunggal dapat berlangsung hingga 120 menit. Koneksi ditutup secara otomatis pada batas ini.
-
Model menyimpan riwayat percakapan hingga batas putaran dan durasi berikut. Ketika terlampaui, riwayat terlama dibuang. Durasi maksimum adalah durasi kumulatif audio atau video (frame gambar) yang disimpan dalam konteks.
Video dimasukkan sebagai frame yang diekstraksi (direkomendasikan: 1 fps). Durasi maksimum video adalah durasi kumulatif frame yang disimpan — misalnya, 240 detik berarti hanya frame dari 240 detik terakhir yang disimpan.
Model
qwen3-omni-flash-realtimememiliki batas 8 putaran dialog (biasanya tercapai lebih dulu). Batas durasinya tergantung pada panjang konteks model dan tidak tercantum secara terpisah.Model
Putaran audio maksimum
Putaran video maksimum
Durasi audio maksimum
Durasi video maksimum
qwen3.5-omni-plus-realtime
100 putaran
50 putaran
600 detik
240 detik
qwen3.5-omni-flash-realtime
80 putaran
50 putaran
480 detik
120 detik
qwen3-omni-flash-realtime
8 putaran
8 putaran
—
—
Memulai
Dapatkan API key dan tetapkan sebagai variabel lingkungan.
Pilih bahasa pemrograman dan ikuti langkah-langkah untuk memulai obrolan real-time.
WebSocket
DashScope Python SDK
- Lingkungan runtime
Pastikan Python 3.10 atau lebih baru telah terinstal.
Instal PyAudio untuk sistem operasi Anda.
macOS
brew install portaudio && pip install pyaudio
Debian/Ubuntu
- Jika Anda tidak menggunakan virtual environment, Anda dapat menginstalnya langsung menggunakan manajer paket sistem:
sudo apt-get install python3-pyaudio
- Jika Anda menggunakan virtual environment, pertama instal dependensi build:
sudo apt update
sudo apt install -y python3-dev portaudio19-dev
Kemudian, instal dengan pip di virtual environment yang diaktifkan:
pip install pyaudio
CentOS
sudo yum install -y portaudio portaudio-devel && pip install pyaudio
Windows
pip install pyaudio
Instal dependensi lainnya:
pip install websocket-client dashscope
-
Mode interaksi
-
Mode VAD (Deteksi Aktivitas Suara, secara otomatis mendeteksi awal dan akhir ucapan)
Server merespons setelah mendeteksi akhir ucapan pengguna.
-
Mode manual (tekan untuk berbicara, lepas untuk mengirim)
Klien mengontrol awal dan akhir ucapan. Setelah berbicara, aplikasi Anda harus memberi tahu server.
Mode VAD
Buat file Python bernama vad_dash.py dan salin kode berikut ke dalam file tersebut:
vad_dash.py
# Dependensi: 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 # Konfigurasi: URL, API key, suara, model, peran model # Tentukan wilayah. 'intl' untuk wilayah Singapura, 'cn' untuk wilayah China (Beijing). Ganti {WorkspaceId} dengan workspace ID Bailian Anda. 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' # Konfigurasi API key. Jika variabel lingkungan tidak diatur, ganti baris berikut dengan API key Anda: dashscope.api_key = "sk-xxx" dashscope.api_key = os.getenv('DASHSCOPE_API_KEY') # Tentukan suara. voice = 'Ethan' # Tentukan model. model = 'qwen3.5-omni-plus-realtime' # Tentukan peran model. instructions = "Anda adalah Xiaoyun, asisten pribadi. Jawab pertanyaan pengguna dengan cara yang lucu dan cerdas." class SimpleCallback(OmniRealtimeCallback): def __init__(self, pya): self.pya = pya self.out = None def on_open(self): # Inisialisasi aliran output audio. 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': # Putar audio. self.out.write(base64.b64decode(response['delta'])) elif response['type'] == 'conversation.item.input_audio_transcription.delta': # Pratinjau streaming: teks adalah awalan yang dikonfirmasi, stash adalah akhiran yang belum dikonfirmasi. preview = response.get('text', '') + response.get('stash', '') print(f"\r[User] {preview}", end='', flush=True) elif response['type'] == 'conversation.item.input_audio_transcription.completed': # Transkripsi selesai. Cetak teks akhir dan pindah ke baris baru. print(f"\r[User] {response['transcript']}") elif response['type'] == 'response.audio_transcript.done': # Cetak teks respons asisten. print(f"[LLM] {response['transcript']}") # 1. Inisialisasi perangkat audio. pya = pyaudio.PyAudio() # 2. Buat fungsi callback dan percakapan. callback = SimpleCallback(pya) conv = OmniRealtimeConversation(model=model, callback=callback, url=url) # 3. Hubungkan dan konfigurasikan sesi. conv.connect() conv.update_session(output_modalities=[MultiModality.AUDIO, MultiModality.TEXT], voice=voice, instructions=instructions) # 4. Inisialisasi aliran input audio. mic = pya.open(format=pyaudio.paInt16, channels=1, rate=16000, input=True) # 5. Loop utama untuk memproses input audio. 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: # Bersihkan sumber daya. conv.close() mic.close() callback.out.close() pya.terminate() print("\nConversation ended")Jalankan
vad_dash.pyuntuk memulai percakapan real-time melalui mikrofon Anda. Sistem mendeteksi ucapan dan mengalirkan audio ke server.Mode manual
Buat file Python bernama
manual_dash.pydan salin kode berikut ke dalam file tersebut:manual_dash.py
# Dependensi: dashscope >= 1.23.9, pyaudio import os import base64 import sys import threading import pyaudio from dashscope.audio.qwen_omni import * import dashscope # Jika variabel lingkungan tidak diatur, ganti baris berikut dengan API key Anda: dashscope.api_key = "sk-xxx" dashscope.api_key = os.getenv('DASHSCOPE_API_KEY') voice = 'Ethan' class MyCallback(OmniRealtimeCallback): """Callback minimal: menginisialisasi speaker saat koneksi terbuka dan memutar audio yang dikembalikan langsung di penanganan event.""" def __init__(self, ctx): super().__init__() self.ctx = ctx def on_open(self) -> None: # Inisialisasi PyAudio dan speaker (24 kHz, mono, 16-bit) setelah koneksi terbentuk. 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): # Decode data base64 dan tulis langsung ke aliran output untuk pemutaran. 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): # Tandai putaran saat ini sebagai selesai, memungkinkan loop utama dilanjutkan. 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): """Lepaskan sumber daya audio dan PyAudio secara aman.""" 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 ...') # Konteks runtime: menyimpan handle audio dan percakapan. 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 berikut untuk wilayah Singapura. Saat memanggil, ganti {WorkspaceId} dengan workspace ID Anda yang sebenarnya. URL berbeda-beda tergantung wilayah. 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 # Konfigurasi sesi: aktifkan output teks dan audio, dan nonaktifkan VAD sisi server untuk perekaman manual. 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="Anda adalah Xiaoyun, asisten pribadi. Mohon jawab pertanyaan pengguna secara akurat dan ramah, selalu tanggapi dengan sikap membantu." ) 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.")Jalankan
manual_dash.py. Tekan Enter untuk mulai merekam, dan tekan Enter lagi untuk berhenti dan mengirim. Respons audio model diputar secara otomatis. -
DashScope Java SDK
Pilih mode interaksi-
Mode VAD (Deteksi Aktivitas Suara, secara otomatis mendeteksi awal dan akhir ucapan)
API Realtime mendeteksi kapan Anda mulai dan berhenti berbicara serta merespons.
-
Mode manual (tekan untuk berbicara, lepas untuk mengirim)
Klien mengontrol awal dan akhir ucapan. Setelah berbicara, klien harus mengirim pesan ke server.
Mode 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 berikut untuk wilayah Singapura. Saat memanggil, ganti {WorkspaceId} dengan workspace ID Anda yang sebenarnya. URL berbeda-beda tergantung wilayah.
.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",
"Anda adalah agen layanan pelanggan AI untuk hotel bintang lima. Jawab pertanyaan pelanggan tentang jenis kamar, fasilitas, harga, dan kebijakan pemesanan secara akurat dan ramah. Selalu tanggapi dengan sikap profesional dan membantu. Jangan memberikan informasi yang belum dikonfirmasi atau informasi di luar cakupan layanan hotel."))
.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":
// Pratinjau streaming: teks adalah awalan yang dikonfirmasi, stash adalah akhiran yang belum dikonfirmasi
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;
}
}
}
Jalankan OmniServerVad.main() untuk memulai percakapan real-time melalui mikrofon Anda. Sistem mendeteksi ucapan dan mengirim audio ke server.
Mode manual
OmniWithoutServerVad.java
// Diperlukan DashScope Java SDK 2.20.9 atau lebih baru.
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();
}
// Memainkan potongan audio dan memblokir hingga pemutaran selesai.
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);
// Menunggu audio dalam buffer selesai diputar.
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();
}
}
}
// Merekam audio dan mengalirkannya ke percakapan secara real-time.
private static void recordAndSend(TargetDataLine line, OmniRealtimeConversation conversation) {
byte[] buffer = new byte[3200];
AtomicBoolean stopRecording = new AtomicBoolean(false);
// Memulai thread untuk mendengarkan tombol Enter.
Thread enterKeyListener = new Thread(() -> {
try {
System.in.read();
stopRecording.set(true);
} catch (IOException e) {
e.printStackTrace();
}
});
enterKeyListener.start();
// Merekam dan mengirim audio secara real-time.
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 key untuk Singapura dan China (Beijing) berbeda. Dapatkan API key dari: https://www.alibabacloud.com/help/en/model-studio/get-api-key
// Jika variabel lingkungan DASHSCOPE_API_KEY tidak diatur, ganti pemanggilan apikey di bawah ini dengan API key Model Studio Anda, misalnya: .apikey("sk-xxx")
.apikey(System.getenv("DASHSCOPE_API_KEY"))
// URL berikut untuk wilayah Singapura. Saat memanggil, ganti {WorkspaceId} dengan workspace ID Anda yang sebenarnya. URL berbeda-beda tergantung wilayah.
.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":
// Pratinjau streaming: teks adalah awalan yang dikonfirmasi, stash adalah akhiran yang belum dikonfirmasi
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)
// Atur peran model.
.parameters(new HashMap<String, Object>() {{
put("instructions","Anda adalah Xiaoyun, asisten pribadi. Mohon jawab pertanyaan pengguna secara akurat dan ramah, selalu tanggapi dengan sikap membantu.");
}})
.build();
conversation.updateSession(config);
// Siapkan mikrofon untuk perekaman.
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; // Keluar dari loop jika terjadi kesalahan.
}
System.out.println("Recording started. Speak now... Press Enter again to stop recording and send.");
recordAndSend(line, conversation);
conversation.commit();
// Tunggu hingga transkripsi selesai sebelum memicu respons model untuk menghindari output yang tercampur
transcriptionDoneLatch.get().await(10, TimeUnit.SECONDS);
System.out.println("Waiting for model response...");
conversation.createResponse(null, null);
responseDoneLatch.get().await();
// Setel ulang latch untuk putaran berikutnya.
responseDoneLatch.set(new CountDownLatch(1));
transcriptionDoneLatch.set(new CountDownLatch(1));
}
} catch (LineUnavailableException e) {
e.printStackTrace();
} finally {
if (line != null) {
line.stop();
line.close();
}
}
}}
Jalankan OmniWithoutServerVad.main(). Tekan Enter untuk mulai merekam, dan tekan lagi untuk berhenti dan mengirim. Respons model diputar secara otomatis.
WebSocket (Python)
-
Siapkan lingkungan runtime
Pastikan Python 3.10 atau lebih baru telah terinstal.
Instal pyaudio untuk sistem operasi Anda.
macOS
brew install portaudio && pip install pyaudioDebian/Ubuntu
sudo apt-get install python3-pyaudio or pip install pyaudioKami merekomendasikan menggunakan
pip install pyaudio. Jika instalasi gagal, pertama instal dependensiportaudiountuk sistem operasi Anda.CentOS
sudo yum install -y portaudio portaudio-devel && pip install pyaudioWindows
pip install pyaudioInstal dependensi WebSocket:
pip install websockets==15.0.1
-
Buat klien
Buat file bernama
omni_realtime_client.pydan salin kode berikut ke dalamnya: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" # Direkomendasikan untuk model seri 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 {} # Status respons saat ini self._current_response_id = None self._current_item_id = None self._is_responding = False # Status pencetakan transkrip input/output self._print_input_transcript = True self._output_transcript_buffer = "" async def connect(self) -> None: """Membuat koneksi WebSocket dengan API Realtime.""" url = f"{self.base_url}?model={self.model}" headers = { "Authorization": f"Bearer {self.api_key}" } self.ws = await websockets.connect(url, additional_headers=headers) # konfigurasi sesi 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: """Memperbarui konfigurasi sesi.""" event = { "type": "session.update", "session": config } await self.send_event(event) async def stream_audio(self, audio_chunk: bytes) -> None: """Mengalirkan data audio mentah ke API.""" # Hanya PCM 16-bit, 16 kHz, mono yang didukung. 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: """Mengirimkan buffer audio untuk memicu pemrosesan.""" event = { "type": "input_audio_buffer.commit" } await self.send_event(event) async def append_image(self, image_chunk: bytes) -> None: """Menambahkan data gambar ke buffer gambar. Data gambar dapat berasal dari file lokal atau aliran video real-time. Catatan: - Format gambar harus JPG atau JPEG. Kami merekomendasikan resolusi 480p atau 720p. Resolusi maksimum yang didukung adalah 1080p. - Gambar tunggal setelah diencode Base64 tidak boleh melebihi 256 KB. Kami merekomendasikan menjaga ukuran gambar mentah di bawah 190 KB sebelum encoding. - Encode data gambar ke Base64 sebelum mengirim. - Kami merekomendasikan mengirim gambar satu frame per detik. - Anda harus mengirim data audio setidaknya sekali sebelum mengirim data gambar. """ 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: """Meminta API untuk menghasilkan respons. Ini hanya diperlukan dalam Mode Manual.""" event = { "type": "response.create" } await self.send_event(event) async def cancel_response(self) -> None: """Membatalkan respons saat ini.""" event = { "type": "response.cancel" } await self.send_event(event) async def handle_interruption(self): """Menangani interupsi pengguna terhadap respons saat ini.""" if not self._is_responding: return # 1. Batalkan respons saat ini. 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: """Menutup koneksi WebSocket.""" if self.ws: await self.ws.close() -
Pilih mode interaksi
-
Mode VAD (Deteksi Aktivitas Suara, secara otomatis mendeteksi awal dan akhir ucapan)
API Realtime mendeteksi kapan Anda mulai dan berhenti berbicara serta menghasilkan respons.
-
Mode manual (tekan untuk berbicara, lepas untuk mengirim)
Anda mengontrol kapan mulai dan berhenti mengirim audio. Setelah berbicara, klien harus mengirim pesan ke server untuk menghasilkan respons.
Mode VAD
Di direktori yang sama dengan
omni_realtime_client.py, buat file bernamavad_mode.pydan salin kode berikut ke dalamnya:vad_mode.py
# -- coding: utf-8 -- import os, asyncio, pyaudio, queue, threading from omni_realtime_client import OmniRealtimeClient, TurnDetectionMode # Kelas pemutar audio yang menangani interupsi 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() # Rekam dari mikrofon dan kirim audio 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( # Ini adalah base_url untuk wilayah Singapura. Base_url untuk wilayah China (Beijing) adalah 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="You are Xiaoyun, a witty and humorous assistant.", # SEMANTIC_VAD direkomendasikan untuk model seperti qwen3.5-omni-realtime. 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...") # Jalankan secara konkuren await asyncio.gather(client.handle_messages(), record_and_send(client)) if __name__ == "__main__": try: asyncio.run(main()) except KeyboardInterrupt: print("\nProgram exited.")Jalankan
vad_mode.pyuntuk memulai percakapan real-time melalui mikrofon Anda. Sistem mendeteksi ucapan dan mengalirkan audio ke server.Mode manual
Di direktori yang sama dengan
omni_realtime_client.py, buat file bernamamanual_mode.pydan salin kode berikut ke dalamnya: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: """Kelas pemutar audio real-time.""" def __init__(self, sample_rate=24000, channels=1, sample_width=2): self.sample_rate = sample_rate self.channels = channels self.sample_width = sample_width # 2 byte untuk 16-bit self.audio_queue = queue.Queue() self.is_playing = False self.play_thread = None self.pyaudio_instance = None self.stream = None self._lock = threading.Lock() # Tambahkan lock untuk akses sinkronisasi. self._last_data_time = time.time() # Catat waktu saat data terakhir diterima. self._response_done = False # Tambahkan flag untuk menunjukkan bahwa respons telah selesai. self._waiting_for_response = False # Flag untuk menunjukkan apakah klien sedang menunggu respons server. # Catat waktu penulisan terakhir ke aliran audio dan durasi potongan audio terakhir untuk deteksi akhir pemutaran yang lebih akurat. self._last_play_time = time.time() self._last_chunk_duration = 0.0 def start(self): """Mulai pemutar audio.""" with self._lock: if self.is_playing: return self.is_playing = True try: self.pyaudio_instance = pyaudio.PyAudio() # Buat aliran output audio. self.stream = self.pyaudio_instance.open( format=pyaudio.paInt16, # 16-bit channels=self.channels, rate=self.sample_rate, output=True, frames_per_buffer=1024 ) # Mulai thread pemutaran. 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): """Hentikan pemutar audio.""" with self._lock: if not self.is_playing: return self.is_playing = False # Kosongkan antrian. while not self.audio_queue.empty(): try: self.audio_queue.get_nowait() except queue.Empty: break # Tunggu thread pemutaran selesai. Tunggu di luar lock untuk menghindari deadlock. if self.play_thread and self.play_thread.is_alive(): self.play_thread.join(timeout=2.0) # Dapatkan lock lagi untuk membersihkan sumber daya. with self._lock: self._cleanup_resources() print("Audio player stopped") def _cleanup_resources(self): """Bersihkan sumber daya audio. Ini harus dipanggil dalam lock.""" try: # Tutup aliran audio. 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): """Tambahkan data audio ke antrian pemutaran.""" if self.is_playing and audio_data: self.audio_queue.put(audio_data) with self._lock: self._last_data_time = time.time() # Perbarui waktu saat data terakhir diterima. self._waiting_for_response = False # Data diterima, tidak perlu menunggu lagi. def stop_receiving_data(self): """Tandai bahwa tidak ada data audio baru yang akan diterima.""" with self._lock: self._response_done = True self._waiting_for_response = False # Respons berakhir, tidak perlu menunggu lagi. def prepare_for_next_turn(self): """Setel ulang status pemutar untuk putaran percakapan berikutnya.""" 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 # Mulai menunggu respons berikutnya. # Kosongkan data audio yang tersisa dari putaran sebelumnya. while not self.audio_queue.empty(): try: self.audio_queue.get_nowait() except queue.Empty: break def is_finished_playing(self): """Periksa apakah semua data audio telah diputar.""" 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 # ---------------------- Deteksi akhir cerdas ---------------------- # 1. Metode utama: Jika server telah menandai selesai dan antrian pemutaran kosong, # tunggu potongan audio terbaru selesai diputar (durasi potongan + toleransi 0,1 detik). if self._response_done and queue_size == 0: min_wait = max(self._last_chunk_duration + 0.1, 0.5) # Tunggu minimal 0,5 detik. if time_since_last_play >= min_wait: return True # 2. Metode cadangan: Jika tidak ada data baru yang diterima selama lebih dari satu detik dan antrian pemutaran kosong. # Logika ini berfungsi sebagai pengaman jika server tidak secara eksplisit mengirim `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): """Thread pekerja untuk memutar data audio.""" while True: # Periksa apakah harus berhenti. with self._lock: if not self.is_playing: break stream_ref = self.stream # Dapatkan referensi ke aliran. try: # Dapatkan data audio dari antrian, dengan timeout 0,1 detik. audio_data = self.audio_queue.get(timeout=0.1) # Periksa status dan validitas aliran lagi. with self._lock: if self.is_playing and stream_ref and not stream_ref.is_stopped(): try: # Putar data audio. stream_ref.write(audio_data) # Perbarui informasi pemutaran terbaru. 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 # Tandai blok data ini sebagai diproses. self.audio_queue.task_done() except queue.Empty: # Lanjutkan menunggu jika antrian kosong. continue except Exception as e: print(f"Error playing audio: {e}") break class MicrophoneRecorder: """Perekam mikrofon real-time.""" 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): """Thread pekerja perekaman.""" # Terus-menerus membaca data dari aliran audio selama _is_recording bernilai True. while self._is_recording: try: # Gunakan exception_on_overflow=False untuk menghindari crash karena luapan buffer. data = self.stream.read(self.chunk_size, exception_on_overflow=False) self.frames.append(data) except (IOError, OSError) as e: # Saat aliran ditutup, operasi baca mungkin menimbulkan kesalahan. print(f"Error reading from recording stream, it might be closed: {e}") break def start(self): """Mulai perekaman.""" 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): """Hentikan perekaman dan kembalikan data audio.""" if not self._is_recording: return None self._is_recording = False # Tunggu thread perekaman keluar dengan aman. 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): """Bersihkan sumber daya PyAudio dengan aman.""" 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(): """ Uji interaktif untuk percakapan multi-putaran dengan dukungan audio dan gambar. """ # ------------------- 1. Inisialisasi dan koneksi (sekali saja) ------------------- # API key untuk wilayah Singapura dan China (Beijing) berbeda. Untuk mendapatkan API key, lihat https://www.alibabacloud.com/help/en/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( # Ini adalah base_url untuk wilayah Singapura. Jika Anda menggunakan model di wilayah China (Beijing), ganti base_url dengan 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="You are Xiaoyun, a personal assistant. Please answer the user's questions accurately and in a friendly manner, always responding with a helpful attitude.", # Atur peran model. 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. Loop percakapan multi-putaran ------------------- while True: print(f"\n--- Turn {turn_counter} ---") audio_player.prepare_for_next_turn() recorded_audio = None image_paths = [] # --- Dapatkan input pengguna: Rekam dari mikrofon --- loop = asyncio.get_event_loop() recorder = MicrophoneRecorder(sample_rate=16000) # Kami merekomendasikan laju sampel 16k untuk pengenalan suara. 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 # --- Dapatkan input gambar (opsional) --- # Fitur input gambar dinonaktifkan secara default. Hapus komentar kode di bawah ini untuk mengaktifkannya. # 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. Kirim data dan dapatkan respons --- print("\n--- Input Confirmation ---") print(f"Audio to process: 1 (from microphone), Images: {len(image_paths)}") print("------------------") # 3.1 Kirim rekaman audio. 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 Kirim semua file gambar. # Fitur input gambar dinonaktifkan secara default. Hapus komentar kode di bawah ini untuk mengaktifkannya. # 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 Kirim dan tunggu respons. 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. Bersihkan sumber daya ------------------- 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.")Jalankan
manual_mode.py. Tekan Enter untuk mulai merekam, dan tekan Enter lagi untuk berhenti dan mengirim. -
WebRTC
Python
-
Lingkungan runtime
Diperlukan Python 3.10 atau lebih baru. Instal dependensi berikut:
pip install aiortc aiohttp sounddevice numpy certifi av
-
Jalankan demo
Buat file Python bernama
webrtc_demo.pydan tempel kode berikut:webrtc_demo.py
# Dependensi: 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 # Ganti dengan API key Anda, atau atur variabel lingkungan DASHSCOPE_API_KEY API_KEY = os.getenv("DASHSCOPE_API_KEY", "your-api-key") MODEL = "qwen3.5-omni-plus-realtime" # Wilayah Singapura. Ganti {WorkspaceId} dengan workspace ID Bailian Anda. URL berbeda-beda tergantung wilayah. SIGNALING_URL = "https://{WorkspaceId}.ap-southeast-1.maas.aliyuncs.com/api/v1/webrtc/realtime?model=" + MODEL # --------------- Parsing frame audio --------------- def _nb_channels(frame: AudioFrame) -> int: """Dapatkan jumlah channel dalam frame audio, kompatibel dengan berbagai versi 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: """ Frame audio server adalah stereo interleaved. Pengubahan bentuk langsung menyebabkan ketidaksesuaian channel. Susun ulang menjadi (samples, channels) berdasarkan jumlah channel sebenarnya. Versi decoder aiortc yang berbeda mengembalikan bentuk array yang berbeda untuk audio yang sama, sehingga penanganan terpadu diperlukan di sini. """ 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}") # --------------- Pemutar audio latensi rendah --------------- class RemoteAudioPlayer: """ Pemutar audio latensi rendah yang memainkan blok audio 5ms untuk meminimalkan penundaan. Mendukung interupsi suara: mengosongkan buffer saat pengguna mulai berbicara, menghentikan pemutaran respons model lama. Menggabungkan audio stereo server menjadi mono (rata-rata channel kiri dan kanan) untuk pemutaran. """ 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): """Kosongkan buffer pemutaran untuk interupsi suara""" 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): """Terima frame audio, gabungkan channel secara otomatis dan masukkan ke antrian""" 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 # --------------- utama --------------- async def main(): pc = RTCPeerConnection(RTCConfiguration(iceServers=[])) # Inisialisasi pemutar audio (output mono, blocksize 5ms untuk latensi rendah) speaker = RemoteAudioPlayer(samplerate=48000, out_channels=1, blocksize=240, max_seconds=0.2) speaker.start() # Inisialisasi mikrofon (avfoundation macOS; untuk Linux gunakan pulse atau 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) # Klien membuat DataChannel (nama dapat disesuaikan); server mendorong event melalui channel bernama "txt" pc.createDataChannel("oai-events") remote_dc = None got_first_txt_msg = False def make_session_update() -> dict: """Bangun konfigurasi session.update: suara, format audio, strategi VAD, parameter inferensi""" 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, }, } # Tangani event DataChannel yang didorong server @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')}") # Kosongkan buffer pemutaran saat pengguna mulai berbicara (interupsi suara) if isinstance(evt, dict) and evt.get("type") == "input_audio_buffer.speech_started": speaker.clear() print("[Playback] User speech detected, clearing buffer (interruption)") # Kirim session.update setelah menerima pesan pertama di channel txt 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") # Terima audio server dan putar dengan latensi rendah @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() # Pertukaran SDP: buat Offer dan POST ke server signaling, dapatkan 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())Jalankan
webrtc_demo.pyuntuk memulai percakapan real-time dengan model Qwen-Omni-Realtime melalui mikrofon Anda. Sistem mendeteksi awal ucapan Anda dan mengirim audio ke server secara otomatis.
JavaScript
-
Prasyarat
- Gunakan browser modern yang mendukung WebRTC (Chrome, Edge, Firefox, Safari, dll.).
- Browser memerlukan izin mikrofon.
- Karena kebijakan keamanan lintas-origin browser, browser tidak dapat langsung mengirim permintaan koneksi ke server. Anda perlu menjalankan perintah curl di terminal untuk menyelesaikan pengaturan koneksi.
-
Jalankan demo
Buat file HTML bernama
webrtc_demo.htmldan tempel kode berikut:webrtc_demo.html
<!DOCTYPE html> <html lang="en"> <head> <meta charset="UTF-8" /> <title>WebRTC Realtime Voice Chat</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; } /* Bilah atas lengket */ .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; } /* Bilah alat */ .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; } /* Tombol */ 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; } /* Indikator status */ .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; } } /* Kartu 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 */ .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; } /* Panel event */ .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 Realtime Voice Chat</h1> <div class="sticky-top"> <div class="toolbar"> <button id="startBtn" class="btn-primary">Start Session</button> <button id="setAnswerBtn" disabled>Set Answer</button> <button id="endBtn" class="btn-danger" disabled>End Session</button> <button id="downloadBtn" disabled>Download Remote Audio</button> <label> <input id="sendVideoCheckbox" type="checkbox" /> Enable Video </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>Copy Offer SDP</button> </div> <textarea id="offerBox" rows="6" readonly placeholder="Auto-generated after clicking Start Session"></textarea> <div class="card-hint">Auto-generated after ICE gathering completes. Copy and send to the server via curl to get the Answer.</div> </div> <div class="card"> <div class="card-title">curl Command</div> <div style="margin-bottom: 8px;"> <button id="copyCurlBtn" disabled>Copy curl Command</button> </div> <div class="card-hint" style="margin-bottom: 4px;">Wilayah Singapura. Ganti {WorkspaceId} dengan workspace ID Bailian Anda. URL berbeda-beda tergantung wilayah.</div> <textarea id="curlBox" rows="6" placeholder="Auto-filled after Offer SDP is generated"></textarea> <div class="card-hint">Ganti {WorkspaceId} dengan workspace ID Anda, lalu salin perintah ini ke terminal Anda. Tempel Answer SDP yang dikembalikan di bawah ini.</div> </div> <div class="card"> <div class="card-title">Answer SDP</div> <textarea id="answerBox" rows="6" placeholder="Paste the Answer SDP returned by curl here"></textarea> <div class="card-hint">Setelah menempel, klik Set Answer di atas untuk membuat koneksi.</div> </div> <div class="video-section" id="videoSection" style="display:none;"> <div class="video-label">Local Video Preview</div> <video id="localVideo" autoplay playsinline muted></video> </div> <div class="events-title">Events (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)); } // ===== Panel event ===== 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 = "Waiting for events..."; 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 copied"); }; copyCurlBtn.onclick = async () => { const txt = curlBox.value; if (!txt) return; await navigator.clipboard.writeText(txt); alert("curl command copied. Run it in your terminal."); }; downloadBtn.onclick = () => { if (audioBlob) downloadBlob(audioBlob, 'remote-audio.webm'); else alert('No audio recording available'); }; async function startSession() { if (pc) return; pc = new RTCPeerConnection({ iceServers: [] }); clearUIEvents(); setStatus('Requesting microphone access...', '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. Start speaking.', '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, "'\\''"); // Wilayah Singapura. Ganti {WorkspaceId} dengan workspace ID Bailian Anda. URL berbeda-beda tergantung wilayah. 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 generated. Copy the curl command to your terminal to get the 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('Click "Start Session" first to generate the Offer.'); const txt = answerBox.value.trim(); if (!txt) return alert("Please paste the Answer SDP"); const answerSdp = normalizeSdpForSetRemote(txt); try { await pc.setRemoteDescription({ type: 'answer', sdp: answerSdp }); setStatus('Establishing connection...', 'connecting'); } catch (e) { alert("Failed to set 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('Disconnected', ''); 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 create failed:", 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>Buka file ini di browser dan ikuti langkah-langkah berikut:
- Klik Start Session. Halaman secara otomatis menghasilkan Offer SDP dan perintah curl yang sesuai.
- Klik Copy curl Command dan jalankan di terminal Anda. Outputnya adalah Answer SDP.
- Tempel Answer SDP ke kotak teks Answer SDP, lalu klik Set Answer untuk membuat koneksi dan memulai obrolan suara.
Alur interaksi
Mode VAD
Atur session.turn_detection.type dalam session.update ke "server_vad" atau "semantic_vad" untuk mengaktifkan mode VAD. Cocok untuk skenario panggilan suara. Baik WebSocket maupun WebRTC mendukung mode VAD dengan event server yang sama; perbedaannya hanya pada cara transmisi audio dan gambar.
WebRTC hanya mendukung mode VAD dan tidak mendukung Mode Manual. Dengan WebRTC, audio ditransmisikan langsung melalui RTP tanpa mengirim event
input_audio_buffer.append; gambar ditransmisikan melalui track video tanpa dukungan untuk eventinput_image_buffer.append. Perintah kontrol dan event server ditransmisikan melalui DataChannel dengan jenis event yang sama seperti WebSocket.
Alur interaksinya sebagai berikut:
- Klien mengirim data audio. WebSocket mengirimnya melalui event input_audio_buffer.append; WebRTC mentransmisikannya secara otomatis melalui track audio (RTP) tanpa mengirim event secara manual.
- Server mendeteksi awal ucapan dan mengirim event input_audio_buffer.speech_started melalui DataChannel (WebRTC) atau WebSocket.
- Server mendeteksi akhir ucapan dan mengirim event input_audio_buffer.speech_stopped.
- Server mengirimkan buffer audio dan mengirim event input_audio_buffer.committed.
- Server mulai menghasilkan respons, mengirim event conversation.item.created dan event lainnya. Respons audio dikembalikan secara bertahap melalui event WebSocket
response.audio.delta, atau ditransmisikan langsung melalui track audio WebRTC (RTP). - Selama respons, server mengembalikan transkripsi teks bertahap melalui event
response.audio_transcript.delta, dan mengirim eventresponse.donesaat respons selesai.
| Siklus hidup | Event klien | Event server |
|---|---|---|
Inisialisasi sesi |
|
|
Input audio pengguna |
| input_audio_buffer.speech_started
input_audio_buffer.speech_stopped
|
Keluaran Audio Server | Tidak ada |
response.audio_transcript.delta
response.audio_transcript.done
|
Mode manual
Atur session.turn_detection dalam session.update ke null untuk mode manual. Klien mengirim input_audio_buffer.commit dan response.create untuk meminta respons. Cocok untuk skenario tekan-untuk-bicara seperti pesan suara di aplikasi obrolan.
Alur interaksinya sebagai berikut:
-
Klien dapat mengirim event input_audio_buffer.append dan input_image_buffer.append kapan saja untuk menambahkan audio dan gambar ke buffer.
Anda harus mengirim setidaknya satu event
input_audio_buffer.appendsebelum mengirim eventinput_image_buffer.append. -
Klien mengirim event input_audio_buffer.commit untuk mengirimkan buffer audio dan gambar, memberi sinyal ke server bahwa semua input pengguna (audio dan gambar) untuk putaran saat ini telah dikirim.
-
Server merespons dengan event input_audio_buffer.committed.
-
Klien mengirim event response.create dan menunggu server mengembalikan output model.
-
Server merespons dengan event conversation.item.created.
| Siklus hidup | Event klien | Event server |
|---|---|---|
Inisialisasi sesi |
|
|
Input audio pengguna |
|
|
Keluaran Audio Server |
|
response.audio_transcript.delta
response.audio_transcript.done
|
Pencarian web
Pencarian web memungkinkan model menggunakan data real-time untuk menjawab pertanyaan tentang informasi terkini seperti harga saham dan cuaca. Model secara otomatis menentukan apakah pencarian diperlukan.
Hanya
qwen3.5-omni-plus-realtimeyang mendukung pencarian web. Dinonaktifkan secara default; aktifkan dengansession.update.
Untuk penagihan, lihat kebijakan
agentdalam aturan penagihan.
Aktifkan pencarian web
Tambahkan parameter berikut ke event session.update:
enable_search: Atur ketrueuntuk mengaktifkan fitur pencarian web.search_options.enable_source: Atur ketrueuntuk menyertakan sumber hasil pencarian dalam respons.
Untuk parameter lainnya, lihat session.update.
Format respons
Saat pencarian web diaktifkan, objek usage dalam response.done menyertakan field plugins dengan informasi metering pencarian:
{
"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"
}
}
}
}
Contoh kode
Aktifkan pencarian web dalam percakapan real-time:
DashScope Python SDK
Teruskan parameter enable_search dan search_options dalam pemanggilan update_session:
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')
# Wilayah Singapura. Ganti {WorkspaceId} dengan workspace ID Bailian Anda. URL berbeda-beda tergantung wilayah.
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="You are Xiaoyun, a personal assistant.",
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
Dalam metode updateSession, teruskan konfigurasi pencarian web dalam argumen 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"))
// Wilayah Singapura. Ganti {WorkspaceId} dengan workspace ID Bailian Anda. URL berbeda-beda tergantung wilayah.
.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", "You are Xiaoyun, a personal assistant.",
"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)
Tambahkan field enable_search dan search_options ke payload JSON dari event session.update:
import json
import os
import websocket
import base64
import pyaudio
import threading
API_KEY = os.getenv("DASHSCOPE_API_KEY")
# Wilayah Singapura. Ganti {WorkspaceId} dengan workspace ID Bailian Anda. URL berbeda-beda tergantung wilayah.
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": "You are Xiaoyun, a personal assistant.",
"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()
Referensi API
Penagihan dan pembatasan laju
Penagihan
Penagihan berbasis token, diukur berdasarkan modalitas (audio, gambar, teks). Periksa konsol Model Studio untuk harga.
CatatanDalam percakapan real-time multi-putaran, setiap kali model menghasilkan respons, ia memproses semua konten percakapan historis dalam jendela konteks — termasuk audio, gambar, dan teks dari putaran sebelumnya — bersama dengan input baru putaran saat ini sebagai token input. Akibatnya, token input terakumulasi setiap putaran alih-alih hanya dihitung untuk input baru putaran saat ini.
Sebagai contoh, jika input audio 10 detik dikonversi menjadi 70 token (Qwen3.5-Omni-Realtime), dan audio tersebut masih berada dalam jendela konteks pada putaran 3, audio tersebut tetap dihitung terhadap token input untuk putaran 3. Token input yang ditagih sebenarnya = token dari semua putaran historis dalam jendela konteks + token dari input baru putaran saat ini.
Aturan konversi audio dan gambar ke token
Audio
-
Qwen3.5-Omni-Realtime:- Audio input:
total token = Durasi audio (detik) * 7 - Audio output:
total token = Durasi audio (detik) * 12,5
- Audio input:
-
Qwen3-Omni-Flash-Realtime:Audio input dan output menggunakan rumus yang sama:total token = Durasi audio (detik) * 12,5 -
Qwen-Omni-Turbo-Realtime:Audio input dan output menggunakan rumus yang sama:total token = Durasi audio (detik) * 25Durasi audio kurang dari 1 detik ditagih sebagai 1 detik.
Gambar
- Model seri
Qwen3.5-Omni-Realtimemenggunakan 1 token per32x32piksel - Model
Qwen3-Omni-Flash-Realtimemenggunakan 1 token per32x32piksel - Model
Qwen-Omni-Turbo-Realtimemenggunakan 1 token per28x28piksel
Gambar mengonsumsi 4–1.280 token. Gunakan kode berikut untuk memperkirakan konsumsi token berdasarkan dimensi gambar dan durasi sesi:
# Instal library Pillow dengan menjalankan: pip install Pillow
from PIL import Image
import math
# Untuk model Qwen-Omni-Turbo-Realtime, faktor penskalaan adalah 28.
# factor = 28
# Untuk model Qwen3-Omni-Flash-Realtime dan Qwen3.5-Omni-Realtime, faktor penskalaan adalah 32.
factor = 32
def token_calculate(image_path='', duration=10):
"""
:param image_path: Jalur gambar
:param duration: Durasi sesi
:return: Total token untuk gambar berdasarkan durasi sesi
"""
if len(image_path) > 0:
# Buka file gambar.
image = Image.open(image_path)
# Dapatkan dimensi asli gambar.
height = image.height
width = image.width
print(f"Dimensi gambar sebelum penskalaan: height={height}, width={width}")
# Sesuaikan tinggi menjadi kelipatan factor.
h_bar = round(height / factor) * factor
# Sesuaikan lebar menjadi kelipatan factor.
w_bar = round(width / factor) * factor
# Batas bawah untuk token gambar: 4 token.
min_pixels = factor * factor * 4
# Batas atas untuk token gambar: 1.280 token.
max_pixels = 1280 * factor * factor
# Penskalaan gambar agar sesuai dengan batas jumlah piksel.
if h_bar * w_bar > max_pixels:
# Hitung faktor penskalaan beta untuk menghindari melebihi max_pixels.
beta = math.sqrt((height * width) / max_pixels)
# Hitung ulang tinggi yang disesuaikan untuk memastikan merupakan kelipatan bilangan bulat dari factor.
h_bar = math.floor(height / beta / factor) * factor
# Hitung ulang lebar yang disesuaikan untuk memastikan merupakan kelipatan bilangan bulat dari factor.
w_bar = math.floor(width / beta / factor) * factor
elif h_bar * w_bar < min_pixels:
# Hitung faktor penskalaan beta agar jumlah piksel gambar yang diskala tidak kurang dari min_pixels.
beta = math.sqrt(min_pixels / (height * width))
# Hitung ulang tinggi yang disesuaikan untuk memastikan merupakan kelipatan bilangan bulat dari factor.
h_bar = math.ceil(height * beta / factor) * factor
# Hitung ulang lebar yang disesuaikan untuk memastikan merupakan kelipatan bilangan bulat dari factor.
w_bar = math.ceil(width * beta / factor) * factor
print(f"Dimensi gambar setelah penskalaan: height={h_bar}, width={w_bar}")
# Hitung token gambar.
token = int((h_bar * w_bar) / (factor * factor))
print(f"Jumlah token setelah penskalaan: {token}")
total_token = token * math.ceil(duration / 2)
print(f"Total token: {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)
Pembatasan laju
Untuk batas laju model, lihat Pembatasan laju.
Kode kesalahan
Jika pemanggilan model gagal dan mengembalikan pesan kesalahan, lihat Kode kesalahan untuk resolusi.
Daftar suara
Untuk daftar suara yang tersedia untuk model Qwen-Omni-Realtime, lihat Suara.