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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 04:22:09