Flink Table API是否支持row_number()等窗口函数等价实现?
Flink Table API 排序窗口函数实现方案
1. 是否存在等价实现?
存在。Flink 1.15.1版本的Table API完全支持row_number()、rank()、dense_rank()这类排序窗口函数,只是官方文档示例多聚焦于avg、min这类聚合窗口函数,容易让开发者忽略这类排序函数的Table API写法。
2. 具体写法
你只需要调用org.apache.flink.table.api.Expressions工具类中的对应方法,替换代码中的/* ??? */即可。针对你的需求,完整的Table API代码如下:
import org.apache.flink.table.api.Expressions; // 其他必要导入... final Table someTable = tableEnvironment.from("SOME_TABLE") .window(Over.partitionBy($("COLUMN_A")) .orderBy($("EVENT_TIME").desc()) .as("window")) .select( $("COLUMN_A"), $("COLUMN_B"), $("COLUMN_C"), Expressions.row_number().over($("window")).as("ROW_NUM") ) .where($("ROW_NUM").isEqual(1));
补充说明:
- 如果需要使用
rank()或dense_rank(),写法和row_number()完全一致,只需替换为Expressions.rank()或Expressions.dense_rank()即可。 - 可以通过静态导入简化写法:
import static org.apache.flink.table.api.Expressions.*;,这样直接写row_number()就能调用对应方法。
内容的提问来源于stack exchange,提问作者Bart Gerard
相关产品推荐
相关产品推荐

