借助MongoDB实现大数据的流式处理:从原理到实践的深度指南

引言

1.1 流式处理:大数据时代的"实时刚需"

在数字化浪潮下,数据正以"洪流"之势涌现——物联网设备每秒钟产生数百万条传感器读数,电商平台实时处理数十万笔交易,社交媒体每小时生成PB级用户行为数据。传统的批处理模式(如每天夜间运行MapReduce任务)已无法满足现代业务对"实时性"的需求:欺诈检测需要在交易发生时立即响应,物流调度依赖实时路况数据,个性化推荐必须基于用户最新行为调整。

流式处理(Stream Processing) 应运而生:它将数据视为持续流动的"流",在数据产生的瞬间进行实时捕获、处理和分析,最终输出即时结果。与批处理相比,流式处理具有三大核心优势:

  • 低延迟:数据处理延迟从小时级降至毫秒/秒级
  • 高吞吐:支持每秒数十万甚至数百万条记录的处理能力
  • 持续可用:7x24小时不间断运行,不中断数据流动

1.2 MongoDB:被低估的流式处理" Swiss Army Knife"

提到流式处理,多数工程师会首先想到Kafka、Flink、Spark Streaming等专用框架。但鲜为人知的是,MongoDB——这款以灵活文档模型著称的数据库,其实内置了强大的流式处理能力。

MongoDB的流式处理能力源于三大核心特性的深度融合:

  • Change Streams:实时监听数据库变更的"数据水龙头",基于MongoDB复制集的oplog实现
  • 时间序列集合(Time Series Collections):专为时序数据优化的存储引擎,提供高效写入、自动过期和降采样能力
  • 多模态数据支持:原生存储JSON、BSON、地理空间数据、二进制数据(如IoT设备的原始日志),无需预处理即可直接处理

这些特性使MongoDB成为流式处理的理想选择:它既能作为流数据的"着陆区"(Ingestion Layer),也能作为实时处理引擎的"计算平台",还能作为最终结果的"存储中心"。尤其对于需要实时数据访问+灵活处理+低架构复杂度的场景,MongoDB提供了"一站式"解决方案。

1.3 本文脉络:从原理到落地的完整指南

本文将系统讲解MongoDB流式处理的技术栈,包含四大部分:

  1. 基础篇:流式处理核心概念与MongoDB技术底座解析
  2. 原理篇:Change Streams深度剖析(从oplog到断点续传)
  3. 实战篇:完整案例(物联网实时监控系统从0到1实现)
  4. 进阶篇:高可用架构设计与性能优化指南

无论你是需要构建实时数据管道的架构师,还是处理海量时序数据的开发者,本文都将为你提供从理论到代码的全链路指导。

一、基础篇:流式处理与MongoDB技术底座

1.1 流式处理核心概念:从数据特征到架构模型

1.1.1 流数据的三大特征

流数据(Stream Data)是指持续生成、无界、顺序到达的数据序列,其核心特征可概括为"3V+1C":

  • Volume(海量):单条记录小(KB级),但总吞吐量高(如每秒10万条IoT读数)
  • Velocity(高速):数据产生与处理需实时同步,延迟要求通常<1秒
  • Variety(多样):结构灵活(JSON日志、CSV传感器数据、二进制文件)
  • Continuity(连续性):数据无"结束"概念,处理逻辑需长期运行
1.1.2 流式处理的核心挑战

相比批处理(Batch Processing),流式处理面临独特挑战:

  • 无界数据处理:无法等待全部数据到达后处理,需"来一条处理一条"
  • 状态管理:需维护中间计算状态(如滑动窗口内的平均值),且支持故障恢复
  • 数据乱序:网络延迟可能导致数据到达顺序与生成顺序不一致
  • 资源效率:需平衡实时性与计算资源消耗(避免过度分配资源)
1.1.3 主流流式处理架构模型

当前流式处理架构可分为三类:

  • 边缘处理(Edge Processing):数据在产生端(如IoT设备)就近处理,仅上传结果(适合带宽有限场景)
  • 流-批混合(Lambda Architecture):实时层(流处理)+批处理层(校准结果),如Storm+MapReduce
  • 统一处理(Kappa Architecture):用单一流处理引擎处理所有数据,如Flink的"流批一体"

MongoDB的流式处理能力主要服务于流-批混合统一处理架构,尤其擅长作为"实时数据枢纽"连接数据源与计算引擎。

1.2 MongoDB技术底座:为何它天生适合流式处理?

MongoDB作为文档型数据库,其核心特性与流式处理需求高度契合,可概括为"4个原生优势":

