☰
SeaTunnel Sink Options Placeholders 使用指南:通过占位符动态获取上游表元数据
2026/9/27 7:27:29 网站建设 项目流程
  • 数据工程
  • 大数据
  • 批处理
  • 流处理

【免费下载链接】seatunnel

SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.

项目地址:https://gitcode.com/gh_mirrors/sea/seatunnel
点击查看免费下载

导读

SeaTunnel 的 Sink Options Placeholders(Sink 参数占位符)功能允许在 Sink 配置中使用${database_name}、${table_name}、${primary_key}等占位符,在连接器启动前由框架自动替换为上游 Catalog Table 的真实元数据。本文将以 sink-options-placeholders.md 为核心,完整讲解全部占位符的语义、默认值语法、配置前提与实战示例,并结合仓库中 TablePlaceholder.java 的源码实现和单元测试,深入解析替换机制的底层原理、多表写入场景下的工作方式,以及占位符未被替换时的排查思路。

一、功能概述与典型应用场景

在多表同步场景下,同一个 Sink 配置往往需要被多个上游表复用。例如使用 MySQL-CDC 或 Oracle-CDC 捕获多张表的变化,再写入下游 JDBC 目标库时,目标库名、目标表名、主键字段都随上游表而变化。如果每张表都手写一份 Sink 配置,既繁琐又难以维护。

Sink Options Placeholders 正是为解决这一问题而设计:SeaTunnel 允许在 Sink 配置中写入占位符表达式,框架会在连接器启动之前完成替换,将上游 Catalog Table 的元数据注入 Sink 配置,从而以一份配置驱动多表写入。

该功能在以下三类引擎上均受支持:

  • SeaTunnel Zeta(SeaTunnel 自研引擎)
  • Flink(通过 seatunnel-flink-starter 运行)
  • Spark(通过 seatunnel-spark-starter 运行)

二、占位符清单与语义

SeaTunnel 共提供 8 种占位符,主要分为“表标识(Table Identifier)”与“表结构信息”两类。表标识类占位符从上游 Catalog Table 的TableIdentifier中取值,表结构信息类占位符从上游表的TableSchema中取值。

占位符表达式说明取值来源
${database_name}上游 Catalog Table 的数据库名(database)TableIdentifier
${schema_name}上游 Catalog Table 的 schema 名TableIdentifier
${table_name}上游 Catalog Table 的表名(table)TableIdentifier
${schema_full_name}数据库 + schema 的完整路径,以.连接TableIdentifier
${table_full_name}数据库 + schema + 表名的完整路径,以.连接TableIdentifier
${primary_key}上游表的主键字段列表,多个字段以,分隔TableSchema 的 PrimaryKey
${unique_key}上游表的唯一键字段列表,多个字段以,分隔TableSchema 的 ConstraintKeys
${field_names}上游表的全部字段名列表,多个字段以,分隔TableSchema 的 FieldNames

以上占位符常量定义在 TablePlaceholder.java 中,其中NAME_DELIMITER(.)用于拼接完整路径,FIELD_DELIMITER(,)用于拼接多字段列表。

默认值语法

对于${database_name}、${schema_name}、${table_name}这类表标识占位符,可以通过表达式指定默认值。当上游元数据缺失对应字段时,将回退使用默认值:

${database_name:default_my_db} # 数据库缺失时使用 default_my_db ${schema_name:default_my_schema} # schema 缺失时使用 default_my_schema ${table_name:default_my_table} # 表名缺失时使用 default_my_table

默认值语法也适用于${schema_full_name}、${table_full_name}、${primary_key}、${unique_key}、${field_names}——从源码看,替换逻辑replacePlaceholders(String input, String placeholderName, String value, String defaultValue)对所有占位符统一处理形如\$\{name(:[^}]*)?\}的模式,冒号后即为默认值(TablePlaceholder.java)。

注意:原文档与源码中默认值写法的空格不影响解析——正则中默认值部分会经过.trim()处理,${database_name: default_db}与${database_name:default_db}等价。

三、使用前提:Sink 连接器需实现 TableSinkFactory API

