君は春の中にいる、かけがえのない春の中にいる.

你驻足于春色中,于那独一无二的春色之中.

利用Mongo数据库的oplog机制实现准实时数据操作监控

最近有一个需求是要实时获取到新插入到MongoDB的数据,而插入程序本身已经有一套处理逻辑,所以不方便直接在插入程序里写相关程序,传统的数据库大多自带这种触发器机制,但是Mongo没有相关的函数可以用(也可能我了解的太少了,求纠正),当然还有一点是需要python实现,于是收集整理了一个相应的实现方法。

0x01 引子

首先可以想到,这种需求其实很像数据库的主从备份机制,从数据库之所以能够同步主库是因为存在某些指标来做控制,我们知道MongoDB虽然没有现成触发器,但是它能够实现主从备份,所以我们就从它的主从备份机制入手。

0x02 Oplog

首先,需要以master模式来打开mongod守护,命令行使用–master,或者配置文件增加master键为true。

此时,我们可以在Mongo的系统库local里见到新增的collection——oplog,此时oplog.$main里就会存储进oplog信息,如果此时还有充当从数据库的Mongo存在,就会还有一些slaves的信息,由于我们这里并不是主从同步,所以不存在这些集合。

再来看看oplog结构:

"ts" : Timestamp(6417682881216249, 1),  时间戳
"h" : NumberLong(0),  长度
"v" : 2,  
"op" : "n", 操作类型
"ns" : "",  操作的库和集合
"o2" : "_id"  update条件
"o" : {}  操作值,即document

这里需要知道op的几种属性:

insert,'i'
update, 'u'
remove(delete), 'd'
cmd, 'c'
noop, 'n' 空操作

从上面的信息可以看出,我们只要不断读取到ts来做对比,然后根据op即可判断当前出现的是什么操作,相当于使用程序实现了一个从数据库的接收端。

0x03 Code

在Github上找到了别人的实现方式,不过它的函数库太老旧,所以在他的基础上进行修改。

https://github.com/RedBeard0531/mongo-oplog-watcher

首先是实现一个数据库的初始化,设定一个延迟时间(准实时):

self.poll_time = poll_time
self.connection = connection or pymongo.MongoClient()

主要的函数是start(),实现一个时间的比对并进行相应字段的处理:

def start(self):
    oplog = self.connection.local['oplog.$main']
    #读取之前提到的库
    ts = oplog.find().sort('$natural', -1)[0]['ts']
    #获取一个时间边际
    while True:
        if self._ns_filter is None:
            filter = {}
        else:
            filter = {'ns': self._ns_filter}
        filter['ts'] = {'$gt': ts}
        try:
            cursor = oplog.find(filter)
            #对此时间之后的进行处理
            while True:
                for op in cursor:
                    ts = op['ts']
                    id = self.__get_id(op)
                    self.all_with_noop(ns=op['ns'], ts=ts, op=op['op'], id=id, raw=op)
                    #可以指定处理插入监控,更新监控或者删除监控等
                time.sleep(self.poll_time)
                if not cursor.alive:
                    break
        except AutoReconnect:
            time.sleep(self.poll_time)

循环这个start函数,在all_with_noop这里就可以编写相应的监控处理逻辑。

这样就可以实现一个简易的准实时Mongo数据库操作监控器,下一步就可以配合其他操作来对新入库的程序进行相应处理。