1. 项目概述从“监控”到“存储”的自动化桥梁最近在折腾一个物联网数据采集的小项目传感器数据源源不断地进来我需要一个既能实时监控数据流状态又能把有效数据持久化存储的方案。这听起来像是两个独立的任务监控告警和数据库写入。一开始我也打算分开处理用个脚本轮询再写个服务入库但很快就被各种边缘情况搞到头大——脚本挂了怎么办数据格式突变怎么预警写入失败如何重试直到我把Watcher和MongoDB这两个工具组合起来才发现原来有一条如此顺畅的“快速通道”。简单来说Watcher 负责盯梢一旦发现符合条件的数据变化或事件就自动触发动作而 MongoDB 则作为海量、灵活的数据仓库接收并存储这些被“捕获”的数据。这个组合的核心价值在于自动化和解耦。你不再需要编写冗长的轮询逻辑和异常处理代码只需声明“当XX发生时把数据存到YY地方”系统就会自动、可靠地执行。这非常适合哪些场景呢如果你正在处理日志文件的实时分析、物联网设备状态监控、用户行为事件追踪、或者任何需要将特定事件数据落库的场景这个模式都能大幅降低开发复杂度。它把“事件响应”和“数据持久化”这两个关注点分离让每个部分更专注、更健壮。接下来我就结合自己的踩坑经验带你快速上手这套组合拳实现从事件监控到数据入库的无缝衔接。2. 核心组件解析Watcher 与 MongoDB 的角色定位在搭建这条自动化流水线之前我们必须先吃透两个核心组件各自的能力边界和设计哲学。理解它们“为什么”这样设计比记住“怎么用”更重要这能帮助我们在后续配置中做出更合理的选择。2.1 Watcher不只是“监视器”更是“事件路由器”很多人一听到 Watcher就想到简单的文件变化监听。这低估了它。在现代应用架构中Watcher 本质上是一个轻量级的事件驱动引擎。它的核心职责是持续监听一个或多个数据源如文件、目录、API端点、数据库的Oplog、消息队列并在预定义的条件被满足时例如新文件创建、内容匹配特定模式、数据超过阈值触发一个或多个后续动作。我选择 Watcher 类工具通常基于以下几个关键特性低侵入性它通常以独立进程或服务运行不需要你大幅修改现有业务代码。你只需要告诉它“监视什么”和“触发什么”。条件过滤能力这是它的智慧所在。不是所有变化都值得响应。一个强大的 Watcher 允许你通过正则表达式、脚本判断、字段匹配等方式精细过滤事件避免无意义的动作触发和数据洪流。动作的多样性与可扩展性触发动作不应只是记录日志。优秀的 Watcher 支持调用 HTTP API、执行 Shell 命令、发送邮件/消息通知当然也包括写入数据库。这使其成为系统集成的枢纽。可靠性机制比如具备重试策略当动作执行失败时自动重试、死信队列处理反复失败的事件以及状态持久化服务重启后不丢失监控状态这些都是生产环境不可或缺的。在本次快速入门中我们将 Watcher 视为一个高度可配置的事件侦测与转发代理。它填补了数据生产端和消费端MongoDB之间的空白让数据流动变得规则化、自动化。2.2 MongoDB非关系型数据库的灵活性与陷阱MongoDB 作为知名的 NoSQL 数据库以其灵活的文档模型著称。它不像传统关系型数据库那样需要预先严格定义表结构Schema每条记录文档都可以拥有不同的字段。这种灵活性正是应对 Watcher 所捕获的、可能结构多变的事件数据的绝佳选择。但是“灵活”不等于“随意”。这是我用 MongoDB 存储事件数据时最深刻的体会。初期我任由 Watcher 将各种形态的数据直接塞入 MongoDB很快遇到了问题查询性能低下没有索引的字段在数据量增长后查询慢如蜗牛。数据含义模糊同样的userId字段有时是字符串有时是数字导致应用层逻辑复杂且容易出错。存储膨胀随意嵌套的文档结构可能导致存储空间非预期增长。因此即便使用 MongoDB我也强烈建议实施“宽松但有约束的 Schema 设计”。例如为所有文档定义一个基础结构包含如event_type、timestamp、source等通用字段而事件主体数据则放入一个如payload的字段中。同时必须为高频查询字段建立索引比如timestamp和event_type。MongoDB 的灵活性应该用在应对业务字段的增减上而不是在基础数据规范上放飞自我。另一个关键点是写入考量。Watcher 触发写入可能是并发的。我们需要根据数据重要性在 MongoDB 写入关注级别Write Concern上做出选择。对于监控日志可能使用{w: 0}非阻塞写入以追求速度而对于关键业务事件则应使用{w: “majority”}以确保数据已安全写入多数节点。这需要在 Watcher 的 MongoDB 输出动作中进行配置。3. 环境准备与工具选型搭建你的实验舞台理论讲完了我们动手搭建环境。这里我会提供两种主流的实现路径一种是基于成熟开源软件的“组装”方案另一种是编写轻量脚本的“自制”方案。你可以根据自身技术栈和需求选择。3.1 方案一使用 Node-RED 作为可视化 Watcher对于快速原型验证、运维人员或不希望写太多代码的开发者Node-RED是一个神器。它是一个基于流的低代码编程工具通过连接不同的节点Node来创建应用。我们可以用它轻松构建一个功能强大的 Watcher。安装与启动如果你本地有 Node.js 环境一行命令即可安装npm install -g node-red安装后在终端运行node-red浏览器打开http://localhost:1880就能看到流编辑界面。对于 Windows 用户也可以下载免安装的压缩包解压后运行node-red.exe即可。核心节点介绍在 Node-RED 中我们需要关注这几类节点输入节点Watcher 部分inject手动触发或定时触发模拟事件。watch监视文件或目录的变化。可能需要额外安装节点如node-red-node-watcherhttp in监听 HTTP 请求将 API 调用转化为事件。mqtt in订阅 MQTT 主题接收物联网消息。处理节点过滤与转换function编写 JavaScript 代码对消息进行过滤、判断、格式转换。这是实现复杂条件逻辑的核心。switch根据消息属性路由到不同分支。change设置、修改或删除消息属性。输出节点写入 MongoDBmongodb官方或社区提供的 MongoDB 节点用于执行插入、查询等操作。你需要先配置 MongoDB 连接信息。一个简单流程示例从左侧面板拖入一个inject节点配置它每10秒触发一次。拖入一个function节点连接到inject节点后。在函数节点中我们可以生成或处理数据并决定是否继续传递。例如// 模拟一个温度传感器事件 var temperature Math.random() * 30 10; // 10-40度随机数 msg.payload { sensorId: temp_001, value: temperature, unit: °C, timestamp: new Date() }; // 添加一个事件类型 msg.event_type sensor_reading; // 假设只关心高温报警 if (temperature 35) { msg.high_temp_alert true; return msg; // 只有温度35才传递下去 } // 其他情况我们可以选择不传递或者传递到另一条分支 return null; // 丢弃此消息拖入mongodb节点连接到function节点后。双击配置填入你的 MongoDB 连接字符串如mongodb://localhost:27017/mydb并选择操作类型为“insert”。点击右上角“部署”按钮流就开始运行了。每10秒它会生成一个随机温度只有高温数据会被插入 MongoDB。注意Node-RED 的function节点功能强大但也是容易出错的点。务必做好异常捕获避免因为一个消息处理错误导致整个流停止。可以在关键函数外用try...catch包裹。3.2 方案二使用 Python 脚本构建定制化 Watcher如果你需要更精细的控制、复杂的业务逻辑或者希望将监控逻辑深度集成到现有 Python 项目中自己编写脚本是更灵活的选择。Python 生态中有大量优秀的库支持。核心库选择监视库对于文件系统监控watchdog库是行业标准高效且易用。对于监听 HTTP 请求可以使用Flask或FastAPI快速搭建 API 端点。对于监听消息队列有pika(RabbitMQ)、kafka-python等。MongoDB 驱动pymongo是官方驱动稳定且功能全面。异步框架可选如果事件频率很高考虑使用asyncio配合异步版本的库如motor用于异步 MongoDB 操作来提升吞吐量。一个基于watchdog和pymongo的示例脚本骨架import time import json from watchdog.observers import Observer from watchdog.events import FileSystemEventHandler from pymongo import MongoClient from pymongo.errors import ConnectionFailure, DuplicateKeyError import logging # 配置日志 logging.basicConfig(levellogging.INFO, format%(asctime)s - %(levelname)s - %(message)s) logger logging.getLogger(__name__) class MyEventHandler(FileSystemEventHandler): def __init__(self, mongo_collection): self.collection mongo_collection # 可以在这里初始化一些状态或缓存 def on_created(self, event): # 当新文件创建时触发 if not event.is_directory: self.process_file(event.src_path) def on_modified(self, event): # 当文件被修改时触发注意某些编辑器保存文件会多次触发 if not event.is_directory: # 为了避免频繁触发可以加入防抖逻辑 self.process_file(event.src_path) def process_file(self, file_path): try: logger.info(fProcessing file: {file_path}) # 1. 读取并解析文件内容假设是JSON格式的日志 with open(file_path, r, encodingutf-8) as f: # 这里可以按行读取处理大文件 content json.load(f) # 2. 业务逻辑过滤例如只处理特定类型的日志 if content.get(level) ! ERROR: logger.debug(fIgnore non-ERROR log: {content.get(level)}) return # 3. 数据增强添加处理时间、来源等元数据 document { **content, # 原始内容 _source_file: file_path, _processed_at: time.time(), event_type: file_log_error } # 4. 写入 MongoDB result self.collection.insert_one(document) logger.info(fDocument inserted with id: {result.inserted_id}) # 5. 可选处理成功后可以归档或删除源文件 # os.remove(file_path) except json.JSONDecodeError as e: logger.error(fFailed to parse JSON from {file_path}: {e}) except IOError as e: logger.error(fFailed to read file {file_path}: {e}) except ConnectionFailure as e: logger.error(fMongoDB connection failed: {e}) # 这里应实现重试逻辑或放入重试队列 except Exception as e: logger.exception(fUnexpected error processing {file_path}: {e}) def main(): # 连接 MongoDB try: client MongoClient(mongodb://localhost:27017/, serverSelectionTimeoutMS5000) # 快速检查连接 client.admin.command(ismaster) db client[event_db] collection db[error_logs] # 确保有索引 collection.create_index([(timestamp, -1)]) collection.create_index([(event_type, 1)]) logger.info(Connected to MongoDB successfully.) except ConnectionFailure as e: logger.error(fCould not connect to MongoDB: {e}) return # 设置文件监视 event_handler MyEventHandler(collection) observer Observer() path_to_watch /path/to/your/log/directory # 替换为你的目录 observer.schedule(event_handler, path_to_watch, recursiveTrue) observer.start() logger.info(fStarted watching directory: {path_to_watch}) try: while True: time.sleep(1) except KeyboardInterrupt: observer.stop() logger.info(Watcher stopped by user.) observer.join() client.close() if __name__ __main__: main()这个脚本展示了一个完整的生命周期连接数据库、监听文件事件、过滤业务数据、增强元信息、安全写入。你可以在此基础上扩展比如监听多个目录、整合多种事件源、增加更复杂的聚合逻辑等。4. 连接与配置实战打通 Watcher 与 MongoDB环境准备好了现在我们来完成最关键的一步建立可靠的数据通道。这里会涉及连接配置、数据格式设计以及一些提升鲁棒性的高级技巧。4.1 MongoDB 连接配置与安全最佳实践无论使用哪种 Watcher连接 MongoDB 的第一步都是配置连接字符串。一个基础的连接字符串如下mongodb://username:passwordhost:port/database?authSourceadmin关键参数解析与避坑指南认证生产环境务必使用用户名密码。authSource参数指定了用户凭证所在的数据库通常是admin。切勿将连接字符串硬编码在代码中应使用环境变量或配置文件管理。连接选项serverSelectionTimeoutMS5000指定尝试连接服务器的超时时间毫秒。避免网络波动时脚本长时间挂起。socketTimeoutMS60000套接字操作如查询、插入的超时时间。对于 Watcher如果 MongoDB 响应慢这个超时可以防止线程被无限阻塞。maxPoolSize50连接池最大大小。即使你的 Watcher 是单线程的使用连接池也能复用连接提升效率。根据并发量调整。retryWritestrue启用写入重试。在网络瞬断或主节点切换时驱动会自动重试写入操作提高可靠性。TLS/SSL如果 MongoDB 启用了 TLS需要在连接字符串中加入?tlstrue或?ssltrue并可能需要配置证书路径。在 Node-RED 中配置在 mongodb 节点的配置面板中通常有一个“Add new mongodb-config”的选项在那里你可以填入连接字符串和给配置起个名字方便多个流复用。在 Python (pymongo) 中配置from pymongo import MongoClient import os # 从环境变量读取更安全 MONGO_URI os.getenv(MONGO_URI, mongodb://localhost:27017/) client MongoClient( MONGO_URI, serverSelectionTimeoutMS5000, socketTimeoutMS30000, maxPoolSize30, retryWritesTrue )4.2 数据模型设计与写入优化数据怎么存直接决定了以后怎么用。对于 Watcher 捕获的事件数据我推荐以下文档结构{ _id: ObjectId(...), // MongoDB自动生成 event_type: user_login, // 事件类型用于分类和索引 timestamp: ISODate(2023-10-27T08:30:00Z), // 事件发生时间务必用Date类型 source: api_gateway, // 事件来源 severity: info, // 严重级别debug, info, warn, error, fatal payload: { // 事件主体灵活存储 user_id: u123456, ip_address: 192.168.1.1, user_agent: Mozilla/5.0..., metadata: {...} // 其他任意数据 }, processed_at: ISODate(2023-10-27T08:30:05Z), // Watcher处理时间 tags: [authentication, security] // 标签便于多维查询 }设计理由与优化技巧分离元数据与主体数据event_type,timestamp,source等是通用的元数据适合单独字段并建立索引。payload字段容纳可变的主体内容保持了 MongoDB 的灵活性优势。使用正确的日期类型timestamp务必存为 MongoDB 的 Date 类型 (ISODate)而不是字符串或数字时间戳。这样可以利用 MongoDB 丰富的日期查询运算符$gte,$lte,$dayOfMonth等并且日期索引效率最高。索引策略至少为{timestamp: -1}和{event_type: 1, timestamp: -1}创建复合索引。前者用于按时间倒序查看最新事件后者是查询特定类型事件的最常用模式。根据查询需求可能还需要在source、severity或tags上建索引。写入关注级别 (Write Concern)在 Watcher 中配置。对于可丢失的监控数据可以用{w: 0}非确认写入换取极致速度。对于关键业务事件使用{w: “majority”}和{j: true}日志刷盘来保证持久性。在 pymongo 中可以在集合级别或每次插入操作时指定collection.with_options(write_concernWriteConcern(w”majority”))。批量插入如果 Watcher 可能短时间内产生大量事件如日志文件批量解析使用insert_many()比循环insert_one()性能高出几个数量级。但要注意insert_many()默认是顺序执行遇到错误会停止。可以设置orderedFalse来忽略单条错误继续插入但需自行处理错误列表。4.3 错误处理与重试机制构建健壮的流水线网络会抖动数据库可能临时不可用磁盘可能写满。一个生产级的 Watcher 必须能优雅地处理这些故障。1. 连接失败与重连在脚本初始化时如上面 Python 示例使用serverSelectionTimeoutMS进行快速失败检测。一旦连接建立驱动通常有内置的心跳机制维持连接。但你仍需在代码中捕获pymongo.errors.ConnectionFailure等异常。一个简单的重连逻辑可以这样实现import time from pymongo import MongoClient from pymongo.errors import ConnectionFailure, AutoReconnect def get_mongo_client(max_retries5, delay2): retries 0 while retries max_retries: try: client MongoClient(your_uri, serverSelectionTimeoutMS5000) # 发送一个简单命令测试连接 client.admin.command(ismaster) return client except (ConnectionFailure, AutoReconnect) as e: retries 1 logger.warning(fFailed to connect to MongoDB (attempt {retries}/{max_retries}): {e}) if retries max_retries: time.sleep(delay * retries) # 指数退避 else: logger.error(Max retries reached. Could not connect to MongoDB.) raise2. 写入失败与重试队列对于单次写入失败如主节点切换导致的短暂不可写可以利用驱动自带的retryWritestrue。但对于更复杂的业务逻辑失败如数据格式校验不通过需要业务层面的重试。一个常见的模式是引入一个“本地死信队列”。当写入 MongoDB 失败时先将事件数据写入一个本地文件或一个轻量级本地数据库如 SQLite。然后由一个独立的恢复进程定期检查这个队列并尝试重新投递。import sqlite3 import json class DeadLetterQueue: def __init__(self, db_pathdead_letter.db): self.conn sqlite3.connect(db_path) self._init_db() def _init_db(self): cursor self.conn.cursor() cursor.execute( CREATE TABLE IF NOT EXISTS failed_events (id INTEGER PRIMARY KEY AUTOINCREMENT, event_data TEXT NOT NULL, error TEXT, failed_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, retry_count INTEGER DEFAULT 0) ) self.conn.commit() def put(self, event_data, error_msg): cursor self.conn.cursor() cursor.execute(INSERT INTO failed_events (event_data, error) VALUES (?, ?), (json.dumps(event_data), str(error_msg))) self.conn.commit() def get_retry_batch(self, limit10): cursor self.conn.cursor() cursor.execute(SELECT id, event_data FROM failed_events WHERE retry_count 3 ORDER BY failed_at LIMIT ?, (limit,)) return cursor.fetchall() def mark_success(self, event_id): cursor self.conn.cursor() cursor.execute(DELETE FROM failed_events WHERE id?, (event_id,)) self.conn.commit() def mark_retry(self, event_id): cursor self.conn.cursor() cursor.execute(UPDATE failed_events SET retry_count retry_count 1 WHERE id?, (event_id,)) self.conn.commit()在你的主处理逻辑中捕获写入异常后调用dlq.put(event_data, str(e))。然后可以启动一个定时任务如另一个线程或 cron job来调用get_retry_batch()并重新尝试插入 MongoDB。5. 高级模式与性能调优当你的 Watcher-MongoDB 流水线稳定运行后可能会面临数据量增长、性能要求提高等挑战。下面分享几个进阶模式和调优思路。5.1 模式一基于 MongoDB Change Streams 的递归监听这是一个非常强大的模式它让 Watcher 本身也监听 MongoDB 的数据变化实现数据处理的管道化。例如你可以用一个 Watcher 将原始日志存入raw_events集合然后利用MongoDB 的 Change Streams功能监听这个集合的变化。当有新文档插入时触发另一个处理程序可以是另一个 Watcher 或函数对数据进行清洗、丰富再写入processed_events集合。优势解耦了数据接收和数据处理每个阶段可以独立扩展和容错。实现要点使用pymongo的watch()方法监听集合。注意Change Streams 需要副本集或分片集群环境单机 MongoDB 默认不支持。pipeline [{$match: {operationType: insert}}] # 只监听插入操作 with db.raw_events.watch(pipeline) as stream: for change in stream: raw_doc change[fullDocument] # 在这里进行你的数据处理逻辑 processed_doc transform_data(raw_doc) db.processed_events.insert_one(processed_doc)5.2 模式二分片集群下的写入考量如果事件数据量非常庞大单台 MongoDB 服务器可能成为瓶颈。这时需要考虑使用MongoDB 分片集群。分片的核心是选择一个好的分片键。对于时间序列类的事件数据以timestamp字段作为分片键是常见选择这能将新数据均匀分散到不同分片。但是这可能导致“热点”分片所有最新数据都写入同一个分片。一个更好的策略是使用复合分片键例如{source: 1, timestamp: 1}或{event_type: 1, timestamp: 1}。这样数据首先按来源或类型分散再按时间排序既能分散写入压力又能在查询特定类型或来源的数据时利用局部性。在 Watcher 配置上连接到分片集群与连接到单机实例没有区别连接字符串指向mongos路由器即可。但你需要预先在 MongoDB 中启用分片并对数据库、集合进行分片设置。5.3 性能监控与调优要点监控 Watcher 自身记录它处理的事件数量、成功率、延迟。可以在 Watcher 中增加代码将自身的运行指标如events_processed_per_minute,avg_processing_latency也写入一个专门的 MongoDB 监控集合便于用同一套工具查看。监控 MongoDB关注db.serverStatus()中的opcounters操作计数器、connections连接数以及wiredTiger.cache缓存利用率。写入压力大时可能会看到opcounters.insert持续高位连接数增长。此时需要考虑批量插入、优化索引避免过多索引影响写入速度甚至水平扩展。索引优化定期使用db.collection.explain()分析慢查询。使用db.collection.totalIndexSize()查看索引大小避免索引体积超过数据本身。对于只用于归档、很少查询的历史数据集合可以考虑删除非必要的索引来提升写入性能。TTL 索引自动清理很多监控事件数据具有时效性。可以利用 MongoDB 的TTL生存时间索引自动删除旧数据。例如为timestamp字段创建一个 30 天过期的 TTL 索引db.events.createIndex({“timestamp”: 1}, {expireAfterSeconds: 2592000})。这样就不需要额外写清理脚本了。6. 常见问题排查与实战技巧在实际操作中你一定会遇到各种报错和意外情况。下面是我总结的一些典型问题及其解决方法。6.1 连接与认证问题错误pymongo.errors.ServerSelectionTimeoutError可能原因网络不通、MongoDB 服务未启动、防火墙阻挡、连接字符串中的主机名或端口错误。排查在 Watcher 机器上用telnet mongodb_host mongodb_port测试网络连通性。检查 MongoDB 服务状态 (sudo systemctl status mongod或ps aux | grep mongod)。确认连接字符串中的用户名、密码、认证数据库 (authSource) 是否正确。如果是远程连接检查 MongoDB 配置bindIp是否允许了 Watcher 的 IP。错误pymongo.errors.OperationFailure: Authentication failed.可能原因密码错误、用户权限不足。排查用 MongoDB Shell 或 Compass 使用相同凭证尝试连接。确认用户是否在目标数据库上有readWrite角色。创建用户的命令示例db.createUser({user: “watcher”, pwd: “password”, roles: [{role: “readWrite”, db: “event_db”}]})。6.2 数据写入与格式问题错误bson.errors.InvalidDocument: Cannot encode object可能原因试图插入一个包含 MongoDB 不支持数据类型的文档比如 Python 的datetime对象应使用bson.datetime.datetime但pymongo会自动转换或者自定义类的实例。解决确保你要插入的数据是基本的 Python 类型dict, list, str, int, float, bool, None或能被bson编码的类型。对于复杂对象先将其转换为字典。错误插入成功但查询不到或时间不对可能原因时区问题。如果你用new Date()或datetime.now()生成时间它使用的是系统时区。而 MongoDB 默认存储 UTC 时间。解决在生成时间戳时显式使用 UTC 时间。在 Python 中datetime.utcnow()。在 Node-RED 的 Function 节点中new Date().toISOString()会生成 ISO 格式的 UTC 时间字符串MongoDB 可以正确解析。查询时也要注意时区转换。现象写入速度越来越慢可能原因索引过多每次插入都需要更新所有相关索引。评估并删除不必要或使用率极低的索引。磁盘 IO 瓶颈检查磁盘使用率和 IO 等待时间。考虑使用更快的 SSD。锁竞争虽然 MongoDB 的 WiredTiger 存储引擎是文档级锁但某些元数据操作仍可能引发竞争。使用db.currentOp()查看是否有长时间运行的操作阻塞了写入。文档体积过大MongoDB 单个文档限制为 16MB。如果payload字段不断累积数据可能接近或超过此限制。考虑定期归档或拆分文档。6.3 Watcher 自身运行问题现象Node-RED 流意外停止或无响应排查检查 Node-RED 日志通常位于~/.node-red/logs或在启动终端中查看。检查 Function 节点中的代码是否有未捕获的异常导致整个流上下文崩溃。务必用try...catch包裹可能出错的代码。检查是否有节点如 HTTP Request在等待外部服务响应时超时导致消息堆积。合理设置超时时间。现象Python Watchdog 漏掉文件事件或重复触发原因某些编辑器如 VS Code、Vim在保存文件时会先写入临时文件再重命名可能触发on_modified和on_moved等多个事件。频繁保存会导致大量重复事件。解决实现“防抖”逻辑。例如记录每个文件的最后处理时间如果在短时间内如1秒再次收到同一文件的事件则忽略。class DedupeEventHandler(FileSystemEventHandler): def __init__(self): self.last_processed {} self.debounce_interval 1.0 # 秒 def on_modified(self, event): if not event.is_directory: current_time time.time() last_time self.last_processed.get(event.src_path, 0) if current_time - last_time self.debounce_interval: self.last_processed[event.src_path] current_time self.process_file(event.src_path)6.4 一个综合排查案例API 错误日志入库假设你有一个 Watcher 监听应用错误日志文件并将错误插入 MongoDB。突然发现 MongoDB 里没有新数据了。排查步骤检查 Watcher 进程ps aux | grep -E “(node-red|python.*watch)”确认进程还在运行。检查 Watcher 日志查看 Node-RED 管理界面或 Python 脚本的日志输出看是否有报错信息。模拟触发手动在监控目录创建一个新的错误日志文件观察 Watcher 是否有反应日志是否输出。检查权限确认 Watcher 进程用户有权限读取被监控的日志文件和写入 MongoDB。检查 MongoDB登录 MongoDB Shell检查目标集合的db.collection.stats()看是否有插入尝试但失败。检查 MongoDB 日志 (/var/log/mongodb/mongod.log) 是否有身份验证失败、主节点切换等记录。检查网络与连接从 Watcher 服务器尝试连接 MongoDB 端口。使用db.adminCommand({“currentOp”: 1})查看是否有长时间运行的查询占用了锁。检查数据格式如果 Watcher 日志显示正在处理但插入失败很可能是数据格式问题。在插入代码前后增加调试日志打印出准备插入的文档内容检查是否有特殊字符或类型错误。按照这个由近及远、由内到外的顺序排查大部分问题都能定位。最关键的是给你的 Watcher 加上详尽的日志记录记录它每个阶段在做什么这是事后排查的黄金依据。