ARTICLE DETAIL

资讯详情

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

平台流量雷达:基于Filebeat+Kafka+Flink的实时监控告警方案

平台流量雷达:基于Filebeat+Kafka+Flink的实时监控告警方案 如果你是一个平台的开发者或运维一定经历过这样的事大促前夜流量突然暴涨接口延迟飙升你打开监控面板却只能看到一片堆叠的折线根本看不出问题出在哪或者某个晚上服务器被爬虫刷爆短信告警响了一宿你被叫起来三次第二天看日志才发现异常流量来源早就摆在那儿只是没人盯着。PLFM_RADAR这套系统就是用来解决这个问题的。PLFM是Platform的缩写RADAR即雷达合起来就是“平台雷达”。它不是一个具体的开源项目而是一套面向平台流量与运行状态的实时监控告警方案对标我们常说的业务监控系统中的“流量态势感知”模块。它能做的事情很明确实时采集平台入口流量与核心接口的访问日志通过规则引擎识别异常波动再通过可视化大盘和告警通道把问题在用户感知之前暴露出来。适合后端开发、运维工程师、SRE、以及所有被线上问题折腾过但又不想继续被动挨打的同学参考。这套方案我从零搭过两轮第一轮用开源组件硬拼第二轮把规则引擎和告警逻辑做了重构。过程中踩了不少坑也总结出一套可以落地复制的打法。下面我把整体设计、技术选型、核心实现和排障经验完整拆开讲。1. 项目背景与整体设计思路1.1 为什么要给平台装一套“雷达”先讲一个很具体的场景。某个平台上线了新功能第二天PV翻了一倍这本该是好事。但实际结果是后端数据库连接池被打满接口平均响应时间从80ms飙升到2.3秒部分用户开始反馈页面白屏。这时候问题来了监控面板上所有图表都在动但没有人能第一时间判断出“流量翻倍是不是正常业务增长”。如果没有一套基于流量特征的雷达系统你只能等用户投诉、等告警短信响、等自己登录服务器敲命令看日志。信号是有的但缺乏“过滤、放大、定向追踪”的环节。所谓平台雷达做的就是这三件事过滤掉正常波动放大异常信号定位到具体来源和接口。我做的PLFM_RADAR定位非常明确不替代APM应用性能监控不替代链路追踪只关注平台入口的“流量三要素”——速率、来源、质量。速率就是QPS每秒请求数来源就是访问方IP、UA、地理区域质量就是响应时间、错误率、成功率。这三个要素覆盖了大多数平台级故障的第一现场。另一个原因是成本。很多团队一上来就上全套可观测性平台投入大、周期长而雷达系统可以以很轻的方式先跑起来。我自己第一版就是一台4核8G的机器加一套开源组件两周内就跑出了能用的监控大盘。1.2 系统设计的三条核心原则这套系统在设计时有三个原则我建议任何做监控的同学都尽量遵守。第一条实时性优先。雷达的语义就是“尽早发现”所以从日志产生到告警触发的端到端延迟必须控制在30秒以内。这意味着采集端要轻量、传输链路要短、计算环节要尽量少做复杂聚合。我在第一版中用Logstash做采集和解析结果发现它在高吞吐下CPU占用很高后来换成了Filebeat延迟降了一半。第二条可解释性优先。告警不能只说“流量异常”必须能清楚地回答三个问题哪个入口、什么指标、和什么基线相比。比如一条告警应该写成这样“登录接口QPS 5分钟均值8321环比1小时均值上涨312%疑似爬虫活动”。我见过太多监控系统给一堆数字却没有任何上下文值班人员看到告警还要自己去翻图表这种体验非常糟糕。第三条低侵入接入。雷达系统的数据来源主要是访问日志和网关指标不需要在业务代码中埋大量点。这个原则决定了系统能不能快速推广到多个业务线。我的做法是统一从Nginx和网关层采集日志业务方基本无感知。1.3 与现有监控体系的边界划分搭PLFM_RADAR之前团队里已经有Prometheus Grafana的指标监控体系。为什么还要另一套系统听起来重复其实两者解决的问题完全不同。Prometheus擅长监控“机器和进程”层面的指标比如CPU、内存、磁盘IO、请求延迟分位数它采集的是数值型指标按固定频率拉取。但它对“请求来源是哪里”“访问频率的形态是怎样的”“某个UA段是不是突然增多”这类多维分析场景支持较弱日志数据天然是文本要先把标签提取出来才能建模。PLFM_RADAR则专注在“流量行为”层面它的输入是日志文本输出是结构化的流量画像。两者正好互补一个看“机器是否健康”一个看“流量是否正常”。在实际排障中后者的预警时间通常更早因为任何问题最终都会先反映为流量的形态变化。2. 技术选型与架构拆解2.1 数据采集层Filebeat Nginx日志为什么放弃Agent埋点数据采集是整个雷达的“天线”天线收不到信号后面全是白搭。我最终选定的组合是Filebeat采集Nginx访问日志直接发送到Kafka。这套链路最大的特点是稳Filebeat内部有背压机制当Kafka写入变慢时它会自动降速不会丢数据而且它维护了读取文件的offsetFilebeat重启后可以从上次位置继续读不会重复发送。为什么不推荐在业务代码里埋点做Agent采集两条理由。一是侵入性太强每个服务都要引入SDK对于老系统来说改动成本极高二是日志丢失问题业务进程一重启内存中的日志队列就没了可靠性反而不如文件型日志。文件型日志的优势在于只要进程还在写文件Filebeat就能读到就算Filebeat挂了日志还在磁盘上躺着事后可以补。Nginx日志格式一定要包含关键字段我使用的log_format如下log_format radar $remote_addr|$time_iso8601|$request_method|$request_uri| $status|$body_bytes_sent|$request_time|$http_user_agent| $http_referer|$upstream_addr|$upstream_response_time;管道符|做分隔符是个实用细节比用空格不容易出错。有些UA里带有空格用默认的combined格式解析时会错位而管道符在合法请求中几乎不会出现。2.2 消息队列与流处理Kafka做缓冲Flink做窗口计算采集层之后的第一个枢纽是Kafka。Kafka在这里的作用是削峰填谷——流量突发时日志会产生一个尖峰如果直接写入分析引擎会被打垮Kafka可以先把消息堆积起来让下游按自己的节奏消费。这就像雷达信号接收机后面的缓冲存储器保证瞬间的大量信号不会导致系统崩溃。我用的版本是Kafka 3.4Topic按业务入口分了几类默认分区数设为12。分区数不是越大越好分区太多会增加文件句柄开销和Leader切换成本我测试下来12到24个分区在这个数据量级是最稳的。流处理层我选了Flink主要看中三点真正的流式计算毫秒级延迟不像Spark Streaming本质上是微批自带窗口机制滚动窗口、滑动窗口、会话窗口都是现成的状态管理能力强可以做去重、计数、比值这类需要跨事件计算的操作。Flink消费Kafka里的原始日志按分钟粒度做聚合产出QPS、错误率、响应时间均值、P95这些指标再写入下游存储。2.3 存储与可视化ClickHouse承载多维分析Grafana做展示指标数据落库我选了ClickHouse。原因很直接日志聚合后的数据是典型的时序数据带有大量维度标签接口名、IP段、UA分类、区域查询模式基本是“按时间范围 若干维度做聚合”这正是列式存储的强项。我建的表结构大致如下CREATE TABLE radar.metric_1m ( ts DateTime, endpoint String, ip_segment String, ua_type String, qps UInt32, error_count UInt32, avg_response_ms Float32, p95_response_ms Float32 ) ENGINE MergeTree PARTITION BY toYYYYMMDD(ts) ORDER BY (ts, endpoint, ip_segment);排序键的设计很关键我踩过坑一开始ORDER BY只写了ts查询按endpoint过滤时全表扫描速度慢得离谱。改成复合排序键之后相同endpoint的数据在存储上连续查询性能提升了十倍不止。可视化层用的是Grafana告警规则引擎我先用Prometheus风格的表达式后来改成了自定义Python规则服务这个后面详细讲。Grafana的好处是插件生态好、告警通知能直连企业微信和钉钉而且支持变量模板可以做“先选接口、再看趋势”的交互式大盘。3. 核心功能解析与实操要点3.1 流量总览大盘把“平台健康状况”放在一块屏上雷达总览页我设计了九个核心面板每个面板回答一个运营问题面板名称指标回答的问题实时QPSQPS当前请求量多大请求成功率success_rate整体健康吗P95响应时间p95用户体感卡吗TOP10接口qps_by_endpoint哪里流量最集中错误状态码分布status_5xx有服务挂了吗来源IP地域分布geo_map流量从哪里来UA类型占比ua_pie是浏览器还是爬虫环比异常幅度diff_ratio和平时比差了多少上游响应时间upstream_time后端服务慢不慢Grafana面板配置上有一个被我反复使用的技巧所有图表使用同一个时间变量$interval再配合$endpoint做模板过滤。这样看板就从固定图表变成了一个可交互的分析工具点一下接口名全部面板联动刷新排查效率高很多。3.2 异常检测规则引擎从固定阈值到动态基线雷达的“识别异常”能力全部体现在规则引擎上。这个模块我重构过一版第一版全部是固定阈值规则效果不理想——业务低谷期的晚上固定阈值会漏报大促期间的白天固定阈值会误报。后来我改成“动态基线 固定线”双模式才真正稳定下来。动态基线实现方式对每个接口维护一个历史窗口比如过去7天同一时段计算均值和标准差。当前值超过均值 3*标准差时判定为异常。这里有一个统计常识要注意流量数据普遍是长尾分布直接用均值容易偏离我实际采用中位数 3 * MAD绝对中位差来定义基线鲁棒性更好。规则的应用示例def is_anomaly(current, history): median np.median(history) mad np.median(np.abs(history - median)) if mad 0: mad 0.01 # 超过基线的4个MAD判定异常约等于3倍标准差的效果 threshold median 4 * mad return current threshold, threshold规则的另一个关键参数是连续触发次数。只有当连续3个数据点都超过阈值时才告警这个参数直接决定了雷达会不会“草木皆兵”。设置太短容易误报太长又会拖慢发现速度我实测3次是个平衡点对应的时间成本是3分钟。3.3 告警通道与值班降噪推给对的人不轰所有人告警的最终目的是让人采取行动所以每次告警必须包含足够的信息才能做到这一点。我最终格式定型为【平台雷达告警】 级别严重 接口/api/v1/user/login 指标QPS 当前值83215分钟均值 基线值2018 偏离倍数4.1x 时间2024-06-18 21:35 来源地域华东IP占比58%异常UAScrapy 可能原因疑似爬虫活动或热门活动上线这样一条消息推送到企业微信后值班同学不需要打开电脑就能判断优先级。告警接收人按业务线配置接口A的告警只发给业务A的后端群接口B的告警发给业务B的群。避免全局广播产生“狼来了”效应——所有人都在收不相关的告警最终结果就是所有人都不认真看告警。3.4 访问来源画像识别爬虫与恶意流量雷达系统一个很有价值的应用是建立访问来源画像识别爬虫和异常流量。我的实现方式是基于UA用户代理关键字和User-Agent聚类的双重判断。同行对比数据正常浏览器的UA里必然包含Mozilla、Chrome、Safari这些标记而爬虫的UA通常只有一段很简短的字符串或者包含Python-requests、Go-http-client、Scrapy这样的特征。还有一种更隐蔽的伪装了完整UA但请求频率不符合人类行为——人不可能在1秒内请求5个完全不同页面的接口。我设计了一张UA特征表按规则打分特征分值UA长度小于30字符2包含pythoncurl不包含Mozilla/Chrome/Safari/Edg2单IP每秒请求超过20次4累计分值大于6标记为爬虫这个模型不追求100%准确但能筛掉绝大多数的明显异常流量。识别出来之后规则引擎会触发一条“爬虫活动告警”级别的通知同时在雷达大盘上给这条流量路径单独打标方便后续排查。4. 实操过程与关键环节实现4.1 环境准备与组件部署整套系统以Docker Compose方式部署包含四个服务Filebeat、Kafka、Zookeeper、Flink。ClickHouse和Grafana我建议独立部署放在同一台机器上会互相抢资源。我的部署目录结构如下radar/ ├── docker-compose.yml ├── filebeat/ │ ├── filebeat.yml │ └── Dockerfile ├── flink/ │ ├── flink-job.jar │ └── flink-sql.yaml ├── rules/ │ └── anomaly_rules.py └── scripts/ └── kafka_create_topic.shdocker-compose.yml中Kafka和Zookeeper直接使用官方镜像Flink使用Standalone模式带一个TaskManager。这里给个提醒Kafka的ADVERTISED_LISTENERS一定要配置成外部可访问的地址否则Flink容器内部连不上Kafka这个坑我调试了整整一个下午。部署完之后的验证方法很简单用kafka-console-consumer订阅源Topic然后随便请求一下Nginx接口如果能实时看到日志内容整条采集链路就没问题了。4.2 Flink实时指标计算的核心逻辑Flink作业的核心逻辑是从source Topic消费日志解析字段然后开一个1分钟的滑动窗口做聚合运算。我用的是DataStream API关键代码如下DataStreamMetric metrics logs .map(new LogParser()) // 解析日志为LogEvent .keyBy(event - event.getEndpoint()) .window(TumblingEventTimeWindows.of(Time.minutes(1))) .aggregate(new MetricAggregator()); public static class MetricAggregator implements AggregateFunctionLogEvent, MetricAccumulator, Metric { Override public MetricAccumulator createAccumulator() { return new MetricAccumulator(); } Override public MetricAccumulator add(LogEvent event, MetricAccumulator acc) { acc.setTs(event.getTimestamp()); acc.setEndpoint(event.getEndpoint()); acc.setIpSegment(event.getIpSegment()); acc.setQps(acc.getQps() 1); if (event.getStatus() 500) { acc.setErrorCount(acc.getErrorCount() 1); } acc.addRT(event.getRequestTime()); return acc; } Override public Metric getResult(MetricAccumulator acc) { Metric metric new Metric(); metric.setTs(acc.getTs()); metric.setEndpoint(acc.getEndpoint()); metric.setIpSegment(acc.getIpSegment()); metric.setQps(acc.getQps()); metric.setErrorCount(acc.getErrorCount()); metric.setAvgResponseMs(acc.getAvgRT()); // 按分钟窗口时间戳向下取整 metric.setTs(WindowTimestampHelper.roundToMinute(acc.getTs())); return metric; } }时间窗口还有一个注意点使用EventTime而不是ProcessingTime。ProcessingTime指的是Flink所在机器的当前时间一旦Flink作业重启或者消费滞后窗口数据就会错乱。EventTime基于日志里记录的时间戳配合Watermark可以正确处理数据乱序问题。Watermark我这里设了10秒的延迟容忍度超过这个范围的数据直接丢弃避免窗口一直不触发。4.3 告警规则服务的实现与热更新告警规则我单独做成了一个Python服务常驻运行每30秒从ClickHouse读取最近分钟指标再跑一遍规则集。这种设计比在Grafana里配告警更灵活因为规则可以写得非常复杂而且支持热更新——改了规则文件不用重启服务每30秒自动重载一次。核心的规则调度代码如下def evaluate_rules(): rules load_rules_from_yaml(rules.yaml) recent_metrics query_recent_metrics(minutes5) alerts [] for rule in rules: for endpoint in rule.endpoints: history query_history_series(endpoint, rule.window) current recent_metrics.get(endpoint) if current and current.qps rule.threshold_based_on(history): alerts.append(build_alert(rule, endpoint, current, history)) return alertsrules.yaml里定义了一条针对登录接口的动态基线规则- name: login_qps_surge endpoints: [/api/v1/user/login] metric: qps strategy: dynamic_baseline baseline_window: 7d current_window: 5m multiplier: 4.0 trigger_count: 3 level: serious规则引擎里我特别注意了一个问题告警去重与恢复。同一个接口的同一个异常如果没有恢复不能每分钟都发告警否则值班人员的手机就废了。我用了Grafana风格的告警状态机一条告警先进入pending状态连续触发达到trigger_count次后转为firing接口恢复正常后定时扫描任务会把它标记为resolved并发送一条恢复通知。4.4 Grafana大盘配置与参数选择Grafana面板的配置不需要代码但有几个参数选择需要经验。首先是数据源我选了ClickHouse的官方插件版本一定要和ClickHouse Server匹配否则查询时会报协议错误。每个仪表板的查询语句我都会做一次“维度下钻”设计。比如总览面板的QPS图查询语句是SELECT ts, sum(qps) AS qps FROM radar.metric_1m WHERE ts now() - INTERVAL 1 HOUR GROUP BY ts ORDER BY ts点击任意数据点想看详情时Grafana变量$endpoint会替换到WHERE条件中。为了这个联动效果我费了不少工夫搞清楚Grafana模板变量和SQL子句的拼接顺序现在总结出来就一条经验模板变量名在SQL里用${endpoint}占位配合All值的特殊处理联动效果非常自然。面板刷新率我设为30秒查询范围默认显示最近1小时提供7天/1天/6小时/1小时快捷切换按钮同时脸大屏上只用暗色主题视觉疲劳感低很多。5. 常见问题与排查技巧实录5.1 数据链路延迟为什么告警总是慢半拍问题描述Grafana图表延迟明显明明已经进入Kafka的数据Flink计算完写入ClickHouse后图表要过两三分钟才显示。查了一圈发现瓶颈根本不在计算。排查过程首先确认Kafka积压情况用kafka-consumer-groups --describe查看消费者lag结果没有明显积压再看Filebeat的采集速度发现默认配置里Filebeat对单文件有读取限制multiline配置不对也会拖慢最后发现在Flink的JVM参数里默认的Heap大小只有1GB开启Checkpoint后经常Full GC导致计算速度跟不上。解决Filebeat配置调大了harvester的读取缓冲区Flink增加TaskManager的内存并调整了并行度。修改后端到端延迟稳定在15秒以内。5.2 误报风暴动态基线模型为什么总在中秋节当晚触发告警问题描述每个中国传统节日前一晚流量总和平时不一样动态基线模型会触发批量告警一晚上几百条。排查过程我仔细看过数据发现节假日之前的流量曲线有一个“提前回落”的特征——大家都准备下班过节去了流量跌到低谷这时候假设有小批量的晚间活动流量对比基线的偏离倍数很容易超过阈值。基线模型本身没有错错在它不懂“节日效应”。解决给规则引擎加了一个“日历自定义”模块在重要节日前N天自动拉长基线窗口用过去14天而不是7天的数据做对比同时降低偏离倍数的敏感度。效果立竿见影误报率下降了80%。5.3 Kafka消费者组频繁Rebalance问题描述Flink作业运行几个小时后下游数据突然中断几分钟日志里全是“Rebalance in progress”的警告。排查过程Rebalance的典型原因是消费者心跳超时要么是处理耗时太长要么是最大轮询间隔设置太短。我一开始怀疑是某个接口的数据量太大导致处理变慢后来跟踪日志发现其实是ClickHouse写入并发太高偶发阻塞拖长了Flink的算子处理时间。解决调整两个参数——max.poll.interval.ms提高到了5分钟同时把ClickHouse的写入改为批量异步提交。另外还加入了独立的写入线程池避免写入阻塞主处理流程。之后Rebalance没有再出现过。5.4 常见问题速查表症状可能原因解决方式Kafka有数据但Flink不消费Group ID冲突/消费组被占用重置消费者组并清理offset告警重复推送缺少告警状态去重引入pending/firing/resolved状态机时间窗口不触发EventTime与ProcessingTime混淆检查Watermark设置统一用日志时间戳ClickHouse查询超时排序键未包含过滤字段修改ORDER BY为复合排序键Grafana无数据数据源版本不兼容升级ClickHouse数据源插件爬虫封禁误伤正常用户UA特征规则过于激进加入IP历史行为分析后再封禁凌晨流量异常告警当日基线窗口太短延长基线窗口为24小时磁盘快速增长Kafka保存时间太长调整log.retention.hours为12小时大屏刷新卡顿Grafana图表数量过多拆分面板使用Dashboard变量这套雷达系统上线运行后最大的变化是排障模式从“被动响应”变成了“主动发现”。之前凌晨两点被电话叫醒的滋味相信每个运维都懂。后来再出现接口异常我们的值班流程变成了收到告警消息打开雷达大盘看异常来源定位到具体IP或UA然后直接在网关层处理——整个过程不需要登服务器。如果你们的平台也经常出现“流量一涨就慌流量一跌更慌”的情况我的建议是不要先急着上全套可观测平台可以按这个雷达的架构用两周时间先跑起来。数据采集用Filebeat Kafka流式计算用Flink存储查询用ClickHouse可视化用Grafana告警规则单独做一个服务掌握在自己手里。这套组合已经经过生产验证自己改造成本也不高。最后再分享一个小技巧每个接口的监控阈值不要一次性配死先开一周的“观察模式”不告警、只记录。一周后你手里有了每个接口在真实业务下的流量分布数据再回头定基线、配规则你会惊讶地发现原来你以为的“异常”其实只是人家的日常。雷达这东西先学会分辨正常声音你才能听出真正的危险信号。
返回列表