1.2.1 文档模型:天然适配半结构化流数据

流数据常包含动态字段(如IoT设备可能新增传感器类型),MongoDB的BSON文档模型支持:

  • 动态Schema:无需预定义表结构,新增字段自动适配
  • 嵌套结构:可直接存储复杂数据(如一条日志包含设备信息、传感器读数、位置坐标)
  • 多数据类型:原生支持int、string、date、array、binary等,无需数据类型转换

例如,一条IoT设备日志可直接写入MongoDB,无需拆分表:

{
  "deviceId": "sensor-001",
  "timestamp": ISODate("2024-05-20T12:34:56Z"),
  "readings": {
    "temperature": 23.5,
    "humidity": 65.2,
    "pressure": 1013.25  // 新增字段,无需修改表结构
  },
  "location": { "type": "Point", "coordinates": [116.4, 39.9] }  // 地理空间数据
}
1.2.2 复制集:流式处理的高可用基石

MongoDB复制集(Replica Set)由1个主节点(Primary)+多个从节点(Secondary)组成,为流式处理提供两大保障:

  • 数据冗余:主节点故障时,从节点自动切换(RTO<10秒),避免数据丢失
  • 读负载分流:流处理任务可连接从节点消费数据,不影响主节点写入性能

更重要的是,复制集的oplog(操作日志) 是MongoDB流式处理的"源头"——所有写入操作(insert/update/delete)都会记录到oplog,Change Streams正是基于此实现实时数据监听。

1.2.3 时间序列集合:时序流数据的"超级存储"

MongoDB 5.0引入的时间序列集合(Time Series Collections) 专为时序流数据优化,核心特性包括:

  • 自动分区:按时间自动分桶(Bucket)存储,默认每1小时一个桶,大幅提升查询效率
  • 降采样支持:内置$bucketAuto$group等聚合操作,高效计算滑动窗口指标
  • 数据自动过期:通过expireAfterSeconds配置TTL索引,自动删除历史数据(如保留30天)
  • 高压缩比:针对时序数据的重复字段(如设备ID)进行特殊压缩,存储成本降低60%+

创建时间序列集合的示例:

db.createCollection("sensor_readings", {
  timeseries: {
    timeField: "timestamp",  // 时间戳字段
    metaField: "deviceId",   // 元数据字段(用于分组)
    granularity: "minutes"   // 时间粒度(minutes/hours/days)
  },
  expireAfterSeconds: 2592000  // 30天后自动删除数据
})
1.2.4 聚合管道:流数据的实时计算引擎

MongoDB的聚合管道(Aggregation Pipeline) 支持在数据库内部对数据进行实时转换与计算,避免数据传输到应用层的开销。其核心算子包括:

  • 过滤($match):筛选符合条件的流数据(如温度>30℃的异常记录)
  • 分组($group):按设备ID、时间窗口分组计算(如每5分钟平均值)
  • 连接($lookup):关联其他集合数据(如将设备ID映射为设备名称)
  • 窗口函数($setWindowFields):计算滑动窗口/会话窗口指标(如最近10分钟最大温度)

例如,实时计算每个设备最近5分钟的平均温度:

db.sensor_readings.aggregate([
  { $match: { timestamp: { $gte: new Date(Date.now() - 5*60*1000) } } },
  { $group: {
      _id: "$deviceId",
      avgTemp: { $avg: "$readings.temperature" }
    }
  }
])

1.3 MongoDB流式处理技术栈全景图

MongoDB提供从数据接入→实时处理→存储→分析的全链路工具链,核心组件包括:

组件功能适用场景
Change Streams监听数据库变更,实时推送增量数据数据同步、实时告警
时间序列集合优化时序数据存储与查询IoT监控、日志存储
聚合管道实时数据转换与计算实时指标计算、数据清洗
MongoDB Kafka Connector与Kafka集成,作为流数据源/目标构建企业级数据管道
MongoDB Atlas Data Lake查询S3/GCS中的历史数据,与实时数据联合分析流批混合分析
Realm Sync边缘设备数据实时同步到云端边缘-云端协同处理(如移动应用)

后续章节将重点解析Change Streams(核心引擎)、时间序列集合(存储层)与聚合管道(计算层)的协同工作机制。

二、原理篇:Change Streams深度剖析——MongoDB流式处理的"心脏"

2.1 Change Streams本质:从oplog到实时数据流

2.1.1 oplog:MongoDB的"事务日志"

Change Streams的底层依赖MongoDB复制集的oplog(操作日志)。oplog是一个特殊的capped集合(固定大小的循环队列),存储所有数据库写入操作(insert/update/delete)的逻辑日志(而非物理日志)。

