Flink1.12执行Temporal Table Join报join key不能为空问题咨询
问题触发原因
该报错是Flink 1.12.x版本的已知问题,核心原因有两类:
- 1.12版本的Temporal Table Join关联键解析逻辑存在缺陷:
LogicalCorrelateToJoinFromGeneralTemporalTableRule规则在提取ON子句中的关联键时,容易对多条件的ON子句出现识别错误,哪怕已经写了等值关联条件,也会误判为Join Key为空,触发验证异常。 - Temporal Join的右表(即示例中的Customers表)未显式声明主键:Flink 1.12要求Temporal Join的右表必须在DDL中定义主键,否则解析器无法对应关联键,也会抛出该报错。
解决办法
可按优先级选择以下方案处理:
优先升级Flink版本
该关联键解析的bug已经在Flink 1.13.0及之后的稳定版本完全修复,升级到1.13+版本即可直接解决问题,不需要修改原有SQL逻辑。修正右表DDL配置
如果暂时无法升级版本,首先检查Customers表的CREATE TABLE语句,是否显式声明了主键,示例配置如下:
CREATE TABLE Customers ( id INT, country STRING, zip STRING, -- 必须显式声明主键,主键字段要和关联键对应 PRIMARY KEY (id) NOT ENFORCED ) WITH ( -- 实际连接器配置 );
- 调整SQL写法规避bug
如果已经配置了主键还是报错,可以将ON子句中的非等值过滤条件(比如is not null判断)移到WHERE子句中,规避1.12版本解析器的多条件识别缺陷,调整后的第一条SQL示例:
SELECT o.order_id, o.total, c.country, c.zip FROM Orders AS o JOIN Customers FOR SYSTEM_TIME AS OF o.proc_time AS c ON o.customer_id = c.id WHERE o.customer_id IS NOT NULL AND c.id IS NOT NULL;
内容的提问来源于stack exchange,提问作者xiaojin
相关产品推荐
相关产品推荐

