SeaTunnel 多表 Transform(Multi-Table Transform)完整指南:单配置处理多张上游表
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
SeaTunnel 的 Transform 层原生支持多表(Multi-Table)变换能力:当上游插件一次输出多张表(如JDBCSource、MySQL-CDC等)时,你可以在一个 Transform 配置块内为不同表分别配置变换规则,并把多张表的规则合并进同一个 Transform,统一管理。读完本文,你将掌握table_match_regex、table_transform、table_path、rule_match_mode四个核心参数的含义与优先级,能够用 Copy 等任意 Transform 以"正则批量 + 单表单条"的组合方式完成多表变换,并理解其底层匹配与分发实现。
一、什么是多表 Transform
在 SeaTunnel 中,大多数 Source 插件默认输出一张表,但部分插件支持一次输出多张表(multi-catalog / multi-table),例如JDBCSource(多表抽取)和MySQL-CDC(整库/多表变更捕获)。当数据流中存在多张表时,如果仍使用传统单表 Transform 配置,你就需要为每张表分别书写一段 Transform,配置冗长且难以维护。
多表 Transform(Multi-Table Transform)正是为解决这个问题而设计的:它允许在一个 Transform 配置中同时声明多张表的变换规则,由框架按表分发并执行,实现"一处配置、多表生效"。
能力边界(官方明确说明):多表 Transform 对 Transform 的能力没有任何限制,任何 Transform 配置都可以用于多表 Transform。多表 Transform 的本质是:对数据流中的多张表分别独立处理,同时把多张表的 Transform 配置合并到一个 Transform 中,便于统一管理。
这一点在源码结构上得到了印证:seatunnel-transforms-v2中几乎所有 Transform 都提供了对应的*MultiCatalogTransform实现,例如:
- Copy:
CopyFieldMultiCatalogTransform(CopyFieldMultiCatalogTransform.java) - Calcite SQL:
CalciteMultiCatalogTransform - 字段映射:
FieldMapperMultiCatalogTransform - 字段过滤:
FilterFieldMultiCatalogTransform - 字段加密:
FieldEncryptMultiCatalogTransform - 元数据:
MetadataMultiCatalogTransform - LLM / Embedding / Python / DynamicCompile 等均有对应实现
它们统一继承自 AbstractMultiCatalogTransform.java,共享同一套多表配置与匹配机制。
二、多表 Transform 属性(Properties)
多表 Transform 在原有 Transform 参数之上,额外增加以下公共参数。这些参数定义在 TransformCommonOptions.java 中,对所有支持多表的 Transform 通用:
| 名称 | 类型 | 是否必填 | 默认值 | 说明 |
|---|---|---|---|---|
table_match_regex | String | 否 | .* | 用于匹配"需要做变换的表"的正则表达式,默认匹配所有表。注意:这里匹配的是上游真实表名(即 table path),不是plugin_output指定的数据集名称。 |
table_transform | List | 否 | - | 在table_transform中可以用列表方式为单张表指定变换规则。如果某张表在table_transform中配置了专属规则,则外层规则对该表不生效(table_transform中的规则优先)。 |
table_transform.table_path | String | 否 | - | 在table_transform中为某张表配置规则时,必须通过table_path指定表路径。表路径格式为databaseName[.schemaName].tableName,采用精确匹配。 |
rule_match_mode | String | 否 | FIRST_MATCH | 控制当多条table_transform规则指向完全相同的table_path时的求值方式。可选值:FIRST_MATCH与ALL_MATCH。 |
从源码可以看到每个参数的底层定义:
// seatunnel-transforms-v2/src/main/java/org/apache/seatunnel/transform/common/TransformCommonOptions.java public static final Option<List<Map<String, Object>>> MULTI_TABLES = Options.key("table_transform") .type(new TypeReference<List<Map<String, Object>>>() {}) .defaultValue(Collections.emptyList()) .withDescription("The table transform config"); public static final Option<String> TABLE_PATH = Options.key("table_path") .stringType() .noDefaultValue() .withDescription("The table path of catalog table"); public static final Option<String> TABLE_MATCH_REGEX = Options.key("table_match_regex") .stringType() .defaultValue(".*") .withDescription("The regex to match the table path"); public static final Option<RuleMatchMode> RULE_MATCH_MODE = Options.key("rule_match_mode") .enumType(RuleMatchMode.class) .defaultValue(RuleMatchMode.FIRST_MATCH) .withDescription("The rule match mode for table transform config");几点值得注意的源码细节:
table_transform的默认值为Collections.emptyList(),即不配置任何单表规则;table_match_regex默认值为.*,即默认匹配所有表——这意味着如果你不写任何多表参数,多表 Transform 的行为与普通 Transform 一致,对每张表都套用外层规则;rule_match_mode的默认值是FIRST_MATCH(注意与下文"未配置时拒绝重复 table_path"的行为相区分:是否显式配置了该参数,行为不同)。
三、多表匹配逻辑与配置优先级
3.1 匹配优先级
对每一张表,配置的生效优先级为:
table_transform(精确的单表规则) > table_match_regex(正则批量规则)如果某张表既没有命中table_transform中的table_path,也不匹配table_match_regex,那么该表不应用任何变换(数据原样透传)。
3.2 源码中的匹配流程
AbstractMultiCatalogTransform的构造函数完整实现了上述逻辑(见 AbstractMultiCatalogTransform.java):
- 编译
table_match_regex为正则Pattern; - 读取
table_transform列表,并过滤出含table_path的条目; - 遍历每张输入
CatalogTable,取出其 table path(形如database.schema.table或database.table):- 若存在
table_path精确等于该表路径的规则 → 使用该(些)单表规则; - 否则,若
table_match_regex能匹配该表路径 → 使用外层(Transform 块级)配置; - 否则 → 为该表创建 Identity(透传)Transform,不做任何变换。
- 若存在
inputCatalogTables.forEach(inputCatalogTable -> { String tableId = inputCatalogTable.getTableId().toTablePath().toString(); List<ReadonlyConfig> tableConfigs = singleTableConfigs.stream() .filter(c -> tableId.equals(c.get(TransformCommonOptions.TABLE_PATH))) .collect(Collectors.toList()); if (!tableConfigs.isEmpty()) { transformMap.put(tableId, buildTransform(inputCatalogTable, selectTableConfigs(tableConfigs, ruleMatchMode))); } else if (tableMatchRegex.matcher(tableId).matches()) { transformMap.put(tableId, buildTransform(inputCatalogTable, Collections.singletonList(config))); } else { transformMap.put(tableId, createIdentityTransform(inputCatalogTable)); } });在运行时,数据行会根据其所属表 ID 被分发到对应的内部 Transform。以 Map 类 Transform 为例(AbstractMultiCatalogMapTransform.java):
@Override public SeaTunnelRow map(SeaTunnelRow row) { if (transformMap.size() == 1) { return ((SeaTunnelMapTransform<SeaTunnelRow>) transformMap.values().iterator().next()).map(row); } return ((SeaTunnelMapTransform<SeaTunnelRow>) transformMap.get(row.getTableId())).map(row); }即:只有一个内部 Transform 时直接复用;多表时则按row.getTableId()精确路由到该表对应的变换器。
3.3 table_path 精确匹配与 rule_match_mode
table_transform.table_path采用精确匹配,格式为databaseName[.schemaName].tableName(例如test.xyz、mydb.dbo.users)。rule_match_mode只控制"多条table_transform条目使用了同一个精确table_path"这一种情况:
- 未配置
rule_match_mode:配置解析阶段直接拒绝重复的精确table_path条目(抛出异常),避免歧义; FIRST_MATCH:按声明顺序应用第一条匹配的table_transform条目;ALL_MATCH:按声明顺序应用所有匹配的table_transform条目,前一条规则对同表的输出会作为下一条规则的输入(规则链式叠加)。
源码中对应的两个方法直观地展示了这一行为:
// 未配置 rule_match_mode 时,拒绝重复 table_path private void rejectDuplicateTablePaths(List<ReadonlyConfig> tableConfigs) { Map<String, Boolean> tablePaths = new HashMap<>(); for (ReadonlyConfig tableConfig : tableConfigs) { String tablePath = tableConfig.get(TransformCommonOptions.TABLE_PATH); if (tablePaths.put(tablePath, true) != null) { throw new IllegalStateException(String.format( "Duplicate table_transform rules are configured for table_path [%s]", tablePath)); } } } // FIRST_MATCH 取第一条;ALL_MATCH 取全部 private List<ReadonlyConfig> selectTableConfigs( List<ReadonlyConfig> tableConfigs, TransformCommonOptions.RuleMatchMode ruleMatchMode) { if (ruleMatchMode == TransformCommonOptions.RuleMatchMode.FIRST_MATCH) { return Collections.singletonList(tableConfigs.get(0)); } return tableConfigs; }对于ALL_MATCH模式,多个规则会按顺序逐个构建 Transform 并串成链(ChainedMapTransform),前一规则的输出表结构(getProducedCatalogTable())会成为下一规则的输入结构:
for (ReadonlyConfig config : configs) { SeaTunnelTransform<SeaTunnelRow> transform = buildTransform(currentCatalogTable, config); transforms.add(transform); currentCatalogTable = transform.getProducedCatalogTable(); }四、完整示例:一个 Copy Transform 处理五张表
4.1 场景假设
假设上游一次性读取了五张结构相同的表:test.abc、test.abcd、test.xyz、test.xyzxyz、test.www,每张表都有三个字段:id、name、age。
我们希望通过 Copy Transform 复制这些表的数据,具体要求如下:
- 对
test.abc和test.abcd:将name字段复制到新字段name1; - 对
test.xyz:将name字段复制到name2; - 对
test.xyzxyz:将name字段复制到name3; - 对
test.www:不做任何变换。
4.2 一个配置搞定多张表
transform { Copy { plugin_input = "fake" // 可选:指定读取的数据集名称 plugin_output = "fake1" // 可选:指定输出的数据集名称 table_match_regex = "test.a.*" // 1. 匹配需要变换的表,这里命中 test.abc 和 test.abcd src_field = "name" // 源字段 dest_field = "name1" // 目标字段 table_transform = [{ table_path = "test.xyz" // 2. 指定要变换的表名 src_field = "name" // 源字段 dest_field = "name2" // 目标字段 }, { table_path = "test.xyzxyz" src_field = "name" dest_field = "name3" }] } }4.3 配置解读
- 通过正则
test.a.*及对应的 Copy 参数,命中test.abc和test.abcd,将name复制为name1; - 通过
table_transform为test.xyz单独指定规则,将name复制为name2;同理为test.xyzxyz将name复制为name3; test.www既不匹配正则,也没有table_transform规则,因此不应用任何变换。
这样,我们就在一个 Transform 配置内完成了多张表的变换处理。
关于
src_field/dest_field:从 CopyTransformConfig.java 源码看,src_field与dest_field属于已标记@Deprecated的旧式写法;推荐的新写法是使用fields映射(例如fields { name = "name1" }),一次可声明多组字段复制关系。二者在CopyTransformConfig.of()中会统一归一化为LinkedHashMap<String, String> fields。多表示例中为了直观展示"每张表一套规则"的写法仍沿用旧式参数,实际生产配置建议优先使用fields。
4.4 各表实际生效的配置与输出结构
- test.abc 与 test.abcd(外层规则):
transform { Copy { src_field = "name" dest_field = "name1" } }输出结构:
| id | name | age | name1 |
- test.xyz(table_transform 精确规则):
transform { Copy { src_field = "name" dest_field = "name2" } }输出结构:
| id | name | age | name2 |
- test.xyzxyz(table_transform 精确规则):
transform { Copy { src_field = "name" dest_field = "name3" } }输出结构:
| id | name | age | name3 |
- test.www(无规则命中,透传):
transform { // 无需任何变换 }输出结构:
| id | name | age |
4.5 优先级再确认
以上示例再次印证了优先级规则:test.abc/test.abcd虽然也在table_transform规则之外,但它们命中外层table_match_regex,因此应用外层规则;而test.xyz/test.xyzxyz命中了table_transform中的精确table_path,外层规则对它们不生效,各自应用专属规则;test.www两类规则均未命中,原样输出。
五、rule_match_mode 实战:ALL_MATCH 链式变换
当多条table_transform条目指向同一个table_path时,可用rule_match_mode控制求值方式。下面的配置在ALL_MATCH模式下,先对test.xyz把name复制为name2,再对同一张表把name2复制为name3:
transform { Copy { rule_match_mode = "ALL_MATCH" table_transform = [{ table_path = "test.xyz" src_field = "name" dest_field = "name2" }, { table_path = "test.xyz" src_field = "name2" dest_field = "name3" }] } }输出结构:
| id | name | age | name2 | name3 |
这里前一条规则(name→name2)的输出字段name2成为后一条规则(name2→name3)的输入,两条规则按声明顺序链式执行。
需要注意的边界情况:
- 如果不配置
rule_match_mode,上述"两条规则指向同一 table_path"的配置会在配置解析阶段直接报错(Duplicate table_transform rules are configured for table_path [test.xyz]),这是源码 AbstractMultiCatalogTransform.java 中rejectDuplicateTablePaths的强制校验; FIRST_MATCH模式则只取声明顺序中的第一条规则,第二条会被忽略。
六、其他 Transform 的多表用法
多表 Transform 的能力不仅限于 Copy。本文示例使用的是 Copy Transform,但 SeaTunnel 中所有 Transform 都支持多表变换,你只需在对应 Transform 的配置块中按照同样的方式书写table_match_regex、table_transform等参数即可。
例如,对字段过滤(Filter)、字段加密(Encrypt)、SQL(Calcite)、元数据提取(Metadata)等 Transform,多表配置结构与本文完全一致——它们各自的*MultiCatalogTransform实现均继承自AbstractMultiCatalogTransform,共享同一套参数解析与匹配机制。可以参考 seatunnel-transforms-v2 模块下的实现,以及 transforms 文档目录 中各 Transform 的详细参数说明。
七、小结
| 要点 | 说明 |
|---|---|
| 适用场景 | 上游一次输出多张表(如JDBCSource、MySQL-CDC)时,在一个 Transform 内统一配置变换 |
| 批量规则 | table_match_regex用正则匹配表路径(默认.*匹配全部),匹配的是真实上游表名 |
| 单表规则 | table_transform+table_path(database[.schema].table精确匹配)为单表定制规则 |
| 优先级 | table_transform>table_match_regex;均未命中则透传 |
| 同表多规则 | 不配rule_match_mode会拒绝重复table_path;FIRST_MATCH取第一条;ALL_MATCH链式叠加 |
| 适用范围 | 所有 Transform 均支持,能力不受限 |
配置多个多表规则时,建议遵循:能用正则批量覆盖的场景用table_match_regex,需要为个别表定制差异规则时用table_transform,同时注意同一表路径的多规则必须显式声明rule_match_mode为ALL_MATCH(或FIRST_MATCH),避免配置解析失败。
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考