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

发布时间 - 2026-01-11 01:03:58    点击率:

前言

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

一、引子

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

二、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即可判断当前出现的是什么操作,相当于使用程序实现了一个从数据库的接收端。

三、CODE

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

Github地址:https://github.com/RedBeard0531/mongo-oplog-watcher

mongo_oplog_watcher.py如下:

#!/usr/bin/python
import pymongo
import re
import time
from pprint import pprint # pretty printer
from pymongo.errors import AutoReconnect

class OplogWatcher(object):
  def __init__(self, db=None, collection=None, poll_time=1.0, connection=None, start_now=True):
    if collection is not None:
      if db is None:
        raise ValueError('must specify db if you specify a collection')
      self._ns_filter = db + '.' + collection
    elif db is not None:
      self._ns_filter = re.compile(r'^%s\.' % db)
    else:
      self._ns_filter = None

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

    if start_now:
      self.start()

  @staticmethod
  def __get_id(op):
    id = None
    o2 = op.get('o2')
    if o2 is not None:
      id = o2.get('_id')

    if id is None:
      id = op['o'].get('_id')

    return id

  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, tailable=True)
        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)

  def all_with_noop(self, ns, ts, op, id, raw):
    if op == 'n':
      self.noop(ts=ts)
    else:
      self.all(ns=ns, ts=ts, op=op, id=id, raw=raw)

  def all(self, ns, ts, op, id, raw):
    if op == 'i':
      self.insert(ns=ns, ts=ts, id=id, obj=raw['o'], raw=raw)
    elif op == 'u':
      self.update(ns=ns, ts=ts, id=id, mod=raw['o'], raw=raw)
    elif op == 'd':
      self.delete(ns=ns, ts=ts, id=id, raw=raw)
    elif op == 'c':
      self.command(ns=ns, ts=ts, cmd=raw['o'], raw=raw)
    elif op == 'db':
      self.db_declare(ns=ns, ts=ts, raw=raw)

  def noop(self, ts):
    pass

  def insert(self, ns, ts, id, obj, raw, **kw):
    pass

  def update(self, ns, ts, id, mod, raw, **kw):
    pass

  def delete(self, ns, ts, id, raw, **kw):
    pass

  def command(self, ns, ts, cmd, raw, **kw):
    pass

  def db_declare(self, ns, ts, **kw):
    pass

class OplogPrinter(OplogWatcher):
  def all(self, **kw):
    pprint (kw)
    print #newline

if __name__ == '__main__':
  OplogPrinter()

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

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数据库操作监控器,下一步就可以配合其他操作来对新入库的程序进行相应处理。

总结

以上就是这篇文章的全部内容了,希望本文的内容对大家的学习或者工作能带来一定的帮助,如果有疑问大家可以留言交流,谢谢大家对的支持。


# mongodb  # oplog  # 开启  # 查询  # MongoDB利用oplog恢复数据的方法  # 剖析后OpLog订阅MongoDB的数据变更就没那么难了  # mongodb中oplog介绍和格式详析  # Mongodb的oplog详解  # 关于单台MongoDB实例开启Oplog的过程详解  # 分享MongoDB修改oplog大小的4种方法  # 就会  # 就可以  # 来做  # 的是  # 是因为  # 基础上  # 可以用  # 我们可以  # 再来  # 不存在  # 几种  # 可以看出  # 这篇文章  # 自带  # 谢谢大家  # 很像  # 还有一些  # 命令行  # 配置文件  # 有一套 


相关栏目: 【 网站优化151355 】 【 网络推广146373 】 【 网络技术251813 】 【 AI营销90571


相关推荐: 文字头像制作网站推荐软件,醒图能自动配文字吗?  三星网站视频制作教程下载,三星w23网页如何全屏?  如何注册花生壳免费域名并搭建个人网站?  如何用y主机助手快速搭建网站?  Laravel如何使用Laravel Vite编译前端_Laravel10以上版本前端静态资源管理【教程】  html5怎么画眼睛_HT5用Canvas或SVG画眼球瞳孔加JS控制动态【绘制】  微信小程序 require机制详解及实例代码  微信推文制作网站有哪些,怎么做微信推文,急?  高性能网站服务器部署指南:稳定运行与安全配置优化方案  如何在 Go 中优雅地映射具有动态字段的 JSON 对象到结构体  Midjourney怎样加参数调细节_Midjourney参数调整技巧【指南】  历史网站制作软件,华为如何找回被删除的网站?  高端网站建设与定制开发一站式解决方案 中企动力  Linux系统命令中tree命令详解  Laravel怎么实现模型属性转换Casting_Laravel自动将JSON字段转为数组【技巧】  香港服务器租用每月最低只需15元?  Laravel如何设置自定义的日志文件名_Laravel根据日期或用户ID生成动态日志【技巧】  C#如何调用原生C++ COM对象详解  中山网站推广排名,中山信息港登录入口?  详解Android图表 MPAndroidChart折线图  如何彻底卸载建站之星软件?  PythonWeb开发入门教程_Flask快速构建Web应用  如何正确下载安装西数主机建站助手?  Laravel如何使用Service Provider服务提供者_Laravel依赖注入与容器绑定【深度】  Laravel怎么连接多个数据库_Laravel多数据库连接配置  网页设计与网站制作内容,怎样注册网站?  JS中使用new Date(str)创建时间对象不兼容firefox和ie的解决方法(两种)  谷歌Google入口永久地址_Google搜索引擎官网首页永久入口  百度浏览器如何管理插件 百度浏览器插件管理方法  Internet Explorer官网直接进入 IE浏览器在线体验版网址  Laravel如何使用Telescope进行调试?(安装和使用教程)  如何快速生成凡客建站的专业级图册?  php在windows下怎么调试_phpwindows环境调试操作说明【操作】  音响网站制作视频教程,隆霸音响官方网站?  北京网站制作公司哪家好一点,北京租房网站有哪些?  东莞专业网站制作公司有哪些,东莞招聘网站哪个好?  iOS正则表达式验证手机号、邮箱、身份证号等  php增删改查怎么学_零基础入门php数据库操作必知基础【教程】  HTML透明颜色代码怎么让图片透明_给img元素加透明色的技巧【方法】  大同网页,大同瑞慈医院官网?  Win10如何卸载预装Edge扩展_Win10卸载Edge扩展教程【方法】  焦点电影公司作品,电影焦点结局是什么?  Android 常见的图片加载框架详细介绍  edge浏览器无法安装扩展 edge浏览器插件安装失败【解决方法】  如何获取PHP WAP自助建站系统源码?  Laravel怎么为数据库表字段添加索引以优化查询  html5如何设置样式_HTML5样式设置方法与CSS应用技巧【教程】  香港服务器租用费用高吗?如何避免常见误区?  手机网站制作平台,手机靓号代理商怎么制作属于自己的手机靓号网站?  详解阿里云nginx服务器多站点的配置