借助MongoDB实现大数据的流式处理
借助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流式处理的技术栈,包含四大部分:
- 基础篇:流式处理核心概念与MongoDB技术底座解析
- 原理篇:Change Streams深度剖析(从oplog到断点续传)
- 实战篇:完整案例(物联网实时监控系统从0到1实现)
- 进阶篇:高可用架构设计与性能优化指南
无论你是需要构建实时数据管道的架构师,还是处理海量时序数据的开发者,本文都将为你提供从理论到代码的全链路指导。
一、基础篇:流式处理与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,而是对其进行抽象与过滤,提供更友好的变更事件流。其工作流程如下:
- 建立游标:客户端通过
watch()方法创建Change Stream,指定关注的集合与过滤条件 - 监听oplog:MongoDB后台线程监听oplog,当目标集合发生变更时触发事件
- 事件转换:将oplog记录转换为标准化的Change Event(包含操作类型、文档内容等)
- 推送给客户端:通过游标将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),可作为"恢复点"。
断点续传的实现步骤:
- 记录最后处理的事件ID:每次处理完事件后,持久化存储
change["_id"](如写入本地文件或Redis) - 重启时指定恢复点:创建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": { ... } }
]
}
处理事务事件时,需:
- 检查
operationType是否为"transaction" - 遍历
events数组,按顺序处理每个子事件 - 仅在所有子事件处理完成后,才更新恢复点(避免部分处理)
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复制集(生产级配置)
-
安装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 -
初始化复制集:
# 启动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"} ] }) -
创建数据库与集合:
// 切换到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 安装依赖工具与库
-
MQTT Broker(Mosquitto):
sudo apt install -y mosquitto mosquitto-clients -
Python依赖库:
pip install pymongo # MongoDB驱动 paho-mqtt # MQTT客户端 python-dotenv # 环境变量管理 requests # 发送HTTP告警 -
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 关键技术点解析
-
异常检测逻辑:
is_anomaly函数判断温度>30℃或湿度>80%,可根据需求扩展(如添加气压、振动等指标)。 -
断点续传实现:通过
load_resume_token/save_resume_token函数持久化_id,确保服务重启后从上次位置继续监听。 -
错误处理:捕获
PyMongoError(MongoDB连接异常),确保服务不会因临时网络问题崩溃。 -
性能优化:
- 仅监听
insert事件(无需处理update/delete,因原始数据写入后不会修改) - 直接在MongoDB内部过滤(通过pipeline),减少传输到应用层的数据量
- 仅监听
3.5 可视化层:Grafana实时看板
3.5.1 配置MongoDB数据源
- 访问Grafana(默认端口3000,用户名/密码admin/admin)
- 进入Configuration → Data Sources → Add data source
- 选择MongoDB,配置:
- Name: MongoDB IoT
- Connection String:
mongodb://localhost:27017/iot_db - Database:
iot_db
- 点击"Save & Test",确认连接成功
3.5.2 创建实时监控仪表板
- 新建Dashboard,添加Panel,选择"Time Series"图表
- 配置查询(从
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" } } ]) - 设置可视化:X轴为
timestamp,Y轴为temperature,按deviceId分组显示多条线 - 添加异常标记:新建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性能调优:索引与配置
-
索引优化:
- 为
deviceId+timestamp创建复合索引(支持按设备+时间范围查询)db.sensor_readings.createIndex({ "deviceId": 1, "timestamp": -1 }) - 时间序列集合无需为
timeField单独创建索引(MongoDB自动优化)
- 为
-
写入优化:
- 使用批量写入:将多条记录合并为
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))
- 使用批量写入:将多条记录合并为
-
MongoDB配置优化:
- 增大
wiredTigerCacheSizeGB(缓存大小,建议设为物理内存的50%) - 启用
journal(事务日志)确保崩溃后数据可恢复 - 调整
oplogSizeMB(oplog大小,高写入场景设为100GB+)
- 增大
4.2.2 Change Streams性能调优
-
减少事件传输量:
- 通过pipeline过滤不需要的字段(仅保留
deviceId、timestamp、readings)pipeline = [ {"$match": {"operationType": "insert"}}, {"$project": { "fullDocument.deviceId": 1, "fullDocument.timestamp": 1, "fullDocument.readings": 1, "_id": 0 }} ]
- 通过pipeline过滤不需要的字段(仅保留
-
控制批处理大小:
- 通过
batchSize参数调整每次返回的事件数量(默认100,可根据网络带宽调整)change_stream = collection.watch(pipeline, batch_size=500)
- 通过
-
异步处理事件:
- 使用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
- 使用Python的
更多推荐
所有评论(0)