Mengonsumsi Aliran dengan Python SDK
Gunakan klien asinkron OpenAI dengan async for untuk mengonsumsi penyelesaian yang dialirkan, mengumpulkan respons lengkap, dan menangani galat di tengah aliran tanpa kehilangan keluaran parsial.
Mengonsumsi Aliran dengan Python SDK adalah pelajaran AI Engineering Academy gratis di CoddyKit. Ini adalah pelajaran 2 dari 4. Kamu bisa membaca pelajaran lengkapnya di bawah secara gratis — lalu praktikkan langsung di browser dengan editor kode bawaan dan tutor AI 24/7. Ini adalah bagian dari jalur belajar AI Engineering Academy, dan progresmu tersinkronisasi di web dan aplikasi CoddyKit. Kursus AI Engineering Academy mencakup 4 pelajaran total.
Klien Streaming Sinkron vs Asinkron
OpenAI Python SDK menyediakan klien sinkron OpenAI dan klien asinkron AsyncOpenAI. Untuk skrip baris perintah dan aplikasi sederhana, klien sinkron lebih mudah digunakan. Untuk server web, API, dan aplikasi yang menangani beberapa permintaan secara bersamaan, klien asinkron sangat penting—klien ini tidak memblokir loop peristiwa saat menunggu token, sehingga permintaan lain dapat dilayani secara bersamaan.
# Synchronous client (simple scripts)
from openai import OpenAI
client = OpenAI()
# Asynchronous client (web servers, concurrent workloads)
from openai import AsyncOpenAI
async_client = AsyncOpenAI()
# The async client has the same API surface as the sync client
# but all methods are coroutines that must be awaitedStreaming Asinkron dengan AsyncOpenAI
Dengan klien AsyncOpenAI, pemanggilan streaming menjadi coroutine. Anda menggunakan async for untuk melakukan iterasi pada potongan, bukan perulangan for biasa. Loop peristiwa dapat menjadwalkan coroutine lain di antara kedatangan setiap potongan, sehingga server Anda dapat menangani permintaan lain sambil menunggu token berikutnya dari LLM—ini adalah keunggulan utama dibandingkan streaming sinkron dalam konteks web.
import asyncio
from openai import AsyncOpenAI
async_client = AsyncOpenAI()
async def async_stream_completion(prompt: str) -> str:
stream = await async_client.chat.completions.create(
model='gpt-4o-mini',
messages=[{'role': 'user', 'content': prompt}],
stream=True,
)
full_text = ''
async for chunk in stream:
delta = chunk.choices[0].delta.content
if delta:
print(delta, end='', flush=True)
full_text += delta
print()
return full_text
# Run the coroutine
asyncio.run(async_stream_completion('Explain what async/await does in Python'))Menggunakan Pengelola Konteks Stream
OpenAI SDK juga menyediakan pengelola konteks stream melalui client.chat.completions.stream(). Pendekatan ini secara otomatis menutup stream saat konteks berakhir dan menyediakan metode praktis seperti stream.text_stream yang hanya menghasilkan delta teks yang bukan None serta stream.get_final_completion() untuk statistik penggunaan setelah stream selesai, tanpa perlu mengumpulkannya secara manual.
from openai import AsyncOpenAI
import asyncio
async def stream_with_context_manager(prompt: str):
async with async_client.chat.completions.stream(
model='gpt-4o-mini',
messages=[{'role': 'user', 'content': prompt}],
) as stream:
# text_stream filters None deltas automatically
async for text in stream.text_stream:
print(text, end='', flush=True)
# Access final completion after stream ends
completion = await stream.get_final_completion()
print(f'\nUsage: {completion.usage}')
return completion
asyncio.run(stream_with_context_manager('What are the benefits of async I/O?'))Menangani Error di Tengah Stream dengan Baik
Error dapat terjadi kapan saja selama stream berlangsung: saat koneksi awal, setelah token pertama, atau mendekati akhir respons yang panjang. Bungkus iterasi stream dalam blok try/except dan tangani openai.APIConnectionError, openai.RateLimitError, dan openai.APIStatusError secara terpisah, karena masing-masing memerlukan strategi pemulihan yang berbeda (percobaan ulang, jeda bertahap, atau pemberitahuan kepada pengguna).
import openai
async def resilient_stream(prompt: str):
try:
stream = await async_client.chat.completions.create(
model='gpt-4o-mini',
messages=[{'role': 'user', 'content': prompt}],
stream=True,
)
accumulated = ''
async for chunk in stream:
delta = chunk.choices[0].delta.content
if delta:
accumulated += delta
yield delta # async generator
except openai.RateLimitError:
yield '[Rate limit reached — please wait and retry]'
except openai.APIConnectionError:
yield '[Connection error — check your network]'
except openai.APIStatusError as e:
yield f'[API error {e.status_code}]'
except Exception as e:
yield f'[Unexpected error: {type(e).__name__}]'Generator Asinkron untuk Streaming
Pola asinkron yang paling bersih untuk streaming adalah fungsi generator asinkron yang menghasilkan token. Konsumen melakukan iterasi terhadapnya dengan async for. Dengan demikian, logika streaming terpisah dari cara output digunakan—endpoint FastAPI, pengendali WebSocket, dan pengujian semuanya menggunakan generator yang sama tanpa perlu saling mengetahui.
from typing import AsyncGenerator
async def token_stream(
messages: list[dict],
model: str = 'gpt-4o-mini',
) -> AsyncGenerator[str, None]:
stream = await async_client.chat.completions.create(
model=model,
messages=messages,
stream=True,
)
async for chunk in stream:
delta = chunk.choices[0].delta.content
if delta:
yield delta
# Consumer 1: print to terminal
async def print_stream(messages):
async for token in token_stream(messages):
print(token, end='', flush=True)
# Consumer 2: collect to string
async def collect_stream(messages) -> str:
return ''.join([t async for t in token_stream(messages)])Permintaan Streaming Bersamaan
Salah satu manfaat utama streaming asinkron adalah kemampuan menjalankan beberapa stream secara bersamaan dalam satu proses. Dengan menggunakan asyncio.gather, Anda dapat memulai beberapa permintaan streaming LLM secara bersamaan dan memproses tokennya saat token tersebut tiba. Ini berguna untuk pola fan-out ketika Anda ingin membandingkan beberapa variasi prompt atau menjalankan beberapa subtugas secara paralel.
import asyncio
async def run_parallel_streams(queries: list[str]) -> list[str]:
async def collect(query):
messages = [{'role': 'user', 'content': query}]
return ''.join([t async for t in token_stream(messages)])
results = await asyncio.gather(*[collect(q) for q in queries])
return results
queries = [
'What is RAG?',
'What is a vector database?',
'What is BM25?',
]
async def main():
answers = await run_parallel_streams(queries)
for q, a in zip(queries, answers):
print(f'Q: {q}\nA: {a[:100]}\n')
asyncio.run(main())Batas Waktu dan Pembatalan
Stream yang berjalan lama sebaiknya memiliki batas waktu untuk mencegah pemblokiran tanpa batas. Gunakan asyncio.wait_for untuk menerapkan batas waktu pada tingkat coroutine atau httpx.Timeout untuk menetapkan batas waktu koneksi dan pembacaan pada tingkat klien HTTP. Kedua pendekatan ini memastikan stream yang macet tidak menahan permintaan tanpa batas. Selalu batalkan stream secara eksplisit saat pengguna terputus agar komputasi GPU tidak terbuang.
import asyncio
async def stream_with_timeout(messages, timeout_seconds: float = 30.0):
try:
async with asyncio.timeout(timeout_seconds):
stream = await async_client.chat.completions.create(
model='gpt-4o-mini',
messages=messages,
stream=True,
timeout=timeout_seconds, # HTTP-level timeout
)
async for chunk in stream:
delta = chunk.choices[0].delta.content
if delta:
yield delta
except asyncio.TimeoutError:
yield '\n[Stream timed out after {:.0f}s]'.format(timeout_seconds)Menyangga Baris yang Belum Lengkap
Saat melakukan streaming ke klien yang memproses baris lengkap (seperti CLI yang merender markdown), Anda mungkin ingin menyangga token hingga karakter baris baru atau batas kalimat sebelum meneruskannya. Cara ini mencegah tampilan kalimat yang belum lengkap berkedip-kedip. Kumpulkan token dalam penyangga, kosongkan penyangga ke konsumen saat Anda mendeteksi tanda baca akhir kalimat atau karakter baris baru, dan selalu kosongkan sisa penyangga di akhir stream.
async def buffered_line_stream(messages):
buffer = ''
flush_on = {'.', '!', '?', '\n'}
async for token in token_stream(messages):
buffer += token
if any(c in buffer for c in flush_on):
# Find the last sentence-ending position
for i, c in enumerate(reversed(buffer)):
if c in flush_on:
split_pos = len(buffer) - i
yield buffer[:split_pos]
buffer = buffer[split_pos:]
break
if buffer: # flush remainder
yield bufferMencatat Latensi Stream di Lingkungan Produksi
Di lingkungan produksi, lakukan instrumentasi pada setiap stream untuk mencatat TTFT dan total waktu pembuatan guna pemantauan. Simpan metrik ini dalam basis data deret waktu dan buat peringatan saat TTFT melampaui ambang SLA Anda (biasanya 1–2 detik untuk aplikasi interaktif). Hubungkan lonjakan TTFT dengan panjang prompt, beban model, dan waktu dalam sehari untuk mengidentifikasi akar penyebab penurunan latensi.
import time
from dataclasses import dataclass
@dataclass
class StreamMetrics:
prompt_chars: int
ttft_ms: float
total_ms: float
token_count: int
async def instrumented_stream(messages) -> tuple[str, StreamMetrics]:
t_start = time.perf_counter()
t_first = None
token_count = 0
full_text = ''
stream = await async_client.chat.completions.create(
model='gpt-4o-mini', messages=messages, stream=True
)
async for chunk in stream:
delta = chunk.choices[0].delta.content
if delta:
if t_first is None:
t_first = time.perf_counter()
token_count += 1
full_text += delta
t_end = time.perf_counter()
prompt_len = sum(len(m.get('content', '')) for m in messages)
metrics = StreamMetrics(
prompt_chars=prompt_len,
ttft_ms=(t_first - t_start) * 1000 if t_first else 0,
total_ms=(t_end - t_start) * 1000,
token_count=token_count,
)
return full_text, metricsMenguji Kode Streaming Asinkron
Menguji streaming asinkron memerlukan perhatian khusus. Gunakan pytest-asyncio untuk menjalankan fungsi pengujian asinkron, dan buat tiruan klien OpenAI agar tidak terjadi pemanggilan API nyata dalam pengujian unit. Buat stream palsu yang menghasilkan potongan yang telah ditentukan dengan penundaan yang dapat dikonfigurasi untuk menguji pemrosesan token pada jalur normal maupun jalur penanganan error tanpa menghabiskan anggaran API.
# pip install pytest pytest-asyncio
import pytest
from unittest.mock import AsyncMock, MagicMock
async def fake_stream(tokens: list[str]):
for token in tokens:
chunk = MagicMock()
chunk.choices[0].delta.content = token
yield chunk
@pytest.mark.asyncio
async def test_stream_accumulates_correctly(monkeypatch):
mock_create = AsyncMock(return_value=fake_stream(['Hello', ', ', 'world', '!']))
monkeypatch.setattr(async_client.chat.completions, 'create', mock_create)
result = await collect_stream([{'role': 'user', 'content': 'Hi'}])
assert result == 'Hello, world!'Pembantu SDK: stream.text dan stream.final_message
Pengelola konteks stream milik OpenAI Python SDK menyediakan atribut pembantu yang menghindarkan Anda dari pengumpulan manual. stream.text_stream adalah iterable asinkron yang hanya menghasilkan string konten yang bukan None. Setelah stream selesai, await stream.get_final_message() mengembalikan ChatCompletionMessage lengkap dengan teks penuh dan data penggunaan. Pembantu ini mengurangi kode berulang dan secara otomatis menangani kasus khusus seperti delta kosong.
async def clean_streaming_example(prompt: str):
async with async_client.chat.completions.stream(
model='gpt-4o-mini',
messages=[{'role': 'user', 'content': prompt}],
) as stream:
# Iterate only over text tokens, None deltas filtered automatically
async for text in stream.text_stream:
print(text, end='', flush=True)
# After context exit, get accumulated result
final = await stream.get_final_completion()
return final.choices[0].message.contentPemeriksaan Singkat
Uji pemahaman Anda tentang streaming asinkron dengan OpenAI Python SDK dari pelajaran ini.
Ringkasan Pelajaran
Dalam pelajaran ini Anda telah mempelajari: AsyncOpenAI memungkinkan streaming tanpa pemblokiran sehingga server dapat menangani permintaan secara bersamaan; generator asinkron adalah pola paling bersih untuk menghasilkan token streaming bagi konsumen hilir; dan asyncio.wait_for serta parameter timeout mencegah stream yang macet tanpa batas memblokir server Anda. Pengelola konteks stream menyediakan pembantu praktis seperti text_stream dan get_final_completion. Selanjutnya, kita akan mengekspos streaming LLM kepada klien peramban melalui FastAPI dan Server-Sent Events.
Pertanyaan yang Sering Diajukan
Apakah pelajaran “Mengonsumsi Aliran dengan Python SDK” gratis?
Ya — teks lengkap “Mengonsumsi Aliran dengan Python SDK” gratis dibaca di sini di web. Untuk praktiknya secara interaktif (editor kode bawaan dan tutor AI 24/7) dan buka sisa kursus AI Engineering Academy, upgrade ke CoddyKit PRO. Kursus AI Engineering Academy mencakup 4 pelajaran total.
Apa yang akan aku pelajari di “Mengonsumsi Aliran dengan Python SDK”?
Gunakan klien asinkron OpenAI dengan async for untuk mengonsumsi penyelesaian yang dialirkan, mengumpulkan respons lengkap, dan menangani galat di tengah aliran tanpa kehilangan keluaran parsial. Kamu berlatih AI Engineering Academy dengan kode praktik yang langsung kamu jalankan di browser, dan tutor AI 24/7 menjawab pertanyaanmu saat kamu mengerjakan pelajaran ini.
Apakah aku perlu pengalaman untuk memulai AI Engineering Academy?
Tidak diperlukan pengalaman sebelumnya. AI Engineering Academy di CoddyKit dirancang untuk pemula hingga pelajar tingkat lanjut, jadi kamu bisa memulai di sini atau dari awal dan belajar sesuai kecepatan kamu sendiri. Ini adalah pelajaran 2 dari 4.
Berapa lama pelajaran “Mengonsumsi Aliran dengan Python SDK” memakan waktu?
Sebagian besar pelajaran CoddyKit memakan waktu sekitar 5–10 menit. Setiap pelajaran ringkas dan interaktif, jadi kamu membuat kemajuan stabil dan melanjutkan dari tempat kamu tinggalkan di web dan aplikasi.
Bisakah aku menulis dan menjalankan kode dalam pelajaran AI Engineering Academy ini?
Ya. Setiap pelajaran AI Engineering Academy menyertakan editor kode bawaan, jadi kamu menulis dan menjalankan kode nyata langsung di browser dan mendapatkan umpan balik AI instan — tidak diperlukan penyiapan lokal.
Semua pelajaran dalam kursus ini
- Memahami Pengaliran Token
- Mengonsumsi Aliran dengan Python SDK
- Pengaliran di FastAPI dengan Peristiwa yang Dikirim Server
- Menangani Pemanggilan Alat dalam Respons yang Dialirkan