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

如何让Node.js中的MongoDB变更流持续运行?

解决Mongoose变更流停止监听的问题

你的变更流偶尔静默停止,核心原因是未处理连接中断、错误事件,且缺少自动重连与断点恢复逻辑,同时长时间无变更时可能因连接闲置被服务器断开。以下是具体解决方案:

1. 监听变更流的错误与关闭事件

变更流遇到网络中断、服务器重启等情况时会触发error或close事件,若不处理这些事件,流会直接停止且无报错。需在事件触发时重新创建变更流。

2. 持久化Resume Token

你当前将resumeTokenRis存在内存中,应用重启或崩溃后会丢失,导致无法从上次中断的位置恢复监听。建议将其持久化到文件、Redis或MongoDB集合中。

3. 配置心跳与长连接参数

显式配置参数避免长时间无变更时连接被断开:

  • heartbeatFrequencyMS:设置心跳间隔(默认30000ms),缩短间隔确保服务器感知连接存活
  • maxAwaitTimeMS:设置单次等待变更的最长时间(默认10000ms),到期后返回空结果维持连接活跃

4. 封装变更流创建逻辑

把创建变更流的代码封装成函数,方便在错误、关闭事件中调用重试。

修改后的完整代码示例

const fs = require('fs').promises;
const RESUME_TOKEN_PATH = './resume-token-ris.json';

// 从文件读取持久化的resumeToken
async function loadResumeToken() {
  try {
    const data = await fs.readFile(RESUME_TOKEN_PATH, 'utf8');
    return JSON.parse(data);
  } catch (err) {
    // 文件不存在时返回null
    return null;
  }
}

// 保存resumeToken到文件
async function saveResumeToken(token) {
  try {
    await fs.writeFile(RESUME_TOKEN_PATH, JSON.stringify(token), 'utf8');
  } catch (err) {
    console.error('Failed to save resume token:', err);
  }
}

// 创建变更流的函数
async function createChangeStream() {
  let resumeTokenRis = await loadResumeToken();
  const ris = mongoose.connection.db.collection("ris");
  
  try {
    const changeStreamRis = ris.watch([], {
      resumeAfter: resumeTokenRis,
      heartbeatFrequencyMS: 15000, // 15秒发送一次心跳
      maxAwaitTimeMS: 30000 // 30秒无变更则返回空,保持连接
    });

    // 监听变更事件
    changeStreamRis.on("change", async (change) => {
      if (change.operationType === "insert") {
        resumeTokenRis = change._id;
        await saveResumeToken(resumeTokenRis); // 持久化最新的token
        console.log("RI Inserted");

        const ri = change.fullDocument;
        let userIds = ri.issueStat.taggedUser?.map(userId => userId.toString()) || [];
        userIds = [...new Set(userIds)].map(id => new ObjectId(id));

        const userTokens = await NotificationToken.find({ userId: { $in: userIds } });
        const user = await User.findById(ri.user);

        let tokens = userTokens.flatMap(userToken => userToken.tokens);

        if (tokens.length === 0) {
          console.log("No tokens found for this user.");
          return;
        }

        const message = {
          notification: {
            title: `New RI Assignment`,
            body: `Assigned by ${user.name} with ${ri.priority} priority.`,
          },
          data: {
            id: `${ri.refId}`,
            type: "RI",
          },
          tokens: tokens,
        };

        try {
          const response = await admin.messaging().sendEachForMulticast(message);
          console.log(response);
          for (let index = 0; index < response.responses.length; index++) {
            const res = response.responses[index];
            if (!res.success) {
              const failedToken = tokens[index];
              console.log(`Failed token ID: ${failedToken}`);
              try {
                await NotificationToken.updateOne(
                  { userId: { $in: userIds } },
                  { $pull: { tokens: failedToken } }
                );
                console.log(`Token ${failedToken} deleted successfully`);
              } catch (error) {
                console.error(`Failed to delete token ${failedToken}: ${error}`);
              }
            }
          }
        } catch (error) {
          console.error("Failed to send notification:", error);
        }
      }
    });

    // 监听错误事件,出错后重新创建流
    changeStreamRis.on("error", async (err) => {
      console.error("Change stream error:", err);
      // 延迟几秒后重试,避免频繁重试
      setTimeout(() => createChangeStream(), 5000);
    });

    // 监听关闭事件,关闭后重新创建流
    changeStreamRis.on("close", () => {
      console.log("Change stream closed, restarting...");
      setTimeout(() => createChangeStream(), 5000);
    });

    console.log("Change stream started successfully");
  } catch (err) {
    console.error("Failed to create change stream:", err);
    setTimeout(() => createChangeStream(), 5000);
  }
}

// 启动变更流
createChangeStream();

额外注意事项

  • 若使用MongoDB单节点,变更流的断点恢复功能受限,建议改用副本集或分片集群
  • 生产环境中,可将resumeToken持久化到Redis或MongoDB集合,替代文件存储
  • 重试逻辑可增加指数退避策略(如第一次等5秒,第二次10秒,第三次20秒),避免给数据库造成压力

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 06:24:56