- 数据工程
- 大数据
- 批处理
- 流处理
【免费下载链接】seatunnel
SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.
导读
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 调用链全景
- 入口:Sink 创建统一走 FactoryUtil.createAndPrepareSink。它发现
TableSinkFactory后,调用TableSinkFactoryContext.replacePlaceholderAndCreate(...)创建上下文; - 上下文构造:TableSinkFactoryContext.replacePlaceholderAndCreate 内部调用
TablePlaceholder.replaceTablePlaceholder(options, catalogTable, excludeTablePlaceholderReplaceKeys),得到替换后的ReadonlyConfig; - 校验与创建:替换后的配置经
ConfigValidator.validate(factory.optionRule())校验通过后,才交给factory.createSink(context)真正创建连接器实例。
因此,连接器拿到的永远是替换完成的最终配置,占位符不会泄漏到下游连接器的运行时逻辑中。
5.2 替换逻辑的核心实现
TablePlaceholder是这一功能的唯一实现类,其替换流程分四个步骤(TablePlaceholder.java):
- 遍历配置:克隆
ReadonlyConfig的源 Map(copy-on-write),逐键处理,命中excludeKeys(即连接器通过excludeTablePlaceholderReplaceKeys()声明排除的键)则跳过; - 表标识替换:
replaceTableIdentifier依次替换${database_name}、${schema_name}、${table_name},并按“database + schema”(schema_full_name)、“database + schema + table”(table_full_name)的拼接顺序替换完整路径占位符;拼接时跳过为 null 的层级; - 结构信息替换:
replaceTablePrimaryKey、replaceTableUniqueKey、replaceTableFieldNames分别从TableSchema中提取主键列、唯一键列、全部字段名,以,连接后替换对应占位符; - 特殊处理 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}也会保持原样。
两类应对方案:
- 使用默认值语法兜底:如
${database_name:default_my_db},元数据缺失时自动回退; - 按数据源特性选择占位符:先确认上游 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.
相关推荐
SeaTunnel Sink 参数占位符(Sink Options Placeholders)完全指南:动态获取上游表元数据,实现多表自动路由写入
SeaTunnel Sink 参数占位符(Sink Options Placeholders)完全指南:动态获取上游表元数据,实现多表自动路由写入 导读 本文系
数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Sink Options 占位符(Placeholders)深度指南:基于上游表元数据的动态写入配置
SeaTunnel Sink Options 占位符(Placeholders)深度指南:基于上游表元数据的动态写入配置 SeaTunnel 提供了一套 Sin
数据集成ETL大数据批处理流处理变更数据捕获Apache SeaTunnel Sink参数占位符使用详解
Apache SeaTunnel Sink参数占位符使用详解 引言 在数据集成和处理场景中,我们经常需要将数据从源系统抽取后写入到目标系统。Apache Sea
数据集成ETL大数据批处理流处理变更数据捕获
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考