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

Flink 1.16.0中dense_rank()报非时间字段排序不支持问题排查

问题描述

使用Flink 1.16.0版本的Flink SQL调用dense_rank()函数,确认ts为时间属性,但仍报错"Sort on a non time attribute field is not supported"。已设置配置项table.exec.non-temporal-sort.enabled为true,但该配置未生效。需求是按uid分区、ts排序获取排名。

代码实现

public class dws_user_profile_count_type {
    public static void main(String[] args) throws Exception {
        //TODO 1.获取执行环境
        Configuration conf = new Configuration();
        conf.setBoolean("table.exec.non-temporal-sort.enabled", true);
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(conf);
        env.setParallelism(1);
        StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
        // TODO 2. 状态后端设置
//        env.enableCheckpointing(3000L, CheckpointingMode.EXACTLY_ONCE);
//        env.getCheckpointConfig().setCheckpointTimeout(60 * 1000L);
//        env.getCheckpointConfig().setMinPauseBetweenCheckpoints(3000L);
//        env.getCheckpointConfig().enableExternalizedCheckpoints(
//                CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION
//        );
//        env.setRestartStrategy(RestartStrategies.failureRateRestart(
//                3, Time.days(1), Time.minutes(1)
//        ));
//        env.setStateBackend(new HashMapStateBackend());
//        env.getCheckpointConfig().setCheckpointStorage(
//                "hdfs://master/ck"
//        );
//        System.setProperty("HADOOP_USER_NAME", "root");
        //TODO 3.临时表a
        //1.使用flink连接器消费dwd_traffic_page_log,建表
        tableEnv.executeSql("CREATE TABLE dwd_traffic_page_log(" +
                "  `common` MAP<STRING,STRING>,  " +
                "  `page` MAP<STRING,STRING>,  " +
                "  `ts_str` STRING " +
//                "   `pt` AS PROCTIME()  " +
                ")" +
                MyKafkaUtil.getKafkaDDL("dwd_traffic_page_log", "active"));
        //2.对查询需要的字段做成临时视图traffic_page
        Table result = tableEnv.sqlQuery("select " +
                "`common`['uid'] uid, " +
                "`page`['during_time'] during_time, " +
                "`page`['last_page_id'] last_page_id, " +
                "TO_DATE(FROM_UNIXTIME(CAST(`ts_str` AS BIGINT) / 1000, 'yyyy-MM-dd')) `happen_date`, " +
                "TO_TIMESTAMP(FROM_UNIXTIME(CAST(`ts_str` AS BIGINT)/1000, 'yyyy-MM-dd HH:mm:ss')) ts " +
//                " `pt` " +
                "from dwd_traffic_page_log");
        tableEnv.createTemporaryView("traffic_page", result);
//        tableEnv.sqlQuery("select TYPEOF(ts) from traffic_page").execute().print();
        //3.在traffic_page1中筛选近30天的数据          创建临时表tmp_traffic_page
        //创建临时表
        tableEnv.executeSql("CREATE TEMPORARY TABLE tmp_traffic_page(   " +
                "   uid STRING,   " +
                "   during_time BIGINT,   " +
                "   last_page_id STRING,   " +
                "   happen_date DATE,   " +
                "   `ts` TIMESTAMP(3), " +
                "   WATERMARK FOR `ts` AS `ts` - INTERVAL '5' SECOND)" +
                MyKafkaUtil.getKafkaDDL("tmp_traffic_page","tmp"));
        tableEnv.sqlQuery("SELECT   " +
                " uid,  " +
                " CAST(during_time AS BIGINT) during_time,  " +
                " last_page_id,  " +
                " happen_date,  " +
                " `ts`  " +
                "FROM traffic_page   " +
                "WHERE `happen_date` >= CURRENT_DATE - INTERVAL '30' DAY " +
                "ORDER BY uid, `ts` "
        ).executeInsert("tmp_traffic_page");
//
        //4.对临时表tmp_traffic_page进行DENSE_RANK()开窗,用户id分区、ts排序
        tableEnv.executeSql("CREATE TEMPORARY TABLE tmp_traffic_page2(  " +
                "   uid STRING,   " +
                "   during_time BIGINT,   " +
                "   last_page_id STRING,   " +
                "   happen_date DATE,  " +
                "   `ts` TIMESTAMP(3), " +
                "   `rank` BIGINT " +
                "   )  " +
                MyKafkaUtil.getKafkaDDL("tmp_traffic_page2","tmp2"));
        tableEnv.sqlQuery("SELECT " +
                "uid, " +
                "during_time, " +
                "last_page_id, " +
                "happen_date, " +
                "`ts`, " +
                "DENSE_RANK() OVER (PARTITION BY uid ORDER BY `ts` ASC) AS `rank` " +
                "FROM tmp_traffic_page " +
                "ORDER BY happen , `ts` ASC"
        ).executeInsert("tmp_traffic_page2");
    }
}

错误信息