oplog记录的结构示例(简化版):

{
  "ts": Timestamp(1621507200, 1),  // 操作时间戳(秒+计数器)
  "t": 1,                          // 事务ID(用于多文档事务)
  "h": NumberLong("123456789"),    // 操作哈希值(唯一标识)
  "v": 2,                          // oplog版本
  "op": "i",                       // 操作类型(i:插入, u:更新, d:删除, c:命令)
  "ns": "iot_db.sensor_readings",  // 命名空间(数据库.集合)
  "o": { "_id": "123", "temp": 25 } // 操作内容(插入的文档)
}

关键特性

  • 顺序写入:oplog按操作发生顺序追加写入,保证时间线一致性
  • 循环覆盖:达到最大容量后,旧记录会被新记录覆盖(需合理配置大小)
  • 复制同步:从节点通过拉取主节点的oplog实现数据同步
2.1.2 Change Streams如何翻译oplog?

Change Streams并非直接暴露oplog,而是对其进行抽象与过滤,提供更友好的变更事件流。其工作流程如下:

  1. 建立游标:客户端通过watch()方法创建Change Stream,指定关注的集合与过滤条件
  2. 监听oplog:MongoDB后台线程监听oplog,当目标集合发生变更时触发事件
  3. 事件转换:将oplog记录转换为标准化的Change Event(包含操作类型、文档内容等)
  4. 推送给客户端:通过游标将Change Event异步推送给客户端

Change Event的标准结构:

{
  "_id": { "_data": "8260A..." },  // 事件唯一标识(包含oplog时间戳)
  "operationType": "insert",       // 操作类型(insert/update/delete/replace)
  "clusterTime": Timestamp(1621507200, 1),  // 集群时间戳
  "fullDocument": { "_id": "123", "temp": 25 },  // 完整文档(仅insert/replace有)
  "ns": { "db": "iot_db", "coll": "sensor_readings" },  // 命名空间
  "documentKey": { "_id": "123" }  // 被修改文档的主键
}
2.1.3 Change Streams vs 轮询:为何前者更高效?

传统的实时数据获取方式是轮询(Polling):客户端定期查询数据库(如每秒一次),检查是否有新数据。但这种方式存在两大问题:

  • 延迟高:数据产生到被查询到的延迟=轮询间隔(如1秒轮询则平均延迟0.5秒)
  • 资源浪费:大部分轮询请求返回空结果(无新数据时),浪费CPU/IO资源

Change Streams采用推模式(Push),仅在数据变更时才推送事件,优势显著:

  • 低延迟:数据变更后毫秒级响应(延迟取决于oplog同步速度)
  • 零空查询:仅在有数据时才通信,节省网络带宽与数据库资源
  • 断点续传:支持从上次断开的位置继续监听,避免数据丢失

测试数据显示:在1000 TPS写入场景下,Change Streams的平均延迟比1秒轮询低80%,数据库CPU占用降低60%

2.2 Change Streams核心能力:从基础监听到高级特性

2.2.1 基础监听:单集合数据变更捕获

创建Change Stream的基础语法(以Python为例,使用pymongo驱动):

from pymongo import MongoClient
from bson.json_util import dumps

client = MongoClient("mongodb://localhost:27017/")
db = client["iot_db"]
collection = db["sensor_readings"]

# 创建Change Stream,监听所有插入事件
change_stream = collection.watch([{"$match": {"operationType": "insert"}}])

# 循环接收事件
for change in change_stream:
    print(dumps(change))  # 处理新插入的文档

支持的操作类型过滤:

  • insert:新文档插入
  • update:文档更新(返回更新前后的字段)
  • delete:文档删除(仅返回documentKey)
  • replace:文档替换(返回新文档)
  • drop/rename/dropDatabase:集合/数据库级事件(需管理员权限)
2.2.2 高级过滤:按文档内容实时筛选

Change Streams支持通过聚合管道对变更事件进行过滤与转换,仅推送符合条件的事件。例如,仅监听温度>30℃的异常数据:

pipeline = [
    {"$match": {
        "operationType": "insert",
        "fullDocument.readings.temperature": {"$gt": 30}  # 过滤文档内容
    }}
]

change_stream = collection.watch(pipeline)

常用过滤场景:

  • 字段存在性"fullDocument.location": {"$exists": True}(仅含位置信息的记录)
  • 数组匹配"fullDocument.tags": {"$in": ["critical", "high-priority"]}
  • 多条件组合:结合$and/$or实现复杂逻辑(如温度>30℃且湿度<40%)
