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

การสตรีมผลลัพธ์ใน LangChain

ใช้การสตรีมโทเคนผ่านสายงาน LCEL เพื่อให้แอปพลิเคชันแสดงแต่ละคำทันทีที่มาถึง แทนที่จะรอการตอบกลับทั้งหมด ซึ่งช่วยให้ผู้ใช้รับรู้ว่าระบบตอบสนองเร็วขึ้น

การสตรีมผลลัพธ์ใน LangChain เป็นบทเรียน AI Engineering Academy ฟรีบน CoddyKit นี่คือบทเรียนที่ 4 จากทั้งหมด 4 บทเรียน คุณสามารถอ่านบทเรียนทั้งหมดด้านล่างฟรี — จากนั้นลองปฏิบัติด้วยตัวคุณเองในเบราว์เซอร์พร้อมตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 24/7 บทเรียนนี้เป็นส่วนหนึ่งของเส้นทางการเรียน AI Engineering Academy และความก้าวหน้าของคุณจะซิงค์ข้ามเว็บและแอป CoddyKit คอร์ส AI Engineering Academy มีบทเรียนทั้งหมด 4 บทเรียน

เหตุใดการส่งแบบสตรีมจึงสำคัญ

หากไม่มีการส่งแบบสตรีม ผู้ใช้จะต้องมองหน้าจอว่างเปล่าระหว่างรอให้ LLM สร้างผลลัพธ์เสร็จ ซึ่งอาจใช้เวลา 5–30 วินาทีสำหรับคำตอบยาว ๆ เมื่อใช้ การส่งแบบสตรีม โทเค็นจะแสดงขึ้นทันทีที่ถูกสร้าง ทำให้ผู้ใช้ได้รับการตอบสนองในทันทีและรู้สึกว่าระบบตอบสนองได้รวดเร็วขึ้นอย่างมาก LCEL ของ LangChain จะส่งต่อการทำงานแบบสตรีมผ่านสายโซ่ทั้งหมดโดยอัตโนมัติเมื่อคุณเรียกใช้ .stream()

การส่งแบบสตรีมพื้นฐานด้วย .stream()

สายโซ่ LCEL ทุกสายมีเมธอด .stream() ซึ่งส่งคืนตัววนซ้ำของชิ้นส่วน สำหรับสายโซ่ที่ลงท้ายด้วย StrOutputParser แต่ละชิ้นส่วนจะเป็นส่วนหนึ่งของสตริง คุณวนซ้ำผ่านชิ้นส่วนเหล่านี้แล้วพิมพ์หรือส่งต่อทันทีที่ได้รับ การส่งแบบสตรีมเกิดขึ้นในระดับ HTTP โดยโทเค็นแต่ละรายการจาก API ของ OpenAI จะถูกส่งผ่านตัวแยกวิเคราะห์ทันทีที่มาถึง

from langchain_openai import ChatOpenAI
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.output_parsers import StrOutputParser

chain = (
    ChatPromptTemplate.from_template('Explain {topic} in detail.')
    | ChatOpenAI(model='gpt-4o-mini')
    | StrOutputParser()
)

# Stream tokens to stdout
for chunk in chain.stream({'topic': 'quantum entanglement'}):
    print(chunk, end='', flush=True)
print()  # final newline

การส่งแบบสตรีมชนิดอะซิงโครนัสด้วย .astream()

.astream() เป็นเวอร์ชันอะซิงโครนัสของ .stream() โดยส่งคืนตัววนซ้ำแบบอะซิงโครนัสที่คุณใช้ด้วย async for วิธีนี้เป็นแนวทางที่ถูกต้องใน FastAPI, Starlette และกรอบงานเว็บแบบอะซิงโครนัสอื่น ๆ ซึ่งตัวจัดการคำขอเป็นโครูทีน การใช้การสตรีมแบบซิงโครนัสภายในตัวจัดการแบบอะซิงโครนัสจะทำให้ลูปเหตุการณ์ถูกบล็อก

