Стриминг и SSE
Стриминг позволяет отправлять данные по частям, не дожидаясь полной готовности.
ReadableStream
export async function GET() {
const encoder = new TextEncoder();
const stream = new ReadableStream({
async start(controller) {
for (let i = 0; i < 10; i++) {
controller.enqueue(encoder.encode(`chunk ${i}\n`));
await new Promise((r) => setTimeout(r, 500));
}
controller.close();
},
});
return new Response(stream, {
headers: { 'Content-Type': 'text/plain' },
});
}
SSE (Server-Sent Events)
Формат: data: <json>\n\n.
export async function GET() {
const encoder = new TextEncoder();
const stream = new ReadableStream({
async start(controller) {
for (let i = 0; i < 10; i++) {
const data = JSON.stringify({ count: i });
controller.enqueue(encoder.encode(`data: ${data}\n\n`));
await new Promise((r) => setTimeout(r, 1000));
}
controller.close();
},
});
return new Response(stream, {
headers: {
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache, no-transform',
'Connection': 'keep-alive',
},
});
}
Клиент для SSE
'use client';
import { useEffect, useState } from 'react';
export function Counter() {
const [count, setCount] = useState(0);
useEffect(() => {
const source = new EventSource('/api/stream');
source.onmessage = (event) => {
const data = JSON.parse(event.data);
setCount(data.count);
};
source.onerror = () => source.close();
return () => source.close();
}, []);
return <p>{count}</p>;
}
EventSource автоматически переподключается.
AI-стриминг
Vercel AI SDK:
import { streamText } from 'ai';
import { openai } from '@ai-sdk/openai';
export async function POST(request: Request) {
const { messages } = await request.json();
const result = streamText({
model: openai('gpt-4'),
messages,
});
return result.toDataStreamResponse();
}
Клиент:
'use client';
import { useChat } from 'ai/react';
export function Chat() {
const { messages, input, handleInputChange, handleSubmit } = useChat();
return (
<>
{messages.map((m) => (
<div key={m.id}>{m.role}: {m.content}</div>
))}
<form onSubmit={handleSubmit}>
<input value={input} onChange={handleInputChange} />
<button>Отправить</button>
</form>
</>
);
}
Токены приходят по мере генерации.
Стриминг OpenAI напрямую
import OpenAI from 'openai';
const openai = new OpenAI({ apiKey: process.env.OPENAI_API_KEY });
export async function POST(request: Request) {
const { prompt } = await request.json();
const completion = await openai.chat.completions.create({
model: 'gpt-4',
messages: [{ role: 'user', content: prompt }],
stream: true,
});
const encoder = new TextEncoder();
const stream = new ReadableStream({
async start(controller) {
for await (const chunk of completion) {
const content = chunk.choices[0]?.delta?.content ?? '';
if (content) {
controller.enqueue(encoder.encode(content));
}
}
controller.close();
},
});
return new Response(stream, {
headers: { 'Content-Type': 'text/plain' },
});
}
Прогресс долгой задачи
export async function POST() {
const encoder = new TextEncoder();
const stream = new ReadableStream({
async start(controller) {
const steps = ['fetch', 'process', 'save', 'done'];
for (const step of steps) {
await doStep(step);
controller.enqueue(encoder.encode(`data: ${JSON.stringify({ step })}\n\n`));
}
controller.close();
},
});
return new Response(stream, {
headers: { 'Content-Type': 'text/event-stream' },
});
}
Отмена
Клиент закрывает соединение — стрим останавливается:
export async function GET(request: Request) {
const stream = new ReadableStream({
async start(controller) {
while (true) {
if (request.signal.aborted) break;
controller.enqueue('data\n');
await new Promise((r) => setTimeout(r, 1000));
}
controller.close();
},
});
return new Response(stream);
}
Когда использовать
- AI-ответы — токен за токеном
- Долгие задачи — прогресс
- Real-time — обновления
- Большие данные — не ждать всё
Итоги
ReadableStreamдля стриминга- SSE через
text/event-stream EventSourceна клиенте- AI SDK упрощает стриминг
- Прогресс через шаги
request.signal.abortedдля отмены