2.2.3 断点续传:从故障中恢复的"生命线"

网络中断或应用重启时,Change Streams需要支持断点续传(Resuming),避免重新处理历史数据。实现原理基于事件的_id字段——该字段包含oplog的时间戳(ts)与计数器(inc),可作为"恢复点"。

断点续传的实现步骤:

  1. 记录最后处理的事件ID:每次处理完事件后,持久化存储change["_id"](如写入本地文件或Redis)
  2. 重启时指定恢复点:创建Change Stream时通过resume_after参数指定上次记录的事件ID

代码示例:

import json

# 从文件加载上次的恢复点
def load_resume_token():
    try:
        with open("resume_token.json", "r") as f:
            return json.load(f)
    except FileNotFoundError:
        return None

# 保存当前事件ID作为恢复点
def save_resume_token(token):
    with open("resume_token.json", "w") as f:
        json.dump(token, f)

resume_token = load_resume_token()

# 创建支持断点续传的Change Stream
if resume_token:
    change_stream = collection.watch(resume_after=resume_token)
else:
    change_stream = collection.watch()

for change in change_stream:
    process_change(change)  # 处理事件
    save_resume_token(change["_id"])  # 保存恢复点

注意:MongoDB仅保留oplog的历史记录(默认配置下约24小时),若断点超过oplog保留期,需使用startAtOperationTime指定一个更早的时间点(需确保该时间点在oplog有效期内)。

2.2.4 多文档事务支持:确保数据一致性

MongoDB 4.0支持多文档事务,Change Streams会将事务中的所有变更事件合并为一个事务事件,避免部分事件被单独处理导致的数据不一致。

事务事件的结构示例:

{
  "_id": { ... },
  "operationType": "transaction",  // 事务类型标识
  "ts": Timestamp(1621507200, 1),  // 事务提交时间戳
  "events": [  // 事务包含的所有变更事件
    { "operationType": "insert", "fullDocument": { ... } },
    { "operationType": "update", "fullDocument": { ... } }
  ]
}

处理事务事件时,需:

  1. 检查operationType是否为"transaction"
  2. 遍历events数组,按顺序处理每个子事件
  3. 仅在所有子事件处理完成后,才更新恢复点(避免部分处理)

2.3 复制集环境:Change Streams的"基石"

2.3.1 为何Change Streams必须依赖复制集?

Change Streams的底层依赖oplog,而oplog仅存在于复制集环境(单节点实例无oplog)。因此,使用Change Streams前必须部署MongoDB复制集(至少1主1从1仲裁节点,测试环境可简化为单主1从)。

复制集的作用:

  • 数据冗余:主节点故障时,从节点可切换为主节点,确保oplog不丢失
  • 读写分离:Change Streams可连接从节点监听oplog,减轻主节点负载
  • 故障恢复:主节点宕机后,新主节点的oplog会继续记录写入操作
2.3.2 复制集配置:最小化延迟的关键参数

为确保Change Streams的低延迟,复制集需优化以下参数:

  • oplog大小:默认 oplog 大小为磁盘空间的5%,建议调大至10-20%(避免高频写入时 oplog 过快覆盖)
    # 启动时指定oplog大小(单位MB)
    mongod --replSet rs0 --oplogSize 102400  # 100GB
    
  • 从节点同步延迟:通过replSetSyncFrom指定从节点同步源(优先同步延迟低的节点)
  • 写关注(Write Concern):客户端写入时指定w: "majority",确保操作被大多数节点确认后才返回,避免主节点宕机导致oplog记录丢失
2.3.3 从节点读取:减轻主节点负载的最佳实践

默认情况下,Change Streams连接主节点监听oplog。但在高写入场景下,主节点负载较高,可将Change Streams路由到从节点(secondary),实现读写分离。

配置方式(Python示例):

client = MongoClient("mongodb://primary:27017,secondary:27017/?replicaSet=rs0")
db = client.get_database("iot_db", read_preference=ReadPreference.SECONDARY_PREFERRED)
# SECONDARY_PREFERRED:优先从从节点读取,若从节点不可用则从主节点读取
change_stream = db.sensor_readings.watch()

注意:从节点同步存在延迟(通常毫秒级,但极端情况可能达秒级),因此从节点Change Streams的事件推送会比主节点略晚。需根据业务对延迟的敏感度选择读取节点。

三、实战篇:物联网实时监控系统从0到1实现

3.1 项目背景与架构设计

3.1.1 业务需求:实时监控10万台设备的运行状态

