如何获取Flink各算子维护的State大小并实现自动监控告警?
Flink 1.8 RocksDB状态后端算子状态大小获取与告警方案
能否获取每个算子的State大小
可以获取,Flink 1.8 已原生支持RocksDB状态后端的状态大小指标采集,无需额外二次开发即可拿到对应数值。
具体实现步骤
1. 开启RocksDB状态指标采集
在flink-conf.yaml中添加如下配置,开启RocksDB原生指标暴露:
# 开启RocksDB状态指标 state.backend.rocksdb.metrics: true
开启后会自动暴露两个核心相关指标(subtask粒度):
rocksdb.estimate-live-data-size:单subtask维护的状态总大小估算值,单位为字节rocksdb.estimate-num-keys:单subtask维护的状态key总数估算值
将同一算子下所有subtask的rocksdb.estimate-live-data-size值求和,即可得到该算子当前维护的总状态大小。
2. 对接指标上报链路
推荐使用Flink内置的指标Reporter将数据上报到通用监控系统,对作业性能影响极小,无需侵入业务代码:
- 若使用Prometheus作为监控存储,配置示例如下:
metrics.reporter.prom.class: org.apache.flink.metrics.prometheus.PrometheusReporter metrics.reporter.prom.port: 9250-9260
配置完成后重启作业,Prometheus即可从配置的端口抓取到状态指标,所有指标默认携带operator_name、job_id等标签,可直接按算子维度聚合数值。
- 若需要上报到自定义监控平台,可自行实现Flink的
MetricReporter接口,在接口实现中按固定周期拉取状态指标聚合后上报即可。
3. 告警规则配置
在监控系统中配置对应告警规则即可实现状态异动告警,常见规则参考:
- 同比1小时前状态大小涨幅/跌幅超过30%触发告警
- 状态大小超过预设的业务阈值(比如单算子状态超过100G)触发告警
- 状态大小连续3个采集周期持续上涨触发告警
注意事项
- RocksDB暴露的状态大小为估算值,误差范围通常在10%以内,完全满足趋势监控和告警需求,不需要追求绝对精确。
- 指标统计的是算子全量状态大小,和是否开启增量检查点无关,符合总状态大小统计要求。
- 尽量不要在业务代码中直接嵌入状态大小采集逻辑,会增加作业负担,通过外部指标采集的方式稳定性更高。
内容的提问来源于stack exchange,提问作者user239874
相关产品推荐
相关产品推荐

