Next.js 14 App Router中MongoDB Change Stream推送及关闭问题
Next.js 14 App Router 中 MongoDB Change Stream 实时推送与资源释放方案
你的核心问题在于:普通的Next.js API GET请求是一次性响应模型,无法持续向客户端推送数据;同时没有监听客户端连接断开事件来清理Change Stream资源。下面是具体的解决方案:
一、用Server-Sent Events(SSE)实现实时数据推送
SSE是浏览器原生支持的单向服务器推送技术,适合这种需要持续监听数据库变更并推送到前端的场景。我们需要在API路由中创建一个可读流,将Change Stream的变更数据以SSE格式发送给客户端。
修改后的API路由代码
import { NextResponse } from 'next/server'; import { ObjectId } from 'mongodb'; import { getMongoDb } from '@/lib/mongodb'; // 替换为你的MongoDB连接工具路径 export async function GET(request, { params: { eventID } }) { try { const db = await getMongoDb(); const collection = db.collection("payments"); const pipeline = [ { $match: { $and: [ { $or: [ { operationType: "insert" }, { operationType: "update" }, ], }, { "fullDocument.eventID": new ObjectId(eventID) }, ], }, }, ]; const changeStream = collection.watch(pipeline, { fullDocument: "updateLookup", }); // 创建SSE可读流 const stream = new ReadableStream({ async start(controller) { // 监听Change Stream变更事件 const onChange = (change) => { const updatedData = change.fullDocument; console.log("updatedData", updatedData); // 按SSE格式写入数据(必须以data: 内容\n\n结尾) controller.enqueue(`data: ${JSON.stringify(updatedData)}\n\n`); }; changeStream.on("change", onChange); // 监听客户端断开连接,自动关闭Change Stream request.signal.addEventListener('abort', async () => { await changeStream.close(); controller.close(); }); // 处理Change Stream异常 changeStream.on("error", (err) => { controller.error(err); console.error("Change Stream异常:", err); }); }, }); // 返回SSE响应,设置正确的响应头 return new NextResponse(stream, { headers: { 'Content-Type': 'text/event-stream', 'Cache-Control': 'no-cache', 'Connection': 'keep-alive', }, }); } catch (e) { return NextResponse.json( { error: { message: e.message } }, { status: 500 } ); } }
二、前端接收SSE数据
在前端组件中,使用EventSource连接到API路由,监听数据变更并更新页面:
'use client'; import { useEffect, useState } from 'react'; export default function PaymentUpdates({ eventID }) { const [payments, setPayments] = useState([]); useEffect(() => { // 建立SSE连接 const eventSource = new EventSource(`/api/payments/watch/${eventID}`); // 接收服务器推送的数据 eventSource.onmessage = (event) => { const updatedData = JSON.parse(event.data); // 根据业务需求更新状态,比如追加或替换现有数据 setPayments(prev => [...prev, updatedData]); }; // 处理SSE连接错误 eventSource.onerror = (error) => { console.error("SSE连接错误:", error); eventSource.close(); }; // 组件卸载时关闭SSE连接 return () => { eventSource.close(); }; }, [eventID]); return ( <div> <h3>实时支付更新</h3> {payments.map(payment => ( <div key={payment._id} className="border p-2 mb-2"> <p>金额: {payment.amount}</p> <p>状态: {payment.status}</p> </div> ))} </div> ); }
关键细节说明
- 资源清理:通过
request.signal.addEventListener('abort')监听客户端断开事件,自动关闭MongoDB Change Stream,避免资源泄漏。 - SSE格式规范:必须以
data: 内容\n\n的格式发送数据,浏览器才能正确解析为事件消息。 - 客户端连接管理:前端组件卸载时务必调用
eventSource.close(),避免无效连接占用服务器资源。
内容的提问来源于stack exchange,提问作者Khalid
相关产品推荐
相关产品推荐

