การรับข้อมูลแบบต่อเนื่องด้วย Python SDK
ใช้ไคลเอนต์แบบอะซิงโครนัสของ OpenAI ร่วมกับ async for เพื่อรับผลลัพธ์ที่ส่งต่อเนื่อง สะสมการตอบกลับทั้งหมด และจัดการข้อผิดพลาดระหว่างการส่งโดยไม่สูญเสียผลลัพธ์บางส่วน
การรับข้อมูลแบบต่อเนื่องด้วย Python SDK เป็นบทเรียน AI Engineering Academy ฟรีบน CoddyKit นี่คือบทเรียนที่ 2 จากทั้งหมด 4 บทเรียน คุณสามารถอ่านบทเรียนทั้งหมดด้านล่างฟรี — จากนั้นลองปฏิบัติด้วยตัวคุณเองในเบราว์เซอร์พร้อมตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 24/7 บทเรียนนี้เป็นส่วนหนึ่งของเส้นทางการเรียน AI Engineering Academy และความก้าวหน้าของคุณจะซิงค์ข้ามเว็บและแอป CoddyKit คอร์ส AI Engineering Academy มีบทเรียนทั้งหมด 4 บทเรียน
ไคลเอ็นต์สตรีมแบบซิงโครนัสเทียบกับอะซิงโครนัส
OpenAI Python SDK มีทั้งไคลเอ็นต์แบบซิงโครนัส OpenAI และไคลเอ็นต์แบบอะซิงโครนัส AsyncOpenAI สำหรับสคริปต์บรรทัดคำสั่งและแอปพลิเคชันทั่วไป ไคลเอ็นต์แบบซิงโครนัสจะใช้งานได้ง่ายกว่า ส่วนเซิร์ฟเวอร์เว็บ บริการเอพีไอ และแอปพลิเคชันที่จัดการคำขอพร้อมกันหลายรายการ จำเป็นต้องใช้ ไคลเอ็นต์แบบอะซิงโครนัส เพราะจะไม่บล็อกลูปเหตุการณ์ขณะรอโทเค็น ทำให้สามารถให้บริการคำขออื่นไปพร้อมกันได้
# 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 awaitedการสตรีมแบบอะซิงโครนัสด้วย AsyncOpenAI
เมื่อใช้ไคลเอ็นต์ AsyncOpenAI การเรียกสตรีมจะกลายเป็นโครูทีน คุณใช้ async for เพื่อวนซ้ำผ่านชังก์แทนการใช้ลูป for ปกติ ลูปเหตุการณ์สามารถจัดตารางเวลาให้โครูทีนอื่นทำงานระหว่างที่ชังก์แต่ละรายการมาถึงได้ ทำให้เซิร์ฟเวอร์จัดการคำขออื่นระหว่างรอโทเค็นถัดไปจาก LLM ได้ นี่คือข้อได้เปรียบสำคัญเหนือการสตรีมแบบซิงโครนัสในบริบทเว็บ
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'))การใช้ตัวจัดการบริบทของสตรีม
OpenAI SDK ยังมี ตัวจัดการบริบทของสตรีมผ่าน client.chat.completions.stream() วิธีนี้จะปิดสตรีมโดยอัตโนมัติเมื่อออกจากบริบท และมีเมธอดอำนวยความสะดวก เช่น stream.text_stream ซึ่งให้ผลลัพธ์เฉพาะเดลตาข้อความที่ไม่ใช่ None และ stream.get_final_completion() สำหรับสถิติการใช้งานหลังจบสตรีม โดยไม่ต้องสะสมข้อมูลด้วยตนเอง
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?'))การจัดการข้อผิดพลาดระหว่างสตรีมอย่างเหมาะสม
ข้อผิดพลาดอาจเกิดขึ้นได้ทุกช่วงระหว่างการสตรีม ไม่ว่าจะเป็นขณะเชื่อมต่อครั้งแรก หลังได้รับโทเค็นแรก หรือใกล้สิ้นสุดการตอบกลับที่ยาวนาน ให้ห่อการวนซ้ำของสตรีมด้วยบล็อก try/except และจัดการ openai.APIConnectionError, openai.RateLimitError และ openai.APIStatusError แยกกัน เนื่องจากแต่ละประเภทต้องใช้แนวทางกู้คืนต่างกัน เช่น การลองใหม่ การหน่วงเวลาแบบเพิ่มขึ้น หรือการแจ้งเตือนผู้ใช้
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__}]'ตัวสร้างแบบอะซิงโครนัสสำหรับการสตรีม
รูปแบบอะซิงโครนัสที่สะอาดที่สุดสำหรับการสตรีมคือ ฟังก์ชันตัวสร้างแบบอะซิงโครนัสที่ให้ผลลัพธ์เป็นโทเค็น ผู้ใช้จะวนซ้ำผ่านฟังก์ชันนี้ด้วย async for วิธีนี้แยกตรรกะการสตรีมออกจากวิธีนำผลลัพธ์ไปใช้ โดยปลายทาง FastAPI ตัวจัดการ WebSocket และการทดสอบสามารถใช้ตัวสร้างเดียวกันได้โดยไม่ต้องรู้จักกัน
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)])คำขอสตรีมพร้อมกัน
ข้อดีสำคัญของการสตรีมแบบอะซิงโครนัสคือ ความสามารถในการเรียกใช้ สตรีมหลายรายการพร้อมกันภายในโพรเซสเดียว เมื่อใช้ asyncio.gather คุณสามารถเริ่มคำขอสตรีมไปยัง LLM หลายรายการพร้อมกัน และประมวลผลโทเค็นทันทีที่มาถึง วิธีนี้มีประโยชน์สำหรับรูปแบบการกระจายงาน เมื่อคุณต้องการเปรียบเทียบรูปแบบของพรอมต์หลายแบบ หรือเรียกใช้งานย่อยแบบขนาน
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())การหมดเวลาและการยกเลิก
สตรีมที่ทำงานเป็นเวลานานควรมี การหมดเวลา เพื่อป้องกันการบล็อกไม่สิ้นสุด ใช้ asyncio.wait_for เพื่อกำหนดการหมดเวลาในระดับโครูทีน หรือใช้ httpx.Timeout เพื่อตั้งค่าการหมดเวลาสำหรับการเชื่อมต่อและการอ่านข้อมูลในระดับไคลเอ็นต์ HTTP ทั้งสองวิธีช่วยให้สตรีมที่หยุดชะงักไม่ครองคำขอไว้อย่างไม่มีกำหนด ควรยกเลิกสตรีมอย่างชัดเจนเสมอเมื่อผู้ใช้ตัดการเชื่อมต่อ เพื่อไม่ให้สิ้นเปลืองการประมวลผลของ GPU
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)การบัฟเฟอร์บรรทัดบางส่วน
เมื่อสตรีมไปยังไคลเอ็นต์ที่ประมวลผลทีละบรรทัดที่สมบูรณ์ เช่น CLI ที่แสดงผลมาร์กดาวน์ คุณอาจต้องการ บัฟเฟอร์โทเค็นไว้จนกว่าจะพบอักขระขึ้นบรรทัดใหม่หรือขอบเขตประโยคก่อนส่งต่อ วิธีนี้ช่วยป้องกันการแสดงผลประโยคที่ยังไม่สมบูรณ์และกะพริบสะดุด ให้สะสมโทเค็นไว้ในบัฟเฟอร์ ล้างบัฟเฟอร์ส่งให้ผู้ใช้เมื่อพบเครื่องหมายวรรคตอนปิดประโยคหรืออักขระขึ้นบรรทัดใหม่ และล้างบัฟเฟอร์ส่วนที่เหลือเสมอเมื่อสตรีมสิ้นสุด
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 bufferการบันทึกความหน่วงของสตรีมในระบบจริง
ในระบบจริง ให้เพิ่มเครื่องมือวัดในทุกสตรีมเพื่อบันทึก TTFT และเวลาสร้างทั้งหมดสำหรับการตรวจติดตาม จัดเก็บเมตริกเหล่านี้ในฐานข้อมูลอนุกรมเวลา และแจ้งเตือนเมื่อ TTFT เกินค่าขีดจำกัด SLA ซึ่งโดยทั่วไปอยู่ที่ 1–2 วินาทีสำหรับแอปพลิเคชันแบบโต้ตอบ เปรียบเทียบจุดที่ TTFT พุ่งสูงกับความยาวของพรอมต์ ภาระของโมเดล และช่วงเวลาของวัน เพื่อระบุสาเหตุรากของความหน่วงที่แย่ลง
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, metricsการทดสอบโค้ดสตรีมแบบอะซิงโครนัส
การทดสอบสตรีมแบบอะซิงโครนัสต้องใช้ความระมัดระวังเป็นพิเศษ ใช้ pytest-asyncio เพื่อเรียกใช้ฟังก์ชันทดสอบแบบอะซิงโครนัส และจำลองไคลเอ็นต์ OpenAI เพื่อหลีกเลี่ยงการเรียกใช้เอพีไอจริงในการทดสอบหน่วย สร้างสตรีมจำลองที่ให้ผลลัพธ์เป็นชังก์ที่กำหนดไว้ล่วงหน้า พร้อมความล่าช้าที่ตั้งค่าได้ เพื่อทดสอบทั้งการประมวลผลโทเค็นในกรณีปกติและเส้นทางจัดการข้อผิดพลาด โดยไม่ต้องใช้โควตาเอพีไอ
# 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!'ตัวช่วย SDK: stream.text และ stream.final_message
ตัวจัดการบริบทของสตรีมใน OpenAI Python SDK มีแอตทริบิวต์ตัวช่วยที่ไม่ต้องสะสมข้อมูลด้วยตนเอง stream.text_stream เป็นตัววนซ้ำแบบอะซิงโครนัสที่ให้ผลลัพธ์เฉพาะสตริงเนื้อหาที่ไม่ใช่ None หลังสตรีมเสร็จสิ้น await stream.get_final_message() จะส่งคืน ChatCompletionMessage ฉบับเต็มพร้อมข้อความและข้อมูลการใช้งานทั้งหมด ตัวช่วยเหล่านี้ลดโค้ดที่ต้องเขียนซ้ำ และจัดการกรณีพิเศษ เช่น เดลตาว่าง โดยอัตโนมัติ
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.contentตรวจสอบความเข้าใจ
ทดสอบความเข้าใจเกี่ยวกับการสตรีมแบบอะซิงโครนัสด้วย OpenAI Python SDK จากบทเรียนนี้
สรุปบทเรียน
ในบทเรียนนี้ คุณได้เรียนรู้ว่า AsyncOpenAI ช่วยให้สตรีมทำงานโดยไม่บล็อก ทำให้เซิร์ฟเวอร์จัดการคำขอพร้อมกันได้ ตัวสร้างแบบอะซิงโครนัสเป็นรูปแบบที่สะอาดที่สุดสำหรับการส่งโทเค็นสตรีมไปยังผู้ใช้ปลายทาง และ asyncio.wait_for กับพารามิเตอร์ timeoutช่วยป้องกันไม่ให้สตรีมที่หยุดชะงักอย่างไม่มีกำหนดบล็อกเซิร์ฟเวอร์ ตัวจัดการบริบทของสตรีมมีตัวช่วยอำนวยความสะดวก เช่น text_stream และ get_final_completion บทถัดไป เราจะเปิดให้ไคลเอ็นต์เบราว์เซอร์เข้าถึงการสตรีม LLM ผ่าน FastAPI และ Server-Sent Events
คำถามที่พบบ่อย
บทเรียน “การรับข้อมูลแบบต่อเนื่องด้วย Python SDK” ฟรีหรือไม่
ใช่ — ข้อความเต็มของ “การรับข้อมูลแบบต่อเนื่องด้วย Python SDK” ฟรีให้อ่านที่นี่บนเว็บ เพื่อปฏิบัติแบบโต้ตอบ (ตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 24/7) และปลดล็อคส่วนที่เหลือของคอร์ส AI Engineering Academy ให้อัปเกรดเป็น CoddyKit PRO คอร์ส AI Engineering Academy มีบทเรียนทั้งหมด 4 บทเรียน
คุณจะเรียนรู้อะไรในบทเรียน “การรับข้อมูลแบบต่อเนื่องด้วย Python SDK”
ใช้ไคลเอนต์แบบอะซิงโครนัสของ OpenAI ร่วมกับ async for เพื่อรับผลลัพธ์ที่ส่งต่อเนื่อง สะสมการตอบกลับทั้งหมด และจัดการข้อผิดพลาดระหว่างการส่งโดยไม่สูญเสียผลลัพธ์บางส่วน คุณปฏิบัติ AI Engineering Academy ด้วยโค้ดที่ใช้งานได้จริงที่คุณเรียกใช้โดยตรงในเบราว์เซอร์ และติวเตอร์ AI ตลอด 24/7 ตอบคำถามของคุณขณะที่คุณไปผ่านบทเรียน
คุณต้องมีประสบการณ์ก่อนที่จะเริ่มเรียน AI Engineering Academy หรือไม่
ไม่จำเป็นต้องมีประสบการณ์มาก่อน AI Engineering Academy บน CoddyKit ออกแบบมาสำหรับผู้เริ่มต้นไปจนถึงผู้เรียนขั้นสูง คุณสามารถเริ่มต้นที่นี่หรือเริ่มจากตัวแรกและเรียนด้วยความเร็วของคุณเอง นี่คือบทเรียนที่ 2 จากทั้งหมด 4 บทเรียน
บทเรียน “การรับข้อมูลแบบต่อเนื่องด้วย Python SDK” ใช้เวลานานแค่ไหน
บทเรียน CoddyKit ส่วนใหญ่ใช้เวลาประมาณ 5–10 นาที แต่ละบทเรียนจึงสั้นและเป็นแบบโต้ตอบ คุณสามารถก้าวหน้าอย่างต่อเนื่องและกลับมาเรียนต่อจากตรงที่เพิ่งหยุดบนเว็บและแอปได้เลย
ฉันเขียนและรันโค้ดในบทเรียน AI Engineering Academy นี้ได้ไหม
ได้ บทเรียน AI Engineering Academy ทุกบทมีตัวแก้ไขโค้ดในตัว คุณจึงเขียนและรันโค้ดจริงได้เลยในเบราว์เซอร์ และได้รับข้อเสนอแนะจาก AI ในทันที — ไม่ต้องติดตั้งในเครื่องของคุณ
บทเรียนทั้งหมดในหลักสูตรนี้
- ทำความเข้าใจการส่งโทเค็นแบบต่อเนื่อง
- การรับข้อมูลแบบต่อเนื่องด้วย Python SDK
- การส่งข้อมูลแบบต่อเนื่องใน FastAPI ด้วยเหตุการณ์ที่ส่งจากเซิร์ฟเวอร์
- การจัดการการเรียกใช้เครื่องมือในการตอบกลับแบบต่อเนื่อง