Урок 07

Стриминг и 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 для отмены