1. 项目定位:PLFM_RADAR 到底在做什么
PLFM_RADAR 这个名字是我自己起的,PLFM 取 Platform 的缩写,RADAR 不是蹭军事概念,而是想表达这套系统的核心工作方式:像雷达一样周期性扫描目标平台,捕捉变化、滤除噪声、输出信号。这不是一个爬虫项目,也不是一个报表工具,而是一套面向多平台公开数据的主动监测与预警系统,我把它叫作“平台雷达”。
做这个东西的起因很实际。之前团队做内容运营,每天早晨要人工打开四五个后台,挨个看热榜、看竞品动态、看关键词排名,再手动记录到表格里。数据一多,漏看、错看、看不完就是常态。更难受的是,某些关键词的热度突然起飞,往往是半夜或者周末,等人发现的时候流量红利已经过去了。PLFM_RADAR 就是为了解决这类问题:统一采集各平台的公开指标数据,集中存储和计算,自动识别异常波动,并推送预警消息。
这个项目适合谁参考?一类是负责内容运营、新媒体投放的同学,需要一个低成本的舆情/热点监控手段;另一类是后端或数据工程师,想了解如何用 Kafka、ClickHouse、Redis 搭一套轻量级的实时数据处理链路;还有一类是产品经理,想给团队搭一个竞品动态雷达。它不依赖任何商业 SaaS 产品,所有模块都是开源组件,数据源层面也只使用各平台公开的开放接口和授权数据,完全合规,跑起来之后维护成本很低。
先说清楚一件事:雷达和摄像头的区别。摄像头是全天候记录,雷达是周期扫描、快速识别“异常目标”。PLFM_RADAR 的设计哲学也是这个——它不是把所有平台数据都无限存下来,而是围绕你关心的关键词、榜单、指标做定向扫描,在数据里捕捉“信号”,其余噪声直接丢弃。这一点贯穿了整个架构设计,后面每个模块的选择都是围绕“轻量、聚焦、可预警”展开的。
2. 整体架构设计与选型思路
2.1 模块划分与数据流转
整个系统分成五个层面:采集层、缓冲层、计算层、存储层、通知层。采集层负责按固定节奏调用各平台的公开接口,解析返回数据并转换成统一结构;缓冲层用 Kafka 解耦采集与消费,避免某个环节抖动拖垮整条链路;计算层负责滑动窗口聚合、异常检测、评分排序;存储层用 ClickHouse 保存明细数据,Redis 保存实时状态;通知层把预警信息推送到企业微信群、钉钉群或者邮件,同时保留一个查询面板给运营同学按需看数。
数据流转大概是这样的:采集器定时抓取目标数据,转换成 JSON 消息发送到 Kafka;消费服务从 Kafka 拉取消息,做清洗和增益处理后批量写入 ClickHouse;另一条链路从 ClickHouse 或 Redis 读取最近 N 个窗口的数据,执行异常检测逻辑,命中规则就生成预警事件,再通过 Webhook 发送到群。Grafana 直接对接 ClickHouse,提供趋势曲线和看板。
这套链路最核心的取舍是:采集和计算完全解耦。早期我做第一版的时候,采集器直接往 ClickHouse 里写数据,看着简单,实际上很被动。一旦 ClickHouse 写入抖动或者消费逻辑需要重跑,采集器就会阻塞,最终数据丢失。改成 Kafka 之后,采集器只负责生产消息,完全不关心下游是否健康。下游哪怕停了半小时,消息还在 Kafka 里,恢复后继续消费,不会丢数据,这就是缓冲区的意义。
2.2 组件选型:为什么是 Kafka、ClickHouse、Redis 这三件套
很多朋友一开始会问:为什么不用 MySQL?为什么不用 RabbitMQ?我逐个说一下选型理由。
先说存储。雷达系统的典型特征是写多读少、按时间范围扫描、实时聚合。MySQL 在这种场景下有两个硬伤:第一,高频写入容易产生锁竞争和主从延迟;第二,对“最近一小时窗口聚合”这类查询支持很差,GROUP BY 跑一次要扫全表。ClickHouse 是列式存储,写入吞吐极高,单体单机能到每秒几十万行,而且它对时间范围聚合有原生优化,跑一条SELECT keyword, avg(heat) ... GROUP BY keyword比 MySQL 快一到两个数量级。代价是事务和更新能力弱,但雷达数据基本都是 append-only,正好踩在它的舒适区。
再说消息队列。我们用 Kafka 而不是 Redis Stream,主要是考虑两点:一是 Kafka 的消息堆积能力和回溯消费能力更强,出问题的时候可以重置 offset 重新处理一段历史数据;二是 Kafka 的分区机制天然支持按关键词或平台维度做并发扩展。Redis Stream 在数据量小、链路短的情况下也能用,一旦你要做多消费者组、按时间回溯,Kafka 的维护成本并没有想象中的高。
Redis 在这里不是缓存主力,它的角色更像“状态寄存器”。存放每个关键词最近一次预警时间、当前评分、检测窗口等状态信息。为什么要单独放 Redis?因为这些状态需要高频读写,而且要求低延迟。如果每次预警判断都去 ClickHouse 查一遍历史,查询压力和延迟都不可控。Redis 的 TTL 机制还可以天然处理“冷却期”的状态过期,非常契合预警去重这个场景。
2.3 为什么没有直接上一套商业方案或开源全家桶
市场上其实有开源的舆情系统,也有商业的舆情监测产品。没有直接拿现成方案,不是不相信它们,而是 PLFM_RADAR 的定位决定了它不适合大而全。商业舆情平台通常按月收费,价格不低,而且数据源、更新频率、预警规则都是黑盒,运营想加一个自定义指标,流程长到你忘记当初的需求。开源全家桶,比如某个完整的数据采集分析平台,功能确实全,但部署起来组件十几个,对一台小服务器来说太重了。
PLFM_RADAR 刻意保持精简,核心组件就三个,加上采集和消费两个自研服务,总共五个模块。这样做的好处是:新人花半天就能看懂数据流,出了问题能直接定位到某个组件,服务器配置要求也不高,4 核 8G 的机器跑起来绰绰有余。雷达系统最关键的不是功能多,而是稳定和响应快,少一个中间环节,就少一个故障点。
3. 核心实现:从采集到预警的关键细节
3.1 采集层:并发、限流与重试
采集层是整个系统的起点,也是坑最多的地方。因为要同时对接多个平台,每个平台的接口频率限制、返回结构、鉴权方式都不一样,所以我在采集层做了几个统一封装。
第一是统一限流。每个平台单独配一个令牌桶,比如平台 A 每秒最多 5 次请求,平台 B 每秒 2 次。为什么必须单独限流?因为各平台对 QPS 的容忍度完全不同,统一限流会导致某个平台跑不满配额,另一个平台又频繁触发 429。令牌桶可以用 Python 的asyncio.Semaphore加时间窗口来实现,也可以用现成的aiolimiter库。
第二是统一重试。接口请求失败是常态,超时、网络抖动、返回 5xx 都会遇到。重试策略我用的是指数退避加抖动,最多重试 3 次。注意一定要加抖动,否则多个任务同时失败后同时重试,会对平台造成一波集中的请求尖峰。
下面是采集器里一个典型的请求封装:
import asyncio import aiohttp from aiolimiter import AsyncLimiter class PlatformClient: def __init__(self, name, qps): self.name = name self.limiter = AsyncLimiter(qps, 1) async def fetch_json(self, session, url, params=None): async with self.limiter: for attempt in range(3): try: async with session.get(url, params=params, timeout=10) as resp: if resp.status == 429: wait = 2 ** attempt + 0.5 * attempt await asyncio.sleep(wait) continue resp.raise_for_status() return await resp.json() except asyncio.TimeoutError: await asyncio.sleep(1.0) except aiohttp.ClientError as e: await asyncio.sleep(2 ** attempt) return None这里有几个容易被忽略的细节。AsyncLimiter(qps, 1)是每秒允许 qps 个请求,间隔均匀,不会出现第一毫秒打满、剩下 999 毫秒空转的情况。超时时间 10 秒是调过几次之后的经验值,太短了容易误判,太长了会拖累整体采集周期。重试时对 429 和普通异常的退避策略分开处理,429 的等待时间更长,因为平台明确告诉你被限流了,此时继续硬闯只会加重封禁风险。
采集调度采用固定间隔任务,每个平台一个循环:拉取配置中心里的目标关键词列表,逐个调用接口,最后组装消息发送到 Kafka。整体并发度控制在一个平台一个采集协程,不用每个关键词都开协程,主要原因是大部分接口都支持批量查询,没必要拆太细。
3.2 统一数据模型与 ClickHouse 建表
不同平台的返回结构千差万别,但落到雷达系统里,抽象出来的核心字段其实是一致的:哪个平台、哪个关键词、什么时间、热度值多少、排在第几位。统一数据模型之后,后续的存储、聚合、检测逻辑全部只认这一套结构,平台差异被隔离在采集层。
采集器输出的消息结构大致如下:
{ "platform": "content_site_a", "keyword": "某话题", "rank": 3, "heat": 286000, "source_type": "hot_search", "record_time": "2024-06-15T10:30:00+08:00", "version": 1 }source_type用来区分数据是来自热榜、关键词趋势还是竞品监测,后续可以按这个维度拆分不同的检测策略。version是数据版本号,用于应对平台改版导致字段结构变化,便于重放历史数据时区分新旧格式。
ClickHouse 建表语句如下:
CREATE TABLE IF NOT EXISTS plfm_radar.raw_event ( platform LowCardinality(String), keyword String, rank Int32, heat Float64, source_type LowCardinality(String), record_time DateTime('Asia/Shanghai'), version UInt32 ) ENGINE = MergeTree ORDER BY (record_time, platform, keyword) PARTITION BY toYYYYMMDD(record_time) TTL record_time + INTERVAL 90 DAY;建表用了LowCardinality修饰平台和来源类型,因为这两个字段的重复度极高,字典编码之后能显著降低存储空间和查询扫描量。排序键把时间放在第一位,因为绝大多数查询都是按时间范围过滤;TTL 设了 90 天,雷达明细数据保留三个月足够,过期自动删除,不用写定时任务。
写入采用批量方式,消费端攒够 5000 条或者每隔 3 秒刷一次,使用 ClickHouse 的 JSONEachRow 格式批量插入。实测在 4 核 8G 的机器上,单批 5000 条写入耗时在 50 毫秒左右,完全不是瓶颈。如果数据量继续增长,可以考虑用Buffer表引擎做缓冲,但对雷达这个量级来说没有必要。
3.3 异常检测算法与 RADAR 评分
采集只是过程,预警才是雷达的输出。怎么判断一个关键词的热度波动是异常?最简单的方案是固定阈值,但实际运行几天就发现不行:头部关键词日均热度几十万,尾部关键词才几百,一个统一的阈值要么对头部无效,要么对尾部天天误报。所以我改用了一套组合评分,名字就叫 RADAR Score,由三个因子构成。
首先是归一化热度因子,用来衡量“当前热度在它自己历史水平中的位置”。不要直接用绝对热度跨关键词比较,而是跟该关键词自身近期分布比。我取过去 7 天同一时段的热度分布,用 P90 分位数做分母,当前热度取对数后除以分位数的对数,这样避免极端值主导。
其次是加速度因子,衡量“上升得有多快”。计算当前时间窗口内热度的线性回归斜率,和历史窗口斜率做比值。如果当前斜率是历史平均斜率的 5 倍以上,说明正在急速拉升。这个因子对“突然蹿升”特别敏感,能捕捉到那些基础热度不高、但有明显启动迹象的长尾关键词。
最后是持续性因子,衡量“异常状态是否稳定”。如果只是单次采样尖刺,很可能是接口抖动或者数据异常,不值得报警。持续性因子统计最近 6 个采样点中有几个超过历史阈值,比例越高,信号越可信。
最终评分公式是这样的:
RADAR_Score = heat_norm * accel_factor * persist_factor三个因子相乘而非相加,是因为雷达信号需要同时满足“绝对位置高、上升快、能持续”三个条件。任何一个因子接近 0,总分就会被压下来,有效过滤掉单点噪声。会看公式的朋友可以试算一个例子:某关键词热度在过去 7 天的 P90 是 8000,当前热度 24000,heat_norm 约等于 1.36;当前窗口斜率 1200,历史窗口斜率 180,accel_factor 为 6.67;6 个采样点中有 5 个超过阈值,persist_factor 为 0.83。三项相乘,RADAR Score 约为 7.52,高于阈值 3.0,触发预警。
阈值本身也不是拍脑袋定的。我把历史数据按评分倒序排列,观察 Top 100 的事件里哪些是真实的运营事件、哪些是数据噪声,然后选定 3.0 作为默认阈值。这里我建议团队根据自己数据做一次校准,不要直接抄任何人的阈值,因为数据源的噪声水平不一样。
3.4 预警调度与防重复机制
评分计算出来之后,最忌讳的就是“同一件事反复报警”。关键词热度一旦起飞,会在好几个小时内持续维持高位,如果检测任务每个窗口都触发,群里就会被刷屏,运营同学最后直接把群屏蔽了,雷达就变成噪声制造机。
我的做法是在 Redis 里维护一个关键词维度的冷却记录:当一个关键词触发预警后,写入plfm:alert:last_sent:{keyword},值为当前时间戳,TTL 设为 4 小时。在发送预警前先检查这个 key 是否存在,存在就直接跳过。4 小时不是固定的,对于热点生命周期短的内容平台可以缩短到 2 小时,对于趋势缓慢的行业指数可以延长到 12 小时,按source_type分别配置。
还需要做聚合发送。单条推送到群里容易刷屏,我让检测服务把每轮命中的关键词缓存到内存队列,每 5 分钟批量发送一次。推送格式包含平台、关键词、当前评分、环比变化、热度值,以及一个指向面板的趋势链接。这样既保证及时性,又控制了消息数量,群里不会变成刷屏现场。
4. 部署实操与参数调优
4.1 用 Docker Compose 快速拉起整套环境
整个项目我全部用 Docker Compose 编排,一台 Linux 服务器就够了。下面是 docker-compose.yml 里几个关键服务的配置,完整版还会包含健康检查、日志轮转和网络配置:
services: clickhouse: image: clickhouse/clickhouse-server:24.8 container_name: plfm_ch volumes: - ./data/ch:/var/lib/clickhouse - ./schema/init.sql:/docker-entrypoint-initdb.d/init.sql:ro ports: - "8123:8123" restart: unless-stopped kafka: image: bitnami/kafka:3.7 container_name: plfm_kafka environment: - KAFKA_CFG_NODE_ID=0 - KAFKA_CFG_PROCESS_ROLES=controller,broker - KAFKA_CFG_LISTENERS=PLAINTEXT://:9092,CONTROLLER://:9093 - KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092 - KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE=true ports: - "9092:9092" restart: unless-stopped redis: image: redis:7.2 container_name: plfm_redis command: redis-server --appendonly yes ports: - "6379:6379" restart: unless-stopped collector: build: ./collector container_name: plfm_collector depends_on: - kafka environment: - KAFKA_BOOTSTRAP=plfm_kafka:9092 restart: unless-stopped detector: build: ./detector container_name: plfm_detector depends_on: - clickhouse - redis - kafka restart: unless-stopped grafana: image: grafana/grafana:11.1 container_name: plfm_grafana ports: - "3000:3000" environment: - GF_PATHS_PROVISIONING=/etc/grafana/provisioning volumes: - ./grafana/provisioning:/etc/grafana/provisioning:ro restart: unless-stopped第一次启动的流程是:先拉起 ClickHouse、Kafka、Redis,确认它们健康后再启动采集和检测服务。我踩过一个坑,Compose 里所有服务一起启动,采集器连不上 Kafka,虽然应用层有重试机制,但日志里全是连接报错,排查起来眼花缭乱。后来我显式配置了depends_on条件,并且给服务加了健康检查,只有上游真正就绪才启动下游。
4.2 关键参数与调优建议
运行一段时间之后,有四个参数值得反复调。
Kafka 的分区数。topicradar.raw.event我建了 6 个分区,主要考虑是消费端需要按照关键词维度并行处理,分区数太少会限制消费者并发,太多则增加 Broker 元数据开销。如果你监控的关键词在 1000 个以内,6 个分区完全够用。
ClickHouse 的写入批量大小。刚才提到攒够 5000 条或 3 秒刷一次,这个值来自一个简单的权衡:批量太小,插入次数多,ClickHouse 每次插入都要落盘,吞吐上不去;批量太大,内存占用和单次失败的影响范围都变大。5000 条在这个数据规模下是甜点值。
预警检测的窗口长度。我用的是 10 分钟窗口,每 5 分钟扫描一次,也就是每次检测覆盖最近两个采样点。窗口太长会让预警变钝,热点起飞后要等很久才触发;窗口太短又容易受噪声干扰。10 分钟对于大多数内容平台的热点节奏是合适的,你可以按自己业务的实时性要求调整。
Redis 里冷却期的 TTL。这个前面讲过,按 source_type 分开配置,热榜类 2 小时,趋势类 6 小时,竞品类 12 小时。设置的原则是:信号生命周期短,冷却短一点,避免漏报二次爆发;信号生命周期长,冷却长一点,避免重复轰炸。
4.3 效果验证:跑了一周,我看到了什么
部署之后的头七天,我接入了两个内容平台和一个电商平台的公开数据,监控约 300 个关键词,每 5 分钟采集一次。七天下来的明细数据量在 120 万条左右,ClickHouse 占用磁盘不到 1.5GB,整体运行非常轻量。
异常检测方面,7 天共触发预警 34 次,手动核对后确认其中 5 次是数据源字段漂移导致的误报,有效预警 29 次。在这 29 次里,有两次确实对应了平台热搜榜的突发放量,运营同学在预警推送后 10 分钟内就做了跟进,这就是雷达系统最直接的价值。误报的那 5 次也很有代表性,最后定位到是平台改版导致热度字段的单位从“万”变成了“个”,量级跳变本身不真实。后来我在采集层加了一个单位一致性校验,用环比跳变超 50 倍的记录自动标记为可疑数据,误报率立刻降下来了。
Grafana 面板我做了三个大类:总览页看整体采集延迟和消息吞吐,趋势页按关键词看热度曲线,异常页列出历史预警事件及评分构成。运营同学日常只看总览和异常页,技术细节收起来,避免信息过载。
5. 常见问题与排查实录
5.1 采集侧:字段漂移与限流封禁
运行期间遇到最多的问题就是平台接口返回结构变化,也就是字段漂移。平台的开放接口不受你控制,偶尔加个字段、改个单位、调整 JSON 层级,采集器的解析逻辑就得跟着变。我的对策是两层:第一,解析时用容错取值,拿不到关键字段就用默认值,不让单条数据解析失败导致整个批次丢弃;第二,采集器里记录原始响应的 schema hash,发现和上次不一致就告警出来,运营同学看到告警就知道平台改接口了,而不是看到数据突然全为零。
限流封禁的问题也遇到过。某平台某天突然把 QPS 限制从 10 降到了 2,采集器没有及时感知,连续触发 429,最后被封了一段时间。教训是:不要只看单次请求是否成功,要统计限流占比,如果某平台连续 5 分钟内 429 的比例超过 20%,自动降低该平台的 QPS 配额,同时推一条通知出来。加了这一步之后,再也没出现过封禁事故。
5.2 存储侧:分区膨胀与查询变慢
ClickHouse 的一个隐藏坑是分区数量过多。刚开始我按小时分区,一天 24 个分区,数据量小的时候看着没什么,但跑了几天之后分区数上千,查询时 Meta 扫描开销明显变大。后来改回按天分区,查询性能立刻改善。如果你以后单日数据量大到需要小时级分区,建议配合TTL把历史分区合并,别让陈旧分区无限累积。
另一个和查询相关的问题是 ORDER BY 使用不当。我第一版排序键是(platform, record_time, keyword),导致按时间范围查询时不高效,因为排序键的首列是 platform,等于把所有平台的数据按平台分组排列,而不是按时间排列。后来改成(record_time, platform, keyword),同样的查询快了近三倍。排序键的选择顺序,一定是查询过滤最频繁的字段在最前面,这条经验对任何 ClickHouse 表都适用。
5.3 预警侧:误报与消息风暴
误报有两个来源,一个是外部数据源的问题,前面说的字段漂移和单位变化都属于这一类;另一个是算法本身的问题。RADAR Score 里加速度因子对斜率比值太敏感,当历史窗口的斜率趋近于 0 时,比值会爆炸性放大。比如某关键词过去 7 天热度一直平稳,偶尔有两个采样点波动,计算出的历史斜率几乎为 0,当前只要有一个正常波动,accel_factor 就可能冲到几十。后来我给分子分母都加了一个极小值平滑项,并给 accel_factor 加了上限,超过 20 一律按 20 算,误报才真正控制住。
消息风暴的问题在加了冷却机制后基本解决,但还有一个漏网之鱼:多个平台同时命中同一关键词。比如某关键词在内容平台和电商平台同时起飞,会发出两条独立预警,运营同学看着像重复消息。我的处理是在发送队列里增加一个“同关键词跨平台合并”的规则,5 分钟窗口内同一个关键词命中多个平台,合并成一条消息,展示各平台的数据对比,信息量反而更大。
5.4 运维侧:数据延迟与消费积压
雷达系统对时效性要求高,数据从采集到入库再到检测,整体延迟超过 10 分钟,预警的意义就打了折扣。排查延迟的思路,我一般是从上往下逐层看:先看 Kafka 消费组 Lag 是不是在增长,再看 ClickHouse 写入耗时,最后看检测任务的执行间隔。
有一次发现消费端 Lag 越来越高,排查后发现是消费服务里做数据清洗时调了一个外部接口做关键词分类,这个接口偶尔会卡几秒钟,导致单个消息处理时间从 20 毫秒飙升到 3 秒。这里的设计教训是:消费链路里尽量不要同步调用外部服务。后面我把关键词分类改成本地规则引擎,处理时间稳定在 5 毫秒以内。简单地说,消费端每一毫秒的耗时都会被放大,特别是在数据高峰期,积压一旦形成,追 Lag 的过程里还会堆积更多新数据。
还有一个运维小技巧:Kafka 的消息保留时间我设成 7 天,日常用不到这么多,但万一 ClickHouse 出了故障,可以把消费组的 offset 重置到一天前,把丢失窗口的数据重新处理一遍。这个能力在数据链路上相当于一道安全网,出事故的时候能救命。
最后再分享一点实际体会:做这种雷达系统,真正难的不是把链路跑通,而是让预警结果持续可信。刚上线那周大家很兴奋,每天都看预警,但一周之后误报和重复消息多了,群里的人就开始免疫。后来我花了大量时间在降误报和收敛消息上,而不是继续加数据源。一个让运营愿意每天打开的预警系统,一定是一个很少打扰他们的系统——雷达扫到目标才响,扫不到的时候就应该安安静静地转。PLFM_RADAR 现在的状态基本达到了这个标准,后续如果继续扩展,我建议优先做预警反馈闭环:运营同学对每次预警点“有用/无用”,把这些反馈喂回阈值校准里,让系统越用越准。