Mục lục
- Mục tiêu bài học
- Streaming — định nghĩa và cơ chế
- Vì sao streaming gần như bắt buộc cho chat
- SSE — Server-Sent Events ngắn gọn
- OpenAI streaming — stream=True
- OpenAI — lấy usage cuối stream
- Anthropic streaming — messages.stream()
- Anthropic — event-based stream
- HuggingFace — TextStreamer local
- Async streaming — AsyncOpenAI / AsyncAnthropic
- FastAPI — StreamingResponse SSE
- Streamlit — st.write_stream
- Gradio chat — built-in streaming
- Frontend JS — response.body.getReader()
- Cancel stream khi đủ
- Structured Output + streaming
- Tool use + streaming
- Error giữa chừng stream
- Buffering và reverse proxy
- LangChain streaming
- TTFT — Time To First Token
- Code Python tổng hợp
- Bài tập
Mục tiêu bài học
Sau bài này, bạn cần:
- Mô tả streaming LLM là gì, vì sao nó khác blocking và khi nào nên dùng.
- Đọc và viết được loop streaming cho OpenAI (
stream=True,chunk.choices[0].delta.content) và Anthropic (messages.stream(),text_stream, event-based). - Stream local model bằng HuggingFace
TextStreamer/TextIteratorStreamer. - Dựng FastAPI endpoint trả về
StreamingResponseSSE và consume từ frontend bằngresponse.body.getReader(). - Dùng Streamlit
st.write_streamvà Gradio chat built-in. - Lấy được token usage sau stream với
stream_options={"include_usage": True}(OpenAI) và eventmessage_delta/message_stop(Anthropic). - Cancel stream giữa chừng, xử lý error mid-stream, biết về buffering ở reverse proxy.
- Đo và báo cáo metric TTFT cho ứng dụng.
Streaming — định nghĩa và cơ chế
Inference của LLM autoregressive sinh token tuần tự — mỗi forward pass cho ra một token mới. Ở chế độ blocking, server đợi sinh đủ rồi gửi cả response trong một HTTP body. Ở chế độ streaming, server đẩy từng chunk (thường vài token, đôi khi từng token) ngay khi có, qua một HTTP connection giữ mở.
Cơ chế phía dưới của các provider lớn (OpenAI, Anthropic, Mistral, Cohere, Together, Groq...) là Server-Sent Events (SSE) — chuẩn HTTP one-way push từ server sang client. SDK Python ẩn SSE đi và biến stream thành iterator / async iterator cho dev gọi như loop bình thường.
Nội dung mỗi chunk thường gồm:
- Delta text: vài token mới sinh thêm (không phải full response tích lũy).
- Index: vị trí trong
choices(OpenAI), index của content block (Anthropic). - Metadata:
finish_reason,usage,stop_reason— chỉ xuất hiện ở chunk cuối.
Client phải tự nối delta để ra full text. Phần lớn SDK đã có helper concat sẵn.
Vì sao streaming gần như bắt buộc cho chat
- Giảm perceived latency: với prompt cho response 500-2000 token, blocking phải đợi 5-30 giây trước khi user thấy gì. Streaming hiển thị chữ đầu tiên sau ~0.3-1s — UX khác hẳn dù tổng thời gian như nhau.
- Cancellable: user thấy hướng đi sai → bấm stop, tiết kiệm token và tiền. Blocking không có cơ hội cancel sớm.
- Tỷ lệ abandon thấp hơn: chat UI không có chữ trong vài giây dễ khiến user nghĩ lỗi và rời.
- Đồng bộ tốt với thinking-out-loud: chain-of-thought, agent reasoning dài rất khó UX nếu chờ trọn.
- Server có thể release sớm: với app multi-user, streaming giúp connection ngắn hơn về mặt thực chờ.
Khi nào blocking phù hợp hơn:
- Job batch không có user xem trực tiếp.
- Output cần parse trọn JSON / Pydantic schema trước khi xử lý (xem bước 16).
- Cần đo độ dài chính xác trước khi dùng.
SSE — Server-Sent Events ngắn gọn
SSE là một chuẩn HTML5 trên HTTP/1.1, content-type text/event-stream. Frame text-based đơn giản, không cần WebSocket. Format mỗi event:
event: message
data: {"id": "...", "choices": [{"delta": {"content": "Hello"}}]}
data: {"id": "...", "choices": [{"delta": {"content": " world"}}]}
data: [DONE]
Quy ước OpenAI dùng:
- Mỗi event = một dòng
data: <JSON>, kết thúc bằng dòng trống. - Chunk cuối là
data: [DONE]để client biết stream xong. - Header response:
content-type: text/event-stream,cache-control: no-cache, thường kèmx-accel-buffering: no(nginx).
Anthropic dùng SSE có event: field rõ ràng (message_start, content_block_delta...), kèm data: JSON. SDK Anthropic parse hộ.
So với WebSocket: SSE one-way (server → client), tự reconnect, không cần upgrade handshake — gọn cho use case streaming text. WebSocket chỉ cần khi client cũng đẩy data realtime (voice, multi-agent).
OpenAI streaming — stream=True
from openai import OpenAI
client = OpenAI()
stream = client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": "Viết hàm Python sàng Eratosthenes."}],
stream=True,
)
for chunk in stream:
delta = chunk.choices[0].delta.content
if delta:
print(delta, end="", flush=True)
print()
Bốn điểm cần để ý:
stream=True→ return về một iterator chunk thay vì object response một lần.chunk.choices[0].delta.contentchứa text mới của chunk; có thể làNone(chunk đầu chứa role"assistant", chunk cuối chứafinish_reason).flush=Trueđể Python in ngay, không buffer dòng.- SDK đã ẩn parse SSE — dev chỉ lo loop và concat.
chunk.choices[0].finish_reason sẽ có ở chunk cuối: "stop", "length", "content_filter", "tool_calls". Tích lũy text để xử lý sau:
collected = []
for chunk in stream:
if chunk.choices[0].delta.content:
collected.append(chunk.choices[0].delta.content)
full_text = "".join(collected)
OpenAI — lấy usage cuối stream
Mặc định stream không trả về usage (vì chunk cuối là [DONE]). Để lấy được:
stream = client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": "Hello"}],
stream=True,
stream_options={"include_usage": True},
)
usage = None
for chunk in stream:
if chunk.choices and chunk.choices[0].delta.content:
print(chunk.choices[0].delta.content, end="", flush=True)
if chunk.usage:
usage = chunk.usage # chunk cuối có usage
print("\n", usage)
Khi bật, OpenAI gửi thêm một chunk cuối có choices=[] và usage đầy đủ (prompt_tokens, completion_tokens, total_tokens). Cần guard if chunk.choices and ... vì chunk usage không có choices.
Mọi tính cost / monitoring production nên bật include_usage; thiếu nó phải đoán cost bằng cách đếm token output qua tiktoken — sai số nhỏ nhưng không cần thiết khi đã có sẵn cờ.
Anthropic streaming — messages.stream()
from anthropic import Anthropic
client = Anthropic()
with client.messages.stream(
model="claude-opus-4-7",
max_tokens=1024,
messages=[{"role": "user", "content": "Viết hàm Python fibonacci memo."}],
) as stream:
for text in stream.text_stream:
print(text, end="", flush=True)
final = stream.get_final_message()
print("\nusage:", final.usage)
SDK Anthropic cung cấp context manager messages.stream(...):
stream.text_stream— iterator chỉ trả về text delta (không cần lọc thủ công).stream.get_final_message()— sau khi stream xong, dựng lạiMessageđầy đủ (content, stop_reason, usage) như gọi blocking.- Context manager tự đóng connection khi exit; nếu raise exception cũng cleanup an toàn.
So với stream=True kiểu OpenAI: cách của Anthropic gọn hơn cho 80% use case — chỉ cần text. Khi muốn xem từng event raw (cho tool use, thinking), dùng cách event-based ở bước 8.
Anthropic — event-based stream
with client.messages.stream(
model="claude-opus-4-7",
max_tokens=1024,
messages=[{"role": "user", "content": "Hello"}],
) as stream:
for event in stream:
if event.type == "message_start":
print("[start] id:", event.message.id)
elif event.type == "content_block_delta":
if event.delta.type == "text_delta":
print(event.delta.text, end="", flush=True)
elif event.type == "message_delta":
# usage update (output_tokens cumulative), stop_reason
print("\n[delta]", event.usage, event.delta.stop_reason)
elif event.type == "message_stop":
print("[stop]")
Các event Anthropic định nghĩa trong protocol:
message_start— bắt đầu, kèm metadataid,model,usage.input_tokens.content_block_start— bắt đầu một block (text, tool_use, thinking).content_block_delta— delta của block hiện tại;delta.typecó thể làtext_delta,input_json_delta(cho tool argument streaming),thinking_delta.content_block_stop— đóng block.message_delta— update metadata,stop_reason,usage.output_tokenstích lũy.message_stop— kết thúc stream.
Dạng event-based bắt buộc khi muốn handle tool use (input_json_delta) hoặc extended thinking realtime.
HuggingFace — TextStreamer local
Với model chạy local qua transformers, streaming thực hiện bằng TextStreamer (in trực tiếp) hoặc TextIteratorStreamer (cho iterator). Cả hai chạy model.generate() trên thread phụ và đẩy token ra theo từng bước.
from transformers import AutoModelForCausalLM, AutoTokenizer, TextStreamer
name = "Qwen/Qwen2.5-3B-Instruct"
tok = AutoTokenizer.from_pretrained(name)
model = AutoModelForCausalLM.from_pretrained(name, device_map="auto")
prompt = tok.apply_chat_template(
[{"role": "user", "content": "Một câu chuyện ngắn về AI."}],
tokenize=False, add_generation_prompt=True,
)
inputs = tok(prompt, return_tensors="pt").to(model.device)
streamer = TextStreamer(tok, skip_prompt=True, skip_special_tokens=True)
_ = model.generate(**inputs, max_new_tokens=256, streamer=streamer)
Dùng TextIteratorStreamer khi cần đẩy ra HTTP / WebSocket:
from threading import Thread
from transformers import TextIteratorStreamer
streamer = TextIteratorStreamer(tok, skip_prompt=True, skip_special_tokens=True)
gen_kwargs = dict(**inputs, max_new_tokens=256, streamer=streamer)
Thread(target=model.generate, kwargs=gen_kwargs).start()
for text in streamer:
print(text, end="", flush=True)
Nâng cấp cho serving production: vLLM, Text Generation Inference (TGI), hoặc llama.cpp server có sẵn endpoint OpenAI-compatible streaming, dùng được luôn cùng SDK OpenAI bằng cách đổi base_url.
Async streaming — AsyncOpenAI / AsyncAnthropic
FastAPI, Starlette, aiohttp đều chạy event loop async — code stream phải async để không block worker. SDK đều có biến thể Async:
from openai import AsyncOpenAI
client = AsyncOpenAI()
async def stream_openai(prompt: str):
stream = await client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": prompt}],
stream=True,
)
async for chunk in stream:
delta = chunk.choices[0].delta.content
if delta:
yield delta
from anthropic import AsyncAnthropic
aclient = AsyncAnthropic()
async def stream_anthropic(prompt: str):
async with aclient.messages.stream(
model="claude-opus-4-7",
max_tokens=1024,
messages=[{"role": "user", "content": prompt}],
) as stream:
async for text in stream.text_stream:
yield text
Hai generator này dùng được trực tiếp trong FastAPI StreamingResponse (bước 11). Lưu ý: không mix sync OpenAI() trong endpoint async def — sẽ block toàn bộ event loop suốt thời gian sinh token.
FastAPI — StreamingResponse SSE
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
from openai import AsyncOpenAI
app = FastAPI()
client = AsyncOpenAI()
async def sse_openai(prompt: str):
stream = await client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": prompt}],
stream=True,
)
async for chunk in stream:
delta = chunk.choices[0].delta.content
if delta:
yield f"data: {delta}\n\n"
yield "data: [DONE]\n\n"
@app.get("/chat")
async def chat(q: str):
return StreamingResponse(
sse_openai(q),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"X-Accel-Buffering": "no", # tắt buffer của nginx
"Connection": "keep-alive",
},
)
Ba điểm quan trọng:
- Mỗi event phải đúng format
data: ...\n\n(hai newline kết thúc). BrowserEventSourcedựa vào đó để cắt event. - Header
X-Accel-Buffering: no(nginx) vàCache-Control: no-cacheđể reverse proxy không gom chunk trước khi đẩy. - Nếu cần đẩy JSON cho từng event (để client parse rõ hơn):
yield f"data: {json.dumps({'text': delta})}\n\n".
Production thường wrap thêm: per-request id, heartbeat : keepalive\n\n mỗi 15s tránh proxy đóng idle, send event: error khi exception (xem bước 18).
Streamlit — st.write_stream
import streamlit as st
from openai import OpenAI
client = OpenAI()
def stream_text(prompt: str):
stream = client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": prompt}],
stream=True,
)
for chunk in stream:
if chunk.choices[0].delta.content:
yield chunk.choices[0].delta.content
prompt = st.chat_input("Hỏi gì đó...")
if prompt:
with st.chat_message("user"):
st.write(prompt)
with st.chat_message("assistant"):
st.write_stream(stream_text(prompt))
Streamlit (1.30+) có st.write_stream nhận trực tiếp một generator string và render progressive, hỗ trợ markdown realtime. Không cần xử lý SSE thủ công. Với multi-turn, lưu history trong st.session_state.messages như các template chat mẫu.
Lưu ý: Streamlit chạy script lại từ đầu mỗi tương tác — generator phải gọi lại từ OpenAI() mỗi lần (đã đúng trong ví dụ trên).
Gradio chat — built-in streaming
import gradio as gr
from openai import OpenAI
client = OpenAI()
def chat_fn(message, history):
history_msgs = []
for u, a in history:
history_msgs += [{"role": "user", "content": u},
{"role": "assistant", "content": a}]
history_msgs.append({"role": "user", "content": message})
stream = client.chat.completions.create(
model="gpt-4o-mini",
messages=history_msgs,
stream=True,
)
acc = ""
for chunk in stream:
delta = chunk.choices[0].delta.content
if delta:
acc += delta
yield acc # gradio chat: yield accumulated string
gr.ChatInterface(chat_fn).launch()
Gradio ChatInterface nhận function generator yield ra string tích lũy (không phải delta). Mỗi lần yield, UI replace nội dung message hiện tại. Đây là khác biệt với Streamlit (yield delta).
Gradio cũng dùng được với gr.Blocks để custom hoàn toàn UI; pattern stream tương tự.
Frontend JS — response.body.getReader()
Trong app web tự build (Next.js, Vite, vanilla), consume SSE bằng Fetch + ReadableStream gọn hơn EventSource (vì EventSource chỉ hỗ trợ GET, không gửi body POST):
async function streamChat(prompt) {
const res = await fetch("/api/chat", {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({ prompt }),
});
const reader = res.body.getReader();
const decoder = new TextDecoder();
let buffer = "";
while (true) {
const { value, done } = await reader.read();
if (done) break;
buffer += decoder.decode(value, { stream: true });
// tách theo "\n\n" — biên SSE event
const parts = buffer.split("\n\n");
buffer = parts.pop(); // phần dư chưa hoàn chỉnh
for (const part of parts) {
const line = part.replace(/^data:\s*/, "");
if (line === "[DONE]") return;
// append vào UI
document.getElementById("out").textContent += line;
}
}
}
Ba bẫy hay gặp:
- Quên
{ stream: true }trongdecoder.decode→ byte UTF-8 multi-byte (tiếng Việt) bị cắt đôi giữa chunk → ký tự lỗi. - Tách theo
\nthay vì\n\n→ một event SSE bị parse thành nhiều phần rời. - Quên giữ
bufferdư — event cuối hay bị mất.
Với Vercel AI SDK, LangChain JS, các thư viện cấp cao đã wrap parse SSE, dev chỉ cần subscribe.
Cancel stream khi đủ
User bấm Stop, hoặc app đã đủ thông tin (đã thấy đủ JSON key), nên cancel sớm để khỏi sinh thêm token tốn tiền. Phía client cần đóng connection.
Python OpenAI: stream là context manager, dùng break hoặc stream.close().
with client.chat.completions.create(
model="gpt-4o-mini", messages=msgs, stream=True,
) as stream:
for chunk in stream:
delta = chunk.choices[0].delta.content
if delta and "STOP" in delta:
break # SDK tự đóng HTTP connection khi thoát with
Python Anthropic: context manager messages.stream() cũng tự đóng khi thoát with.
Frontend: dùng AbortController.
const ctrl = new AbortController();
fetch("/api/chat", { signal: ctrl.signal, ... });
// khi user bấm Stop:
ctrl.abort();
FastAPI server: nếu client disconnect, ASGI server (uvicorn) raise asyncio.CancelledError trong generator — bắt và cleanup, hoặc để propagate. Cẩn thận: provider chỉ thực sự dừng tính tiền khi server đóng connection upstream, không chỉ là client disconnect.
Structured Output + streaming
Bài 27 (Structured Output) đã đi vào ràng buộc schema JSON / Pydantic. Khi kết hợp với streaming có hai vấn đề:
- JSON chỉ valid khi nhận đủ — parse giữa stream sẽ lỗi
JSONDecodeError. - Hiển thị JSON nửa chừng cho user không hữu ích.
Hai cách xử lý:
- Buffer toàn bộ, parse cuối — giữ streaming chỉ để giảm TTFB-server-to-client, không show UI; chỉ render UI khi nhận đủ và parse OK. Phù hợp khi UI hiển thị field rời (table, card).
- Partial JSON parser — dùng thư viện như
json-stream,partial-json-parser(npm),ijson(Python) chấp nhận JSON chưa đóng}. OpenAI có helperclient.beta.chat.completions.stream(...).response_format=...trả về object có thuộc tínhparsedcập nhật dần.
from pydantic import BaseModel
class Movie(BaseModel):
title: str
year: int
with client.beta.chat.completions.stream(
model="gpt-4o-mini",
messages=[{"role": "user",
"content": "Cho 1 movie format JSON {title, year}."}],
response_format=Movie,
) as stream:
for event in stream:
if event.type == "content.delta":
if event.parsed:
print(event.parsed) # Movie object update dần (field bổ sung từng bước)
final = stream.get_final_completion()
print(final.choices[0].message.parsed)
Nếu output dài (list nhiều object), partial parsing có ích cho UI table progressive. Output ngắn (<500 token) → buffer rồi parse cuối là đủ.
Tool use + streaming
Khi model quyết định gọi tool, argument là một JSON object. Trong stream, JSON đến từng phần (vài ký tự một).
OpenAI: mỗi chunk có chunk.choices[0].delta.tool_calls là list partial; phải accumulate theo index:
tool_calls = {} # index -> dict {id, name, arguments}
for chunk in stream:
for tc in (chunk.choices[0].delta.tool_calls or []):
slot = tool_calls.setdefault(tc.index, {"args": ""})
if tc.id:
slot["id"] = tc.id
if tc.function and tc.function.name:
slot["name"] = tc.function.name
if tc.function and tc.function.arguments:
slot["args"] += tc.function.arguments
# cuối stream: tool_calls[0]["args"] là chuỗi JSON đầy đủ
import json
args = json.loads(tool_calls[0]["args"])
Anthropic: trong event-based stream, các event content_block_delta có delta.type == "input_json_delta" với field partial_json; SDK đã cộng dồn sẵn trong get_final_message():
final = stream.get_final_message()
for block in final.content:
if block.type == "tool_use":
print(block.name, block.input) # input là dict đã parse
Với app stream tool argument ra UI realtime (hiển thị Claude đang gõ tool name), parse input_json_delta tay; còn lại để SDK cộng dồn.
Error giữa chừng stream
Stream có thể lỗi sau khi đã đẩy vài chunk:
RateLimitErrorở giữa do TPM vượt ngưỡng.- Connection drop (mạng client / load balancer timeout).
- Model overload — provider trả
errorevent và đóng stream. - Content filter — OpenAI có thể stop với
finish_reason="content_filter"giữa response.
Xử lý chuẩn:
from openai import APIError
collected = ""
try:
stream = client.chat.completions.create(
model="gpt-4o-mini", messages=msgs, stream=True,
)
for chunk in stream:
delta = chunk.choices[0].delta.content
if delta:
collected += delta
print(delta, end="", flush=True)
except APIError as e:
# collected có phần đã nhận; có thể save partial + thông báo lỗi cho user
print(f"\n[stream error] {e}; partial length: {len(collected)}")
Trong endpoint FastAPI, gửi event error SSE và cleanup:
async def safe_stream(prompt):
try:
async for delta in stream_openai(prompt):
yield f"data: {delta}\n\n"
except Exception as e:
yield f"event: error\ndata: {str(e)}\n\n"
finally:
yield "data: [DONE]\n\n"
Retry: streaming khó retry transparent vì đã đẩy chunk cho client. Pattern hay dùng: cho user thấy nút "Regenerate" thay vì retry ngầm.
Buffering và reverse proxy
Triệu chứng quen thuộc: code Python in từng chunk OK ở terminal, nhưng deploy lên server với nginx / Cloudflare thì client thấy text dồn cục, hoặc đợi đến cuối mới nhận.
Nguyên nhân tầng buffer:
- Python stdout: thiếu
flush=True. - Uvicorn/Hypercorn: thường không buffer SSE; OK mặc định.
- Nginx: bật
proxy_bufferingmặc định → gom response. Tắt bằng headerX-Accel-Buffering: notừ app, hoặcproxy_buffering off;trong nginx config cho route. - Cloudflare: tự gom với route HTTP. Cần dùng route streaming (SSE qua một domain riêng / bật streaming setting), hoặc bỏ proxy cho path đó.
- Compression: gzip/brotli buffer để nén → tắt cho
text/event-stream. - Browser dev tool: Network tab cũng có thể không update realtime với UI; xem actual chunk arrival qua
console.logtrong reader.
Heartbeat: gửi : ping\n\n mỗi 15s để LB không kill idle connection.
LangChain streaming
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage
llm = ChatOpenAI(model="gpt-4o-mini")
for chunk in llm.stream([HumanMessage("Viết haiku về Python.")]):
print(chunk.content, end="", flush=True)
LangChain (0.3+) chuẩn hoá interface stream cho mọi provider:
llm.stream(...)sync iterator.llm.astream(...)async iterator.chain.astream_events(...)trả event chi tiết của toàn pipeline (LCEL chain) — hữu ích cho agent observability.
Dùng khi: pipeline đã viết bằng LangChain, cần stream chain composite (retriever + LLM + parser). Nếu chỉ gọi 1 LLM thuần, dùng SDK gốc gọn hơn.
TTFT — Time To First Token
TTFT là khoảng từ lúc gửi request đến lúc nhận token đầu tiên. Đây là metric chính của perceived latency — quan trọng hơn cả tổng thời gian với chat UI.
import time
start = time.perf_counter()
first_token_time = None
stream = client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": "Hello"}],
stream=True,
)
n_tokens = 0
for chunk in stream:
delta = chunk.choices[0].delta.content
if delta:
if first_token_time is None:
first_token_time = time.perf_counter()
n_tokens += 1
end = time.perf_counter()
ttft = first_token_time - start
total = end - start
tps = n_tokens / (end - first_token_time) if first_token_time else 0
print(f"TTFT: {ttft*1000:.0f} ms | total: {total:.2f}s | output tps: {tps:.1f}")
Tham chiếu tốt cho user experience: TTFT < 1s là OK cho chat, < 0.5s là tốt. Yếu tố ảnh hưởng:
- Kích thước prompt input (prefill cost, tỉ lệ thuận với context dài).
- Model size — Haiku / mini có TTFT thấp hơn Opus / o-series nhiều lần.
- Region — gọi từ Việt Nam đến endpoint US thêm 150-250 ms RTT; có thể dùng Azure OpenAI region gần (Japan East) hoặc Anthropic qua Bedrock region phù hợp.
- Prompt caching (Anthropic, OpenAI cached input): cache hit giảm TTFT vài lần với prompt 10K+ token.
- Extended thinking / reasoning models có TTFT cao hơn nhiều (reasoning phase chạy trước token đầu của response visible).
Đo TTFT P50 / P95 đều đặn cho prod; report kèm cost và quality khi A/B test model.
Code Python tổng hợp
(a) OpenAI stream — print + usage:
from openai import OpenAI
client = OpenAI()
stream = client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": "Giải thích GIL ngắn gọn."}],
stream=True,
stream_options={"include_usage": True},
)
usage = None
for chunk in stream:
if chunk.choices and chunk.choices[0].delta.content:
print(chunk.choices[0].delta.content, end="", flush=True)
if chunk.usage:
usage = chunk.usage
print("\n", usage)
(b) Anthropic stream — context manager:
from anthropic import Anthropic
ac = Anthropic()
with ac.messages.stream(
model="claude-opus-4-7",
max_tokens=512,
messages=[{"role": "user", "content": "Giải thích GIL ngắn gọn."}],
) as stream:
for text in stream.text_stream:
print(text, end="", flush=True)
print("\n", stream.get_final_message().usage)
(c) FastAPI SSE endpoint:
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
from openai import AsyncOpenAI
import json
app = FastAPI()
client = AsyncOpenAI()
async def sse(prompt: str):
stream = await client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": prompt}],
stream=True,
)
async for chunk in stream:
d = chunk.choices[0].delta.content
if d:
yield f"data: {json.dumps({'text': d})}\n\n"
yield "data: [DONE]\n\n"
@app.get("/chat")
async def chat(q: str):
return StreamingResponse(
sse(q),
media_type="text/event-stream",
headers={"X-Accel-Buffering": "no", "Cache-Control": "no-cache"},
)
(d) HuggingFace local stream:
from transformers import AutoModelForCausalLM, AutoTokenizer, TextStreamer
name = "Qwen/Qwen2.5-3B-Instruct"
tok = AutoTokenizer.from_pretrained(name)
model = AutoModelForCausalLM.from_pretrained(name, device_map="auto")
prompt = tok.apply_chat_template(
[{"role": "user", "content": "Một câu chuyện ngắn về AI."}],
tokenize=False, add_generation_prompt=True,
)
inputs = tok(prompt, return_tensors="pt").to(model.device)
streamer = TextStreamer(tok, skip_prompt=True, skip_special_tokens=True)
model.generate(**inputs, max_new_tokens=200, streamer=streamer)
Bốn snippet bao trùm phần lớn use case streaming cho dev app: hai provider chính, một server endpoint, một local model.
Bài tập
- Viết script
stream_openai.pystream response cho prompt"Viết bài thơ 8 câu về Hà Nội"vớigpt-4o-mini. In từng delta, bậtstream_options.include_usage, inusagecuối. - Viết script
stream_anthropic.pyvớiclaude-haiku-4cùng prompt trên dùngmessages.stream(), inget_final_message().usage. So sánh output và usage với bài 1. - Dựng
chat.pyStreamlit chat dùngst.write_stream, giữ history multi-turn trongst.session_state, cho phép chọn provider (OpenAI / Anthropic) qua sidebar. - Viết
ttft.pyđo TTFT cho 4 model (gpt-4o-mini,gpt-4o,claude-haiku-4,claude-opus-4-7) trên 10 prompt giống nhau (mix short / long). Báo cáo bảng TTFT P50, total P50, output tokens/s. - (Tùy chọn) Dựng FastAPI endpoint
/chattrả SSE, viết HTML một file dùngfetch + getReaderrender progressive. Test qua nginx local; verifyX-Accel-Buffering: nogiải quyết hiện tượng dồn chunk.
- OpenAI — Chat Completions streaming reference
- OpenAI — Streaming responses guide
- OpenAI Cookbook — How to stream completions
- Anthropic — Streaming Messages
- Anthropic — Messages streaming events reference
- anthropic-sdk-python — GitHub
- openai-python — GitHub
- HuggingFace — TextStreamer
- HuggingFace — TextIteratorStreamer
- FastAPI — StreamingResponse
- Streamlit — st.write_stream
- Gradio — Creating a Chatbot fast
- MDN — Using server-sent events
- MDN — Using readable streams
- nginx — proxy_buffering
- LangChain — Streaming
- vLLM — OpenAI-compatible server