假设我们需要构建一个物联网实时监控系统,核心需求:

  • 数据接入:接收10万台设备的传感器数据(每设备每30秒发送一条记录,约3333 TPS)
  • 实时存储:持久化存储原始数据(保留30天),支持按设备/时间范围查询
  • 异常检测:实时识别温度>30℃、湿度>80%的异常数据,触发告警
  • 实时看板:展示各设备的实时状态(当前温度、湿度)与最近1小时趋势
  • 历史分析:支持查询设备某天的温度变化曲线,计算日平均值/峰值
3.1.2 系统架构:MongoDB+Change Streams+Grafana

基于上述需求,设计系统架构如下:

[设备层] → [接入层] → [存储层] → [处理层] → [展示层]
  │          │          │          │          │
传感器 → MQTT Broker → MongoDB → Change Streams → Grafana
                                   ↓
                                 告警服务

各组件职责:

  • MQTT Broker:接收设备发送的MQTT协议数据(轻量级IoT协议)
  • MongoDB:存储原始数据(时间序列集合)、异常记录(普通集合)
  • Change Streams:监听原始数据变更,实时推送至处理层
  • 处理层:Python服务,接收Change Streams事件,检测异常并写入异常集合
  • Grafana:连接MongoDB,可视化实时/历史数据(通过MongoDB Data Source插件)

3.2 环境准备:从MongoDB部署到依赖安装

3.2.1 部署MongoDB复制集(生产级配置)
  1. 安装MongoDB(以Linux为例):

    # 添加MongoDB源
    echo "deb [ arch=amd64,arm64 ] https://repo.mongodb.org/apt/ubuntu focal/mongodb-org/6.0 multiverse" | sudo tee /etc/apt/sources.list.d/mongodb-org-6.0.list
    # 安装
    sudo apt update && sudo apt install -y mongodb-org
    
  2. 初始化复制集

    # 启动MongoDB(指定复制集名称rs0)
    mongod --replSet rs0 --dbpath /data/db --bind_ip_all &
    # 连接MongoDB Shell
    mongosh
    # 初始化复制集(主节点)
    rs.initiate({
      _id: "rs0",
      members: [
        {_id: 0, host: "mongodb-primary:27017"},
        {_id: 1, host: "mongodb-secondary:27017"}
      ]
    })
    
  3. 创建数据库与集合

    // 切换到iot_db数据库
    use iot_db
    
    // 创建时间序列集合(存储原始传感器数据)
    db.createCollection("sensor_readings", {
      timeseries: {
        timeField: "timestamp",
        metaField: "deviceId",
        granularity: "seconds"
      },
      expireAfterSeconds: 2592000  // 30天过期
    })
    
    // 创建普通集合(存储异常记录)
    db.createCollection("anomaly_records")
    
    // 创建索引(优化查询)
    db.sensor_readings.createIndex({ "deviceId": 1, "timestamp": -1 })
    db.anomaly_records.createIndex({ "deviceId": 1, "timestamp": -1 })
    
3.2.2 安装依赖工具与库
  1. MQTT Broker(Mosquitto):

    sudo apt install -y mosquitto mosquitto-clients
    
  2. Python依赖库

    pip install pymongo  # MongoDB驱动
    paho-mqtt           # MQTT客户端
    python-dotenv       # 环境变量管理
    requests            # 发送HTTP告警
    
  3. Grafana与MongoDB插件

    # 安装Grafana
    sudo apt install -y grafana
    # 启动Grafana
    sudo systemctl start grafana-server
    # 安装MongoDB Data Source插件(需重启Grafana)
    grafana-cli plugins install mongodbatlas-mongodb-datasource
    

3.3 数据接入层:从MQTT到MongoDB

3.3.1 模拟设备数据:Python脚本发送MQTT消息

编写Python脚本,模拟10万台设备发送数据:

import json
import time
import random
import paho.mqtt.client as mqtt
from datetime import datetime

# MQTT连接配置
MQTT_BROKER = "localhost"
MQTT_PORT = 1883
TOPIC = "sensor/data"

# 设备ID列表(模拟10万台设备)
DEVICE_IDS = [f"sensor-{i:06d}" for i in range(100000)]

def on_connect(client, userdata, flags, rc):
    if rc == 0:
        print("Connected to MQTT Broker")
    else:
        print(f"Failed to connect, rc={rc}")

client = mqtt.Client()
client.on_connect = on_connect
client.connect(MQTT_BROKER, MQTT_PORT)
client.loop_start()

