Flink Table API内置函数使用求助:ROW_NUMBER()调用失败
在Flink Table API中实现ROW_NUMBER()的正确方式
直接使用Table API内置的row_number()函数
Flink Table API本身提供了对应SQL中ROW_NUMBER()的内置方法,无需通过Expression.callSql(),直接结合窗口(Over Window)即可实现。以下是Java和Scala的示例:
Java示例
import org.apache.flink.table.api.*; import static org.apache.flink.table.api.Expressions.*; // 假设已有输入表 inputTable,包含user_id、ts、value字段 Table resultTable = inputTable .select( $("user_id"), $("ts"), $("value"), // 按user_id分区,ts降序排序生成行号 row_number().over( PartitionBy.partitionBy($("user_id")) .orderBy($("ts").desc()) .rangeBetween(UNBOUNDED_RANGE, CURRENT_RANGE) ).as("row_num") );
Scala示例
import org.apache.flink.table.api._ import org.apache.flink.table.api.Expressions._ val resultTable = inputTable .select( $"user_id", $"ts", $"value", row_number().over( PartitionBy.partitionBy($"user_id") .orderBy($"ts".desc) .rangeBetween(UNBOUNDED_RANGE, CURRENT_RANGE) ).as("row_num") )
若要使用Expression.callSql()的正确写法
如果坚持用callSql()实现,需确保SQL语句完整且字段与表结构匹配,示例如下:
Table resultTable = inputTable .select( $("user_id"), $("ts"), $("value"), callSql("ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY ts DESC)").as("row_num") );
注意事项
- 确保使用的Flink版本在1.12及以上,该版本后Table API的窗口函数API趋于稳定
- 分区(
partitionBy)和排序(orderBy)规则需根据业务需求调整,rangeBetween指定窗口范围,此处的UNBOUNDED_RANGE到CURRENT_RANGE对应SQL中的默认窗口范围
内容的提问来源于stack exchange,提问作者JoeHills
相关产品推荐
相关产品推荐