占位符替换由 SeaTunnel 框架在创建 Sink 时统一完成,前提是所使用的 Sink 连接器实现了TableSinkFactoryAPI。这是所有基于新版 API 的连接器的通用要求:

  • 接口定义位于 TableSinkFactory.java,其中还提供了excludeTablePlaceholderReplaceKeys()默认方法,用于声明不需要做占位符替换的配置键;
  • 当前仓库中,JDBC、Doris、StarRocks、Iceberg、ClickHouse 等绝大多数 V2 连接器的XxxSinkFactory均实现了该接口(可参考 DorisSinkFactory.java 等实现)。

若使用的连接器仍为旧版 API(未实现TableSinkFactory),则不会触发占位符替换流程。

四、配置示例:MySQL-CDC 多表写入 JDBC

以下两个示例来自官方文档,可直接作为实战模板。它们演示了通过 CDC 捕获上游表元数据后,将占位符嵌入 JDBC Sink 的目标库名、目标表名与主键配置。

示例 1:使用${database_name}映射目标库

env { // ignore... } source { MySQL-CDC { // ignore... } } transform { // ignore... } sink { jdbc { url = "jdbc:mysql://localhost:3306" driver = "com.mysql.cj.jdbc.Driver" user = "root" password = "123456" database = "${database_name}_test" table = "${table_name}_test" primary_keys = ["${primary_key}"] } }

当上游表为shop.orders时,替换后 JDBC Sink 实际写入的目标为:database = shop_test、table = orders_test、primary_keys = [<orders 表主键字段>]。

示例 2:使用${schema_name}映射目标库

env { // ignore... } source { Oracle-CDC { // ignore... } } transform { // ignore... } sink { jdbc { url = "jdbc:mysql://localhost:3306" driver = "com.mysql.cj.jdbc.Driver" user = "root" password = "123456" database = "${schema_name}_test" table = "${table_name}_test" primary_keys = ["${primary_key}"] } }

对于 Oracle 这类以 schema 作为命名空间组织表结构的数据库,使用${schema_name}能正确获取上游 schema 并映射到目标库名。

占位符可以嵌入任意字符串中(如"${database_name}_test"),也可以作为完整值使用(如"${primary_key}")。当作为 List 类型配置(如primary_keys)的完整值时,替换结果会按,拆分成多个元素,这正是 TablePlaceholder.java 中对 List 类型值的特殊处理。

五、源码级原理:替换是如何在连接器启动前完成的

官方文档明确指出:“We will complete the placeholder replacement before the connector is started, ensuring that the sink options is ready before use.”下面结合源码还原这条调用链。

5.1 调用链全景

  1. 入口:Sink 创建统一走 FactoryUtil.createAndPrepareSink。它发现TableSinkFactory后,调用TableSinkFactoryContext.replacePlaceholderAndCreate(...)创建上下文;
  2. 上下文构造:TableSinkFactoryContext.replacePlaceholderAndCreate 内部调用TablePlaceholder.replaceTablePlaceholder(options, catalogTable, excludeTablePlaceholderReplaceKeys),得到替换后的ReadonlyConfig;
  3. 校验与创建:替换后的配置经ConfigValidator.validate(factory.optionRule())校验通过后,才交给factory.createSink(context)真正创建连接器实例。

因此,连接器拿到的永远是替换完成的最终配置,占位符不会泄漏到下游连接器的运行时逻辑中。

5.2 替换逻辑的核心实现

TablePlaceholder是这一功能的唯一实现类,其替换流程分四个步骤(TablePlaceholder.java):

  1. 遍历配置:克隆ReadonlyConfig的源 Map(copy-on-write),逐键处理,命中excludeKeys(即连接器通过excludeTablePlaceholderReplaceKeys()声明排除的键)则跳过;
  2. 表标识替换:replaceTableIdentifier依次替换${database_name}、${schema_name}、${table_name},并按“database + schema”(schema_full_name)、“database + schema + table”(table_full_name)的拼接顺序替换完整路径占位符;拼接时跳过为 null 的层级;
  3. 结构信息替换:replaceTablePrimaryKey、replaceTableUniqueKey、replaceTableFieldNames分别从TableSchema中提取主键列、唯一键列、全部字段名,以,连接后替换对应占位符;
  4. 特殊处理 List 值:若某配置键的值是长度为 1 的字符串列表且恰好等于${primary_key}/${unique_key}/${field_names},则替换后按,拆分为真正的 List 返回,保证primary_keys这类数组型配置的语义正确。

从测试 TablePlaceholderTest.java 可以看到框架对该行为的完整验证:字符串型与数组型配置均被正确替换,例如"xyz_${database_name: default_db}_test"在缺失 database 时被替换为"xyz_default_db_test","${primary_key}"被替换为主键列表["f1", "f2"]。

5.3 多表场景:一份配置,逐表替换

多表写入时,框架会为每个上游CatalogTable独立执行一次替换。测试用例testSinkOptionsWithMultiTable(TablePlaceholderTest.java)验证了同一份配置在table1(含完整 database/schema/table 路径)与table2(路径全空)下分别被替换为不同的结果,前者得到my-database/my-schema/my-table,后者回退到默认值default_db/default_schema/default_table。这说明占位符机制天然支持多表动态路由。

5.4 排除指定键:excludeTablePlaceholderReplaceKeys

某些连接器的配置项可能本身包含类似${...}的字符串且不希望被替换,此时可覆写TableSinkFactory.excludeTablePlaceholderReplaceKeys()返回需要排除的键列表。replaceTablePlaceholder的excludeKeys参数即为此设计,测试用例testSinkOptionsWithExcludeKeys(TablePlaceholderTest.java)验证了排除database键后该键保留原始占位符文本。

六、占位符未被替换时的原因排查

如果 Sink 参数中仍残留${...}文本,通常意味着上游表元数据中缺少该字段。官方文档给出的典型场景包括:

  • MySQL 等数据源不包含${schema_name}:MySQL 的 Catalog Table 通常只包含 database 与 table 两级,schema 层级为 null,因此${schema_name}无法被替换(此时应改用${database_name});
  • Oracle 等数据源不包含${database_name}:Oracle 以 schema 组织表结构,database 层级可能为 null,因此${database_name}无法被替换(此时应改用${schema_name});
  • 同理,若上游表未声明主键/唯一键/约束信息,${primary_key}、${unique_key}也会保持原样。

两类应对方案:

  1. 使用默认值语法兜底:如${database_name:default_my_db},元数据缺失时自动回退;
  2. 按数据源特性选择占位符:先确认上游 Catalog Table 实际包含哪些层级,再选择对应的占位符表达式。

七、最佳实践小结

  • 多表写入优先使用占位符:CDC 同步多张表到 JDBC/OLAP 目标时,用${database_name}、${table_name}动态映射目标库表,避免逐表编写 Sink 配置;
  • 完整路径占位符用于单层命名空间:当上游只含 database 或只含 schema 单层路径时,${schema_full_name}、${table_full_name}会按实际存在的层级拼接(缺省层级自动跳过),比固定写两层更安全;
  • 结构类占位符用于表结构敏感的目标端:${primary_key}、${unique_key}可直接作为 JDBC Sink 的primary_keys、Doris/StarRocks 建表模型等配置,实现表结构自动对齐;
  • 关键配置键声明排除:若连接器配置中存在不应被替换的${...}文本,通过excludeTablePlaceholderReplaceKeys()显式排除;
  • 默认值兜底 + 按数据源选型:对可能缺失的层级提供默认值,并根据上游数据源(MySQL、Oracle 等)的命名空间特性选择合适的占位符。

延伸阅读

  • 官方英文文档:Sink Options Placeholders、中文版见 docs/zh/concept/sink-options-placeholders.md;
  • 核心实现:TablePlaceholder.java、TableSinkFactoryContext.java、FactoryUtil.java;
  • 单元测试:TablePlaceholderTest.java;
  • 相关功能说明:connector-v2-features.md。
  • 数据工程
  • 大数据
  • 批处理
  • 流处理

【免费下载链接】seatunnel

SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.

项目地址:https://gitcode.com/gh_mirrors/sea/seatunnel
点击查看免费下载

相关推荐

上一篇:YCVideoPlayer与常见播放器对比:为什么选择YC作为你的视频解决方案
下一篇:Sigma.js图层架构全景图:7个WebGL与Canvas图层如何协作渲染

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询