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

如何在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:

  1. 在SQL中创建触发器,将变更事件写入Pub/Sub(以PostgreSQL为例,使用pg_notify或第三方扩展)
  2. 编写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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 19:15:38