import asyncio
from langchain_openai import ChatOpenAI
from langchain_core.output_parsers import StrOutputParser
from langchain_core.prompts import ChatPromptTemplate

chain = (
    ChatPromptTemplate.from_template('Write a poem about {subject}')
    | ChatOpenAI(model='gpt-4o-mini')
    | StrOutputParser()
)

async def stream_response():
    async for chunk in chain.astream({'subject': 'the ocean'}):
        print(chunk, end='', flush=True)

asyncio.run(stream_response())

การส่งแบบสตรีมใน FastAPI ด้วย StreamingResponse

ใน FastAPI ให้ห่อเครื่องกำเนิดแบบอะซิงโครนัสด้วย StreamingResponse และใช้ media_type='text/plain' เพื่อส่งโทเค็นข้อความไปยังเบราว์เซอร์ สำหรับเหตุการณ์ที่ส่งจากเซิร์ฟเวอร์ (SSE) ให้ใช้ media_type='text/event-stream' และจัดรูปแบบแต่ละชิ้นส่วนเป็น data: ...\n\n จากนั้นเบราว์เซอร์จะได้รับโทเค็นทันทีที่ถูกสร้าง โดยไม่ต้องรอคำตอบทั้งหมด

from fastapi import FastAPI
from fastapi.responses import StreamingResponse

app = FastAPI()

async def generate_stream(topic: str):
    async for chunk in chain.astream({'topic': topic}):
        yield chunk

@app.get('/stream')
async def stream_endpoint(topic: str):
    return StreamingResponse(
        generate_stream(topic),
        media_type='text/plain'
    )

# SSE format for frontend EventSource
async def sse_stream(topic: str):
    async for chunk in chain.astream({'topic': topic}):
        yield f'data: {chunk}\n\n'

การส่งแบบสตรีมผ่านขั้นตอนระดับกลาง

สายโซ่ LCEL จะส่งต่อการทำงานแบบสตรีมผ่านแต่ละขั้นตอนที่รองรับ StrOutputParser รองรับการสตรีมและส่งต่อชิ้นส่วนทันที อย่างไรก็ตาม ตัวแยกวิเคราะห์บางชนิด เช่น JsonOutputParser จำเป็นต้องเก็บผลลัพธ์ทั้งหมดไว้ก่อนจึงจะแยกวิเคราะห์ได้ ทำให้การสตรีมหยุดชะงัก LangChain ทำให้พฤติกรรมนี้ชัดเจน กล่าวคือ หากขั้นตอนใดไม่รองรับการสตรีม ขั้นตอนนั้นจะสะสมผลลัพธ์ไว้ก่อนส่งต่อไปยังขั้นตอนถัดไป

from langchain_core.output_parsers import JsonOutputParser

# This chain does NOT stream token by token
# JsonOutputParser must buffer the full response before parsing JSON
json_chain = (
    ChatPromptTemplate.from_template('Return JSON: {task}')
    | ChatOpenAI(model='gpt-4o-mini')
    | JsonOutputParser()  # buffers until complete
)

# But partial JSON streaming IS possible with streaming_json_parser
for partial in json_chain.stream({'task': 'list 3 colors'}):
    print(partial)  # prints partial dict as it fills in

astream_events เพื่อการควบคุมอย่างละเอียด

.astream_events() มี API การสตรีมที่ละเอียดกว่า โดยส่งเหตุการณ์สำหรับทุกขั้นตอนในสายโซ่ ไม่ใช่เฉพาะผลลัพธ์สุดท้าย แต่ละเหตุการณ์มีฟิลด์ kind (on_chain_start, on_llm_stream, on_chain_end) และข้อมูลในฟิลด์ data วิธีนี้ทำให้คุณส่งผลลัพธ์จากการเรียกเครื่องมือ การให้เหตุผลระหว่างทาง และผลลัพธ์สุดท้ายแบบแยกกันไปยังส่วนต่าง ๆ ของส่วนติดต่อผู้ใช้ได้

