【问题标题】:MongoDB change stream timeouts if database is down for some time如果数据库关闭一段时间,MongoDB 更改流超时
【发布时间】:2019-02-10 08:43:33
【问题描述】:

我在 nodejs 中使用 mongoDB 更改流,一切正常,但如果数据库关闭需要超过 10 5 秒才能启动更改流会引发超时错误,这是我的更改流观察器代码

Service.prototype.watcher = function( db ){

let collection = db.collection('tokens');
let changeStream = collection.watch({ fullDocument: 'updateLookup' });
let resumeToken, newChangeStream;

changeStream.on('change', next => {
    resumeToken = next._id;
    console.log('data is ', JSON.stringify(next))
    changeStream.close();
    // console.log('resumeToken is ', JSON.stringify(resumeToken))
    newChangeStream = collection.watch({ resumeAfter : resumeToken });
    newChangeStream.on('change', next => {
        console.log('insert called ', JSON.stringify( next ))
    });
});

但是在数据库端我已经处理了它,即如果数据库关闭或使用此代码重新连接

 this.db.on('reconnected', function () {
    console.info('MongoDB reconnected!');
});
this.db.on('disconnected', function() {
    console.warn('MongoDB disconnected!');
});

但我无法处理更改流观察器以在数据库关闭时停止它并在重新连接数据库时重新启动它或是否有任何其他更好的方法?

【问题讨论】:

    标签: node.js mongodb changestream


    【解决方案1】:

    您要做的是将watch() 调用封装在一个函数中。然后,此函数将在出错时调用自身,以使用先前保存的恢复令牌重新观察集合。您拥有的代码中缺少的是错误处理程序。例如:

    const MongoClient = require('mongodb').MongoClient
    const uri = 'mongodb://localhost:27017/test?replicaSet=replset'
    var resume_token = null
    
    run()
    
    function watch_collection(con, db, coll) {
      console.log(new Date() + ' watching: ' + coll)
      con.db(db).collection(coll).watch({resumeAfter: resume_token})
        .on('change', data => {
          console.log(data)
          resume_token = data._id
        })
        .on('error', err => {
          console.log(new Date() + ' error: ' + err)
          watch_collection(con, coll)
        })
    }
    
    async function run() {
      con = await MongoClient.connect(uri, {"useNewUrlParser": true})
      watch_collection(con, 'test', 'test')
    }
    

    注意watch_collection() 包含watch() 方法及其处理程序。在更改时,它将打印更改并存储恢复令牌。出错时,它会调用自己再次重新观看集合。

    【讨论】:

    • 是的,这就是解决方案,我以相同的方式实现,现在它正在工作,我现在正在将恢复令牌保存在文件中,并再次运行它从该文件中获取,现在即使应用程序崩溃或停止,但当它再次运行时,它将从保存的最后一个成功事件 ID 开始观看,即恢复令牌
    【解决方案2】:

    这是我开发的解决方案,只需添加 stream.on(error) 函数,这样它就不会在出现错误时崩溃,因为在重新连接数据库时重新启动流,还将每个事件的恢复令牌保存在文件中,这是当应用程序崩溃或停止并且您再次运行并且在此期间如果添加了 x 条记录时会很有帮助,因此在应用程序重新启动时只需从文件中获取最后一个恢复令牌并从那里启动观察程序,它将获得之后插入的所有记录,因此没有记录会丢失,下面是代码

    var rsToken ;
        try {
            rsToken = await this.getResumetoken()
        } catch (error) {
            rsToken = null ;
        }
    
        if (!rsToken)
            changeStream = collection.watch({ fullDocument: 'updateLookup' });
        else 
            changeStream = collection.watch({ fullDocument: 'updateLookup', resumeAfter : rsToken  });
    
        changeStream.on('change', next => {
    
            resumeToken = next._id;
            THIS.saveTokenInfile(resumeToken)
    
            cs_processor.process( next )
    
    
        });  
        changeStream.on('error', err => {
            console.log('changestream error ')
        })
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2013-12-14
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多