如何在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; }
问题原因及解决方法
你的代码有两个核心问题导致流传输失效:
- 请求方法不匹配:Supabase Edge Function默认处理
GET请求,但你用了POST请求,函数没有对应的处理逻辑,可能导致请求被阻塞。 - 缓存所有数据再返回:你把流的所有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
相关产品推荐
相关产品推荐

