0Pricing
AI Engineering Academy · บทเรียน

การรับข้อมูลแบบต่อเนื่องด้วย 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 ในทันที — ไม่ต้องติดตั้งในเครื่องของคุณ

บทเรียนทั้งหมดในหลักสูตรนี้

  1. ทำความเข้าใจการส่งโทเค็นแบบต่อเนื่อง
  2. การรับข้อมูลแบบต่อเนื่องด้วย Python SDK
  3. การส่งข้อมูลแบบต่อเนื่องใน FastAPI ด้วยเหตุการณ์ที่ส่งจากเซิร์ฟเวอร์
  4. การจัดการการเรียกใช้เครื่องมือในการตอบกลับแบบต่อเนื่อง
← กลับไปที่ AI Engineering Academy