การสตรีมผลลัพธ์ใน 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 inastream_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 ในทันที — ไม่ต้องติดตั้งในเครื่องของคุณ
บทเรียนทั้งหมดในหลักสูตรนี้
- สถาปัตยกรรม LangChain และนามธรรมหลัก
- การสร้างสายงานด้วย LCEL
- สายงานแบบแยกแขนงและแบบขนาน
- การสตรีมผลลัพธ์ใน LangChain