async def stream_with_events(question: str):
    async for event in chain.astream_events(
        {'question': question},
        version='v2'
    ):
        kind = event['event']
        if kind == 'on_llm_stream':
            chunk = event['data']['chunk'].content
            print(chunk, end='', flush=True)
        elif kind == 'on_chain_end':
            print('\n[Done]')
        elif kind == 'on_tool_start':
            print(f'\n[Tool: {event["name"]}]')

การเก็บบัฟเฟอร์ของผลลัพธ์แบบสตรีม

บางครั้งคุณต้องการทั้งส่งโทเค็นให้ผู้ใช้แบบสตรีม และเก็บคำตอบทั้งหมดไว้สำหรับการบันทึกหรือการประมวลผลเพิ่มเติม ให้ใช้ .astream() ร่วมกับตัวสะสมแบบลิสต์ หลังจบลูปให้รวมชิ้นส่วนเข้าด้วยกันเพื่อรับข้อความทั้งหมด รูปแบบนี้ทำให้คุณแสดงผลลัพธ์แบบสตรีมแบบเรียลไทม์ พร้อมจัดเก็บคำตอบทั้งหมดไว้สำหรับการวิเคราะห์ การแคช หรือการประเมินผลได้

async def stream_and_capture(question: str) -> str:
    full_response = []
    async for chunk in chain.astream({'question': question}):
        print(chunk, end='', flush=True)  # stream to user
        full_response.append(chunk)        # also collect
    print()  # newline
    complete = ''.join(full_response)
    await log_response(question, complete)  # log full text
    return complete

การส่งแบบสตรีมร่วมกับการเรียกเครื่องมือ

เมื่อโมเดลสร้างการเรียกเครื่องมือในคำตอบแบบสตรีม อาร์กิวเมนต์ของฟังก์ชันจะมาถึงเป็นส่วนย่อยของโทเค็น คุณต้องเก็บสตริงอาร์กิวเมนต์ JSON ไว้ในบัฟเฟอร์จนกว่าการเรียกเครื่องมือจะเสร็จสมบูรณ์จึงจะเรียกใช้งานได้ LangChain จัดการเรื่องนี้โดยอัตโนมัติในตัวดำเนินการเอเจนต์ แต่หากคุณกำลังสร้างลูปการสตรีมแบบกำหนดเอง คุณต้องตรวจสอบ finish_reason และสะสมส่วนย่อยของ tool_call.function.arguments

from openai import AsyncOpenAI

client = AsyncOpenAI()

async def stream_with_tools(prompt: str):
    tool_call_buffer = {}
    async with client.chat.completions.stream(
        model='gpt-4o-mini',
        messages=[{'role': 'user', 'content': prompt}],
        tools=[weather_tool_schema]
    ) as stream:
        async for chunk in stream:
            delta = chunk.choices[0].delta
            if delta.tool_calls:
                for tc in delta.tool_calls:
                    idx = tc.index
                    if idx not in tool_call_buffer:
                        tool_call_buffer[idx] = ''
                    if tc.function.arguments:
                        tool_call_buffer[idx] += tc.function.arguments

การยกเลิกและการหมดเวลาในการส่งแบบสตรีม

คำตอบแบบสตรีมที่ใช้เวลานานจำเป็นต้องรองรับการยกเลิก ใน Python แบบอะซิงโครนัส คุณสามารถยกเลิก asyncio.Task ที่ครอบการสตรีมได้ ใน FastAPI กรอบงานจะจัดการการยกเลิกเมื่อไคลเอนต์ตัดการเชื่อมต่อโดยอัตโนมัติเมื่อใช้ StreamingResponse ตั้งค่าระยะหมดเวลาผ่านพารามิเตอร์ timeout ของไคลเอนต์ OpenAI หรือห่อการสตรีมด้วย asyncio.wait_for() เพื่อยุติการทำงานเมื่อเกินระยะเวลาสูงสุด

