如何在Flink中将int列聚合为array<int>?继承AggregateFunction的UDF输出类型无法修改的解决建议咨询
解决Flink中聚合int列到array的问题
我之前也碰到过类似的情况,给你几个可行的方案:
直接使用Flink内置的ARRAY_AGG函数(推荐)
从Flink 1.13版本开始,官方已经提供了内置的ARRAY_AGG聚合函数,专门用来将分组后的列聚合为数组类型。你完全可以不用自己写UDF,直接替换成下面的SQL:SELECT a_tag, ARRAY_AGG(type) AS type_arr FROM a GROUP BY a_tag这个函数默认会将int类型聚合为
ARRAY<INT>,完全符合你的需求,而且不用维护自定义UDF,省心很多。调整自定义AggregateFunction的输出类型
如果你必须使用自己写的NEST()UDF,那问题大概率出在AggregateFunction的类型定义上。你需要确保在UDF类中正确指定输出类型为数组类型,而不是默认的其他类型。比如你的UDF类应该这样写:public class NestIntAgg extends AggregateFunction<Integer[], List<Integer>> { // 初始化累加器 @Override public List<Integer> createAccumulator() { return new ArrayList<>(); } // 累加逻辑 public void accumulate(List<Integer> accumulator, Integer value) { accumulator.add(value); } // 指定输出类型 @Override public TypeInformation<Integer[]> getResultType() { return ArrayType.of(IntType.INSTANCE); } // 将累加器转换为最终结果 @Override public Integer[] getValue(List<Integer> accumulator) { return accumulator.toArray(new Integer[0]); } }重点是
getResultType()方法要返回ArrayType.of(IntType.INSTANCE),这样Flink就能识别输出是ARRAY<INT>类型了。查询时添加类型转换(临时 workaround)
如果暂时没法修改UDF的代码,你可以在SQL查询中通过CAST函数强制转换输出类型:SELECT a_tag, CAST(NEST(type) AS ARRAY<INT>) AS type_arr FROM a GROUP BY a_tag不过这个方法只适用于UDF实际输出的类型和
ARRAY<INT>可以兼容的情况(比如UDF返回的是List<Integer>),如果类型不兼容可能会报错,所以还是优先前两个方案。
内容的提问来源于stack exchange,提问作者Kenyore
相关产品推荐
相关产品推荐

