You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.19 19:35:23