try:
    while True:
        # 随机选择100台设备发送数据(模拟每30秒一轮)
        for device_id in random.sample(DEVICE_IDS, 100):
            # 生成随机温度(20-35℃)、湿度(40-90%)
            temperature = round(random.uniform(20, 35), 2)
            humidity = round(random.uniform(40, 90), 2)
            # 构造数据
            data = {
                "deviceId": device_id,
                "timestamp": datetime.utcnow().isoformat(),
                "readings": {
                    "temperature": temperature,
                    "humidity": humidity
                }
            }
            # 发送MQTT消息(设备ID作为消息键,便于分区)
            client.publish(
                topic=TOPIC,
                payload=json.dumps(data),
                qos=1  # 确保消息至少送达一次
            )
        time.sleep(30)  # 每30秒发送一轮
except KeyboardInterrupt:
    client.loop_stop()
    client.disconnect()
3.3.2 MQTT消费者:将数据写入MongoDB

编写MQTT消费者脚本,接收消息并写入MongoDB时间序列集合:

import json
import paho.mqtt.client as mqtt
from pymongo import MongoClient
from datetime import datetime
from dotenv import load_dotenv
import os

load_dotenv()  # 加载环境变量

# MongoDB配置
MONGO_URI = os.getenv("MONGO_URI", "mongodb://localhost:27017/")
DB_NAME = "iot_db"
COLLECTION_NAME = "sensor_readings"

# MQTT配置
MQTT_BROKER = os.getenv("MQTT_BROKER", "localhost")
MQTT_PORT = int(os.getenv("MQTT_PORT", 1883))
TOPIC = "sensor/data"

# 连接MongoDB
client = MongoClient(MONGO_URI)
db = client[DB_NAME]
collection = db[COLLECTION_NAME]

def on_connect(client, userdata, flags, rc):
    print(f"Connected with result code {rc}")
    client.subscribe(TOPIC)

def on_message(client, userdata, msg):
    try:
        # 解析MQTT消息
        data = json.loads(msg.payload.decode())
        # 转换timestamp为datetime类型(MongoDB时间序列集合要求)
        data["timestamp"] = datetime.fromisoformat(data["timestamp"])
        # 写入MongoDB
        collection.insert_one(data)
        print(f"Written to MongoDB: {data['deviceId']} at {data['timestamp']}")
    except Exception as e:
        print(f"Error processing message: {e}")

mqtt_client = mqtt.Client()
mqtt_client.on_connect = on_connect
mqtt_client.on_message = on_message
mqtt_client.connect(MQTT_BROKER, MQTT_PORT)
mqtt_client.loop_forever()

启动消费者后,数据将持续写入MongoDB的sensor_readings时间序列集合。

3.4 实时处理层:Change Streams实现异常检测

3.4.1 编写Change Streams监听服务

编写Python服务,通过Change Streams监听sensor_readings集合的插入事件,实时检测异常数据:

import json
import time
import requests
from pymongo import MongoClient
from pymongo.errors import PyMongoError
from bson.json_util import dumps
from dotenv import load_dotenv
import os

load_dotenv()

# 配置
MONGO_URI = os.getenv("MONGO_URI", "mongodb://localhost:27017/")
DB_NAME = "iot_db"
SOURCE_COLLECTION = "sensor_readings"
DEST_COLLECTION = "anomaly_records"
ALERT_WEBHOOK = os.getenv("ALERT_WEBHOOK", "http://localhost:5000/alert")
RESUME_TOKEN_FILE = "resume_token.json"

# 连接MongoDB
client = MongoClient(MONGO_URI)
db = client[DB_NAME]
source_collection = db[SOURCE_COLLECTION]
dest_collection = db[DEST_COLLECTION]

def load_resume_token():
    """加载上次的恢复点"""
    try:
        with open(RESUME_TOKEN_FILE, "r") as f:
            return json.load(f)
    except (FileNotFoundError, json.JSONDecodeError):
        return None

def save_resume_token(token):
    """保存当前恢复点"""
    with open(RESUME_TOKEN_FILE, "w") as f:
        json.dump(token, f)

def is_anomaly(data):
    """判断数据是否异常"""
    readings = data.get("readings", {})
    
    # 温度>30℃或湿度>80%为异常
    return (
        readings.get("temperature", 0) > 30 or
        readings.get("humidity", 0) > 80
    )

def send_alert(anomaly):
    """发送告警(调用Webhook)"""
    try:
        response = requests.post(
            ALERT_WEBHOOK,
            json={
                "deviceId": anomaly["deviceId"],
                "timestamp": anomaly["timestamp"].isoformat(),
                "readings": anomaly["readings"],
                reason: "High temperature" if anomaly["readings"]["temperature"] > 30 else "High humidity"
            },
            timeout=5
        )
        response.raise_for_status()
        print(f"Alert sent for {anomaly['deviceId']}")
    except Exception as e:
        print(f"Failed to send alert: {e}")

