You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

NextJS 14中使用TransformStream实现SSE时消息一次性发送的问题

在NextJS 14中实现基于POST请求的SSE实时更新问题

我需要在NextJS 14中实现SSE,用于处理数据时实时更新用户状态。因为要接收用户提交的数据,必须用POST请求,所以没法用EventSource,只能用fetch()。客户端代码已经正常工作,但服务端需要在耗时处理函数完成后发送消息,而不是定时发送。我写了在函数前后发消息的代码,但所有消息会一次性发送,甚至可能提前关闭writer。试过以下方法都没用:

  • 在写入/关闭之间加sleep函数
  • await所有写入和关闭操作
  • 将写入操作封装为Promise

客户端代码

"use client";

import { useState } from "react";

export default function Home() {
    const [message, setMessage] = useState("");

    async function onClick() {
        const response = await fetch("/api/test2", {
            method: "POST",
            headers: {
                "Content-Type": "application/json",
            },
            body: JSON.stringify({}),
        });

        const reader = response.body?.getReader();
        if (!reader) return;

        let decoder = new TextDecoder();
        while (true) {
            const { done, value } = await reader.read();
            if (done) break;
            if (!value) continue;
            const lines = decoder.decode(value);
            const text = lines
                .split("\n")
                .filter((line) => line.startsWith("data:"))[0]
                .replace("data:", "")
                .trim();

            setMessage((prev) => prev + text);
        }
    }

    return (
        <div>
            <button onClick={onClick}>START</button>
            <p>{message}</p>
        </div>
    );
}

服务端可用的定时发送代码

import { NextRequest, NextResponse } from "next/server";

export async function POST(req: NextRequest, res: NextResponse) {
    const { readable, writable } = new TransformStream();
    const writer = writable.getWriter();

    const text = "Some test text";

    let index = 0;
    const interval = setInterval(() => {
        if (index < text.length) {
            writer.write(`event: message\ndata: ${text[index]}\n\n`);
            index++;
        } else {
            writer.write(`event: message\ndata: [DONE]\n\n`);
            clearInterval(interval);
            writer.close();
        }
    }, 1);

    return new NextResponse(readable, {
        headers: {
            "Content-Type": "text/event-stream",
            "Cache-Control": "no-cache",
            Connection: "keep-alive",
        },
    });
}

问题代码

import { NextRequest, NextResponse } from "next/server";

export async function POST(req: NextRequest, res: NextResponse) {
    const { readable, writable } = new TransformStream();
    const writer = writable.getWriter();

    writer.write(`event: "start"\ndata:"Process 1"\n\n`)
    await processThatTakesTime(); //类似用Puppeteer点击按钮的操作
    writer.write(`event: "done"\ndata:"Process 1"\n\n`)

    writer.write(`event: "start"\ndata:"Process 2"\n\n`)
    await anotherProcess(); //类似用Puppeteer点击按钮的操作
    writer.write(`event: "done"\ndata:"Process 2"\n\n`)

    writer.close();

    return new NextResponse(readable, {
        headers: {
            "Content-Type": "text/event-stream",
            "Cache-Control": "no-cache",
            Connection: "keep-alive",
        },
    });
}

尝试的Promise封装代码

await new Promise<void>((resolve)=>{
    setTimeout(()=>{
        writer.write(`event: "start"\ndata:"Process 1"\n\n`);
        resolve();
    },100)
});

解决方案

问题核心是两个:一是writer.write()返回Promise但未await,导致写入操作积压;二是SSE格式错误(event字段带引号),同时Next.js可能缓冲响应导致消息延迟。

修改后的服务端代码

import { NextRequest, NextResponse } from "next/server";

export async function POST(req: NextRequest) {
    const { readable, writable } = new TransformStream();
    const writer = writable.getWriter();

    // 封装安全的SSE发送函数,确保消息实时发送
    const sendSSE = async (event: string, data: string) => {
        // 修正SSE格式:event字段无需引号,data用JSON.stringify保证格式正确
        await writer.write(`event: ${event}\ndata: ${JSON.stringify(data)}\n\n`);
        // 等待writer就绪,强制刷新缓冲区
        await writer.ready;
    };

    try {
        await sendSSE("start", "Process 1");
        await processThatTakesTime();
        await sendSSE("done", "Process 1");

        await sendSSE("start", "Process 2");
        await anotherProcess();
        await sendSSE("done", "Process 2");

        // 发送结束标记
        await sendSSE("message", "[DONE]");
    } catch (err) {
        // 异常时终止流,避免连接挂起
        await writer.abort(err);
    } finally {
        await writer.close();
    }

    return new NextResponse(readable, {
        headers: {
            "Content-Type": "text/event-stream",
            "Cache-Control": "no-cache",
            Connection: "keep-alive",
            // 禁用响应压缩,防止Next.js缓冲SSE流
            "Content-Encoding": "identity",
        },
    });
}

关键修改点

  1. await writer.write():每次写入都等待Promise完成,避免消息积压
  2. 修正SSE格式:event字段不带引号,用JSON.stringify()处理data,保证客户端解析正常
  3. 添加writer.ready:强制刷新缓冲区,确保消息实时发送到客户端
  4. 错误处理:用try/catch包裹,异常时调用writer.abort(),避免无效连接
  5. 禁用响应压缩:添加Content-Encoding: identity,防止Next.js或中间件缓冲流

客户端代码优化(处理多消息场景)

// 客户端onClick函数中修改解析逻辑
let decoder = new TextDecoder();
let buffer = "";
while (true) {
    const { done, value } = await reader.read();
    if (done) break;
    if (!value) continue;
    buffer += decoder.decode(value);
    // 按SSE分隔符(两个换行)拆分消息
    const messages = buffer.split("\n\n");
    buffer = messages.pop() || ""; // 保留未完成的消息片段

    for (const msg of messages) {
        if (!msg.trim()) continue;
        const lines = msg.split("\n");
        let text = "";
        for (const line of lines) {
            if (line.startsWith("data:")) {
                text += line.replace("data:", "").trim();
            }
        }
        if (text) {
            setMessage((prev) => prev + text);
        }
    }
}

内容的提问来源于stack exchange,提问作者Adrian11asd

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.27 00:14:53