Flink 1.14.2查询Hive 2.1.1报Distinct without an aggregation异常
问题描述
- 使用Flink Hive Table API连接器查询Hive时抛出
Distinct without an aggregation语义异常 - 同一条SQL通过Hue直连Hive执行可正常运行
- 环境版本:
- Flink版本:
1.14.2 - Hive版本:
2.1.1
- Flink版本:
触发异常的SQL语句:
select devid as pdevid, count(distinct vtype) as vip_type_trans from events where dt = '20220702' and utype > -1 group by devid having count(distinct vtype) > 1
核心异常栈:
org.apache.flink.client.program.ProgramInvocationException: The main method caused an error: org.apache.hadoop.hive.ql.parse.SemanticException: Distinct without an aggregation. at org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:372) at org.apache.flink.client.program.PackagedProgram.invokeInteractiveModeForExecution(PackagedProgram.java:222) at org.apache.flink.client.ClientUtils.executeProgram(ClientUtils.java:114) at org.apache.flink.client.cli.CliFrontend.executeProgram(CliFrontend.java:812) at org.apache.flink.client.cli.CliFrontend.run(CliFrontend.java:246) at org.apache.flink.client.cli.CliFrontend.parseAndRun(CliFrontend.java:1054) at org.apache.flink.client.cli.CliFrontend.lambda$main$10(CliFrontend.java:1132) at java.security.AccessController.doPrivileged(Native Method) at javax.security.auth.Subject.doAs(Subject.java:422) at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1698) at org.apache.flink.runtime.security.contexts.HadoopSecurityContext.runSecured(HadoopSecurityContext.java:41) at org.apache.flink.client.cli.CliFrontend.main(CliFrontend.java:1132) Caused by: java.lang.RuntimeException: org.apache.hadoop.hive.ql.parse.SemanticException: Distinct without an aggregation. at org.apache.flink.table.planner.delegation.hive.HiveParserCalcitePlanner.logicalPlan(HiveParserCalcitePlanner.java:304) at org.apache.flink.table.planner.delegation.hive.HiveParserCalcitePlanner.genLogicalPlan(HiveParserCalcitePlanner.java:272) at org.apache.flink.table.planner.delegation.hive.HiveParser.analyzeSql(HiveParser.java:290) at org.apache.flink.table.planner.delegation.hive.HiveParser.processCmd(HiveParser.java:238) at org.apache.flink.table.planner.delegation.hive.HiveParser.parse(HiveParser.java:208) at org.apache.flink.table.api.internal.TableEnvironmentImpl.sqlQuery(TableEnvironmentImpl.java:716) Caused by: org.apache.hadoop.hive.ql.parse.SemanticException: Distinct without an aggregation. at org.apache.flink.table.planner.delegation.hive.HiveParserCalcitePlanner.genSelectLogicalPlan(HiveParserCalcitePlanner.java:2275) at org.apache.flink.table.planner.delegation.hive.HiveParserCalcitePlanner.genLogicalPlan(HiveParserCalcitePlanner.java:2749) at org.apache.flink.table.planner.delegation.hive.HiveParserCalcitePlanner.genLogicalPlan(HiveParserCalcitePlanner.java:2647) at org.apache.flink.table.planner.delegation.hive.HiveParserCalcitePlanner.logicalPlan(HiveParserCalcitePlanner.java:284)
根因说明
这是Flink 1.14版本内嵌Hive解析器的已知兼容性缺陷:当SQL使用Hive方言解析,且HAVING子句中直接书写和SELECT列表完全一致的count(distinct 字段)聚合表达式时,解析器会错误判定distinct关键字未被聚合函数包裹,抛出语义异常。该问题属于Flink侧解析逻辑bug,和Hive本身的语法支持无关。
解决方案
无需升级Flink版本,任选以下一种方式即可规避:
- 改写SQL,通过子查询提前完成聚合计算,外层再做条件过滤,避开解析器的bug触发逻辑,改写后示例:
select pdevid, vip_type_trans from ( select devid as pdevid, count(distinct vtype) as vip_type_trans from events where dt = '20220702' and utype > -1 group by devid ) tmp where vip_type_trans > 1
- 执行SQL前切换Flink SQL方言为默认方言,不使用Hive方言做语法解析,执行参数设置:
SET table.sql-dialect = default;,该方式需要确认目标Hive表的定义兼容Flink默认方言的解析规则。
内容的提问来源于stack exchange,提问作者xcrossed xcrossed
相关产品推荐
相关产品推荐

