ConverseStream — respons streaming
ConverseStream mengubah satu respons besar jadi aliran event. Bentuk eventnya perlu dihafal, karena semua UI chat dibangun di atasnya.
Intisari
- Butuh izin
bedrock:InvokeModelWithResponseStream— berbeda dari izin Converse biasa. - Responsnya iterator event, bukan dict tunggal. Teksnya ada di
contentBlockDelta. metadatadatang di akhir aliran — di situlahusageberada.- Error bisa muncul di tengah aliran setelah HTTP 200; tangani
internalServerExceptiondkk. - Streaming memperbaiki waktu ke token pertama, bukan waktu total.
Bentuk dasarnya
resp = runtime.converse_stream(
modelId=MODEL_ID,
system=[{"text": "Jawab ringkas dan berbahasa Indonesia."}],
messages=[{"role": "user", "content": [{"text": "Jelaskan RAG."}]}],
inferenceConfig={"maxTokens": 1024},
)
for peristiwa in resp["stream"]:
if "contentBlockDelta" in peristiwa:
print(peristiwa["contentBlockDelta"]["delta"]["text"], end="", flush=True)
Semua jenis event
| Event | Kapan muncul | Isinya |
|---|---|---|
messageStart | Sekali di awal | role |
contentBlockStart | Tiap blok baru dimulai | Indeks blok; untuk tool, nama tool |
contentBlockDelta | Berkali-kali | Potongan teks — ini yang kamu tampilkan |
contentBlockStop | Blok selesai | Indeks blok |
messageStop | Sekali | stopReason |
metadata | Terakhir | usage dan metrics |
Konsekuensi penting untuk pencatatan biaya: jumlah token tidak diketahui sampai
event metadata tiba di akhir. Kalau kamu memutus koneksi lebih awal atau lupa menghabiskan
iterator, kamu kehilangan angka biayanya — sementara tagihannya tetap jalan.
Penanganan lengkap
import logging
log = logging.getLogger(__name__)
def alirkan(pertanyaan: str, model_id: str):
resp = runtime.converse_stream(
modelId=model_id,
messages=[{"role": "user", "content": [{"text": pertanyaan}]}],
inferenceConfig={"maxTokens": 1024},
)
alasan_berhenti = None
for peristiwa in resp["stream"]:
if "contentBlockDelta" in peristiwa:
delta = peristiwa["contentBlockDelta"]["delta"]
if "text" in delta:
yield delta["text"]
elif "messageStop" in peristiwa:
alasan_berhenti = peristiwa["messageStop"]["stopReason"]
elif "metadata" in peristiwa:
pakai = peristiwa["metadata"]["usage"]
log.info("selesai", extra={
"token_masuk": pakai["inputTokens"],
"token_keluar": pakai["outputTokens"],
"alasan_berhenti": alasan_berhenti,
})
# Error yang tiba DI TENGAH aliran — HTTP-nya sudah terlanjur 200
elif "internalServerException" in peristiwa:
raise RuntimeError(peristiwa["internalServerException"]["message"])
elif "modelStreamErrorException" in peristiwa:
raise RuntimeError(peristiwa["modelStreamErrorException"]["message"])
elif "throttlingException" in peristiwa:
raise RuntimeError("throttled di tengah aliran")
Blok error di bagian bawah itu yang sering dilupakan. Karena status HTTP sudah 200 sejak awal,
try/except ClientError saja tidak menangkap kegagalan yang terjadi setelah aliran
dimulai. Aplikasi yang tidak menanganinya akan menampilkan jawaban terpotong tanpa satu pun pesan error.
Menyambung ke frontend
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
app = FastAPI()
@app.post("/tanya")
def tanya(pertanyaan: str):
def sse():
for potongan in alirkan(pertanyaan, MODEL_ID):
yield f"data: {potongan}\n\n"
yield "data: [SELESAI]\n\n"
return StreamingResponse(sse(), media_type="text/event-stream")
| Jalur ke browser | Catatan |
|---|---|
| Lambda + Lambda Web Adapter | Paling sederhana untuk Python; lihat materi Fase 1 |
| API Gateway WebSocket | Dua arah, cocok untuk chat yang bisa dibatalkan |
| ECS Fargate + FastAPI | Tanpa batasan runtime, tapi ada server yang harus dijaga |
Latihan: ubah fungsi tanya() Fase 3 sebelumnya jadi versi streaming yang mencetak ke
terminal token demi token, dan tetap mencatat jumlah token di akhir. Lalu ukur selisih waktu ke karakter
pertama antara versi biasa dan versi streaming untuk pertanyaan yang jawabannya panjang.
Rangkuman ini sengaja dipangkas ke bagian yang dipakai di roadmap. Buka sumber aslinya saat kamu butuh detail lengkap atau referensi parameter.