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

如何在Flink SQL中使用ROW类型列字段?解决DDL报错

问题重现

执行以下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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 21:20:25