def process_change(change):
    """处理Change Stream事件"""
    # 提取完整文档
    full_doc = change.get("fullDocument", {})
    if not full_doc:
        return
    
    # 判断是否为异常
    if is_anomaly(full_doc):
        # 添加异常标记
        anomaly_doc = full_doc.copy()
        anomaly_doc["isAnomaly"] = True
        # 写入异常集合
        dest_collection.insert_one(anomaly_doc)
        # 发送告警
        send_alert(anomaly_doc)

def main():
    resume_token = load_resume_token()
    
    # 构建Change Stream管道(仅监听insert事件)
    pipeline = [{"$match": {"operationType": "insert"}}]
    
    # 创建Change Stream(支持断点续传)
    change_stream = source_collection.watch(
        pipeline,
        resume_after=resume_token,
        full_document="updateLookup"  # update事件返回完整文档(此处仅insert,可选)
    )
    
    print("Change Stream started. Listening for changes...")
    
    try:
        for change in change_stream:
            # 处理事件
            process_change(change)
            # 保存恢复点
            save_resume_token(change["_id"])
    except PyMongoError as e:
        print(f"MongoDB error: {e}")
    except KeyboardInterrupt:
        print("Stopping Change Stream...")
    finally:
        change_stream.close()

if __name__ == "__main__":
    main()
3.4.2 关键技术点解析
  1. 异常检测逻辑is_anomaly函数判断温度>30℃或湿度>80%,可根据需求扩展(如添加气压、振动等指标)。

  2. 断点续传实现:通过load_resume_token/save_resume_token函数持久化_id,确保服务重启后从上次位置继续监听。

  3. 错误处理:捕获PyMongoError(MongoDB连接异常),确保服务不会因临时网络问题崩溃。

  4. 性能优化

    • 仅监听insert事件(无需处理update/delete,因原始数据写入后不会修改)
    • 直接在MongoDB内部过滤(通过pipeline),减少传输到应用层的数据量

3.5 可视化层:Grafana实时看板

3.5.1 配置MongoDB数据源
  1. 访问Grafana(默认端口3000,用户名/密码admin/admin)
  2. 进入Configuration → Data Sources → Add data source
  3. 选择MongoDB,配置:
    • Name: MongoDB IoT
    • Connection String: mongodb://localhost:27017/iot_db
    • Database: iot_db
  4. 点击"Save & Test",确认连接成功
3.5.2 创建实时监控仪表板
  1. 新建Dashboard,添加Panel,选择"Time Series"图表
  2. 配置查询(从sensor_readings集合获取实时温度):
    // 查询最新温度(按设备分组)
    db.sensor_readings.aggregate([
      { $sort: { timestamp: -1 } },  // 按时间倒序
      { $group: { _id: "$deviceId", latest: { $first: "$$ROOT" } } },  // 每个设备取最新记录
      { $project: {
          deviceId: "$_id",
          timestamp: "$latest.timestamp",
          temperature: "$latest.readings.temperature"
        }
      }
    ])
    
  3. 设置可视化:X轴为timestamp,Y轴为temperature,按deviceId分组显示多条线
  4. 添加异常标记:新建Panel,查询anomaly_records集合,用"Scatter"类型标记异常点
3.5.3 历史趋势查询

创建"Historical Temperature Trend" Panel,查询设备某天的温度变化:

// 查询设备"sensor-000001"在2024-05-20的温度数据
db.sensor_readings.aggregate([
  { $match: {
      deviceId: "sensor-000001",
      timestamp: {
        $gte: ISODate("2024-05-20T00:00:00Z"),
        $lt: ISODate("2024-05-21T00:00:00Z")
      }
    }
  },
  { $sort: { timestamp: 1 } },  // 按时间正序
  { $project: {
      timestamp: 1,
      temperature: "$readings.temperature",
      _id: 0
    }
  }
])

配置图表类型为"Line",即可展示温度随时间的变化曲线。

四、进阶篇:高可用架构与性能优化

4.1 高可用设计:应对大规模与故障场景

4.1.1 分片集群:支撑百万级TPS写入

当单复制集无法满足写入需求(如设备数增至30万台,写入TPS达10000),需部署MongoDB分片集群(Sharded Cluster),将数据分布到多个分片(Shard)。

