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系统:
- 修改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; } } }
- 在作业初始化阶段注册该计数器:
// 获取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
相关产品推荐
相关产品推荐

