如何在Flink SQL窗口TVF中定义表提示?Iceberg流读报错
解决Flink SQL窗口TVF结合Iceberg流读的语法错误
你的语法错误根源在于动态表选项的位置错误:Iceberg的流读参数(streaming、monitor-interval)需要附加在源表的引用上,而不是窗口TVF的结果表上。
正确语法示例
-- create connect table CREATE TABLE default_catalog.default_database.`test_p` ( id BIGINT, t BIGINT, dt STRING, ts AS CAST(TO_TIMESTAMP_LTZ(t, 3) AS TIMESTAMP(3)), WATERMARK FOR ts AS ts - INTERVAL '5' SECOND ) PARTITIONED BY (dt) WITH ( 'connector'='iceberg', ... ); SET execution.runtime-mode = streaming; SET table.dynamic-table-options.enabled = true; -- 正确写法:将OPTIONS提示直接附加在源表test_p的引用后面 SELECT * FROM TABLE ( TUMBLE( TABLE test_p /*+ OPTIONS('streaming'='true', 'monitor-interval'='15s')*/, DESCRIPTOR(ts), INTERVAL '5' MINUTES ) );
关键说明
- 开启
table.dynamic-table-options.enabled = true是使用动态表选项的必要前提,你这部分配置是正确的 - Iceberg的流读配置属于表读取专属参数,必须绑定到实际的Iceberg源表引用上;窗口TVF是基于源表的转换操作,它的结果表不支持这类Iceberg特定参数
- 动态表选项提示
/*+ OPTIONS(...) */需要紧跟在TABLE test_p之后,作为源表引用的一部分传入窗口TVF函数
内容的提问来源于stack exchange,提问作者toien
相关产品推荐
相关产品推荐

