如何在Flink SQL中将ROW类型的UDF输出拆分为多列?
解决Flink SQL中拆分UDF返回ROW类型为多列的问题
正确的Flink SQL写法
要拆分RankDif UDF返回的ROW类型为多列,你有两种可行的写法,推荐第一种(性能更优,仅调用一次UDF):
方式1:子查询+别名访问字段
SELECT dvid, rank_name, rank_type, window_start, window_end, rank_result.`order` AS rank_order, rank_result.diff AS rank_diff, rank_result.pt AS rank_pt FROM ( SELECT dvid, rank_name, rank_type, window_start, window_end, RankDif(rank_order, rank_pt) AS rank_result FROM TABLE( HOP(TABLE UniqueRankTable, DESCRIPTOR(rank_pt), INTERVAL '1' DAY, INTERVAL '2' DAY) ) GROUP BY dvid, rank_name, rank_type, window_start, window_end )
方式2:直接在SELECT中提取ROW字段(不推荐,会多次调用UDF)
SELECT dvid, rank_name, rank_type, window_start, window_end, RankDif(rank_order, rank_pt).`order` AS rank_order, RankDif(rank_order, rank_pt).diff AS rank_diff, RankDif(rank_order, rank_pt).pt AS rank_pt FROM TABLE( HOP(TABLE UniqueRankTable, DESCRIPTOR(rank_pt), INTERVAL '1' DAY, INTERVAL '2' DAY) ) GROUP BY dvid, rank_name, rank_type, window_start, window_end
为什么你的原写法失败?
你尝试的RankDif(rank_order,rank_pt) AS (rank_order, rank_diff, rank_pt)语法,Flink SQL并不支持。这种“将ROW别名直接拆分为多列”的写法在部分传统数据库中有效,但Flink SQL的语法规则里没有这个特性,因此会触发语法验证失败。
对应Table API写法的逻辑
你在Table API中使用的:
table.select(call("RankDif", $("rank_order"), $("rank_pt")).as("rank_order, rank_diff, rank_pt"))
本质是将UDF返回的ROW类型的每个字段,依次映射到指定的列名。对应的SQL逻辑就是方式1中,先将UDF结果命名为rank_result,再逐个提取其内部字段并赋予别名的操作——两者是完全等价的。
内容的提问来源于stack exchange,提问作者Singleton
相关产品推荐
相关产品推荐

