- 数据工程
- 大数据
- 批处理
- 流处理
【免费下载链接】seatunnel
SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.
本文以 SeaTunnel 官方中文文档 field-mapper.md 为核心,结合仓库源码与端到端测试配置,系统讲解 FieldMapper 转换插件(Transform Plugin)的工作原理、配置方法与实战用法。读完本文,你将掌握如何通过一段
field_mapper映射配置,在数据流入 Sink 之前完成字段删减、重命名、输出顺序重排等常见表结构整形操作,并理解其在 SeaTunnel 内部 Schema 转换与行级数据搬运层面的实现机制。
一、插件概述:什么是 FieldMapper
FieldMapper是 SeaTunnel Transform V2 体系中的一张字段映射转换插件,其核心能力用一句话概括就是:
添加输入模式(input schema)和输出模式(output schema)之间的映射关系。
在实际的数据集成作业中,源端表结构与目标端表结构往往并不完全一致:目标表可能不需要源表中的某些列,某些列需要换一个名字,列的输出顺序也可能与源表不同。FieldMapper 正是为这种"表结构整形"需求而生——它允许你以配置的方式声明"源字段 → 输出字段"的映射,从而:
- 删除字段:不写入映射的输入字段,在输出表中自然消失;
- 重命名字段:将映射的 value 指定为新字段名;
- 重排字段顺序:输出表的列顺序严格遵循
field_mapper映射的书写顺序; - 保持字段数据类型:输出列的类型、长度、可空性、默认值、注释等属性继承自输入列(见 FieldMapperTransform.java)。
从源码结构看,FieldMapper继承自AbstractCatalogSupportTransform,即它是一类"同时转换表 Schema 与行数据"的 Catalog 感知型转换插件,会在作业初始化阶段完成输出 Schema 的推导,并在数据流经时按列索引搬运数据,而不是逐字段做字符串级别的复制,因此具备较高的执行效率。
二、配置参数详解
2.1 参数总览
| 名称 | 类型 | 是否必须 | 默认值 | 说明 |
|---|---|---|---|---|
| field_mapper | Object | 是 | 无 | 指定输入与输出之间的字段映射关系 |
| source_table_name | string | 否 | - | 指定要读取的输入数据集(临时表)名称 |
| result_table_name | string | 否 | - | 将转换结果注册为可供下游插件访问的数据集(临时表)名称 |
其中source_table_name与result_table_name属于所有转换插件通用的 common options,详细语义见 Transform 常见选项。
2.2 field_mapper [config](必填)
field_mapper是 FieldMapper 唯一的必填参数,类型为Object(键值对 Map),语义为:
- key:输入表中的字段名(源字段);
- value:输出表中的字段名(目标字段)。当需要重命名时,value 写成新字段名;当保持原名时,value 与 key 相同即可。
该参数的合法性在源码中被强制约束:在 FieldMapperTransformFactory.java 中,optionRule()将field_mapper声明为 required 选项,缺少该参数时作业配置校验会直接失败;而 FieldMapperTransformConfig.java 将其定义为mapType()且noDefaultValue()的Map<String, String>类型选项。
值得特别注意的是,配置类中使用LinkedHashMap承载映射关系(FieldMapperTransformConfig.java),这意味着field_mapper中键值对的书写顺序就是输出表字段的排列顺序,这也是该插件能够"重排字段顺序"的实现基础。
2.3 common options [config]
- source_table_name:不指定时,当前转换插件处理的是配置文件中前一个插件输出的数据集;指定时,则处理与该参数对应的、已注册的数据集(临时表)。
- result_table_name:不指定时,转换结果不会被注册为数据集,下游插件无法直接访问;指定时,结果会注册为一个数据集/临时表,下游插件可通过
source_table_name引用它。
在官方示例中,source_table_name = "fake"表示读取上游FakeSource(或其他 Source)注册的fake表,result_table_name = "fake1"表示将映射结果注册为fake1表供后续 Sink 或 Transform 使用。
三、配置示例与结果验证
3.1 官方示例:删字段 + 重命名 + 重排顺序
假设源端读取到的数据表如下:
| id | name | age | card |
|---|---|---|---|
| 1 | Joy Ding | 20 | 123 |
| 2 | May Ding | 20 | 123 |
| 3 | Kin Dom | 20 | 123 |
| 4 | Joy Dom | 20 | 123 |
我们希望删除age字段,将输出字段顺序调整为id、card、name,并把name重命名为new_name。只需像下面这样在transform块中添加FieldMapper转换:
transform { FieldMapper { source_table_name = "fake" result_table_name = "fake1" field_mapper = { id = id card = card name = new_name } } }经过该转换后,结果表fake1中的数据将变为:
| id | card | new_name |
|---|---|---|
| 1 | 123 | Joy Ding |
| 2 | 123 | May Ding |
| 3 | 123 | Kin Dom |
| 4 | 123 | Joy Dom |
可以看到:未出现在映射中的age字段被删除;name字段被重命名为new_name;输出列严格按映射书写顺序排列为id、card、new_name。
3.2 端到端可运行配置:FakeSource + FieldMapper + Assert
仓库的端到端测试中提供了完整的可运行作业配置 field_mapper_transform.conf,展示了 FieldMapper 与FakeSource、AssertSink 的组合用法,可作为实战模板:
env { job.mode = "BATCH" } source { FakeSource { result_table_name = "fake" row.num = 100 schema = { fields { id = "int" name = "string" age = "int" string1 = "string" int1 = "int" c_bigint = "bigint" c_row = { c_row = { c_int = int } } } } } } transform { FieldMapper { source_table_name = "fake" result_table_name = "fake1" field_mapper = { id = id age = age_as int1 = int1_as name = name c_row = c_row } } } sink { Assert { source_table_name = "fake1" rules = { row_rules = [ { rule_type = MIN_ROW rule_value = 100 } ], field_rules = [ { field_name = id field_type = int field_value = [ { rule_type = NOT_NULL } ] }, { field_name = age_as field_type = int field_value = [ { rule_type = NOT_NULL } ] }, { field_name = int1_as field_type = int field_value = [ { rule_type = NOT_NULL } ] }, { field_name = name field_type = string field_value = [ { rule_type = NOT_NULL } ] } ] } } }该配置中,FakeSource产出包含id/name/age/string1/int1/c_bigint/c_row等字段的 100 行数据,FieldMapper将其映射为id、age_as、int1_as、name、c_row五个字段(age重命名为age_as、int1重命名为int1_as、删除了string1与c_bigint、保留嵌套行类型c_row),随后AssertSink 通过source_table_name = "fake1"引用映射结果,校验行数不少于 100 且各字段非空、类型正确。对应的测试用例 TestFieldMapperIT.java 会在多个引擎容器中执行该配置并断言退出码为 0,这从侧面验证了 FieldMapper 在真实运行链路中的行为。
四、源码实现原理深入
4.1 参数解析与插件注册
- FieldMapperTransformConfig.java:定义
field_mapper选项,并将配置解析结果保存为LinkedHashMap<String, String>; - FieldMapperTransformFactory.java:通过
@AutoService(Factory.class)注册为 SeaTunnel 可发现的转换工厂,factoryIdentifier()返回"FieldMapper"(即配置块中的插件名),optionRule()声明field_mapper为必填项,createTransform()取第一个 CatalogTable 与用户配置构造转换实例。
工厂测试 FieldMapperTransformFactoryTest.java 会断言optionRule()非空,保证参数规则可被框架正确加载。
4.2 输出 Schema 的推导(transformTableSchema)
在 FieldMapperTransform.java 中,transformTableSchema()负责将输入 Schema 转换为输出 Schema,关键逻辑如下:
- 按映射构建输出列:遍历
field_mapper的每个键值对,在输入字段名列表中查找 key 对应的索引,取出输入列后,用 value(新字段名)配合输入列的数据类型、长度、可空性、默认值、注释构造新的PhysicalColumn,即输出列完整继承输入列的元数据,仅名字可能变化; - 记录列索引:同时记录每个输出字段对应的输入列索引
needReaderColIndex,供行数据转换阶段直接按索引取值; - 保留主键与约束键:若输入表存在主键(PrimaryKey)且其所有列都出现在输出字段中,则主键被复制保留;约束键(ConstraintKey)同理——只有完全落在输出字段集合内的约束才会被继承,避免因字段被删除而产生失效约束。
4.3 行数据的搬运(transformRow)
FieldMapperTransform.java 中的transformRow()采用"按索引搬运"的方式处理每行数据:
- 根据
field_mapper的大小创建一个等长的Object[]输出数组; - 依次从输入行按
needReaderColIndex取出对应位置的字段值放入输出数组; - 构造新的
SeaTunnelRow,并保留输入行的 RowKind(插入/更新/删除标记)与 TableId,保证 CDC 语义与表标识在转换后不丢失。
这种方式意味着 FieldMapper 不做任何字段值的加工或类型转换,纯粹是"位置重排 + 名字替换",因此执行开销极低。
4.4 输入字段校验与错误处理
插件在构造函数阶段就会对映射合法性做前置校验(FieldMapperTransform.java):
- 逐个检查
field_mapper的 key 是否存在于输入表的物理字段中; - 若存在找不到的字段,会抛出
TransformCommonError.cannotFindInputFieldsError(...)错误,作业在启动阶段即失败,而不是在跑批中途报错; - 在
transformTableSchema()中同样对每个 key 再次检查索引有效性,双保险地杜绝"映射指向不存在的字段"导致的运行时问题。
五、典型应用场景
- 字段裁剪(列裁剪):下游系统只关心部分列时,仅列出需要的字段即可丢弃多余字段,降低网络传输与存储开销;
- 字段重命名:适配目标表/目标系统的命名规范,例如将源端的
user_name映射为下游约定的name; - 字段顺序重排:对接按固定列顺序写入的文件或数据库表时,通过映射书写顺序精确控制输出列序;
- Schema 对齐:在多源合并写入同一目标表的场景中,用 FieldMapper 将不同源统一到同一目标 Schema,再配合其他转换继续处理。
六、更新日志
新版本
- 添加复制转换连接器(Copy Transform Connector,见 copy.md),可与 FieldMapper 组合使用,在保留原始字段的同时复制出经过重命名/重排的新字段。
- 数据工程
- 大数据
- 批处理
- 流处理
【免费下载链接】seatunnel
SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.
相关推荐
SeaTunnel FieldMapper Transform 插件详解:字段映射、重命名与顺序调整实战指南
SeaTunnel FieldMapper Transform 插件详解:字段映射、重命名与顺序调整实战指南 导读 FieldMapper 是 SeaTunne
数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel FieldMapper 字段映射转换插件详解:重命名、排序与裁剪字段
SeaTunnel FieldMapper 字段映射转换插件详解:重命名、排序与裁剪字段 本文基于 SeaTunnel 开源仓库的 FieldMapper 转换
数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel FieldMapper 转换插件深度指南:字段映射、列重排与重命名的完整实战
SeaTunnel FieldMapper 转换插件深度指南:字段映射、列重排与重命名的完整实战 FieldMapper 是 SeaTunnel( seatun
数据工程大数据批处理流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考