分片集群架构:

  • 分片(Shard):存储部分数据的复制集(每个分片是独立的复制集)
  • mongos:路由服务,客户端通过mongos访问集群,自动路由请求到对应分片
  • 配置服务器(Config Server):存储集群元数据(分片键范围、分片位置)

分片策略:

  • 分片键选择:选择deviceId作为分片键(基数高,分布均匀),按范围分片(Range Sharding)
  • 预分片:提前创建分片键范围(如按设备ID前两位分片),避免热点分片

创建分片集合的示例:

// 连接mongos
mongosh "mongodb://mongos1:27017/"

// 启用分片数据库
sh.enableSharding("iot_db")

// 设置分片键(deviceId为分片键)
sh.shardCollection("iot_db.sensor_readings", { "deviceId": "hashed" })  // 哈希分片(均匀分布)
4.1.2 多实例Change Streams:避免单点故障

处理层的Python服务若为单实例,可能因进程崩溃导致流处理中断。解决方案:

  • 部署多个Change Streams实例:连接同一MongoDB集群,监听相同集合
  • 使用消费者组(Consumer Group):确保同一条数据仅被一个实例处理(避免重复告警)

MongoDB 6.0+支持Change Stream Consumer Groups,通过startAfter+groupId实现:

change_stream = collection.watch(
    pipeline,
    start_after=resume_token,
    group_id="anomaly-detection-group"  # 消费者组ID
)

消费者组的工作原理:

  • 同一groupId的多个Change Stream实例共享一个"游标"
  • MongoDB确保每个事件仅推送给组内一个实例(轮询分配)
  • 实例故障后,其未处理的事件会自动分配给组内其他实例
4.1.3 数据备份:避免历史数据丢失

尽管时间序列集合配置了TTL自动删除,但仍需定期备份关键数据(如异常记录):

  • 逻辑备份:使用mongodump导出数据(适合小数据量)
    mongodump --uri "mongodb://localhost:27017/" --db iot_db --collection anomaly_records --out /backup/$(date +%Y%m%d)
    
  • 物理备份:复制集从节点执行mongod --dumpDbPaths(适合大数据量,速度快)
  • 云服务备份:MongoDB Atlas提供自动备份功能(每日快照,支持时间点恢复)

4.2 性能优化:从毫秒级延迟到万级TPS

4.2.1 MongoDB性能调优:索引与配置
  1. 索引优化

    • deviceId+timestamp创建复合索引(支持按设备+时间范围查询)
      db.sensor_readings.createIndex({ "deviceId": 1, "timestamp": -1 })
      
    • 时间序列集合无需为timeField单独创建索引(MongoDB自动优化)
  2. 写入优化

    • 使用批量写入:将多条记录合并为insert_many,减少网络往返
      # 批量写入示例(每100条一批)
      batch = []
      for data in mqtt_messages:
          batch.append(data)
          if len(batch) >= 100:
              collection.insert_many(batch)
              batch = []
      if batch:
          collection.insert_many(batch)
      
    • 调整写入关注:非关键数据使用w: 1(仅主节点确认),降低延迟
      collection.insert_one(data, write_concern=WriteConcern(w=1))
      
  3. MongoDB配置优化

    • 增大wiredTigerCacheSizeGB(缓存大小,建议设为物理内存的50%)
    • 启用journal(事务日志)确保崩溃后数据可恢复
    • 调整oplogSizeMB(oplog大小,高写入场景设为100GB+)
4.2.2 Change Streams性能调优
  1. 减少事件传输量

    • 通过pipeline过滤不需要的字段(仅保留deviceIdtimestampreadings
      pipeline = [
          {"$match": {"operationType": "insert"}},
          {"$project": {
              "fullDocument.deviceId": 1,
              "fullDocument.timestamp": 1,
              "fullDocument.readings": 1,
              "_id": 0
          }}
      ]
      
  2. 控制批处理大小

    • 通过batchSize参数调整每次返回的事件数量(默认100,可根据网络带宽调整)
      change_stream = collection.watch(pipeline, batch_size=500)
      
  3. 异步处理事件

    • 使用Python的asyncio+motor(异步MongoDB驱动),避免处理事件阻塞流接收
      import asyncio
      from motor.motor_asyncio import AsyncIOMotorClient
      
      async def process_event(event):
          # 异步处理事件(写入异常集合、发送告警)
          await dest_collection.insert_one(anomaly_doc)
      
      async def main():
          client = AsyncIOMotorClient(MONGO_URI)
          collection = client[DB_NAME][SOURCE_COLLECTION]
          async with collection.watch(pipeline) as change_stream:
              async for change in
      

更多推荐