Flink 1.16.0中dense_rank()报非时间字段排序不支持问题排查
Flink SQL 1.16.0中dense_rank()排序报错"Sort on a non time attribute field is not supported"排查
问题描述
使用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
相关产品推荐
相关产品推荐

