使用Temporal Table Function关联流时遇类型不一致错误求解决
问题:Temporal Table Function关联流时NOT NULL属性不匹配错误
使用Flink的Temporal Table Function关联两个流时,遇到集合类型与表达式类型因proctime0的NOT NULL属性不一致导致的断言错误,询问该差异产生的原因及解决方法。
报错信息
Exception in thread "main" java.lang.AssertionError: Cannot add expression of different type to set: set type is RecordType(VARCHAR(2147483647) CHARACTER SET "UTF-16LE" order_id, DECIMAL(32, 2) price, VARCHAR(2147483647) CHARACTER SET "UTF-16LE" currency, TIMESTAMP(3) order_time, TIMESTAMP_LTZ(3) *PROCTIME* NOT NULL proctime, VARCHAR(2147483647) CHARACTER SET "UTF-16LE" currency0, BIGINT conversion_rate, TIMESTAMP(3) update_time, TIMESTAMP_LTZ(3) *PROCTIME* proctime0) NOT NULL expression type is RecordType(VARCHAR(2147483647) CHARACTER SET "UTF-16LE" order_id, DECIMAL(32, 2) price, VARCHAR(2147483647) CHARACTER SET "UTF-16LE" currency, TIMESTAMP(3) order_time, TIMESTAMP_LTZ(3) *PROCTIME* NOT NULL proctime, VARCHAR(2147483647) CHARACTER SET "UTF-16LE" currency0, BIGINT conversion_rate, TIMESTAMP(3) update_time, TIMESTAMP_LTZ(3) *PROCTIME* NOT NULL proctime0) NOT NULL set is rel#61:LogicalCorrelate.NONE.any.None: 0.[NONE].[NONE](left=HepRelVertex#59,right=HepRelVertex#60,correlation=$cor0,joinType=inner,requiredColumns={4}) expression is LogicalJoin(condition=[__TEMPORAL_JOIN_CONDITION($4, $7, __TEMPORAL_JOIN_CONDITION_PRIMARY_KEY($5))], joinType=[inner]) LogicalProject(order_id=[$0], price=[$1], currency=[$2], order_time=[$3], proctime=[PROCTIME()]) LogicalTableScan(table=[[default_catalog, default_database, orders]]) LogicalProject(currency=[$0], conversion_rate=[$1], update_time=[$2], proctime=[PROCTIME()]) LogicalTableScan(table=[[default_catalog, default_database, currency_rates]])
相关表定义与代码
事实表(orders)定义
CREATE TABLE `orders` ( order_id STRING, price DECIMAL(32,2), currency STRING, order_time TIMESTAMP(3), proctime as PROCTIME() ) WITH ( 'properties.bootstrap.servers' = '127.0.0.1:9092', 'properties.group.id' = 'test', 'scan.topic-partition-discovery.interval' = '10000', 'connector' = 'kafka', 'format' = 'json', 'scan.startup.mode' = 'latest-offset', 'topic' = 'test1' )
维表(currency_rates)定义
CREATE TABLE `currency_rates` ( currency STRING, conversion_rate BIGINT, update_time TIMESTAMP(3), proctime as PROCTIME() ) WITH ( 'properties.bootstrap.servers' = '127.0.0.1:9092', 'properties.group.id' = 'test', 'scan.topic-partition-discovery.interval' = '10000', 'connector' = 'kafka', 'format' = 'json', 'scan.startup.mode' = 'latest-offset', 'topic' = 'test3' )
时态表函数创建
TemporalTableFunction table_rate = tEnv.from("currency_rates") .createTemporalTableFunction("update_time", "currency"); tEnv.registerFunction("rates", table_rate);
关联查询逻辑
SELECT order_id, price, s.currency, conversion_rate, order_time FROM orders AS o, LATERAL TABLE (rates(o.proctime)) AS s WHERE o.currency = s.currency
差异产生原因
- 核心问题是Flink优化器在处理时态表函数关联时,对维表生成的
proctime0字段的NOT NULL属性推断出现不一致:在集合类型中标记该字段可为空,而在表达式类型中标记为不可为空,导致类型校验时断言失败。 - 这种不一致源于Flink内部对动态生成的处理时间字段元数据的逻辑冲突——时态表函数关联过程中,不同阶段的元数据处理逻辑没有统一该字段的非空性标记。
解决方法
方法一:显式指定维表proctime字段为NOT NULL
修改维表DDL,给proctime字段加上NOT NULL约束,强制统一非空属性:
CREATE TABLE `currency_rates` ( currency STRING, conversion_rate BIGINT, update_time TIMESTAMP(3), proctime as PROCTIME() NOT NULL ) WITH ( 'properties.bootstrap.servers' = '127.0.0.1:9092', 'properties.group.id' = 'test', 'scan.topic-partition-discovery.interval' = '10000', 'connector' = 'kafka', 'format' = 'json', 'scan.startup.mode' = 'latest-offset', 'topic' = 'test3' )
方法二:删除维表中的proctime字段
时态表函数关联依赖的是事件时间update_time,不需要维表的处理时间字段。直接删除维表DDL中的proctime as PROCTIME()定义,从根源上避免属性不匹配:
CREATE TABLE `currency_rates` ( currency STRING, conversion_rate BIGINT, update_time TIMESTAMP(3) ) WITH ( 'properties.bootstrap.servers' = '127.0.0.1:9092', 'properties.group.id' = 'test', 'scan.topic-partition-discovery.interval' = '10000', 'connector' = 'kafka', 'format' = 'json', 'scan.startup.mode' = 'latest-offset', 'topic' = 'test3' )
方法三:改用SQL语法注册时态表函数
避免使用Java API创建时态表函数,改用SQL语法注册,减少元数据推断问题:
CREATE TEMPORAL TABLE FUNCTION rates AS SELECT * FROM currency_rates FOR SYSTEM_TIME AS OF proctime WITH ( 'time_attribute' = 'update_time', 'primary_key' = 'currency' );
之后直接在关联查询中使用rates函数即可。
内容的提问来源于stack exchange,提问作者Tally
相关产品推荐
相关产品推荐

