如何在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
相关产品推荐
相关产品推荐