import asyncio

async def stream_with_timeout(question: str, timeout: float = 30.0):
    async def _stream():
        async for chunk in chain.astream({'question': question}):
            yield chunk

    try:
        async for chunk in asyncio.timeout(_stream(), timeout):
            print(chunk, end='', flush=True)
    except asyncio.TimeoutError:
        print('\n[Stream timed out after 30 seconds]')
    except asyncio.CancelledError:
        print('\n[Stream cancelled by client disconnect]')

SSE ฝั่งไคลเอนต์ด้วย JavaScript

ที่ส่วนหน้า API EventSource ในตัวของเบราว์เซอร์จะรับเหตุการณ์ที่ส่งจากเซิร์ฟเวอร์ เมื่อปลายทาง FastAPI ส่งชิ้นส่วน data: token\n\n EventSource จะส่งเหตุการณ์ message สำหรับแต่ละชิ้นส่วน ให้เพิ่มโทเค็นแต่ละรายการลงใน DOM ทันทีที่มาถึง เพื่อสร้างเอฟเฟกต์การพิมพ์ สำหรับการควบคุมที่มากขึ้น fetch() ร่วมกับ response.body.getReader() จะให้คุณเข้าถึงการสตรีมได้อย่างเต็มรูปแบบ

// Frontend JavaScript (not Python)
const source = new EventSource('/stream?topic=quantum+computing');
const outputDiv = document.getElementById('output');

source.onmessage = (event) => {
    outputDiv.textContent += event.data;
};

source.onerror = () => {
    source.close();
    outputDiv.textContent += ' [done]';
};

// Alternative: fetch with ReadableStream
const response = await fetch('/stream?topic=ai');
const reader = response.body.getReader();
while (true) {
    const {done, value} = await reader.read();
    if (done) break;
    outputDiv.textContent += new TextDecoder().decode(value);
}

แนวทางปฏิบัติที่ดีสำหรับการส่งแบบสตรีม

ปฏิบัติตามแนวทางที่ดีเหล่านี้เมื่อใช้งานการส่งแบบสตรีม: ใช้ flush=True เมื่อพิมพ์ไปยัง stdout เสมอ เพื่อป้องกันการเก็บบัฟเฟอร์ ตั้งค่า stream_usage=True หากต้องการนับโทเค็นอย่างแม่นยำระหว่างการสตรีม ส่งตัวบ่งชี้สิ้นสุด data: [DONE]\n\n เมื่อจบสตรีม SSE เพื่อให้ไคลเอนต์ทราบว่าควรปิดการเชื่อมต่อเมื่อใด ทดสอบปลายทางการสตรีมด้วย curl --no-buffer เพื่อตรวจสอบว่าโทเค็นมาถึงทีละส่วน

# Complete SSE endpoint with DONE sentinel
async def sse_generator(question: str):
    try:
        async for chunk in chain.astream({'question': question}):
            # Escape any newlines in the chunk
            safe_chunk = chunk.replace('\n', ' ')
            yield f'data: {safe_chunk}\n\n'
    finally:
        yield 'data: [DONE]\n\n'

@app.get('/chat/stream')
async def chat_stream(question: str):
    return StreamingResponse(
        sse_generator(question),
        media_type='text/event-stream',
        headers={'Cache-Control': 'no-cache', 'X-Accel-Buffering': 'no'}
    )

ตรวจสอบความเข้าใจอย่างรวดเร็ว

ทดสอบความเข้าใจเกี่ยวกับผลลัพธ์แบบสตรีมใน LangChain

สรุปบทเรียน

