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

如何在Next.js路由处理器中正确消费Server-Sent Events流?

如何在Next.js路由处理器中正确消费Server-Sent Events流?

我用Deno编写了一个Supabase Edge Function,它每秒流式推送一条消息,代码如下:

/// <reference types="https://esm.sh/@supabase/functions-js/src/edge-runtime.d.ts" />

const msg = new TextEncoder().encode("data: hello\r\n\r\n");

Deno.serve((_req) => {
  let timerId: number | undefined;
  const body = new ReadableStream({
    start(controller) {
      timerId = setInterval(() => {
        console.log(msg);
        controller.enqueue(msg);
      }, 1000);
    },
    cancel() {
      if (typeof timerId === "number") {
        clearInterval(timerId);
      }
    },
  });
  return new Response(body, {
    headers: {
      "Content-Type": "text/event-stream",
      "Content-Encoding": "none",
    },
  });
});

我编写了如下Next.js路由处理器来消费该流,但fetch请求一直挂起,完全没有实现流传输:

export const dynamic = 'force-dynamic'

export async function GET() {
  const params = {
  };

  console.log("Before Fetch");
  const res = await fetch("http://127.0.0.1:54321/functions/v1/stream", {
    method: "POST",
    headers: {
      "Content-Type": "application/json",
      "Authorization":
        "Bearer <JWT>",
    },
    body: JSON.stringify(params),
  });

  if (res.body) {
    const rstream = res.body;
    const reader = rstream.getReader();

    let done = false, value = undefined, chunks = [];
    do {
      ({ done, value } = await reader.read());
      if (done) {
        break;
      } else if (value) {
        chunks.push(value);
      }
    } while (!done);
    return Response.json({ data: chunks });
  }
  return Response;
}

问题原因及解决方法

你的代码有两个核心问题导致流传输失效:

  1. 请求方法不匹配:Supabase Edge Function默认处理GET请求,但你用了POST请求,函数没有对应的处理逻辑,可能导致请求被阻塞。
  2. 缓存所有数据再返回:你把流的所有chunk都收集到数组里,直到流结束才返回JSON,完全违背了SSE实时推送的设计,导致请求一直挂起直到流关闭。

修改后的Next.js路由处理器代码

export const dynamic = 'force-dynamic'

export async function GET() {
  console.log("Before Fetch");
  const res = await fetch("http://127.0.0.1:54321/functions/v1/stream", {
    method: "GET", // 改为GET,匹配Supabase函数的请求方法
    headers: {
      "Authorization": "Bearer <JWT>", // 保留你的JWT令牌
    },
  });

  if (!res.body) {
    return new Response("无法获取数据流", { status: 500 });
  }

  // 创建实时转发的可读流
  const stream = new ReadableStream({
    async start(controller) {
      const reader = res.body.getReader();
      try {
        while (true) {
          const { done, value } = await reader.read();
          if (done) break;
          // 将收到的SSE数据实时推送给客户端
          controller.enqueue(value);
        }
      } catch (err) {
        controller.error(err);
      } finally {
        reader.releaseLock();
        controller.close();
      }
    }
  });

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

额外优化建议

如果是本地开发,为了避免跨域问题,需要在Supabase Edge Function的响应头中添加CORS配置:

// 在Supabase函数的Response headers中添加
"Access-Control-Allow-Origin": "*", // 生产环境请替换为你的域名
"Access-Control-Allow-Methods": "GET, OPTIONS",
"Access-Control-Allow-Headers": "Authorization"

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 06:45:58