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

如何在Node.js中将MongoDB集合数据同步至Typesense?

高效同步MongoDB集合到Typesense的方案

针对你的需求,以下是几种规范且内存友好的同步方案,避免全量加载数据到内存:

1. 流式游标分批同步

利用MongoDB游标流式读取数据,每次批量处理固定条数(比如100条),再批量插入Typesense,内存占用低且扩展性强,适合任意数据量:

const { MongoClient } = require('mongodb');
const Typesense = require('typesense');

async function syncCompanies() {
  const mongoClient = new MongoClient(process.env.MONGODB_URI);
  const typesenseClient = new Typesense.Client({
    nodes: [{ host: 'localhost', port: 8108, protocol: 'http' }],
    apiKey: process.env.TYPESENSE_API_KEY,
    connectionTimeoutSeconds: 2
  });

  try {
    await mongoClient.connect();
    const db = mongoClient.db('your-db-name');
    const companiesCollection = db.collection('companies');

    // 确保Typesense集合存在
    try {
      await typesenseClient.collections('companies').retrieve();
    } catch (err) {
      if (err.httpStatus === 404) {
        await typesenseClient.collections().create({
          name: 'companies',
          fields: [
            { name: '_id', type: 'string' },
            { name: 'name', type: 'string' },
            { name: 'industry', type: 'string' },
            // 按你的实际字段定义补充
          ],
          default_sorting_field: 'name'
        });
      } else throw err;
    }

    // 流式分批读取MongoDB数据
    const batchSize = 100;
    let cursor = companiesCollection.find({}, { batchSize });

    while (await cursor.hasNext()) {
      const batch = await cursor.next();
      // 转换MongoDB ObjectId为字符串(Typesense不支持ObjectId)
      const typesenseDocs = batch.map(doc => ({ ...doc, _id: doc._id.toString() }));
      // 批量插入/更新Typesense
      await typesenseClient.collections('companies').documents().import(typesenseDocs, { action: 'upsert' });
    }

    console.log('同步完成');
  } finally {
    await mongoClient.close();
  }
}

// 应用启动时触发同步
syncCompanies().catch(err => console.error('同步失败:', err));

2. 聚合管道预处理+流式同步

如果需要对MongoDB数据做清洗、字段转换,直接用MongoDB聚合管道完成预处理,减少应用层逻辑:

// 替换上述代码中的游标初始化部分
let cursor = companiesCollection.aggregate([
  { $project: {
    _id: { $toString: '$_id' },
    name: 1,
    industry: 1,
    // 按需添加字段或聚合转换逻辑
  } },
  { $batchSize: 100 }
]);

// 后续分批插入逻辑同方案1
while (await cursor.hasNext()) {
  const batch = await cursor.next();
  await typesenseClient.collections('companies').documents().import(batch, { action: 'upsert' });
}

3. 持续同步(兼顾启动后数据变更)

如果需要在初始同步后,实时同步MongoDB的新增、更新、删除操作,用MongoDB Change Streams实现:

async function startContinuousSync() {
  const mongoClient = new MongoClient(process.env.MONGODB_URI);
  const typesenseClient = new Typesense.Client({ /* 同前配置 */ });

  await mongoClient.connect();
  const db = mongoClient.db('your-db-name');
  const companiesCollection = db.collection('companies');

  // 监听MongoDB变更流
  const changeStream = companiesCollection.watch();

  changeStream.on('change', async (change) => {
    switch (change.operationType) {
      case 'insert':
        const newDoc = { ...change.fullDocument, _id: change.fullDocument._id.toString() };
        await typesenseClient.collections('companies').documents().create(newDoc);
        break;
      case 'update':
        const updateFields = { ...change.updateDescription.updatedFields };
        delete updateFields._id; // 禁止修改Typesense文档ID
        await typesenseClient.collections('companies').documents(change.documentKey._id.toString()).update(updateFields);
        break;
      case 'delete':
        await typesenseClient.collections('companies').documents(change.documentKey._id.toString()).delete();
        break;
    }
  });

  console.log('持续同步已启动');
}

// 先执行全量同步,再启动持续同步
syncCompanies().then(startContinuousSync).catch(err => console.error(err));

4. 增量同步(优化重复启动场景)

如果应用可能多次启动,不想每次全量同步,可记录上次同步时间戳,只同步新增/更新的数据:

async function syncIncremental() {
  const mongoClient = new MongoClient(process.env.MONGODB_URI);
  await mongoClient.connect();
  const db = mongoClient.db('your-db-name');
  const companiesCollection = db.collection('companies');
  const syncMetaColl = db.collection('sync_metadata');

  // 获取上次同步时间
  const lastSync = await syncMetaColl.findOne({ collection: 'companies' });
  // 假设你的文档有updatedAt字段,或用ObjectId的时间戳过滤
  const query = lastSync ? { updatedAt: { $gt: lastSync.timestamp } } : {};

  const cursor = companiesCollection.find(query, { batchSize: 100 });
  // 分批插入逻辑同方案1...

  // 更新同步时间戳
  await syncMetaColl.updateOne(
    { collection: 'companies' },
    { $set: { timestamp: new Date() } },
    { upsert: true }
  );
}

以上方案均无需全量加载500条数据到内存,且可无缝扩展到更大数据量,同时支持持续数据同步。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 11:45:35