如何在Firestore与GCP其他数据库间实现100条数据双向实时同步?
双向实时同步指定Firestore记录至次级数据库
针对你需要同步Firestore中100条指定记录到次级数据库,且双向变更保持一致的需求,以下是分场景的实现方案:
核心逻辑
实现双向同步的关键是:分别监听两端数据库的变更事件,触发对方的同步操作,同时通过同步标记字段(如lastSyncedFrom)避免循环触发更新。
场景1:次级数据库为Firestore(同/跨项目)
步骤1:标记需同步的记录
在主Firestore的目标集合中,给需要同步的100条记录添加统一标记,比如sync: true,用于后续精准监听。
步骤2:主→次同步
通过Firestore实时监听API onSnapshot 监听标记记录的变更,同步至次级Firestore:
const admin = require('firebase-admin'); admin.initializeApp(); // 初始化主库和次级库实例 const mainDb = admin.firestore(); const secondaryDb = admin.initializeApp( { projectId: 'your-secondary-project-id' }, 'secondary' ).firestore(); // 监听主库指定记录的变更 mainDb.collection('target-collection').where('sync', '==', true).onSnapshot(querySnapshot => { querySnapshot.docChanges().forEach(async change => { const doc = change.doc; const docData = doc.data(); // 跳过来自次级库的同步变更,避免循环 if (docData.lastSyncedFrom === 'secondary') return; const secondaryRef = secondaryDb.collection('synced-collection').doc(doc.id); switch(change.type) { case 'added': case 'modified': await secondaryRef.set({ ...docData, lastSyncedFrom: 'primary' }, { merge: true }); break; case 'removed': await secondaryRef.delete(); break; } }); });
步骤3:次→主同步
在次级Firestore上监听对应集合的变更,同步回主库,逻辑与主→次一致,只需调整标记判断:
// 次级库监听逻辑 secondaryDb.collection('synced-collection').onSnapshot(querySnapshot => { querySnapshot.docChanges().forEach(async change => { const doc = change.doc; const docData = doc.data(); // 跳过来自主库的同步变更 if (docData.lastSyncedFrom === 'primary') return; const mainRef = mainDb.collection('target-collection').doc(doc.id); switch(change.type) { case 'added': case 'modified': await mainRef.set({ ...docData, lastSyncedFrom: 'secondary' }, { merge: true }); break; case 'removed': await mainRef.delete(); break; } }); });
场景2:次级数据库为Cloud SQL(MySQL/PostgreSQL)
主→次同步
用Firestore实时监听标记记录,变更时通过SQL客户端写入对应表,注意数据类型映射(如Firestore Timestamp转SQL DATETIME):
const mysql = require('mysql2/promise'); // 初始化SQL连接池 const pool = mysql.createPool({ host: 'your-sql-host', user: 'sql-user', password: 'sql-password', database: 'db-name' }); // 主库监听同步逻辑 mainDb.collection('target-collection').where('sync', '==', true).onSnapshot(querySnapshot => { querySnapshot.docChanges().forEach(async change => { const doc = change.doc; const docData = doc.data(); if (docData.lastSyncedFrom === 'secondary') return; const conn = await pool.getConnection(); try { if (change.type === 'added' || change.type === 'modified') { await conn.execute( 'REPLACE INTO synced_table (id, field1, field2, updated_at, last_synced_from) VALUES (?, ?, ?, ?, ?)', [doc.id, docData.field1, docData.field2, docData.updatedAt.toDate(), 'primary'] ); } else if (change.type === 'removed') { await conn.execute('DELETE FROM synced_table WHERE id = ?', [doc.id]); } } finally { conn.release(); } }); });
次→主同步
通过Cloud SQL触发器捕获表变更,将事件发送至Pub/Sub,再触发Cloud Function同步回Firestore:
- 在SQL中创建触发器,将变更事件写入Pub/Sub(以PostgreSQL为例,使用
pg_notify或第三方扩展) - 编写Cloud Function处理Pub/Sub消息:
exports.syncSqlToFirestore = functions.pubsub.topic('sql-table-changes').onPublish(async (message) => { const eventData = JSON.parse(Buffer.from(message.data, 'base64').toString()); const { id, newData, changeType } = eventData; // 跳过主库同步来的变更 if (newData.last_synced_from === 'primary') return; const mainRef = mainDb.collection('target-collection').doc(id); if (changeType === 'insert' || changeType === 'update') { await mainRef.set({ ...newData, updatedAt: admin.firestore.Timestamp.fromDate(new Date(newData.updated_at)), lastSyncedFrom: 'secondary' }, { merge: true }); } else if (changeType === 'delete') { await mainRef.delete(); } });
场景3:次级数据库为MongoDB(Atlas/自托管)
主→次同步
监听Firestore变更,通过MongoDB驱动写入对应集合:
const { MongoClient } = require('mongodb'); const mongoClient = new MongoClient('your-mongodb-connection-string'); await mongoClient.connect(); const mongoCollection = mongoClient.db('db-name').collection('synced-collection'); // 主库监听同步逻辑 mainDb.collection('target-collection').where('sync', '==', true).onSnapshot(querySnapshot => { querySnapshot.docChanges().forEach(async change => { const doc = change.doc; const docData = doc.data(); if (docData.lastSyncedFrom === 'secondary') return; switch(change.type) { case 'added': case 'modified': await mongoCollection.replaceOne( { _id: doc.id }, { ...docData, _id: doc.id, lastSyncedFrom: 'primary' }, { upsert: true } ); break; case 'removed': await mongoCollection.deleteOne({ _id: doc.id }); break; } }); });
次→主同步
使用MongoDB Change Streams监听集合变更,同步回Firestore:
const changeStream = mongoCollection.watch(); changeStream.on('change', async (change) => { const docId = change.documentKey._id.toString(); const mainRef = mainDb.collection('target-collection').doc(docId); // 跳过主库同步来的变更 if (change.fullDocument?.lastSyncedFrom === 'primary') return; switch(change.operationType) { case 'insert': case 'update': await mainRef.set({ ...change.fullDocument, lastSyncedFrom: 'secondary', _id: undefined // 移除MongoDB的_id字段 }, { merge: true }); break; case 'delete': await mainRef.delete(); break; } });
关键注意事项
- 初始同步:先批量将主库的100条标记记录写入次级数据库,再开启实时监听,避免遗漏初始数据。
- 冲突处理:添加
updatedAt字段,同步时以最新更新时间的记录为准,解决并发变更冲突。 - 错误重试:同步逻辑中加入重试机制(如
p-retry库),处理网络或数据库连接错误。 - 权限控制:确保同步服务(如Cloud Function)拥有两端数据库的读写权限。
内容的提问来源于stack exchange,提问作者david jeba
相关产品推荐
相关产品推荐

