Flink SQL血缘解析如何获取Insert语句的目标表插入字段名?
Flink SQL获取INSERT语句目标写入字段的实现方法
以下方案适配Flink 1.12 ~ 1.18版本,更高版本可对照对应源码微调。
首先说结论:CatalogSinkModifyOperation的公开API本来就没暴露INSERT语句里显式写的目标字段,不是你调用方式错了,是这个类从设计上就没把这个属性开放出来。
要拿到dest_f1、dest_f2这类字段,有两种稳定可用的方案:
方案1:解析阶段直接拿原始SqlNode,性能开销最低
SQL刚解析完还没做校验优化的时候,CatalogSinkModifyOperation内部会存着Calcite原始生成的SqlInsert节点,这个节点上的getTargetColumnList()方法存的就是INSERT子句括号里写的所有字段,反射取出来就行:
import org.apache.calcite.sql.SqlIdentifier; import org.apache.calcite.sql.SqlInsert; import org.apache.flink.table.operations.CatalogSinkModifyOperation; import java.util.stream.Collectors; // 解析输入SQL List<Operation> parsedOps = tableEnv.getParser().parse("insert into target_table(dest_f1, dest_f2) select source_f1, source_f2 from source_table"); CatalogSinkModifyOperation sinkOp = (CatalogSinkModifyOperation) parsedOps.get(0); // 反射拿内部持有的SqlInsert实例 Field insertField = CatalogSinkModifyOperation.class.getDeclaredField("sqlInsert"); insertField.setAccessible(true); SqlInsert rawInsert = (SqlInsert) insertField.get(sinkOp); // 提取显式声明的写入列 List<String> explicitColumns = rawInsert.getTargetColumnList() .getList() .stream() .map(node -> ((SqlIdentifier) node).getSimple()) .collect(Collectors.toList()); // 上述示例SQL运行后explicitColumns的值就是["dest_f1", "dest_f2"]
这个方案的注意点:
- 如果SQL没显式写插入列,比如直接写
insert into target_table select xxx from source_table,getTargetColumnList()会返回空列表,这种场景默认写入列是目标表的全量字段,顺序和SELECT输出的字段顺序一一对应 - 带静态分区的INSERT语句,比如
INSERT INTO t PARTITION(dt='2024-01-01') (col1,col2) SELECT ...,这个方法拿到的列表会自动排除静态分区字段,和SQL里实际写的写入列完全一致
方案2:从校验后的RelNode拿结果,准确性最高
做血缘解析优先用这个方案:把Operation转成Planner校验优化后的RelNode树,从根节点的FlinkLogicalTableModify里拿目标字段列表,不需要反射,而且结果是经过语义校验的,不存在字段写错、字段不存在的问题:
import org.apache.flink.table.planner.plan.nodes.logical.FlinkLogicalTableModify; import org.apache.flink.table.planner.delegation.PlannerBase; import org.apache.calcite.rel.RelNode; import java.util.Collections; PlannerBase planner = (PlannerBase) tableEnv.getPlanner(); RelNode relRoot = planner.translateToRel(Collections.singletonList(sinkOp)).build(); FlinkLogicalTableModify modifyNode = (FlinkLogicalTableModify) relRoot; List<String> targetColumns = modifyNode.getUpdateColumnList();
这个方案不管SQL里有没有显式指定写入列,都会返回完整的、和SELECT输出字段顺序严格一一映射的目标字段列表,直接拿来做字段级血缘映射不会出错。
踩坑提醒
- 别直接调用
CatalogSinkModifyOperation.getTableSchema()拿字段,这个方法返回的是目标表的全量表结构,不是INSERT语句实际要写入的列,字段顺序也不一定和SELECT输出顺序匹配,拿来做血缘很容易映射错 - 反射拿SqlInsert的方案跨版本兼容性很好,从1.12到1.18版本,
CatalogSinkModifyOperation里存SqlInsert的字段名一直是sqlInsert,如果后续大版本升级报字段不存在,直接翻对应版本的这个类源码找SqlInsert类型的成员变量就行,逻辑不会变。
内容的提问来源于stack exchange,提问作者slo
相关产品推荐
相关产品推荐

