ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

从零搭建一套 Flink 监控系统:API 采集、Metrics Reporter 与 InfluxDB + Grafana 可视化实战

从零搭建一套 Flink 监控系统:API 采集、Metrics Reporter 与 InfluxDB + Grafana 可视化实战 示例工程大数据【免费下载链接】flink-learningflink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API SQL 等内容的学习案例还有 Flink 落地应用的大型项目案例PVUV、日志存储、百亿数据实时去重、监控告警分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》项目地址https://gitcode.com/gh_mirrors/fl/flink-learning点击查看免费下载导读本篇文章源自仓库 books/flink-in-action-8.2.md是《Flink 实战与性能优化》一书的 8.2 节内容完整讲解如何在缺少公司级监控平台的情况下利用 Flink 自带的 Rest API 与 Metrics Reporter 体系自行搭建一套覆盖 JobManager、TaskManager 与 Job 三个层级的 Flink 监控系统。读完本文你将掌握通过 Chrome 开发者工具定位 Flink Web UI 背后的 REST API 并定时拉取监控数据的方法Counter、Gauge、Histogram、Meter 四种 Metrics 类型的注册与使用JMX、Prometheus、Prometheus PushGateway、InfluxDB 四类 Reporter 的配置方法以及基于 InfluxDB 时序数据库 Grafana 可视化面板落地一套完整监控系统的全过程。本文还结合本仓库 flink-learning-monitor 监控模块与 flink-learning-metrics 自定义 Metrics 示例源码给出可直接复用、可验证的配置与代码片段。8.2.1 利用 Rest API 获取监控数据熟悉 Flink 的读者都知道Flink 的 Web UI 已经详细展示了大量重要监控指标。因此如果暂时不想额外搭建监控系统直接使用 Flink 自带的 UI 即可获取大量监控信息。但更值得关注的是这些 UI 数据本质上都来自 Flink 自身的 Rest API。所以即使要搭建一套粗糙的监控平台也可以直接复用现有接口定时拉取数据把指标写入某种时序数据库再用可视化图表展示——一套完整的监控系统就成型了。下面以 Chrome 浏览器控制台为例演示如何定位这些提供监控数据的 REST API 与对应的页面形态。打开http://localhost:8081/overview可获取整个 Flink 集群的资源信息TaskManager 个数TaskManagers、Slot 总个数Total Task Slots、可用 Slot 个数Available Task Slots、Job 运行个数Running Jobs、Job 运行状态Finished 0 / Canceled 0 / Failed 0等。通过http://localhost:8081/taskmanagers页面查看 TaskManager 列表可以得到该集群下所有 TaskManager 的信息数据端口号Data Port、上一次心跳时间Last Heartbeat、总共的 Slot 个数All Slots、空闲的 Slot 个数Free Slots以及 CPU 与内存的分配使用情况。通过http://localhost:8081/taskmanagers/tm_id查看某个 TaskManager 的具体情况其中tm_id是一个随机的 UUID 值。在该页面除了上一条的监控信息外还可以查看该 TaskManager 的 JVM堆和非堆、Direct 内存、网络、GC 次数和时间。内存与 GC 指标非常重要——很多时候 TaskManager 频繁重启正是因为 JVM 内存设置不合理导致频繁 GC最终 OOM 崩溃被迫重启。在/taskmanagers/tm_id接口后追加/log即可查看该 TaskManager 的日志。需要特别说明Flink 中的日志与平常应用自己打的日志不同——Flink 日志以 TaskManager 为粒度打印而不是以单个 Job 为粒度。如果一个 Job 运行在多个 TaskManager 上日志就会散落在多个 TaskManager 中如果一个 TaskManager 同时运行多个 Job其日志会混在一起看起来既有这个 Job 的日志又有那个 Job 的日志。理解了这一点之前关于日志混乱的疑问就会解开。关于这种设计是否合理不同人有不同看法Flink 的 Issue 中有人提出希望日志能做到 Job 与 Job 之间隔离以便采集、查看与排查问题更快国内也有公司对这一部分做过改进。通过http://localhost:8081/jobmanager/config可查看 JobManager 的配置信息通过http://localhost:8081/jobmanager/log可查看 JobManager 日志详情。通过http://localhost:8081/jobs/job_id页面可查看 Job 的监控数据包含 Job 的 Task 数据、Operator 数据、Exception 数据、Checkpoint 数据等大量指标读者可在本地自行测试查看。以上列举的只是部分 REST API并非全部其核心价值在于既然这些接口是已知的我们就可以定时拉取对应监控数据绘制更直观酷炫的图表从而更好地掌控集群与作业状态。除利用 Flink UI 的接口定时获取监控数据外Flink 还提供多种 reporter 主动上报监控数据例如 JMXReporter、PrometheusReporter、PrometheusPushGatewayReporter、InfluxDBReporter、StatsDReporter 等可按需定制采集方案下文将逐一演示几个常用 reporter。相关 Rest API 细节可参考 Flink 官方文档的 Rest API Integration 章节monitoring/metrics.html#rest-api-integration。仓库配套实践定时采集 Job 概览指标本仓库的 flink-learning-monitor-collector 模块正好提供了一个利用 API 采集指标的最小示例——FlinkJobMetricCollect.javapublic class FlinkJobMetricCollect { public static void main(String[] args) { String jobManagerHost PropertiesUtil.defaultProp.get(flink.jobmanager.host).toString(); String jobOverviewResult HttpUtil.doGet(http:// jobManagerHost /jobs/overview); // 将返回的 JSON 解析、清洗后写入时序数据库即可作为监控数据的输入源 } }它从配置中读取flink.jobmanager.host调用HttpUtil.doGet请求http://jobmanager/jobs/overview获取所有 Job 的概览信息这正是用 REST API 定时采集监控数据这一思路的落地雏形——实际生产中只需把doGet放进定时任务并把 JSON 解析后写入 InfluxDB 等时序库即可。8.2.2 Metrics 类型简介在继承自RichFunction的函数中可以通过getRuntimeContext().getMetricGroup()获取 MetricGroup 并注册自定义指标。Flink 常见的 Metrics 类型有四种Counter计数器、Gauge瞬时值、Histogram直方图、Meter速率。本节结合仓库源码 flink-learning-metrics/src/main/java/com/zhisheng/metrics/custom 中的四个示例类逐一说明每种类型的使用方式。Counter计数器Counter 用于统计累计值例如处理过的记录数。常用操作是inc()自增通过getCount()读取当前值。仓库示例 CustomCounterMetrics.java 在RichMapFunction.open()中注册了三个 Counterindex getRuntimeContext().getIndexOfThisSubtask(); counter1 getRuntimeContext().getMetricGroup() .addGroup(flink-metrics-test) .counter(mapTest index); counter2 getRuntimeContext().getMetricGroup() .addGroup(flink-metrics-test) .counter(filterTest index); counter3 getRuntimeContext().getMetricGroup() .addGroup(flink-metrics-test) .counter(mapCounter, new SimpleCounter());addGroup(flink-metrics-test)用于给指标分组命名对应 Prometheus 等系统里的指标前缀.counter(名称)默认使用SimpleCounter也可显式传入自定义 Counter 实现。在map()中对counter1.inc()无条件自增仅当数据等于50或20时counter2.inc()借此可以验证不同数据条件下的计数差异。Gauge瞬时值Gauge 表示某个可变的瞬时值例如缓存大小、当前队列长度。核心方法是getValue()每次被读取时返回最新状态。仓库示例 CustomGaugeMetrics.java 用一个被 map 累加的value作为指标值private transient int value 0; Override public void open(Configuration parameters) throws Exception { super.open(parameters); getRuntimeContext().getMetricGroup() .addGroup(flink-metrics-test) .gauge(gaugeTest, new GaugeInteger() { Override public Integer getValue() { return value; } }); } Override public String map(String s) throws Exception { value; return s; }注意value被声明为transient且是任务内部维护的状态——Gauge 只在被 reporter 拉取时计算当前值因此非常适合暴露内存使用量这类快照型数据。Histogram直方图Histogram 用于统计数据分布例如延迟的分位数p75/p99。Flink 自身不提供实现需要借助 Dropwizard Metrics 的Histogram与Reservoir再包装成DropwizardHistogramWrapper。仓库示例 CustomHistogramMetrics.java 使用了SlidingWindowReservoir(500)滑动窗口保留最近 500 个样本com.codahale.metrics.Histogram dropwizardHistogram new com.codahale.metrics.Histogram(new SlidingWindowReservoir(500)); histogram getRuntimeContext().getMetricGroup() .addGroup(flink-metrics-test) .histogram(histogramTest, new DropwizardHistogramWrapper(dropwizardHistogram));调用histogram.update(s)更新样本后可读取统计信息getCount()获取样本数getStatistics().getMax()/getMin()/getMean()获取最大、最小、均值getStatistics().getQuantile(0.75)获取 75 分位数。分位数对判断延迟水位非常有价值。Meter速率Meter 用于统计每秒速率例如每秒处理多少条记录。同样基于 Dropwizard Metrics 实现包装成DropwizardMeterWrapper。仓库示例 CustomMeterMetrics.java 的核心逻辑com.codahale.metrics.Meter dropwizardMeter new com.codahale.metrics.Meter(); meter getRuntimeContext().getMetricGroup() .addGroup(flink-metrics-test) .meter(meterTest, new DropwizardMeterWrapper(dropwizardMeter));每条数据经过 map 时调用meter.markEvent()通过meter.getRate()读取当前速率events/secondmeter.getCount()读取累计次数。监控吞吐量时 Meter 是最直观的指标类型。8.2.3 利用 JMXReporter 获取监控数据JMXReporter 是 Flink 默认开启的 reportermetrics.reporter.jmx.factory.class默认配置即 JMX它会将 Flink 的全部指标暴露为 JMX MBean可通过jconsole等 JMX 客户端直接查看。其优势是零成本、无需额外依赖适合本地调试缺点是需要客户端主动连接 JVM规模化采集不太方便。在生产中它通常配合 Jolokia 这类 HTTP-JMX 桥接组件来暴露指标。开启方式在flink-conf.yaml中显式声明若自定义了 JMX 域可配置metrics.reporter.jmx.portmetrics.reporter.jmx.factory.class: org.apache.flink.metrics.jmx.JMXReporterFactory本仓库 flink-learning-extends 下还有自定义 Kafka/Prometheus 扩展指标的相关实现可作为参考了解 reporter 体系的可扩展性。8.2.4 利用 PrometheusReporter 获取监控数据PrometheusReporter 会把指标暴露为 Prometheus 格式的 HTTP 端点默认端口 9249Prometheus Server 通过scrape方式定时拉取。它适合Prometheus Server 主动抓取的架构配置示例flink-conf.yamlmetrics.reporter.prom.class: org.apache.flink.metrics.prometheus.PrometheusReporter metrics.reporter.prom.port: 9249使用前提是把flink-metrics-prometheus.jar放入 Flink 的lib目录。启动 JobManager/TaskManager 后访问http://host:9249/metrics即可看到形如flink_taskmanager_job_task_operator_numRecordsIn的指标。本仓库的 flink-learning-extends/flink-metrics/flink-metrics-prometheus 模块给出了基于 Flink 源码扩展的 Prometheus 实现其 README 说明了编译与部署方式。8.2.5 利用 PrometheusPushGatewayReporter 获取监控数据Prometheus PushGateway 架构适合指标主动推送的场景Flink 各节点定期把指标 POST 到 PushGateway再由 Prometheus Server 从 PushGateway 拉取。这在 Flink 运行于 Kubernetes、需要临时节点场景下尤其常见JobManager/TaskManager 地址不固定不适合主动 scrape。结合仓库 flink-metrics-prometheus/README.md 给出的完整配置可看到全部关键参数# Metrics Reporter metrics.reporter.promgateway.class: org.apache.flink.metrics.prometheus.PrometheusPushGatewayReporter metrics.reporter.promgateway.host: k8s # PushGateway 地址可替换为 localhost metrics.reporter.promgateway.port: 9091 # PushGateway 端口 metrics.reporter.promgateway.clusterMode: k8s # k8s 模式下自动识别集群信息 metrics.reporter.promgateway.jobName: flink-job # 推送到 PushGateway 的 job 标签 metrics.reporter.promgateway.randomJobNameSuffix: false # 是否随机追加 job 后缀多实例部署建议 true metrics.reporter.promgateway.deleteOnShutdown: true # 关闭时是否删除已推送的指标 metrics.reporter.promgateway.interval: 5 SECONDS # 推送周期其中interval: 5 SECONDS表示每 5 秒推送一次deleteOnShutdown: true可避免 Job 停止后旧指标残留在 PushGateway 上多个 Flink 实例共用 PushGateway 时randomJobNameSuffix: true能防止指标互相覆盖。8.2.6 利用 InfluxDBReporter 获取监控数据InfluxDBReporter 直接将 Flink 指标写入 InfluxDB 时序数据库省去中间采集组件。配置示例metrics.reporter.influxdb.class: org.apache.flink.metrics.influxdb.InfluxdbReporter metrics.reporter.influxdb.host: localhost metrics.reporter.influxdb.port: 8086 metrics.reporter.influxdb.db: flink metrics.reporter.influxdb.username: admin metrics.reporter.influxdb.password: admin metrics.reporter.influxdb.interval: 60 SECONDS注意reporter 的 class 需要对应版本的flink-metrics-influxdb依赖该 reporter 在部分 Flink 版本中是flink-metrics内置模块具体以所用 Flink 发行版为准interval控制写入周期默认 60 秒一次。写入 InfluxDB 后即可用 Grafana 直接查询展示。8.2.7 安装 InfluxDB 和 Grafana安装 InfluxDBInfluxDB 是开源时序数据库适合存储监控指标。安装完成后启动服务然后创建数据库与用户不同版本 CLI 语法略有差异以 1.x 为例influx CREATE DATABASE flink CREATE USER admin WITH PASSWORD admin WITH ALL PRIVILEGES对应 InfluxDBReporter 配置中的db: flink、username: admin、password: admin即指向这里创建的库与账号。启动后可用curl http://localhost:8086/ping验证服务是否正常。安装 GrafanaGrafana 是开源的可视化平台安装启动后默认监听 3000 端口首次访问用默认账号admin/admin登录并修改密码。登录后进入 Data Sources 页面添加 InfluxDB 数据源填写 URL如http://localhost:8086、数据库名与账号密码保存后即可基于该数据源创建 Dashboard 和 Panel。8.2.8 配置 Grafana 展示监控数据数据链路整体为Flink JobManager/TaskManager内置指标 │ 通过 Rest API 定时拉取 / 通过 Reporter 主动上报 ▼ 时序数据库InfluxDB │ Grafana 查询 ▼ Grafana Dashboard 可视化面板在 Grafana 中新建 Panel 时使用 InfluxDB 查询语法按 measurement 筛选指标例如查询 JobManager 的 JVM Heap 使用量、TaskManager 的 GC 次数、Job 的 Checkpoint 完成数等。本仓库 flink_monitor_measurements.md 完整列出了可直接用于建面板的核心指标名分为三大类JobManager 监控指标如jobmanager_Status_JVM_Memory_Heap_Used堆内存使用、jobmanager_Status_JVM_GarbageCollector_*GC 次数与耗时、jobmanager_job_numberOfCompletedCheckpoints已完成 Checkpoint 数、jobmanager_job_lastCheckpointDuration最近一次 Checkpoint 耗时、jobmanager_job_uptime/downtime作业运行/宕机时间、jobmanager_numRegisteredTaskManagers注册的 TM 数、jobmanager_taskSlotsAvailable/Total可用/总 Slot 数等TaskManager 监控指标如taskmanager_Status_JVM_Memory_Heap_Used、taskmanager_Status_JVM_GarbageCollector_G1_Old_Generation_Count/Time、taskmanager_Status_Network_AvailableMemorySegments、taskmanager_Status_Shuffle_Netty_*网络缓冲等Job 监控指标如taskmanager_job_task_operator_numRecordsIn/Out及*PerSecond各算子输入输出与速率、taskmanager_job_task_numBytesIn/Out*字节吞吐、taskmanager_job_task_currentInputWatermark当前水位线、taskmanager_job_task_operator_numLateRecordsDropped迟到丢弃记录数、taskmanager_job_task_checkpointAlignmentTime对齐耗时等。仓库配套实践监控数据的 SQL 消费与持久化除了 InfluxDB Grafana 这一方案本仓库还给出了Flink SQL 消费监控指标写入其他存储的替代思路。flink_metrics_2es.sql 演示了用 Flink SQL 从 Kafka 读取metrics-flink-jobs主题中的 Flink 任务消费延迟指标name、fields、tags结构再insert into打印或写入下游的用法——即监控数据也可以先进消息队列再通过 Flink SQL 做 ETL 后落入 ES/其他存储再交给可视化平台展示。同时 flink-learning-monitor 模块还涵盖了监控告警flink-learning-monitor-alert、支持钉钉/短信/邮件通知、日志处理flink-learning-monitor-log、PV/UV 统计flink-learning-monitor-pvuv与可视化展示flink-learning-monitor-dashboard共同构成一套完整的监控生态。8.2.9 小结与反思本节系统讲解了如何利用 Flink Rest API 获取监控数据、四种 Metrics 类型Counter/Gauge/Histogram/Meter的注册与使用以及 JMXReporter、PrometheusReporter、PrometheusPushGatewayReporter、InfluxDBReporter 四类常用 reporter 的配置方法最后通过 InfluxDB Grafana 落地了一套完整的 Flink 监控系统。作业部署上线后的监控尤其重要。虽然 Flink UI 自身提供了不少监控信息但整体功能相对较弱还是应当搭建一套完整的监控系统对 JobManager、TaskManager 和作业本身进行持续监控。本章的方案可以归纳为三条主线采集既可通过 Web UI 背后的 Rest API如/overview、/jobs/overview、/taskmanagers、/jobs/job_id定时拉取也可通过 Metrics Reporter 主动上报存储监控数据写入时序数据库InfluxDB 是常用选择也可以按公司情况换成其他时序库或消息队列 存储的组合展示与运维用 Grafana 这类可视化工具呈现指标变化达到真正的监控运维效果。另外需要反思的是整套方案并不要求绑定某个具体存储或可视化组件——Grafana 支持多种数据源InfluxDB 也可以替换为 Prometheus、Elasticsearch 等。你们公司的监控系统架构是怎样的是否可以复用、改造本节这套方案直接接入你们现有的监控体系答案应当是肯定的——只要抓住API/Reporter 采集 → 时序存储 → 可视化这条主线就能以最小成本为 Flink 作业建立起可用的监控运维能力为作业出问题后的排查节省大量时间。赞分享示例工程大数据【免费下载链接】flink-learningflink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API SQL 等内容的学习案例还有 Flink 落地应用的大型项目案例PVUV、日志存储、百亿数据实时去重、监控告警分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》项目地址https://gitcode.com/gh_mirrors/fl/flink-learning点击查看免费下载相关推荐CloudExplorer Lite故障排查手册常见问题与解决方案汇总CloudExplorer Lite故障排查手册常见问题与解决方案汇总 CloudExplorer Lite作为一款开源的轻量级云管平台在实际部署和使用过程3步搭建Apache Flink Metrics可视化监控Grafana仪表盘配置指南3步搭建Apache Flink Metrics可视化监控Grafana仪表盘配置指南 你还在为Flink集群性能问题排查头疼本文将通过3个步骤教你从零后端大数据流处理批处理EVCC数据可视化InfluxDBGrafana监控面板搭建EVCC数据可视化InfluxDBGrafana监控面板搭建 引言为什么需要专业的充电数据监控 随着电动汽车的普及家庭充电管理变得越来越重要。EVCC后端前端智能硬件物联网能源管理上一篇Apache PLC4X构建工业物联网统一接口的技术架构深度解析下一篇VMware Unlocker 4.2.7打破硬件壁垒让普通PC也能运行macOS虚拟机的革命性方案创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进