ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

从被动看板到主动雷达:实时监控系统的Flink+ClickHouse实战

从被动看板到主动雷达:实时监控系统的Flink+ClickHouse实战 PLFM_RADAR这个名字是我给一套内部监控系统起的代号。PLFM是我们对业务平台的简称RADAR则意味着这玩意儿不是普通的数据看板——它要像雷达一样主动扫描、实时探测在异常信号刚冒头的时候就把警报甩到值班群里。断断续续做了两个月从采集器到告警链路基本跑通了这篇文章把整体设计、核心代码和踩过的坑一次说透。如果你是做平台监控、可观测性建设、负责值班看板的同学应该能从这套方案里直接抄到作业。1. 项目整体设计与思路拆解1.1 需求背景为什么监控系统需要“雷达”而不是“镜子”先把话说在前面以前我们平台对线上状态的感知方式本质上是“镜子模式”——所有监控数据都存在那儿谁想知道系统什么状态主动去看一眼仪表盘。它会告诉你前十分钟发生了什么但绝大多数时候等到有人去看的时候用户已经在群里开骂了。我当时定PLFM_RADAR这个项目时目标非常明确就是要把监控从一面被动的镜子变成一部主动扫描的雷达。它要解决三个具体问题第一业务异常要在用户发现问题之前先被系统发现。比如凌晨三点某个核心服务的P99响应时间突然从200毫秒涨到4秒影响范围已经覆盖了大量请求这时候不应该等第二天有人打开看板才发现。第二异常发生时要能在3分钟内定位到大致的模块和方向。传统思路是“告警了然后去查日志”运气好十分钟能定位到具体服务运气不好翻半小时日志还理不出头绪。雷达体系要求告警消息本身就要携带可定位的信号。第三各个维度的信号要能互相印证。接口响应时间上升、日志中NPE报错增多、JVM的GC暂停时间拉长这三件事单独看都可能被当成“偶发抖动”忽略掉但放到一起看往往就是一次明显的事故前兆。雷达要能把这些信号汇聚起来做交叉判断。PLFM_RADAR这个项目名里的RADAR并不只是名字好听——它本质上是一套“扫描、聚焦、判读、上报”的完整流程。扫描是持续采集各个信号源聚焦是针对某个服务或某个指标做窗口统计判读是判断当前状态是否偏离正常基线上报是把警报推送给该看的人。每一个环节都要求确定性和低延迟这也是后面整个技术选型的基础。1.2 架构选型为什么是FilebeatKafkaFlinkClickHouse直说结论整套链路是Filebeat采集 - Kafka缓冲 - Flink实时计算 - ClickHouse存储 - Grafana展示。为什么这么选而不去用更省事的ELK一条龙原因有三个。第一ELK里Logstash的定位是处理、解析和过滤但信号量大之后它很快就会成为瓶颈。尤其是需要在日志字段里做正则提取时消耗特别大。实测过一组数据单机Logstash在4核8G的容器里跑到每秒2万条日志时正则解析CPU占用已经接近90%而且还没有算后端的索引压力。换成Filebeat只负责轻量采集和发送Kafka采集端几乎零压力Kafka天然能削峰链路稳定性直接上了一个台阶。第二监控查询的口径并不是什么场景都适合ES。ES擅长的是全文检索像“查某个服务在某个时间窗口内所有包含OrderTimeout的日志”这种事ES在行。但做指标聚合、窗口对比、基线计算这类OLAP分析时Elasticsearch的聚合查询写起来既别扭又慢。ClickHouse是列式存储加向量化执行对时间范围扫描加分组聚合有量级上的优势。我们经常要回看某个指标在两周前的同期表现一条SQL两秒出结果这种体验是ES给不了的。第三Flink才是这套系统里真正实现“雷达”的部分。ELK方案本质上是被动的——日志存下来人主动去查。Flink可以常驻在数据流上持续维护滑动窗口状态自己完成异常判定和告警触发。数据从产生到被判定是否异常中间延迟可以控制在秒级。这是“主动扫描”和“被动查询”的实质区别。1.3 为什么要用状态化流计算而不是定时扫表有人可能会问我把信号写入ClickHouse然后每分钟跑一次聚合查询看最近五分钟的错误率和QPS不一样能发现异常吗这也是我踩过的坑之一早期版本确实就是这么做的。当时的实现是每分钟用ScheduledExecutor定时跑一批SQL扫一遍最近的聚合值判断有没有越过阈值。结果暴露出两个硬伤。第一个硬伤是扫描间隔与业务尖峰不匹配。假设某个服务在10点0分30秒到10点0分40秒这10秒内疯狂报错一共报了300条。由于这300条在那一分钟的聚合里会被平均到60秒折算下来每分钟错误数只有180错误率看着并不高根本不触发阈值。等下一分钟扫描时窗口已经滚过去了这个异常信号就被彻底平均成了一条平稳曲线。第二个硬伤是查询会跟正常的业务查询抢资源。监控聚合SQL本身对ClickHouse的压力不算小尤其在高频扫描时CPU和IO的毛刺很明显甚至会拖累真正要查数据的业务看板。Flink的滑动窗口就没有这个问题。窗口粒度可以做到5秒滑动一次每次计算最近5分钟的整体状态等于每5秒就对全平台做一轮主动“扫描”。计算看似频繁但窗口状态本身是缓存在内存里的不需要反复读存储成本几乎是零。真正让这套逻辑跑得稳的关键是“状态化”——Flink维护的窗口状态允许我们做更高级的关联判断。比如拿当前10分钟内某个接口的平均RT和历史30天同时间段的数据做对比看看当前值偏离基线多少个标准差。这种跨时间维度的动态推断用定时扫表几乎没法优雅实现。2. 核心指标设计与信号采集体系2.1 信号分层基础设施、业务接口、日志关键字雷达要探测什么必须先定义清楚。我在PLFM_RADAR里把信号源分成三层每一层独立采集、独立存储但分析时会做交叉关联。第一层是基础设施信号。包括机器CPU使用率、内存使用率、磁盘IO等待时间、网络连接数、JVM堆内存使用量和GC暂停时间。这层来自Prometheus体系Node Exporter采集物理机指标JVM通过Micrometer暴露指标按固定周期拉取。但要注意PLFM_RADAR不会把Prometheus所有指标都搬进来只挑与稳定性强相关的核心信号其余全部保持在Prometheus侧供按需查询避免数据量无意义膨胀。第二层是业务接口信号。包括核心业务接口的QPS、P95和P99响应时间、5xx和4xx错误数、超时数、熔断触发次数。业务框架埋点把这些指标输出为结构化JSON日志由Filebeat采集上报。这一层的信号直接反映“用户看到的系统表现”是雷达判读的主体。第三层是日志信号。应用日志中ERROR、Exception、Timeout、Connection refused这类关键字的出现频率队列积压告警、任务失败的相关日志片段。这一层处理的是非结构化文本本身没有数值需要先做匹配提取再转成结构化计数。三层信号分开管理有讲究。接口信号告诉你“什么接口出问题了”日志信号告诉你“哪类代码层面的错误在上升”基础设施信号告诉你“底层环境有没有发生变化”。单独看任何一层都容易误判但三层交叉对比之后往往能很迅速地把根因范围缩到很小。比如接口超时告警日志层显示大量Druid连接池获取超时基础设施层显示MySQL所在机器CPU打满一条链路就清晰了。2.2 采集端Filebeat的配置细节采集端配置我直接贴出来这版是经过实际环境验证的可以作为模板直接用。filebeat.inputs: - type: filestream enabled: true id: app-error-log paths: - /data/logs/app/*/*.json fields: signal_type: app_signal parsers: - ndjson: target: json - dissect: tokenizer: %{timestamp} %{service_name} %{level} %{message} %{metadata.source} processors: - add_host_metadata: when.equals: fields.signal_type: app_signal output.kafka: hosts: [kafka1:9092, kafka2:9092, kafka3:9092] topic: plfm-radar-signals partition.hash: hash: [service_name] required_acks: 1 compression: lz4 max_message_bytes: 1048576有三个细节值得特别说明。第一paths配置直接匹配JSON日志文件而不是把所有日志文件都采进来。控制台输出和无用日志宁可丢掉也不浪费采集和传输资源。第二partition.hash按service_name做分区这能保证同一个服务的数据始终落到同一个Kafka分区。后面Flink处理时天然按服务分区有序省去了数据重排的动作也简化了状态管理的复杂度。第三max_message_bytes主动设到1MB防止个别超大日志的默认配置把链路撑爆。2.3 指标窗口与阈值设计既要敏感又不能乱报警指标定义好了阈值怎么定才是真正的技术活。不好的阈值很容易落入两个极端阈值设得太低一遍遍告警吵得人想关掉系统阈值设得太高真出问题时半天等不来一个告警。PLFM_RADAR采用的是一套组合式判据。对于QPS、RT这类周期性平稳的指标用“基线偏差法”。系统定期把过去14天同一时段的数据训练成基线计算出各自的均值和标准差。当前值如果超过基线均值加2.5倍标准差判定为异常。这样做的核心好处是阈值跟着业务节奏走凌晨2点的低流量和上午11点的高峰期各自的“正常”标准都不同不需要人工填死值省下大量调参精力。对于错误数、超时数这类低频但紧急的信号就用“绝对值持续时间”的双条件。比如错误数超过每分钟50次并且持续超过2个滑动窗口也就是10秒左右才触发告警。绝对值保证告警的确定性持续时间过滤掉瞬时抖动防止定时任务触发时那个短暂峰值把人骗过去。日志关键字的匹配要非常小心。泛化的关键字比如“Exception”几乎天天出现会带来大量误报。实际经验是宁可让关键字精确到带类名或错误描述比如“NullPointerException: ”或者“Connection refused: connect”配合正则把调用位置提取出来告警消息里直接带上出错类名和方法名值班同学可以不做任何二次查询就判断优先级和影响面。这个细节在关键时刻能省5到10分钟定位时间。3. 实操过程与核心环节实现3.1 部署形态与基础环境准备先说一个踩坑后的最终部署形态这样后续步骤才有参照。推荐生产规模是三台Kafka节点8C16G、两台Flink TaskManager8C16G、三台ClickHouse节点16C32G其中一个作为副本、每个业务节点各挂一个采集代理。如果信号量不算大这些配置都可以往下打折但ClickHouse的存储盘建议至少给到30天以上的容量因为监控数据的价值就在于能回溯对比存太短等于自断一臂。部署方式直接跳过手工装包全走容器编排。Kafka用官方镜像开KRaft模式无需再单独维护ZooKeeper集群管理和故障恢复都简单很多。Flink集群用Session模式来跑作业——因为PLFM_RADAR的处理作业本身是常驻流式的用Session模式比独立应用模式省去大量提交和资源管理开销。ClickHouse用官方镜像数据目录挂两个外部卷一个放热数据盘一个放冷数据归档盘。Filebeat就做成DaemonSet铺到每个业务节点上。这套组合从零到全部跑通花半天到一天时间很正常。大头的时间消耗在Kafka的JVM参数调优和ClickHouse的MergeTree分区配置上这两个东西别指望一把梭哈多少都要根据实际信号量调两三轮。3.2 Flink实时计算作业滑动窗口异常检测这是整个雷达系统的大脑部分先给核心代码再逐步解释为什么这么写。public class RadarSignalJob { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime); env.enableCheckpointing(60000); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000); DataStreamString source env .addSource(new FlinkKafkaConsumer( plfm-radar-signals, new SimpleStringSchema(), kafkaProps() )) .setParallelism(3); DataStreamSignalRecord records source .map(new SignalParser()) .assignTimestampsAndWatermarks( WatermarkStrategy .SignalRecordforBoundedOutOfOrderness(Duration.ofSeconds(10)) .withTimestampAssigner((record, ts) - record.timestamp) ); records .keyBy(record - record.serviceName | record.metricKey) .window(SlidingEventTimeWindows.of(Time.minutes(5), Time.seconds(5))) .aggregate(new MetricAggregator(), new AnomalyWindowFunction()) .filter(new AnomalyFilter(2.5d, 50L)) .addSink(new AlertSink()); env.execute(plfm-radar-signal-job); } }这段代码做了五件关键的事。第一开启Checkpoint并设60秒做一次保证任务重启时窗口状态不丢。监控系统最怕数据断档窗口状态一旦丢失恢复期间所有异常窗口都会变成“裸奔”。这里有个细节setMinPauseBetweenCheckpoints设置了30秒防止过于频繁的checkpoint拖垮正常数据处理。第二从Kafka订阅原始信号流map成SignalRecord对象后再指定EventTime时间戳。信号流里的时间用的是业务日志本身的产生时间而不是采集端时间这一点极其关键。如果错用采集时间所有延迟到达的数据都会被算进错误的窗口曲线全乱。第三Watermark的乱序容忍设了10秒。Kafka传输、采集端积压都可能造成日志时间轻微乱序给10秒容忍度能过滤大部分偶发乱序而不影响实时性。第四keyBy按serviceName加metricKey分组相同服务的相同指标永远进入同一个窗口聚合实例。这一步保证了后续的滑动窗口统计和基线对比能维护正确的局部状态。第五滑动窗口5分钟长、5秒滑一次。这相当于每5秒对最近5分钟的平台信号做一整轮完整重算这保证了异常即使只发生了十几秒也不会像定时扫表那样被平均淹没。AnomalyWindowFunction的内部逻辑是当前窗口聚合完成后拿这个值与状态里维护的历史近2.5分钟数值做平滑同时从ClickHouse读最近14天同小时段的基线值这部分加了2分钟本地缓存算出均值和标准差。如果当前值超过均值加2.5倍标准差或者错误数连续两个窗口超过50条就输出一条异常记录。对外输出的是结构化告警信号而不是直接发消息方便后面做去重和按级别路由。3.3 ClickHouse表结构与写入优化ClickHouse这里承担的是信号明细存储和历史回溯查询。核心表结构如下CREATE TABLE radar.signal_metrics ( event_time DateTime(Asia/Shanghai) DEFAULT now(), service_name LowCardinality(String), metric_key LowCardinality(String), metric_type String, metric_value Float64, tags Map(String, String) ) ENGINE MergeTree() PARTITION BY toYYYYMMDD(event_time) ORDER BY (service_name, metric_key, event_time) TTL event_time INTERVAL 30 DAY;几个设计点值得讲透。按天分区的意义在于数据清理时就是直接DROP分区完全不需要传统DELETE那种逐行扫描这是列式存储体系里最合理的过期数据管理办法。ORDER BY的设计特意把service_name和metric_key放在event_time前面因为查询条件基本都是“某服务、某指标、一段时间范围”这样的排序能让每次查询只扫到极少量的分区数据块而不是全表扫描。metric_key用LowCardinality压缩指标名的取值其实只有几十个压缩后在存储和查询时都有数量级优势。tags用Map保存实例IP、环境等附加属性既保留灵活性又能应对未来增维。写入端用Flink官方JdbcSink持久化但加了攒批逻辑每凑满一万条或积攒2秒就Flush一次。这里有一个非常实用的优化ClickHouse的async_insert模式虽然支持异步写入但高峰期会让MergeTree分区的后台合并变得很频繁反而拖慢整个集群。攒批写入后实测写入速率从每秒三千条提升到两万条附近磁盘IO和后台合并的毛刺都大幅下降。3.4 Grafana看板与告警接入数据进了ClickHouse之后剩下的展示和通知由Grafana负责。我的做法是搭三块看板总览看板、服务详情看板、信号联动看板。总览看板的核心是“雷达靶盘”面板——每个扇区代表一个服务颜色代表最近5分钟的健康状态绿色正常、黄色可疑、红色异常。这个面板看起来像是定制的可视化组件其实就是一个环形饼图数据源SQL里用CASE WHEN做健康评分再映射到颜色上。这玩意是整个界面里最有雷达感的地方整个团队对这个监控系统的第一印象基本都停留在靶盘上。服务详情则是靶盘的纵深页面点击扇区利用Grafana的模板变量下钻到对应服务接口QPS分布、P99响应时间、异常日志条数、JVM GC暂停时间、以及最近5分钟的告警事件流。每个Panel的查询都引用$service变量Grafana原生支持不需要写插件。告警接入统一走WebhookFlink的AlertSink把告警消息POST到一个告警网关网关负责按级别路由到企业微信、钉钉和邮件。告警消息体必须包含五个字段告警发生时间、服务名、指标名、当前值、基线值或标准差倍数。这样值班同学收到告警的瞬间就能做初步判断不用再登录平台二次查询。去重逻辑放在Flink内部而不是网关实现方式是每个服务维护一个告警指纹状态同一指纹在恢复前只发一次恢复后再次异常就重新计算新指纹。这个方案上线以后告警量降了大概80%值班群终于从“刷屏”回归到“偶尔响一声”的平静状态。4. 常见问题与排查技巧实录4.1 时间窗口错位同样的数据曲线就是对不上这个坑是流式监控必修课。起初我们监控面板上某个服务的P99曲线和网关日志里实际记录的对不上大概整体偏移了8小时排查了一圈才发现是日志时间用了UTC时间戳而Flink恢复出来的EventTime直接按UTC处理落库到ClickHouse后Grafana又按北京时间展示偏移就产生了。如果只是显示问题还好更严重的是按天分区查询时统计口径会整体错一天。排查的思路是倒着找先看Flink作业的Watermark是否正常推进确认事件时间戳没有被错误解析再看ClickHouse里event_time字段的时区和真实值最后看Grafana的查询时区配置。最终我做了一个统一规则日志时间进Flink之前全部转成Unix毫秒时间戳落库时建表语句显式指定时区为Asia/ShanghaiGrafana查询时统一用浏览器时区。日志端的时间转换虽然多点工作量但换来的是全链路只有一个时区口径后面再也不需要做各种隐式换算。4.2 告警风暴阈值太敏感值班群变成告警刷屏第一版阈值用的是“超过均值一个标准差就告警”上线当晚值班群就炸了——几百条告警刷下来基本没人能看清到底哪里出了问题。事后分析业务接口天然存在秒级扰动尤其是定时任务分钟级触发时QPS会短暂冲高单点涌动根本代表不了真实故障。解决办法分两层。第一层是把告警门槛抬高需要同时满足“超过基线均值加2.5倍标准差”并且“持续超过两个滑动窗口”也就是持续10秒以上才告警。单轮抖动直接放过去持续异常才会被捕捉。第二层是引入告警抑制规则同一个服务的同一告警在5分钟内最多重发一次。这两层落地后假告警率降到原来的大约五分之一同时真正故障的告警一个没丢。团队对告警的信任度直接回升了这是监控项目最重要的事。4.3 Kafka rebalance导致的重复统计有一段时间升级完Kafka客户端的版本后Flink算出来的QPS偶发偏高每次持续几分钟又自己恢复正常。排查后发现是Flink的Kafka Source在作业恢复或Broker变更时发生了consumer group的rebalance在rebalance过程中部分分区被重复消费同一批信号被统计了两次。虽然Flink的Checkpoint机制按道理能保证精确一次语义但Source端在rebalance瞬时的消息重复仍会偶发穿透。解决方式是在聚合函数里加幂等控制给每一条信号记录生成一个指纹键状态里维护最近收到的指纹集合指纹存在就跳过。代价是状态会稍微膨胀一些但换来的是统计结果的确定性值。另有一个辅助手段是把KafkaSource的session.timeout.ms从默认的10秒调到60秒明显降低了心跳超时导致的rebalance频率。4.4 ClickHouse高基数查询把服务器资源打空雷达系统刚上线时有一面板展示每台机器的累计错误数量按实例IP拆分高峰时这个查询会把ClickHouse整个节点CPU打到60%以上。分析后发现根因是tags字段用的是Map(String, String)在做WHERE过滤时完全没法利用主键索引按IP查几十亿条记录基本等于全表扫描。解决方式是两招第一把高频查询字段从tags里拆出来单独建一个低基数字段instance_ip并加入ORDER BY组合键第二给tags字段建立BloomFilter索引用于非高频的偶发KV过滤。实测下来同一个查询从8秒多降到0.3秒左右。这个教训很现实设计表结构时凡是“查询条件里一定会出现的字段”一开始就应该提升为独立列不要偷懒全塞进KV里否则后面优化代价更大。4.5 Filebeat日志丢失与背压日志采集端的常见故障之一是Filebeat在Kafka不可达时进入背压消息积压在本地内存队列里队列溢出后的日志直接被丢弃。默认queue.mem.events只有2048条换算下来也就4MB左右在高流量业务下几十秒就能灌满队列开始丢数据。我调整成如下配置queue.mem: events: 65536 flush.min_events: 512 flush.timeout: 2s并且给采集端加了一个本地备份目录。Kafka长时间不可达时Filebeat先写本地文件缓存恢复连接后再补发。代价是Flink端需要有一定的数据积压容忍能力和我们前面设置的10秒乱序容忍度正好配合。这个措施基本保证了极端故障下日志信号不丢代价只体现在恢复后的短暂延迟上。5. 后续扩展方向与个人体会5.1 未来可以做的增强方向PLFM_RADAR目前的形态是常见信号的实时扫描但往“智能雷达”方向演进可以做的事情相当多。一个是根因关联。把告警事件、服务依赖图谱、链路追踪数据三者打通让告警消息直接携带影响拓扑。例如订单接口超时告警触发时自动关联它依赖的MySQL慢查询和下游Redis超时把相关性最高的线索排在最前面。另一个是异常形态分类区分趋势型异常、突变型异常和周期性异常不同形态套用不同检测方法线性回归盯趋势、CUSUM盯突变、季节性分解盯周期。再往下走可以接自动恢复比如确认无外部依赖故障的缓存服务过载直接触发限流动作而不是等在线值班的人来点按钮。但所有这些扩展都要守住一个原则监控系统的核心价值是把人从重复劳动中解放出来而不是制造一个更复杂的系统。每加一个功能都反问一句“这个功能让值班同学少做了什么操作”如果答案说不清楚宁可砍掉不做。5.2 拆掉重做几次之后的几句实在话断断续续做了两个月我最大的体会是监控系统的价值根本不在图表有多花哨而在“异常信号出现到值班同学定位原因”这段间隔。PLFM_RADAR把原本需要几十分钟的定位过程压缩到了3到5分钟。这里最值钱的其实是信号分层的设计而不是那套复杂的实时链路。基础设施、业务接口、日志三层信号放一起看很多异常单看任何一层都会被忽略交叉对比之后的指向性就跟探照灯似的。如果你也要搭一套类似的雷达系统我真心建议从最小闭环起步先接业务日志和接口信号两条数据源先用最简单的固定阈值告警跑稳两周再逐步引入滑动窗口、基线和多信号交叉。系统越简单越容易被人信任越复杂越容易在真正的故障来临时失灵。监控体系的最高境界不是告警又快又多而是每一条告警都能让人信服、每条告警都能导向一个值得处理的行动。最后再分享一个小技巧雷达看板上一定要给“最近一条恢复通知”留出位置。我最初的目的只是想展示告警已经恢复结果慢慢发现团队值班的人都会下意识去瞄那一个区域——他们需要的是“当前是不是安全”的确定性而不仅仅是“刚才哪里出了错”。这种确定性其实比一百个炫酷图表都更能给值班的同事吃一颗定心丸。
返回列表