You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flink中无getRuntimeContext()时如何在自定义ScalarFunction中跟踪指标

问题描述

我实现了一个处理JSON字符串的ScalarFunction类ArraySizeUdf,用于返回JSON数组的长度,代码如下:

public class ArraySizeUdf extends ScalarFunction {
    private static final Logger LOG = LoggerFactory.getLogger(ArraySizeUdf.class);
    private final static ObjectMapper mapper = new ObjectMapper();

    public int eval(String stringJsonArray) {
        try {

            if (stringJsonArray == null) {
                return 0;
            }

            JsonNode actualObj = mapper.readTree(stringJsonArray);
            ArrayNode aa = (ArrayNode) actualObj;
            return aa.size();
        } catch (Exception e) {
            LOG.error("Error deserializing json to find size : {}", stringJsonArray);
            return -1;
        }
    }
}

已通过streamTableEnvironment.createTemporaryFunction("array_size", ArraySizeUdf.class);注册该UDF,并在SQL中以select array_size('["some", "other"]') from table;方式使用。现在希望统计异常发生次数,但该类无法访问getRuntimeContext(),请问如何在此场景下跟踪自定义指标?

解决方案

方法1:改用RichScalarFunction基类

Flink的RichScalarFunction提供了访问RuntimeContext的能力,可直接注册和更新自定义指标,这是最标准的实现方式:

public class ArraySizeUdf extends RichScalarFunction {
    private static final Logger LOG = LoggerFactory.getLogger(ArraySizeUdf.class);
    private final static ObjectMapper mapper = new ObjectMapper();
    private Counter exceptionCounter;

    @Override
    public void open(FunctionContext context) throws Exception {
        super.open(context);
        // 注册自定义计数器指标
        exceptionCounter = context.getMetricGroup().counter("json_parse_exception_count");
    }

    public int eval(String stringJsonArray) {
        try {
            if (stringJsonArray == null) {
                return 0;
            }
            JsonNode actualObj = mapper.readTree(stringJsonArray);
            ArrayNode aa = (ArrayNode) actualObj;
            return aa.size();
        } catch (Exception e) {
            LOG.error("Error deserializing json to find size : {}", stringJsonArray);
            exceptionCounter.inc(); // 异常发生时计数器自增
            return -1;
        }
    }
}

UDF注册方式保持不变,之后可在Flink Web UI、Prometheus等监控组件中查看json_parse_exception_count指标的实时数值。

方法2:静态计数器结合手动Metric注册(兼容原基类)

若无法切换基类,可使用静态变量作为计数器,再手动将其注册到Flink Metric系统:

  1. 修改UDF类,添加静态计数器:
public class ArraySizeUdf extends ScalarFunction {
    private static final Logger LOG = LoggerFactory.getLogger(ArraySizeUdf.class);
    private final static ObjectMapper mapper = new ObjectMapper();
    public static final Counter exceptionCounter = new SimpleCounter(); // 静态全局计数器

    public int eval(String stringJsonArray) {
        try {
            if (stringJsonArray == null) {
                return 0;
            }
            JsonNode actualObj = mapper.readTree(stringJsonArray);
            ArrayNode aa = (ArrayNode) actualObj;
            return aa.size();
        } catch (Exception e) {
            LOG.error("Error deserializing json to find size : {}", stringJsonArray);
            exceptionCounter.inc();
            return -1;
        }
    }
}
  1. 在作业初始化阶段注册该计数器:
// 获取StreamExecutionEnvironment
StreamExecutionEnvironment env = streamTableEnvironment.getExecutionEnvironment();
// 将静态计数器注册到Flink Metric组
env.getMetricGroup().addGroup("udf_metrics")
   .counter("json_parse_exception_count", ArraySizeUdf.exceptionCounter);

// 正常注册UDF
streamTableEnvironment.createTemporaryFunction("array_size", ArraySizeUdf.class);

这种方式可绕过UDF内部无法访问RuntimeContext的限制,将计数纳入Flink原生监控体系。

方法3:日志聚合统计(临时场景)

若无需接入Flink原生Metric系统,可在异常捕获时输出含标记的日志,后续通过ELK、Loki等日志采集工具聚合统计异常次数。该方式实时性和精准度不及前两种,适合临时排查使用。

内容的提问来源于stack exchange,提问作者hitesh

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.14 09:22:17