Exception in thread "main" org.apache.flink.table.api.TableException: **Sort on a non-time-attribute field is not supported.**
    at org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecSort.translateToPlanInternal(StreamExecSort.java:75)
    at org.apache.flink.table.planner.plan.nodes.exec.ExecNodeBase.translateToPlan(ExecNodeBase.java:158)
    at org.apache.flink.table.planner.plan.nodes.exec.ExecEdge.translateToPlan(ExecEdge.java:257)
    at org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecSink.translateToPlanInternal(StreamExecSink.java:145)
    at org.apache.flink.table.planner.plan.nodes.exec.ExecNodeBase.translateToPlan(ExecNodeBase.java:158)
    at org.apache.flink.table.planner.delegation.StreamPlanner.$anonfun$translateToPlan$1(StreamPlanner.scala:85)
    at scala.collection.TraversableLike.$anonfun$map$1(TraversableLike.scala:233)
    at scala.collection.Iterator.foreach(Iterator.scala:937)
    at scala.collection.Iterator.foreach$(Iterator.scala:937)
    at scala.collection.AbstractIterator.foreach(Iterator.scala:1425)
    at scala.collection.IterableLike.foreach(IterableLike.scala:70)
    at scala.collection.IterableLike.foreach$(IterableLike.scala:69)
    at scala.collection.AbstractIterable.foreach(Iterable.scala:54)
    at scala.collection.TraversableLike.map(TraversableLike.scala:233)
    at scala.collection.TraversableLike.map$(TraversableLike.scala:226)
    at scala.collection.AbstractTraversable.map(Traversable.scala:104)
    at org.apache.flink.table.planner.delegation.StreamPlanner.translateToPlan(StreamPlanner.scala:84)
    at org.apache.flink.table.planner.delegation.PlannerBase.translate(PlannerBase.scala:197)
    at org.apache.flink.table.api.internal.TableEnvironmentImpl.translate(TableEnvironmentImpl.java:1733)
    at org.apache.flink.table.api.internal.TableEnvironmentImpl.executeInternal(TableEnvironmentImpl.java:825)
    at org.apache.flink.table.api.internal.TableEnvironmentImpl.executeInternal(TableEnvironmentImpl.java:918)
    at org.apache.flink.table.api.internal.TablePipelineImpl.execute(TablePipelineImpl.java:56)
    at org.apache.flink.table.api.Table.executeInsert(Table.java:1064)
    at com.tancong.app.dws.dws_user_profile_count_type.main(dws_user_profile_count_type.java:73)

Process finished with exit code 1

问题排查与修复方案

1. 全局ORDER BY是报错直接原因

最后一段SQL中的ORDER BY happen , ts ASC存在两个问题:

  • 字段名错误:happen应为happen_date,属于笔误
  • 操作多余:你的需求是按uid分区、ts排序获取排名,开窗函数DENSE_RANK() OVER (PARTITION BY uid ORDER BY ts ASC)已经完成了分区内的排序排名,全局ORDER BY在流处理场景中属于不必要操作,且会触发全局排序检查。

2. 配置项设置方式错误

table.exec.non-temporal-sort.enabled是TableEnvironment专属配置项,你通过StreamExecutionEnvironment的Configuration设置不会生效,正确设置方式为:

StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
// 正确设置非时间排序允许配置
tableEnv.getConfig().set("table.exec.non-temporal-sort.enabled", "true");

3. 时间属性的有效性说明

虽然tmp_traffic_page中给ts定义了Watermark,但happen_date是DATE类型,不属于带Watermark的时间属性,无法作为流全局排序的合法时间依据。

修复后的核心代码片段

去掉多余的全局ORDER BY,修正配置设置:

public class dws_user_profile_count_type {
    public static void main(String[] args) throws Exception {
        //TODO 1.获取执行环境
        Configuration conf = new Configuration();
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(conf);
        env.setParallelism(1);
        StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
        // 修正配置项设置方式
        tableEnv.getConfig().set("table.exec.non-temporal-sort.enabled", "true");

        // ... 其他原有代码保持不变 ...

        // 4.对临时表tmp_traffic_page进行DENSE_RANK()开窗,用户id分区、ts排序
        tableEnv.executeSql("CREATE TEMPORARY TABLE tmp_traffic_page2(  " +
                "   uid STRING,   " +
                "   during_time BIGINT,   " +
                "   last_page_id STRING,   " +
                "   happen_date DATE,  " +
                "   `ts` TIMESTAMP(3), " +
                "   `rank` BIGINT " +
                "   )  " +
                MyKafkaUtil.getKafkaDDL("tmp_traffic_page2","tmp2"));
        // 移除多余的全局ORDER BY
        tableEnv.sqlQuery("SELECT " +
                "uid, " +
                "during_time, " +
                "last_page_id, " +
                "happen_date, " +
                "`ts`, " +
                "DENSE_RANK() OVER (PARTITION BY uid ORDER BY `ts` ASC) AS `rank` " +
                "FROM tmp_traffic_page "
        ).executeInsert("tmp_traffic_page2");
    }
}

内容的提问来源于stack exchange,提问作者cong tan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 21:52:00