PLFM_RADAR:一个多平台数据监测预警系统的完整实现复盘
做开发这几年,我一直对"数据雷达"类的系统特别感兴趣。所谓PLFM_RADAR(Platform Radar,平台雷达),本质上就是一套能够同时盯住多个平台、周期性拉取公开数据、自动识别变化并推送预警的监测框架。去年因为业务需要,我完整地把它从零到一搭了一遍,从需求梳理、技术选型到部署上线踩了不少坑,今天把整个过程沉淀下来,包括架构设计、核心模块的实现细节、参数怎么定、线上踩过的坑,一次性讲透。适合正在做数据采集、竞品监控、舆情追踪或者任何"周期性监测类"项目的朋友参考,无论你是刚要起步还是已经做了一半遇到瓶颈,这篇应该都能给你一些可落地的思路。
1. 项目整体设计与思路拆解
1.1 核心需求:我们到底要监控什么
PLFM_RADAR这个项目最开始的需求其实并不复杂,就是要盯住几个平台的公开页面,一旦发现目标数据发生变化就通知我们。听起来简单,但拆开来看有三个绕不开的问题:盯哪些平台、盯什么数据、变化了之后怎么办。
平台侧,需要兼容不同站点结构、不同接口规范、以及不同频率限制策略。数据侧,我们关心的不只是"拿到数据",而是"拿到能对比的数据"——比如同一件商品的价格、同一个账号的粉丝数、同一篇内容的阅读量,这些数据必须有稳定标识(ID)才能做时间序列上的对比。预警侧,必须区分"变化"和"有效变化",不是所有波动都值得打扰人,比如价格波动5%以内可能是促销抖动,但波动20%可能就是竞品调策略了。
这个项目名称其实隐藏了一个关键设计倾向:雷达不是简单的抓取工具,而是"数据管道 + 变化检测 + 决策触发"的组合体。我见过很多类似的系统死在第一个阶段——只做了抓取和存储,没有做检测层,结果数据堆成山但没人看,最终项目被砍掉。所以在需求阶段就要明确:我们最终要的是"信号"而不是"数据"。
1.2 技术选型:为什么是Python全家桶
选型时考虑过几个方向:Java系的重型框架,Go的并发模型,最后选了Python,原因很直接——生态成熟、开发效率高、团队维护成本低。
采集层用了Scrapy和httpx的组合。Scrapy适合大而全的站点级抓取,自带去重、调度、中间件体系,但调试起来偏重;httpx轻量灵活,适合接口型采集,直接同步异步通吃,写起来像写普通函数一样自然。实际项目里大约70%的采集任务走的httpx + asyncio,剩下30%(比如需要对复杂页面做深度解析的)才用Scrapy。
调度层选了Celery + Redis。Celery的beat调度器可以精确控制每个任务的执行周期,配合Redis作为broker和result backend,这套组合在中小规模场景下最成熟也最稳。
存储层是MySQL + MongoDB双轨并行。结构化数据(监测目标、用户表、预警记录)放MySQL,采集到的半结构化快照数据放MongoDB。原因很简单:同一类数据不同平台的字段差异很大,用关系型数据库硬对齐会非常痛苦,MongoDB的文档模型天然适合存这种"形态相似但细节不同"的数据。
预警通知走的是Webhook + 邮件双通道,Webhook接企业微信和钉钉机器人,邮件作为兜底。这个设计在实战中很重要——企业微信机器人的token过期或者被限流是常有的事,必须有第二通道保证预警不丢。
1.3 架构分层:从数据源到用户的完整链路
整个系统分成四层,每一层职责单一,层与层之间通过接口解耦:
- 接入层:管理平台连接配置、登录态(Cookie/Session)、代理策略。每个平台对应一个独立的Adapter(适配器),新平台接入时只需要实现统一的采集接口,不用动上层逻辑。
- 处理层:负责数据清洗、字段标准化、去重、数据落库。这里的关键操作是"实体识别 + 快照对比",判断这条数据是新实体还是已有实体的更新。
- 分析层:执行变化检测、趋势计算、阈值判断。这一层是PLFM_RADAR最核心的部分,预警质量完全取决于这一层的判断逻辑。
- 通知层:根据分析结果决定是否触发通知、通知给谁、走什么通道,并记录通知历史,防止重复打扰。
分层带来的最大好处是排障效率高。线上出问题先看是接入层没拉到数据,还是处理层清洗报错,再或者分析层误报,每层都有独立日志和计数器,几分钟能定位到问题源头。
2. 核心模块解析与实操要点
2.1 数据接入层:Adapter模式是被验证过的正确选择
在设计接入层时,最关键的决策是引入Adapter适配器模式。每个平台实现一个类,继承同一个BaseAdapter基类,基类定义了fetch_data()、parse_data()、validate_data()三个抽象方法,所有平台都遵循这套接口。
拿电商平台举例,一个适配器的骨架长这样:
class BaseAdapter: def __init__(self, platform_name: str, config: dict): self.platform_name = platform_name self.config = config self.session = self._init_session() def fetch_data(self, target: dict) -> dict: raise NotImplementedError def parse_data(self, raw: dict) -> dict: raise NotImplementedError def validate_data(self, parsed: dict) -> bool: raise NotImplementedError class EcommerceAdapter(BaseAdapter): def fetch_data(self, target: dict): # 每家平台的请求逻辑都不一样 resp = self.session.get( f"https://api.example.com/item?id={target['target_id']}", headers={"User-Agent": self.config.get("ua")} ) resp.raise_for_status() return resp.json() def parse_data(self, raw: dict): return { "target_id": raw["data"]["item_id"], "title": raw["data"]["title"], "price": float(raw["data"]["price"]), "stock": int(raw["data"]["stock"]), "sales": int(raw["data"]["sales"]), "collected_at": datetime.utcnow().isoformat(), }这里有个非常容易被忽略的坑:不同平台的"变化"语义完全不同。同样一个字段名"status",电商平台可能是商品的上下架状态,内容平台可能是文章的发布/删除状态,社交平台可能是账号的封禁状态。所以我要求每个Adapter的parse_data必须输出标准化的快照结构,里面至少包含三类字段:标识字段(target_id)、业务字段(具体的关注指标)、时间字段(collected_at)。
时间字段必须用UTC时间。这个看起来是个小细节,但在凌晨排查预警记录时,如果混用了平台本地时间和服务器UTC时间,整个时间序列全乱套,后续debug成本极高。
2.2 增量更新与实体识别
增量更新的核心是"如何判断两条记录指向同一个实体"。我用的方案是"平台Code + 平台实体ID"组成全局唯一标识,即unique_key = f"{platform_code}:{target_id}"。每次采集到数据先计算unique_key,然后查MySQL的monitor_targets表,存在就更新,不存在就新增。
这个方案的坑在于:很多平台页面上的ID和接口返回的ID不是同一个。比如某个平台的页面URL里是字符串ID,但接口返回的是一串数字ID,如果不做映射就会出现大量重复数据。所以每个Adapter里必须显式声明extract_id()方法,专门负责从原始数据中提取稳定的实体标识。
去重我用了两层策略。第一层是Redis的布隆过滤器,用来快速判断一个unique_key是否已经处理过;第二层是MySQL的唯一索引兜底。布隆过滤器会有一点点误判率(我设置的是1%),但处理量大的时候性能优势非常明显。唯一索引兜底保证即便过滤器误判了也不会插入重复数据,而是捕获到DuplicateKeyError后走更新逻辑。
2.3 数据存储模型
MySQL侧的核心表设计:
CREATE TABLE monitor_targets ( id BIGINT AUTO_INCREMENT PRIMARY KEY, platform_code VARCHAR(32) NOT NULL, target_id VARCHAR(128) NOT NULL, unique_key VARCHAR(192) NOT NULL, title VARCHAR(512), status TINYINT DEFAULT 1 COMMENT '1:监控中 0:已停止', created_at DATETIME DEFAULT CURRENT_TIMESTAMP, updated_at DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, UNIQUE KEY uk_unique_key (unique_key), KEY idx_platform (platform_code) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; CREATE TABLE alert_logs ( id BIGINT AUTO_INCREMENT PRIMARY KEY, target_id BIGINT NOT NULL, alert_type VARCHAR(64), alert_level TINYINT, content JSON, created_at DATETIME DEFAULT CURRENT_TIMESTAMP, KEY idx_target_time (target_id, created_at) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;MongoDB侧存储的是快照集合,一个文档就是一次采集的完整快照:
snapshot = { "unique_key": "ecom:123456", "collected_at": "2024-01-15T08:30:00Z", "data": { "price": 199.0, "stock": 86, "sales": 1240, "title": "某品牌无线耳机" } }快照在MongoDB里只增不改,每次采集都是insert一个新的文档,配合TTL索引定期清理超过90天的历史快照。为什么要只增不改?因为后续分析要看趋势曲线,只有不可变的历史快照才能画出准确的时间序列。这也是我在早期版本踩过的坑——当时用了update覆盖,结果分析趋势时数据全是断层,复盘时才明白快照的不可变性是这类系统的基石。
2.4 分析层:变化检测与预警判断
分析层是整个PLFM_RADAR的"大脑",也是最容易写烂的部分。我经历了三个阶段:
第一阶段是简单阈值判断:比如价格低于X就预警。缺点很明显,缺乏上下文,误报率高,促销季一天能发几十条无用预警。
第二阶段是差值百分比判断:比如相对上一次快照变化超过10%就预警。这个比第一阶段好,但仍然没解决"趋势"问题——有些数据每天缓慢上涨5%,连续涨了10天,系统每天都不触发阈值,等于漏掉了重大变化。
第三阶段才真正靠谱:时间窗口综合判断。对每个实体维护一个滑动窗口(默认取最近24小时内的快照序列),计算当前值与窗口内均值的偏离度,结合变化方向、变化速率、是否跨越关键阈值(比如价格跌破历史最低价)三个维度综合打分,超过设定阈值才触发预警。
判断逻辑大致如下:
def should_alert(entity_id: str, new_snapshot: dict, window_hours: int = 24) -> dict: history = get_snapshots(entity_id, window_hours) if len(history) < 3: return {"alert": False, "reason": "insufficient_data"} baseline = np.mean([h["data"].get("price", 0) for h in history]) current_value = new_snapshot["data"].get("price", 0) change_ratio = abs(current_value - baseline) / max(baseline, 0.01) # 跨阈值检测:是否存在历史最低/最高价的突破 min_price = min([h["data"].get("price", float("inf")) for h in history]) is_breaking_low = current_value < min_price # 打分规则 score = 0 if change_ratio > 0.03: score += 2 if is_breaking_low: score += 6 if new_snapshot["data"].get("stock", 1) == 0: score += 8 return { "alert": score >= 6, "score": score, "change_ratio": round(change_ratio, 4), "base_value": baseline, "current_value": current_value, }这里有个经验之谈:预警阈值的设定必须基于一段时间的历史数据反推,不能拍脑袋。最稳妥的做法是先开启"影子模式"——系统正常运行但不实际发预警,只记录"如果当时发了会是什么内容",跑一周,统计触发频率,再根据业务能接受的打扰程度确定阈值。
3. 实操过程与核心环节实现
3.1 环境搭建与项目结构
从零搭建时我建议用Docker Compose做本地开发环境,把MySQL、MongoDB、Redis都跑起来,一条命令完成环境初始化。配置大概这样:
version: "3.8" services: mysql: image: mysql:8.0 environment: MYSQL_ROOT_PASSWORD: root MYSQL_DATABASE: plfm_radar ports: - "3306:3306" volumes: - ./data/mysql:/var/lib/mysql mongo: image: mongo:6.0 ports: - "27017:27017" volumes: - ./data/mongo:/data/db redis: image: redis:7.0 ports: - "6379:6379"项目目录结构如下:
plfm_radar/ ├── adapters/ # 平台适配器,一个平台一个模块 │ ├── base.py │ ├── ecommerce_adapter.py │ └── content_adapter.py ├── collectors/ # 采集任务入口,供Celery调度的task ├── analyzers/ # 分析层:变化检测、打分规则 ├── notifiers/ # 通知层:Webhook、邮件等 ├── schemas/ # 数据模型定义 ├── tests/ # 单元测试与集成测试 ├── configs/ # 配置文件(按环境拆分) │ ├── config.dev.yaml │ └── config.prod.yaml └── main.py # Celery实例与应用入口这个结构用了几个月下来整体还是很舒服的,最关键的收益是"新平台接入成本极低"——业务提一个需求,开发写一个Adapter,改一下配置文件里的监控任务,测试通过就能上,完全不需要改动其他模块。
3.2 Celery任务的定义与调度
Celery的beat调度是整个系统的"心脏"。我定义了三个周期任务:高频任务每5分钟跑一次(监控核心指标),中频任务每30分钟跑一次(监控普通数据),低频任务每2小时跑一次(做深度快照或者全量对账)。
任务定义的核心代码:
from celery import Celery from celery.schedules import crontab app = Celery("plfm_radar", broker="redis://localhost:6379/0") app.config_from_object({ "timezone": "UTC", "beat_schedule": { "collect-high-frequency": { "task": "collectors.core.collect_high_freq", "schedule": 300.0, }, "collect-medium-frequency": { "task": "collectors.core.collect_medium_freq", "schedule": 1800.0, }, "run-change-detection": { "task": "analyzers.core.run_detection", "schedule": crontab(minute="*/5"), } } })这里我遇到过的一个实际问题是:Celery beat任务的时间漂移。默认情况下Celery的beat进程会按"上次任务开始时间 + 调度间隔"来排下一个任务,如果上一个任务执行时间超过调度间隔,后续任务会不断累积延迟。解决办法是给任务加上soft_time_limit(软超时)和acks_late(任务完成后才确认),配合task_time_limit做硬超时,保证一个任务最多跑2分钟,超时强制杀掉,让beat调度回到正常节奏。
3.3 调度表格与频率设计原则
采集频率这块我总结了几个踩坑经验,直接做成速查表:
| 场景 | 建议频率 | 理由 |
|---|---|---|
| 商品价格、库存 | 5-15分钟 | 价格变动频繁,但太高的频率容易触发平台风控 |
| 账号粉丝数、内容阅读量 | 30-60分钟 | 这些指标变化相对平滑,太快没有实际意义 |
| 关键词搜索结果、热门榜单 | 10-30分钟 | 时效性要求高,但结果列表可能包含不稳定因素 |
| 全量对账(扫描全站目标) | 每天1次 | 全部目标逐个验证,正常应放在低峰时段 |
频率选择的核心逻辑是"按数据特征定周期,而不是按我们想看多久定周期"。降价监控5分钟拉一次是合理的,因为用户就是盯折扣;但粉丝数5分钟拉一次就是在浪费资源,波动完全可以用30分钟粒度的曲线描述。
3.4 通知层实现与降噪处理
通知层我踩的最大的坑是"预警风暴"——某个平台的接口突然挂掉,所有实体同时报错,然后每个实体都触发一条预警,结果一个早上把所有人的手机都打爆了。后来加了三层降噪策略:
第一层是同类合并:把同一平台同一批次出现相同类型预警的实体合并成一条通知,信息汇总成列表。
第二层是静默窗口:同一个实体触发预警后进入30分钟静默期,期内不再重复通知。
第三层是分级通道:普通级别走Webhook推送(不@人),重要级别走Webhook且@指定负责人,严重级别(比如核心商品库存清零)才会同时触发邮件和电话(通过短信平台转发)。
通知模板的设计也有讲究,必须包含三要素:实体是谁、发生了什么变化、变化有多大。一条合格的通知长这样:
【价格预警】商品ID: 123456(某品牌无线耳机) 当前价格: ¥159.00 | 历史基线(24h): ¥199.00 跌幅: 20.1% | 已突破历史最低价 链接: https://... | 时间: 2024-01-15 21:30:00 UTC不要推送那种只有"数据有变化"没有具体内容的消息,这种通知收到几次之后就会被人设成免打扰,等于没有预警。
4. 常见问题与排查技巧实录
4.1 采集频率过高导致IP被平台风控
这是所有做采集监测的人都会遇到的第一道坎。某平台在连续高频请求后会返回验证码页,或者直接拒绝访问。排查过程一般是先看返回码,再看请求头,最后看频率分布。
我最终用了一套组合拳才解决:
第一,请求头随机化。User-Agent从维护好的池子里随机取,数量至少在50个以上,还要随机加入Accept-Language等常规头。
第二,代理池调度。代理的获取和分配走独立模块,每次请求从池中按权重取一个可用IP,用三次失败后自动标记过期。
第三,单IP频率限制。每个IP的请求频率上限是每10秒最多3次,超出就排队等下一轮。
Feeling检查一下,代理池的质量其实比数量更重要。有些免费代理响应很慢还大量丢包,浪费了重试次数,后来换成商业代理的按量套餐后稳定性才有保障。
4.2 数据重复率和漏采率双高
数据出现重复通常只有一个原因:实体标识不稳定。排查方法很简单,写一段脚本检查MongoDB快照里同一个unique_key在同一个collected_at是否有多条记录。如果有,回到Adapter看提取ID的地方——很多时候是页面结构与预期不符,解析规则没匹配上,抓了个空值当ID。
漏采的核心检查点是"解析规则与页面变化脱节"。平台前端改版后,CSS选择器或者XPath就失效了,抓回来一堆空数据。我的处理办法是给每个Adapter的validate_data方法加一个"字段完整度校验",如果关键字段(如标题、价格)为空,直接标记这次采集失败并告警,而不是把脏数据写进数据库。
4.3 预警延迟排查
预警延迟排查,大部分问题出在Celery任务队列堆积。用celery inspect active看当前队列里的任务数,用flower看任务执行耗时分布。还有一个容易被忽视的原因是MongoDB慢查询——分析层要查24小时快照,如果索引没建好,每次查询都要全表扫描,直接拖慢整个分析流程。所以(unique_key, collected_at)的复合索引必须在一开始就建好。
4.4 常见问题速查表
| 症状 | 可能原因 | 排查手段 | 解决办法 |
|---|---|---|---|
| 采集数据全部为空 | 页面结构变更/接口参数变化 | 查看Adapter日志及返回原始报文 | 更新解析规则,补齐字段校验 |
| 同一实体重复入库 | 实体ID提取失败 | 检查快照里的unique_key是否乱码 | 修复ID提取逻辑,增加唯一索引 |
| 预警频繁误报 | 阈值过低/基线窗口不合适 | 查看分析日志中的评分记录 | 调整阈值,增大或缩小滑动窗口 |
| 预警风暴 | 平台接口批量异常 | 查看预警日志,确认是否同批次同类型 | 启用同类合并+静默窗口 |
| 任务堆积延迟 | Celery worker资源不足 | 查看worker CPU/内存指标 | 扩容worker数量,增加超时限制 |
| MongoDB查询卡顿 | 缺少复合索引/数据量膨胀 | 查看慢查询日志 | 建索引,清理超过保留期的快照 |
4.5 容易被忽略的运维细节
有几个小坑在开发自测时完全暴露不出来,只有线上跑起来才会触发:
第一个是日志切割。采集系统一天产生的日志量非常大,不配置logrotate的话,一两个月磁盘就满了。建议日志按天滚动,保留15天足够。
第二个是时区问题。Celery的beat调度、MongoDB里的时间字段、预警通知里的时间展示,全部统一用UTC存储,展示的时候再转本地时区。如果存储时混入本地时间,分析跨天的数据趋势时会凭空多出或减少1小时的数据点。
第三个是配置管理。平台接口地址、Cookie、代理信息这些敏感配置千万不能写死在代码里,用环境变量或者独立的secrets文件挂在部署之外,不然换个人接手或者代码泄漏就是安全事故。
5. 运行效果与经验沉淀
系统上线到目前为止运行了将近两个月,监控目标覆盖二十多个平台实体,每天采集快照量大概在40万条左右。预警准确率通过"影子模式"校准后的反馈来看,达到了85%以上——也就是说每100条预警里大约有85条是业务方认为"值得处理"的。这个数字没法跟那些调参精细的商业系统比,但作为自研项目已经达到了预期。
真正让我觉得这套系统的价值不在于抓了多少数据,而在于把"人工盯数据"变成了"系统盯信号"。以前业务同学每天要花大半个小时翻阅各个平台的后台,现在只需要在预警通知里扫一眼就能知道今天哪些目标有异动。这种变化在执行层带来的效率提升非常直观——人工盯屏的目标数量下降了70%,但异常发现的速度反而快了一倍,因为系统是7x24小时在跑,而人总有睡觉的时候。
最后分享一个自己做这类系统最大的心得体会:永远不要相信任何一次采集的数据是"绝对正确"的。比如某个平台接口偶发返回0或者null,如果系统不加校验直接把这个异常值写入快照,那么基线计算会被污染,趋势判断会出现大坑。所以每个Adapter的validate_data方法一定不能只是return True走个形式,要认真写校验逻辑——这是整个系统里成本最低但收益最高的防呆设计。
PLFM_RADAR这个项目后续还可以扩展的方向很多:比如给分析层加更细粒度的趋势预测(用简单线性回归都能跑出不错的效果),或者增加多维度的关联分析(多个实体同涨同跌时单独预警会变成噪音,但联合分析可能就是一条重要信号)。如果你也正在做类似的数据监测系统,希望这篇记录能帮你少踩几个坑。