LLM Streaming 背压工程实践:慢客户端、断流重连与服务端缓冲区的生产设计
LLM Streaming èåå·¥ç¨å®è·µï¼æ ¢å®¢æ·ç«¯ãææµéè¿ä¸æå¡ç«¯ç¼å²åºçç产设计
AINativeè½¯ä»¶å·¥ç¨ 2026-08-07 0 é 读8åéä½ ç LLM åºç¨æ¯ç§å¾å®¢æ·ç«¯æ¨ 50 个 tokenï¼ä½ç¨æ·çææºæµè§å¨å·²ç»ç¡çäºââä½ çæå¡å¨è¿å¨å»å»å°æè¿äº token å¾ä¸ä¸ªæ²¡äººæ¶è´¹çç¼å²åºéå¡ãè¿ç¯æç« 讲çå°±æ¯è¿ä¸ªé®é¢ï¼ä»¥åæä¹å¨ç产ä¸çæ£è§£å³å®ã
èæ¯ï¼SSE Streaming çç产ç°å®
2025 年以åï¼å¤§å¤æ° LLM åºç¨çæ¶ææ¯è¿æ ·çï¼
ç¨æ· â ä½ çå端 â LLM APIï¼çå¾
宿´åå¤ï¼â 䏿¬¡æ§è¿å
è¿ç§æ¨¡å¼æä¸¤ä¸ªææ¾é®é¢ï¼é¦å±çå¾ æ¶é´ï¼TTFTï¼é¿ï¼ç¨æ·ä½éªå·®ï¼èä¸å½åå¤å¾é¿æ¶ï¼æ´ä½å»¶è¿ä¼çº¿æ§å¢é¿ã
äºæ¯å¤§å®¶é½åå°äºæµå¼ï¼
ç¨æ· â ä½ çå端 â LLM APIï¼SSE/streamï¼â é token 转å
ä½åå®ä¹åï¼ä¸æ¹æ°çç产é®é¢åºç°äºï¼å å卿¶¨ï¼ææ¶å涨å¾å¾å¿«ï¼å¶åç客æ·ç«¯æè¿ä¹åæå¡ç«¯è¿å¨ç»§ç»æ¶è tokenï¼é¨åç¨æ·åæ ææ¶ååå¤çªç¶ä¸æç¶åéè¿ï¼Nginx ä¹åçé¨ç½²çè³å¯¼è´æµå¼å®å ¨å¤±æï¼éåæäºå®æ´ååºä¸æ¬¡æ§è¿åã
è¿äºé®é¢ç»ä¸å«å SSE Streaming çèåé®é¢ï¼backpressureï¼ãæ¬æä»å·¥ç¨å®è·µåºåï¼éä¸è§£æèåçåå åç产解æ³ã
第ä¸ä¸ªé·é±ï¼Node.js ç res.write() 没æèå
å¦æä½ ç¨ Node.js + Express 转å LLM çæµå¼ååºï¼æèªç¶çåæ³æ¯è¿æ ·ï¼
// â çèµ·æ¥æ²¡é®é¢ï¼å®é
䏿²¡æèåä¿æ¤
app.post('/chat', async (req, res) => {
res.setHeader('Content-Type', 'text/event-stream');
res.setHeader('Cache-Control', 'no-cache');
res.setHeader('Connection', 'keep-alive');
const stream = await client.chat.completions.create({
model: 'deepseek-chat',
messages: req.body.messages,
stream: true,
});
for await (const chunk of stream) {
const content = chunk.choices[0]?.delta?.content || '';
if (content) {
res.write(`data: ${JSON.stringify({ text: content })}\n\n`);
}
}
res.end();
});
é®é¢åºå¨åªéï¼
res.write() å¨ Node.js éçè¡ä¸ºæ¯ï¼ææ°æ®åå
¥å
é¨ socket çåç¼å²åºï¼ç¶åç«å³è¿åãå³ä½¿å®¢æ·ç«¯æ²¡æå¨æ¶è´¹ï¼å®ä¹ä¼ç»§ç»å¾ç¼å²åºéå ãNode.js ç http.ServerResponse ç»§æ¿èª stream.Writableï¼ä½å¨ç´æ¥è°ç¨ res.write() æ¶ï¼å¹¶ä¸ä¼èªå¨è§¦åèåæåæºå¶ã
ä½ å¯ä»¥ç¨ res.socket.bufferSize æ¥è§å¯è¿ä¸ç¹ï¼
// å®éªï¼è§å¯æ
¢å®¢æ·ç«¯æ¶ç¼å²åºçå¢é¿
setInterval(() => {
if (res.socket) {
console.log(`Socket buffer: ${res.socket.bufferSize} bytes`);
}
}, 500);
å¨ç产ä¸ï¼å½å®¢æ·ç«¯æ¯ä¸ä¸ªæ ¢é 4G è¿æ¥æè ææºè¿å ¥åå°æ¶ï¼è¿ä¸ªç¼å²åºä¼ä»¥ LLM ç token çæé度ï¼é常 50-80 tokens/secï¼æ¯ä¸ª event 约 100-200 bytesï¼æç»å¢é¿ã妿䏿¬¡ä¼è¯æç» 30 ç§ï¼ç´¯ç§¯çº¦ 150KB-240KBãå¦æä½ çæå¡æ 1000 ä¸ªå¹¶åæ ¢å®¢æ·ç«¯ï¼è¿å°±æ¯ 150MB-240MB çé¢å¤å åååï¼èä¸å ¨æ¯"ææä½æ²¡äººæ¶è´¹"çæ°æ®ã
æ£ç¡®çæ£æµæ¹å¼
const BUFFER_THRESHOLD = 256 * 1024; // 256KB
async function streamWithBackpressure(
stream: AsyncIterable<any>,
res: Response
): Promise<void> {
for await (const chunk of stream) {
const content = chunk.choices[0]?.delta?.content || '';
if (!content) continue;
const data = `data: ${JSON.stringify({ text: content })}\n\n`;
// æ£æ¥å®¢æ·ç«¯æ¯å¦è¿å¨æ¶è´¹
const bufferSize = (res as any).socket?.bufferSize ?? 0;
if (bufferSize > BUFFER_THRESHOLD) {
// æ¹æ¡ Aï¼çå¾
drain äºä»¶ï¼ææ¸©åï¼
await waitForDrain(res);
// æ¹æ¡ Bï¼ç´æ¥æå¼ï¼ææ¿è¿ï¼éåæéè¿æºå¶æ¶ï¼
// res.end(); break;
}
res.write(data);
}
res.end();
}
function waitForDrain(res: Response): Promise<void> {
return new Promise((resolve) => {
const socket = (res as any).socket;
if (!socket || socket.bufferSize === 0) {
resolve();
return;
}
socket.once('drain', resolve);
// è¶
æ¶ä¿æ¤ï¼æå¤ç 5 ç§
setTimeout(resolve, 5000);
});
}
第äºä¸ªé·é±ï¼å®¢æ·ç«¯æå¼ï¼æå¡ç«¯è¿å¨ç§ Token
è¿æ¯æå®¹æè¢«å¿½è§çææ¬é®é¢ãå½ç¨æ·å ³é页颿巿°ï¼SSE è¿æ¥æå¼ï¼ä½ä½ çæå¡ç«¯å¯è½è¿å¨ï¼
- æç»ä» LLM API æå tokenï¼è±é±ï¼
- å¾ä¸ä¸ªå·²ç»å ³éç socket é writeï¼Node.js 伿é误æåºæéé»ä¸¢å¼ï¼
- 妿䏿¸¸æ¯é¿æè模åï¼å¦ DeepSeek-R1ãQwen-Plusï¼å¼å¯æ·±åº¦æè模å¼ï¼ï¼ï¼è¿ä¸ªæµªè´¹å¯ä»¥é«è¾¾æ´ä¸ª token é¢ç®
æ£ç¡®çåæ³æ¯ï¼æ£æµå®¢æ·ç«¯æå¼ â ç«å³ abort 䏿¸¸ LLM 请æ±ã
app.post('/chat', async (req, res) => {
res.setHeader('Content-Type', 'text/event-stream');
res.setHeader('Cache-Control', 'no-cache');
res.setHeader('Connection', 'keep-alive');
// â
å建 AbortControllerï¼å
³èå°å®¢æ·ç«¯è¿æ¥
const abortController = new AbortController();
// 客æ·ç«¯æå¼æ¶è§¦å
req.on('close', () => {
console.log('[stream] client disconnected, aborting upstream');
abortController.abort();
});
try {
const stream = await client.chat.completions.create({
model: 'deepseek-chat',
messages: req.body.messages,
stream: true,
}, {
signal: abortController.signal, // â
ä¼ å
¥ abort signal
});
for await (const chunk of stream) {
if (abortController.signal.aborted) break;
const content = chunk.choices[0]?.delta?.content || '';
if (content) {
const ok = res.write(`data: ${JSON.stringify({ text: content })}\n\n`);
if (!ok) {
// write è¿å false 说æå
é¨ç¼å²åºå·²æ»¡ï¼çå¾
drain
await new Promise(resolve => res.once('drain', resolve));
}
}
}
} catch (err: any) {
if (err.name !== 'AbortError') {
console.error('[stream] error:', err.message);
// åªæå¨è¿æ¥è¿æ´»çæ¶æåé误äºä»¶
if (!res.writableEnded) {
res.write(`data: ${JSON.stringify({ error: err.message })}\n\n`);
}
}
} finally {
if (!res.writableEnded) {
res.end();
}
}
});
注æ res.write() çè¿åå¼ï¼å½å
é¨ç¼å²åºæ»¡æ¶ï¼å®è¿å falseï¼å¹¶ä¼å¨ socket æ¸
空å触å drain äºä»¶ãè¿æ¯ Node.js stream.Writable çæ åèååè®®ï¼ä½å¾å¤äººå SSE æ¶ç´æ¥å¿½ç¥äºè¿ä¸ªè¿åå¼ã
宿µæ°æ®
卿们çç产系ç»éï¼å¼å ¥ AbortController ä¹åï¼
- ç¨æ·æåå ³é页é¢çåºæ¯ï¼çº¦å ææä¼è¯ç 12%ï¼ï¼ä¸æ¸¸ token æ¶èåå°äº 47%
- æåº¦ API è´¦åéä½çº¦ 8%ï¼ä¸å¤¸å¼ ï¼å 为 12% à 47% à é¿åå¤çæ¯ä¾å èµ·æ¥å¾å¯è§ï¼
第ä¸ä¸ªé·é±ï¼Nginx æ SSE åæäºæ¹éååº
è¿æ¯æéè½çåãä½ å¨æ¬å°å¼åæ¶æµå¼å·¥ä½æ£å¸¸ï¼ä¸ä¸ Nginx 代çå°±åæäº"çææ token çæå®æè¿å"ï¼æè åºç°å¥æªçåæ®µç°è±¡ã
åå ï¼Nginx ç proxy_buffering é»è®¤æ¯ onã
å½ proxy_buffering on æ¶ï¼Nginx 伿䏿¸¸ï¼ä½ ç Node.jsï¼çååºç¼å²å°å
å/ç£çï¼çå°ååºå®æåå䏿¬¡æ§åç»å®¢æ·ç«¯ãå¯¹äº SSE èè¨ï¼è¿æå³çææ token é½å¾çå° LLM 说 [DONE]ï¼æä¼æå
ååºå»ã
# â é»è®¤é
ç½®ï¼ç ´å SSE 宿¶æ§
location /chat {
proxy_pass http://backend:3000;
# proxy_buffering é»è®¤ on
}
# â
æ£ç¡®é
ç½®
location /chat {
proxy_pass http://backend:3000;
# å
³é®ï¼å
³é代çç¼å²
proxy_buffering off;
proxy_cache off;
# 让 Nginx ç«å³è½¬åï¼ä¸çæ¢è¡
proxy_read_timeout 600s; # LLM 请æ±å¯è½å¾é¿
proxy_connect_timeout 60s;
proxy_send_timeout 600s;
# å
³é Nginx ç gzip å缩ï¼ä¼ç¼å²ï¼
gzip off;
# SSE å¿
éç headers
add_header Cache-Control no-cache;
add_header X-Accel-Buffering no; # åè¯ Nginxï¼å CDNï¼ä¸è¦ç¼å²
# æ¯æ HTTP/1.1ï¼SSE å¿
é¡»ï¼
proxy_http_version 1.1;
proxy_set_header Connection '';
}
X-Accel-Buffering: no è¿ä¸ª header ä¸åªå¯¹ Nginx ææï¼å¾å¤ CDNï¼å¦ CloudflareãAWS CloudFrontï¼ä¹ä¼è¯»å®æ¥å³å®æ¯å¦ç¼å²ååºãå¨ä½ ç Node.js æå¡éä¹åºè¯¥ä¸»å¨è®¾ç½®ï¼
res.setHeader('X-Accel-Buffering', 'no');
res.setHeader('Content-Type', 'text/event-stream');
res.setHeader('Cache-Control', 'no-cache, no-store');
第å个é·é±ï¼SSE éè¿ä¸ last-event-id çå
SSE æä¸ä¸ªå
ç½®çéè¿æºå¶ï¼å½è¿æ¥æå¼æ¶ï¼æµè§å¨ä¼èªå¨å°è¯éè¿ï¼é»è®¤ 3 ç§åï¼ï¼å¹¶å¸¦ä¸ Last-Event-ID headerï¼å¼æ¯ä¸ä¸æ¬¡æ¶å°çæåä¸ä¸ªäºä»¶ç id åæ®µã
è¿çèµ·æ¥å¾ç¾å¥½ï¼ä½æå 个åï¼
å 1ï¼ä½ 没æç» event å id
// â æ²¡æ idï¼æµè§å¨éè¿æ¶ Last-Event-ID 为空
res.write(`data: ${JSON.stringify({ text: token })}\n\n`);
// â
å ä¸éå¢ id
let eventId = 0;
res.write(`id: ${eventId++}\ndata: ${JSON.stringify({ text: token })}\n\n`);
å 2ï¼æå¡ç«¯æ²¡æå¤ç Last-Event-ID
å³ä½¿ä½ å äº idï¼å¦ææå¡ç«¯ä¸å¤çéè¿è¯·æ±ï¼ç¨æ·è¿æ¯ä¼çå°å 容ä»å¤´å¼å§éæ¾ï¼æè 丢失ä¸é´çå 容ã
app.post('/chat', async (req, res) => {
const lastEventId = parseInt(req.headers['last-event-id'] as string) || 0;
// ä»ç¼å䏿¾å°å¯¹åºçåæ¾èµ·ç¹
const sessionId = req.headers['x-session-id'] as string;
const cached = await tokenCache.get(sessionId);
if (cached && lastEventId > 0) {
// éæ¾ lastEventId ä¹åç tokens
const replayTokens = cached.tokens.slice(lastEventId);
for (const token of replayTokens) {
res.write(`id: ${cached.startId + replayTokens.indexOf(token)}\ndata: ${JSON.stringify({ text: token })}\n\n`);
}
// 妿åå§æµå·²å®æï¼ç´æ¥ç»æ
if (cached.done) {
res.end();
return;
}
}
// ... ç»§ç»æ£å¸¸æµå¼
});
å 3ï¼éè¿æºå¶ä¸ AbortController ç交äº
å½ä½ ç¨ AbortController å¨å®¢æ·ç«¯æå¼æ¶ä¸æ¢ä¸æ¸¸è¯·æ±ï¼æµè§å¨çèªå¨éè¿ä¼è§¦åä¸ä¸ªæ°è¯·æ±ãè¿ä¸ªæ°è¯·æ±éè¦ä¸ä¸ªæ°ç LLM 请æ±ï¼ä»æç¹æé头å¼å§ï¼ï¼ä¸è½å¤ç¨å·²ç» abort çæµã
æ£ç¡®çæ¶ææ¯ï¼çä¼è¯ token ç¼å + éè¿çªå£ã
// ç产级å«ç SSE ä¼è¯ç®¡ç
interface StreamSession {
tokens: string[]; // å·²çæçææ tokens
done: boolean; // æ¯å¦å·²å®æ
createdAt: number; // å建æ¶é´æ³
abortController: AbortController;
}
const sessions = new Map<string, StreamSession>();
// æ¯ 5 å鿏
çè¿æä¼è¯
setInterval(() => {
const now = Date.now();
for (const [id, session] of sessions) {
if (now - session.createdAt > 5 * 60 * 1000) {
sessions.delete(id);
}
}
}, 60 * 1000);
第äºä¸ªé·é±ï¼HTTP/2 vs HTTP/1.1 çè¡ä¸ºå·®å¼
å¦æä½ çæå¡åæ¶æ¯æ HTTP/2 å HTTP/1.1ï¼å¾å¤ CDN ä¼å¼ºå¶å级ï¼ï¼è¦æ³¨æï¼
- HTTP/1.1 SSEï¼è¿æ¥æ¯æä¹ çï¼æµè§å¨æ 6 个并åè¿æ¥éå¶ï¼per originï¼
- HTTP/2 SSEï¼ç论ä¸å¯ä»¥å¤è·¯å¤ç¨ï¼ä½å®é ä¸å¤§å¤æ°æµè§å¨å¯¹ SSE è¿æ¯ç¨ç¬ç«æµ
æ´éè¦çæ¯ HTTP/2 ç flow controlï¼HTTP/2 å¨åè®®å±å°±æçæ£çèåæºå¶ï¼WINDOW_UPDATE 帧ï¼ã彿¥æ¶çªå£æ»¡æ¶ï¼åéæ¹å¿
é¡»æåãè¿æå³çå¨ HTTP/2 ä¸ï¼ä½ ç res.write() å®é
ä¸å¯ä»¥çæ£é»å¡ï¼è䏿¯åªæ¯å¨æ¬å°ç¼å²ã
// æ£æµå½åè¿æ¥åè®®çæ¬
const httpVersion = req.httpVersion; // '1.1' æ '2.0'
if (httpVersion === '2.0') {
// HTTP/2 æå
ç½®èåï¼write() ä¼å¨çªå£æ»¡æ¶çå¾
// ä¸éè¦é¢å¤ç bufferSize æ£æµ
} else {
// HTTP/1.1 éè¦æå¨æ£æµ bufferSize
}
çäº§æ¶æï¼ä¸ä¸ªå®æ´çèåæç¥ SSE æå¡
æä¸é¢ææé·é±çè§£æ³æ´åèµ·æ¥ï¼ä¸ä¸ªç产级ç LLM SSE æå¡åºè¯¥æ¯è¿æ ·çï¼
import express from 'express';
// OpenAI å
¼å®¹ SDKï¼åæ¶æ¯æ DeepSeekãQwenãæºè°± GLM çå½äº§å¤§æ¨¡åç compatible API
import OpenAI from 'openai';
const app = express();
// å½äº§æ¨¡åå¦ DeepSeekãQwen åæ ·æ¯ææ¤ streaming åæ³
const client = new OpenAI({
baseURL: 'https://api.deepseek.com/v1', // æéæ¿æ¢ä¸ºå¯¹åºå½äº§æ¨¡å端ç¹
apiKey: process.env.DEEPSEEK_API_KEY,
});
// ä¼è¯åå¨ï¼çäº§ç¨ Redisï¼
const sessionStore = new Map<string, {
tokens: string[];
done: boolean;
createdAt: number;
}>();
const BUFFER_HIGH_WATERMARK = 128 * 1024; // 128KB
const SESSION_TTL_MS = 5 * 60 * 1000; // 5 åé
app.post('/v1/chat/stream', async (req, res) => {
// 1. SSE å¿
é headers
res.setHeader('Content-Type', 'text/event-stream');
res.setHeader('Cache-Control', 'no-cache, no-store, must-revalidate');
res.setHeader('Connection', 'keep-alive');
res.setHeader('X-Accel-Buffering', 'no'); // ç¦æ¢ä»£çç¼å²
// 2. 设置éè¿é´éï¼æ¯«ç§ï¼
res.write('retry: 3000\n\n');
const sessionId = req.headers['x-session-id'] as string || crypto.randomUUID();
const lastEventId = parseInt(req.headers['last-event-id'] as string) || -1;
// 3. å¤çéè¿ï¼å¦ææç¼åï¼å
åæ¾
const cached = sessionStore.get(sessionId);
if (cached && lastEventId >= 0) {
const replayStart = lastEventId + 1;
for (let i = replayStart; i < cached.tokens.length; i++) {
res.write(`id: ${i}\ndata: ${JSON.stringify({ text: cached.tokens[i] })}\n\n`);
}
if (cached.done) {
res.write('event: done\ndata: {}\n\n');
res.end();
return;
}
}
// 4. AbortControllerï¼å®¢æ·ç«¯æå¼æ¶ abort 䏿¸¸
const abortController = new AbortController();
let clientDisconnected = false;
req.on('close', () => {
clientDisconnected = true;
abortController.abort();
});
// 5. åå§åæè·åä¼è¯
if (!sessionStore.has(sessionId)) {
sessionStore.set(sessionId, {
tokens: [],
done: false,
createdAt: Date.now(),
});
}
const session = sessionStore.get(sessionId)!;
let eventId = session.tokens.length; // ä»å·²æ token æ°å¼å§è®¡æ°
try {
const stream = await client.chat.completions.create({
model: 'deepseek-chat',
messages: req.body.messages,
stream: true,
}, {
signal: abortController.signal,
});
for await (const chunk of stream) {
if (clientDisconnected) break;
const content = chunk.choices[0]?.delta?.content || '';
if (!content) continue;
// 6. ç¼å tokenï¼ç¨äºéè¿åæ¾ï¼
session.tokens.push(content);
// 7. æ£æ¥èå
const bufferSize = (res as any).socket?.bufferSize ?? 0;
if (bufferSize > BUFFER_HIGH_WATERMARK) {
// çå¾
drainï¼æå¤ 10 ç§
await Promise.race([
new Promise<void>(resolve => (res as any).socket?.once('drain', resolve)),
new Promise<void>(resolve => setTimeout(resolve, 10000)),
]);
// drain è¶
æ¶å仿ªæ¶è´¹ï¼æå¼è¿æ¥ï¼é¿å
å
åæç»å¢é¿ï¼
if (((res as any).socket?.bufferSize ?? 0) > BUFFER_HIGH_WATERMARK) {
console.warn(`[stream] slow client ${sessionId}, force closing`);
abortController.abort();
break;
}
}
// 8. åå
¥å¸¦ id ç SSE event
const ok = res.write(`id: ${eventId++}\ndata: ${JSON.stringify({ text: content })}\n\n`);
if (!ok) {
// write è¿å falseï¼çå¾
drainï¼è¿æ¯ Node.js æ åèååè®®ï¼
await new Promise<void>(resolve => res.once('drain', resolve));
}
}
// 9. 宿
session.done = true;
res.write('event: done\ndata: {}\n\n');
} catch (err: any) {
if (err.name === 'AbortError') {
// 客æ·ç«¯ä¸»å¨æå¼ï¼ä¸æ¯é误
} else {
console.error('[stream] upstream error:', err.message);
if (!res.writableEnded) {
res.write(`event: error\ndata: ${JSON.stringify({ message: err.message })}\n\n`);
}
}
} finally {
if (!res.writableEnded) {
res.end();
}
// 10. æ¸
çè¿æä¼è¯ï¼ç®å TTLï¼
setTimeout(() => {
sessionStore.delete(sessionId);
}, SESSION_TTL_MS);
}
});
çæ§ï¼ä½ éè¦æ´é²åªäºææ
èåé®é¢å¨æ²¡æçæ§çæ åµä¸å¾é¾åç°ââå®é常以å åç¼æ ¢å¢é¿æå¶åçæ ¢ååºå½¢å¼åºç°ï¼è䏿¯ç´æ¥çé误ã以䏿¯ç产ä¸å¿ é¡»çæ§çææ ï¼
// Prometheus ææ ç¤ºä¾
import { Gauge, Counter, Histogram } from 'prom-client';
const activeStreams = new Gauge({
name: 'llm_active_streams_total',
help: 'Number of active SSE streaming connections',
});
const slowClientDrops = new Counter({
name: 'llm_slow_client_drops_total',
help: 'Number of streams dropped due to slow client backpressure',
});
const clientDisconnects = new Counter({
name: 'llm_client_disconnect_aborts_total',
help: 'Number of upstream LLM requests aborted due to client disconnect',
});
const streamBufferSize = new Histogram({
name: 'llm_stream_buffer_bytes',
help: 'Distribution of socket buffer sizes during streaming',
buckets: [1024, 8192, 32768, 65536, 131072, 262144, 524288],
});
// æ¯ 500ms éæ · socket ç¼å²åºå¤§å°
function monitorBuffer(res: Response, sessionId: string) {
const interval = setInterval(() => {
const size = (res as any).socket?.bufferSize ?? 0;
streamBufferSize.observe(size);
if (size > 256 * 1024) {
console.warn(`[backpressure] session=${sessionId} buffer=${size} bytes`);
}
}, 500);
res.on('finish', () => clearInterval(interval));
req.on('close', () => clearInterval(interval));
}
å ³é®åè¦éå¼
| ææ | è¦åéå¼ | å±é©éå¼ |
|---|---|---|
| å¹³å socket bufferSize | > 64KB | > 256KB |
| æ ¢å®¢æ·ç«¯ drop ç | > 2% | > 10% |
| 客æ·ç«¯æè¿ abort ç | > 15% | > 30% |
| æ´»è·æµæ°é | è§å®¹é | è¶ è¿è®¾è®¡ä¸é |
横å对æ¯ï¼å ç§èåçç¥çæè¡¡
| çç¥ | å åæ§å¶ | ç¨æ·ä½éª | å®ç°å¤æåº¦ | éç¨åºæ¯ |
|---|---|---|---|---|
| çå¾ drainï¼æ¸©åï¼ | ä¸ | 好ï¼èªå¨åéï¼ | ä½ | å¶åæ ¢å®¢æ·ç«¯ |
| è¶ æ¶åæå¼ | 好 | å·®ï¼ééè¿ï¼ | ä¸ | æéè¿æºå¶æ¶ |
| åºå®ç¼å²ä¸é+drop | æå¥½ | æå·® | ä½ | æ CDN/è¾¹ç¼ç¼å |
| æå¡ç«¯ token ç¼å+éè¿ | 好 | æå¥½ | é« | ä»è´¹/æ ¸å¿ç¨æ·åºæ¯ |
| HTTP/2 flow control | æå¥½ | æå¥½ | ä½ï¼åè®®å ç½®ï¼ | å·²è¿ç§» HTTP/2 |
宿 Checklist
å¨ä½ ç LLM æµå¼æå¡ä¸çº¿åï¼éé¡¹æ£æ¥ï¼
æå¡ç«¯ï¼
-
res.write()è¿åå¼å·²å¤çï¼false æ¶çå¾drain - çå¬
req.on('close')å¹¶ abort 䏿¸¸ LLM è¯·æ± - socket bufferSize çæ§å·²æ¥å ¥ææ ç³»ç»
- ä¼è¯ token ç¼åæ TTL æ¸ çæºå¶
- ææ SSE event 齿éå¢
idåæ®µ
Nginx/代çå±ï¼
-
proxy_buffering off -
gzip offï¼æå¯¹ SSE path æé¤ï¼ -
proxy_read_timeout已设置足å¤é¿ï¼â¥ 300sï¼ -
X-Accel-Buffering: noheader 已设置
客æ·ç«¯ï¼
- å¤ç
Last-Event-IDéè¿é»è¾ - ææå¤§éè¿æ¬¡æ°éå¶ï¼é²æ¢ retry stormï¼
- å¤ç
event: erroråevent: doneèªå®ä¹äºä»¶
çæ§ï¼
- æ´»è·æµæ°é
- 客æ·ç«¯æè¿ abort ç
- æ ¢å®¢æ·ç«¯ drop ç
- P99 TTFTï¼é¦ token å»¶è¿ï¼
å°ç»
LLM æµå¼ååºä¸æ¯"ä¼ write SSE å°±è¡"ï¼å®æ¯ä¸ä¸ªéè¦è®¤ç设计èåçç¥çç³»ç»å·¥ç¨é®é¢ãæ ¸å¿ç»è®ºï¼
- Node.js
res.write()没æå ç½®èåï¼å¿ é¡»æå¨æ£æµ socket bufferSize å¹¶å¤ç drain äºä»¶ - 客æ·ç«¯æå¼å¿ é¡» abort 䏿¸¸ï¼ä¸ç¶ä½ å¨ä¸ºä¸ä¸ªæ¶å¤±çç¨æ·ç§ token é±
- Nginx é»è®¤ä¼ç ´å SSEï¼
proxy_buffering offæ¯å¿ é¡»é 置项ï¼ä¸æ¯å¯é项 - éè¿éè¦ token ç¼åï¼
last-event-id好çä½è¦æå¡ç«¯çæ£é åæææä¹ - HTTP/2 æ¯ç»æè§£æ³ï¼åè®®å±èåï¼é¿å æå¨ç®¡çç¼å²åº
线ä¸é®é¢å¾å¾ä¸æ¯"æµå¼ä¸å·¥ä½"ï¼èæ¯"æµå¼å¨å¤§é¨åæ åµä¸å·¥ä½ï¼å¨å°æ°è¾¹ç¼æ åµä¸æææèèµæº"ãå»ºå¥½çæ§ï¼æè¿äºè¾¹ç¼æ 嵿´é²åºæ¥ï¼ææ¯é¿æç¨³å®è¿è¡çå ³é®ã
æ¬æåºäº Node.js 20+ / Express 4.x çç产å®è·µï¼ä»£ç 示ä¾å·²å¨ DeepSeek-V3ãéä¹åé®ï¼Qwen-Maxï¼ãæºè°± GLM ç主æµå¤§æ¨¡åç streaming åºæ¯ä¸éªè¯ã
Aitishiku.com