Flink SQL处理时间时态左连接实现报错及解决方案咨询
问题场景与报错
需要将Kafka数据流与Hadoop上Parquet格式的静态数据做关联增强,最终写入文件系统。尝试两种时态连接方式均报错:
尝试1:Processing-Time时态左连接
SQL语句:
SELECT t1.*,t2.enrichment_data_col from source_stream_table AS t1 LEFT JOIN lookupTable FOR SYSTEM_TIME AS OF t1.proctime AS t2 ON t1.lookup_type = t2.lookup_type
报错信息:
org.apache.flink.table.api.TableException: Processing-time temporal join is not supported yet. at org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecTemporalJoin.createJoinOperator(StreamExecTemporalJoin.java:292) at org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecTemporalJoin.getJoinOperator(StreamExecTemporalJoin.java:254) at org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecTemporalJoin.translateToPlanInternal(StreamExecTemporalJoin.java:179)
尝试2:Temporal Table Function左连接
代码实现:
TableDescriptor lookupDescriptor = TableDescriptor.forConnector("filesystem").format(FormatDescriptor.forFormat("parquet").build()) .option("path",lookupFileLocation) .schema(Schema.newBuilder() .column("lookup_type", DataTypes.STRING().notNull()) .column("enrichment_data_col",DataTypes.INT()) .columnByExpression("proc_time","PROCTIME()") .primaryKey("lookup_type") .build()) .build(); tenv.createTable("lookupTable",lookupDescriptor); TemporalTableFunction tmpLookup = tenv.from("lookupTable").createTemporalTableFunction($("proc_time"),$("lookup_type")); tenv.createTemporarySystemFunction("lookupTableFunc",tmpLookup);
SQL语句:
SELECT t1.*,t2.enrichment_data_col from source_stream_table AS t1 LEFT OUTER JOIN LATERAL TABLE(lookupTableFunc(t1.proctime)) AS t2 ON TRUE
报错信息:
org.apache.flink.table.api.ValidationException: Only single column join key is supported. Found [I@6f8667bb in [Temporal Table Function] at org.apache.flink.table.planner.plan.utils.TemporalJoinUtil$.validateTemporalFunctionPrimaryKey(TemporalJoinUtil.scala:383) at org.apache.flink.table.planner.plan.utils.TemporalJoinUtil$.validateTemporalFunctionCondition(TemporalJoinUtil.scala:365)
注:Temporal Table Function的INNER JOIN可正常运行,但业务需要左连接,使用Flink 1.15.1和1.16.0版本。
解决方案
方案1:改用Lookup Join(推荐,适配静态数据场景)
由于Parquet文件是静态数据,无需时态版本管理,直接使用Flink的Lookup Join(同步查找)即可,这是静态数据关联的最优方案:
- 定义Lookup表时配置缓存参数(可选,提升性能):
TableDescriptor lookupDescriptor = TableDescriptor.forConnector("filesystem") .format(FormatDescriptor.forFormat("parquet").build()) .option("path", lookupFileLocation) .schema(Schema.newBuilder() .column("lookup_type", DataTypes.STRING().notNull()) .column("enrichment_data_col", DataTypes.INT()) .build()) // 开启缓存,避免频繁读取Parquet文件 .option("lookup.cache.max-rows", "10000") .option("lookup.cache.ttl", "1h") .build(); tenv.createTable("lookupTable", lookupDescriptor);
- 使用普通左连接SQL:
SELECT t1.*, t2.enrichment_data_col FROM source_stream_table AS t1 LEFT JOIN lookupTable AS t2 ON t1.lookup_type = t2.lookup_type
方案2:修正Temporal Table Function左连接语法
如果必须使用时态表函数,需修正SQL的关联条件:
- 时态表函数调用时需传入处理时间+关联键两个参数
- ON条件需明确匹配关联键
正确SQL:
SELECT t1.*, t2.enrichment_data_col FROM source_stream_table AS t1 LEFT OUTER JOIN LATERAL TABLE(lookupTableFunc(t1.proctime, t1.lookup_type)) AS t2 ON t1.lookup_type = t2.lookup_type
错误原因说明
- 第一个错误:Flink 1.15/1.16版本不支持基于处理时间的时态左连接,仅支持事件时间的时态左连接,或处理时间的时态内连接。
- 第二个错误:原SQL中
ON TRUE不符合时态表函数的关联规则,必须明确指定关联键匹配,且函数调用需传入关联键参数,否则会被判定为多键关联导致报错。
内容的提问来源于stack exchange,提问作者rony
相关产品推荐
相关产品推荐

