You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.24 05:22:05