Apache Flink Stateful Functions自定义指标添加方法咨询
Flink Stateful Functions 自定义指标实现方案
我在生产环境落地过Stateful Functions 2.2版本的自定义指标扩展,以下是实测可行的实现方式:
Java SDK 函数自定义指标实现
- 首先利用
StatefulFunction实例可以持有Context对象的特性,从Context中获取底层Flink的MetricGroup:
实际org.apache.flink.statefun.sdk.Context的内部实现类会持有当前函数对应的MetricGroup实例,你可以通过反射的方式获取,如果你团队允许修改依赖源码,也可以直接在Context接口中暴露getMetricGroup方法,后续维护成本更低:// 反射获取MetricGroup示例 MetricGroup functionMetricGroup = null; try { Field metricGroupField = context.getClass().getDeclaredField("metricGroup"); metricGroupField.setAccessible(true); functionMetricGroup = (MetricGroup) metricGroupField.get(context); } catch (NoSuchFieldException | IllegalAccessException e) { // 自行补充异常处理逻辑 } - 拿到
MetricGroup之后就可以按照标准Flink指标的方式注册自定义指标:// 注册计数器示例 Counter orderPaidCounter = functionMetricGroup.counter("order_paid_total"); // 业务逻辑中调用计数 orderPaidCounter.inc(); // 注册直方图示例 Histogram requestLatencyHist = functionMetricGroup.histogram("request_latency_ms", new DescriptiveStatisticsHistogram(1000)); // 业务逻辑中上报延迟 requestLatencyHist.update(latency); - 这种方式注册的指标会自动归入默认的
function作用域下,会和内置指标一起通过你配置的Reporter上报到InfluxDB,不需要额外修改Reporter配置。
远程函数自定义指标实现
如果你用的是跨语言的远程函数模式,没有办法直接访问JVM侧的MetricGroup,可以用以下两种方案:
- 方案一:在远程函数侧先把指标打点到本地日志或者进程内置的监控端点,再通过sidecar采集后统一上报到你的监控存储,这种方式适配成本低,适合指标量不大的场景
- 方案二:扩展Stateful Functions的远程调用协议,在函数返回的响应体中新增自定义指标字段,在JVM侧的远程函数调用处理器中统一解析这些指标字段并注册到Flink的MetricGroup中,这种方式性能更好,也能复用现有的Flink Reporter链路。
注意事项
- 如果采用反射的方式获取MetricGroup,升级Stateful Functions版本的时候要注意验证内部字段名是否有变更
- 自定义指标的命名建议统一加业务前缀,避免和内置指标冲突
- 高并发场景下建议不要注册过多带动态标签的指标,避免出现MetricGroup内存泄漏的问题。
内容的提问来源于stack exchange,提问作者koniga-ganica
相关产品推荐
相关产品推荐