ในบทเรียนนี้ คุณได้เรียนรู้ว่า stream() และ astream() ช่วยให้คุณวนซ้ำผ่านชิ้นส่วนโทเค็นทันทีที่ถูกสร้าง จึงไม่ต้องรอคำตอบทั้งหมดเป็นเวลานาน StreamingResponse ใน FastAPI ที่ใช้รูปแบบ SSE จะส่งโทเค็นไปยังไคลเอนต์เบราว์เซอร์แบบเรียลไทม์ และ astream_events() มีจุดเชื่อมต่อเหตุการณ์แบบละเอียดสำหรับแต่ละขั้นตอนในสายโซ่ รวมถึงการเรียกเครื่องมือและผลลัพธ์ระหว่างทาง บทถัดไปเราจะสำรวจการจัดการหน่วยความจำสำหรับการสนทนาหลายรอบ

คำถามที่พบบ่อย

บทเรียน “การสตรีมผลลัพธ์ใน LangChain” ฟรีหรือไม่

ใช่ — ข้อความเต็มของ “การสตรีมผลลัพธ์ใน LangChain” ฟรีให้อ่านที่นี่บนเว็บ เพื่อปฏิบัติแบบโต้ตอบ (ตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 24/7) และปลดล็อคส่วนที่เหลือของคอร์ส AI Engineering Academy ให้อัปเกรดเป็น CoddyKit PRO คอร์ส AI Engineering Academy มีบทเรียนทั้งหมด 4 บทเรียน

คุณจะเรียนรู้อะไรในบทเรียน “การสตรีมผลลัพธ์ใน LangChain”

ใช้การสตรีมโทเคนผ่านสายงาน LCEL เพื่อให้แอปพลิเคชันแสดงแต่ละคำทันทีที่มาถึง แทนที่จะรอการตอบกลับทั้งหมด ซึ่งช่วยให้ผู้ใช้รับรู้ว่าระบบตอบสนองเร็วขึ้น คุณปฏิบัติ AI Engineering Academy ด้วยโค้ดที่ใช้งานได้จริงที่คุณเรียกใช้โดยตรงในเบราว์เซอร์ และติวเตอร์ AI ตลอด 24/7 ตอบคำถามของคุณขณะที่คุณไปผ่านบทเรียน

คุณต้องมีประสบการณ์ก่อนที่จะเริ่มเรียน AI Engineering Academy หรือไม่

ไม่จำเป็นต้องมีประสบการณ์มาก่อน AI Engineering Academy บน CoddyKit ออกแบบมาสำหรับผู้เริ่มต้นไปจนถึงผู้เรียนขั้นสูง คุณสามารถเริ่มต้นที่นี่หรือเริ่มจากตัวแรกและเรียนด้วยความเร็วของคุณเอง นี่คือบทเรียนที่ 4 จากทั้งหมด 4 บทเรียน

บทเรียน “การสตรีมผลลัพธ์ใน LangChain” ใช้เวลานานแค่ไหน

บทเรียน CoddyKit ส่วนใหญ่ใช้เวลาประมาณ 5–10 นาที แต่ละบทเรียนจึงสั้นและเป็นแบบโต้ตอบ คุณสามารถก้าวหน้าอย่างต่อเนื่องและกลับมาเรียนต่อจากตรงที่เพิ่งหยุดบนเว็บและแอปได้เลย

ฉันเขียนและรันโค้ดในบทเรียน AI Engineering Academy นี้ได้ไหม

ได้ บทเรียน AI Engineering Academy ทุกบทมีตัวแก้ไขโค้ดในตัว คุณจึงเขียนและรันโค้ดจริงได้เลยในเบราว์เซอร์ และได้รับข้อเสนอแนะจาก AI ในทันที — ไม่ต้องติดตั้งในเครื่องของคุณ

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

  1. สถาปัตยกรรม LangChain และนามธรรมหลัก
  2. การสร้างสายงานด้วย LCEL
  3. สายงานแบบแยกแขนงและแบบขนาน
  4. การสตรีมผลลัพธ์ใน LangChain
← กลับไปที่ AI Engineering Academy