如何在Flink SQL中使用ROW类型列字段?解决DDL报错
解决Flink SQL中嵌套ROW类型字段作为主键的语法错误
问题重现
执行以下Flink SQL时触发解析错误:
create table team_config_source ( `payload` ROW( `before` ROW( team_config_id int, ... ), `after` ROW( team_config_id int, ... ) ), PRIMARY KEY (`payload`.`after`.`team_config_id`) NOT ENFORCED ) WITH ( 'connector' = 'kafka', 'topic' = 'xxx', 'properties.bootstrap.servers' = 'xxx', 'properties.group.id' = 'xxx', 'scan.startup.mode' = 'earliest-offset', 'format' = 'json', 'key.format' = 'json' )
抛出错误:
org.apache.flink.table.api.SqlParserException: SQL parse failed. Encountered "." at line 51, column 29. Was expecting one of: ")" ... "," ...
尝试去掉反引号直接写payload.after.team_config_id时,又提示column payload.after.team_config_id was not defined。
解决方案
Flink SQL中,当主键是嵌套ROW类型的字段时,需要将整个嵌套字段访问表达式用双重括号包裹,让解析器将其识别为一个完整的表达式而非列名。修正后的DDL如下:
create table team_config_source ( `payload` ROW( `before` ROW( team_config_id int, ... ), `after` ROW( team_config_id int, ... ) ), -- 用双重括号包裹嵌套字段表达式 PRIMARY KEY ((`payload`.`after`.`team_config_id`)) NOT ENFORCED ) WITH ( 'connector' = 'kafka', 'topic' = 'xxx', 'properties.bootstrap.servers' = 'xxx', 'properties.group.id' = 'xxx', 'scan.startup.mode' = 'earliest-offset', 'format' = 'json', 'key.format' = 'json' )
原理说明
Flink SQL的主键定义默认只接受直接的列名,而嵌套ROW的字段访问(如payload.after.team_config_id)属于表达式,需要用双重括号明确告知解析器:这是一个需要计算的表达式,而非单一列名。
内容的提问来源于stack exchange,提问作者Rinze
相关产品推荐
相关产品推荐

