ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

SeaTunnel RowKindExtractor 转换插件详解:将 CDC 数据流转换为 Append-Only 模式并保留变更类型

SeaTunnel RowKindExtractor 转换插件详解:将 CDC 数据流转换为 Append-Only 模式并保留变更类型 SeaTunnel RowKindExtractor 转换插件详解将 CDC 数据流转换为 Append-Only 模式并保留变更类型【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel导读RowKindExtractor 是 SeaTunnel 提供的一个转换Transform插件核心解决 CDCChange Data Capture变更数据捕获同步场景中的一类经典难题当数据湖、分析型系统等下游只支持 Append-Only仅追加写入、不支持 UPDATE/DELETE 操作时如何既能把 CDC 数据降级为纯插入流又不丢失每一行数据原本的变更语义。本文基于仓库中的插件文档与源码实现完整讲解该插件的工作原理、两个配置参数custom_field_name与transform_type、三种可复制的实战配置示例并结合RowKind枚举与单测用例剖析其底层实现机制。读完本文你将掌握在 SeaTunnel 中把 MySQL CDC 数据安全写入 Iceberg、数据湖等只追加存储的标准姿势。插件是什么CDC 数据流与 Append-Only 下游的桥梁在 CDC 数据同步场景中每条数据行都会携带一个 RowKind行变更类型标记用来描述这条数据发生的变更操作。SeaTunnel 中 RowKind 的完整定义位于 RowKind.java共四种取值RowKind 类型短格式shortString字节值byte含义INSERTI0插入操作UPDATE_BEFORE-U1更新前的旧值用于先撤回旧行UPDATE_AFTERU2更新后的新值DELETE-D3删除操作然而许多下游系统如 Iceberg 等数据湖、部分分析型数据库仅支持 Append-Only 写入模式无法直接消费带-U、U、-D语义的变更流。此时需要做两件事将所有数据的 RowKind 统一转换为IINSERT使其变成 Append-Only 的纯插入流把原始的变更类型保存为一个普通字段供后续按变更类型进行统计、审计或回放分析。RowKindExtractor 插件正是为这两个需求而生的它将数据流中所有行的 RowKind 强制改写为I同时把原始 RowKind 信息提取出来写入一个新追加的字段。从源码看该功能由 RowKindExtractorTransform.java 实现它继承自SingleFieldOutputTransform单字段输出型转换的抽象基类见 SingleFieldOutputTransform.java天然支持只新增一个字段的典型转换形态。转换效果示意输入CDC 数据 RowKind: -D (DELETE) Data: id1, nametest1, age20 输出Append-Only 数据 RowKind: I (INSERT) Data: id1, nametest1, age20, row_kindDELETE典型应用场景将 CDC 数据写入仅支持 Append 模式的数据湖在数据仓库中保留完整的变更历史每次变更都成为一行记录对不同类型变更增/改前/改后/删进行统计分析。配置参数详解RowKindExtractor 插件只有两个业务配置参数定义在 RowKindExtractorTransformConfig.java 中nametyperequireddefault valuedescriptioncustom_field_namestringnorow_kind用于存储原始 RowKind 信息的新字段名transform_typeenumnoSHORTRowKind 输出格式可选 SHORT短格式或 FULL完整格式此外插件工厂 RowKindExtractorTransformFactory.java 还注册了三个通用的多表转换选项multi_tables、table_match_regex与rule_match_mode使其可以配合多表场景如多库多表 CDC使用。custom_field_name [string]指定用于存储原始 RowKind 信息的新字段名。默认值row_kind注意事项字段名不能与上游已有字段重名否则插件会直接抛出异常。该校验逻辑在getOutputColumn()中实现RowKindExtractorTransform.java插件会先取出输入表的全部字段名并检查是否包含自定义字段名若已存在则抛出IllegalArgumentException(field name %s already exists)建议使用语义清晰的名称如operation_type、change_type、cdc_op等。示例custom_field_name operation_type # 使用自定义字段名transform_type [enum]指定新字段中 RowKind 值的输出格式。枚举定义在 RowKindExtractorTransformType.java 中仅含SHORT与FULL两个取值。可用选项格式说明输出值SHORT短格式符号表示I、-U、U、-DFULL完整格式英文名称INSERT、UPDATE_BEFORE、UPDATE_AFTER、DELETE默认值SHORT两种格式与 RowKind 的对应关系RowKind 类型SHORT 格式FULL 格式说明INSERTIINSERT插入操作UPDATE_BEFORE-UUPDATE_BEFORE更新前的值UPDATE_AFTERUUPDATE_AFTER更新后的值DELETE-DDELETE删除操作选型建议SHORT 格式节省存储空间符号仅占 2 个字符适合对存储成本敏感的场景FULL 格式可读性更好适合需要人工排查或直接面向分析师的场景。示例transform_type FULL # 使用完整格式底层实现原理从源码看数据是怎么被改写的理解了参数后再深入源码可以更清楚地看到插件的完整执行链路。核心逻辑集中在 RowKindExtractorTransform.java 的transformRow方法中L54-L61protected SeaTunnelRow transformRow(SeaTunnelRow inputRow) { Object fieldValue getOutputFieldValue(new SeaTunnelRowAccessor(inputRow)); inputRow.setRowKind(RowKind.INSERT); // ① 强制改写为 INSERT SeaTunnelRow outputRow getRowContainerGenerator().apply(inputRow); // ② 扩展行容器 outputRow.setField(getFieldIndex(), fieldValue); // ③ 写入原始 RowKind 值 return outputRow; }整个转换可拆解为三个步骤提取原始 RowKind通过getOutputFieldValue()读取输入行的 RowKind 并格式化为目标字符串。SHORT 模式调用RowKind.shortString()返回I/-U/U/-DFULL 模式调用RowKind.name()返回INSERT/UPDATE_BEFORE/UPDATE_AFTER/DELETE见 L63-L74改写 RowKind将输入行的 RowKind 直接置为RowKind.INSERT实现 Append-Only 化追加字段行容器生成器rowContainerGenerator由SingleFieldOutputTransform在表结构转换时构建负责把原字段复制到更长的输出数组再把提取到的值写入新字段的索引位置。关于新字段的 Schema 定义L77-L92源码显示它是一个字符串类型字段PhysicalColumn.of(customFieldName, BasicType.STRING_TYPE, 13L, false, I, Output column of RowKind)即类型为string、长度 13、默认值I。多表场景的支持插件通过 RowKindExtractorMultiCatalogTransform.java 支持多 Catalog 表输入它会为每个输入CatalogTable分别构建一个独立的RowKindExtractorTransform实例未匹配到的表则走IdentityMapTransform透传因此可以无缝对接多表 CDC 同步任务。测试验证仓库中的单元测试 RowKindExtractorTransformTest.java 直接验证了上述行为testCdcRowTransformShort在默认配置下分别将 INSERT、UPDATE_BEFORE、UPDATE_AFTER、DELETE 四种行输入断言输出均为kindI且末尾追加字段依次为I、-U、U、-DL99-L124testCdcRowTransformFull设置transform_typeFULL后断言末尾字段依次为INSERT、UPDATE_BEFORE、UPDATE_AFTER、DELETEL126-L152。此外 RowKindExtractorTransformFactoryTest.java 验证了工厂的optionRule()能正常构建插件可被 SeaTunnel 工厂机制正确发现与装配。完整示例一默认配置SHORT 格式使用默认配置将 CDC 数据转换为 Append-Only 模式RowKind 以短格式保存到row_kind字段。env { parallelism 1 job.mode STREAMING } source { MySQL-CDC { plugin_output cdc_source server-id 5652 username root password your_password table-names [mydb.users] url jdbc:mysql://localhost:3306/mydb } } transform { RowKindExtractor { plugin_input cdc_source plugin_output append_only_data # 使用默认配置 # custom_field_name row_kind # transform_type SHORT } } sink { Console { plugin_input append_only_data } }数据转换过程输入数据CDC 格式 1. RowKindI, id1, nameJohn, age25 2. RowKind-U, id1, nameJohn, age25 3. RowKindU, id1, nameJohn, age26 4. RowKind-D, id1, nameJohn, age26 输出数据Append-Only 格式 1. RowKindI, id1, nameJohn, age25, row_kindI 2. RowKindI, id1, nameJohn, age25, row_kind-U 3. RowKindI, id1, nameJohn, age26, row_kindU 4. RowKindI, id1, nameJohn, age26, row_kind-D完整示例二FULL 格式 自定义字段名使用完整格式输出 RowKind并自定义新字段名适合写入 Iceberg 等数据湖做历史留存与变更分析。env { parallelism 1 job.mode STREAMING } source { MySQL-CDC { plugin_output cdc_source server-id 5652 username root password your_password table-names [mydb.orders] url jdbc:mysql://localhost:3306/mydb } } transform { RowKindExtractor { plugin_input cdc_source plugin_output append_only_data custom_field_name operation_type # 自定义字段名 transform_type FULL # 使用完整格式 } } sink { Iceberg { plugin_input append_only_data catalog_name iceberg_catalog database mydb table orders_history # Iceberg 表中将包含 operation_type 字段记录每条数据的变更类型 } }数据转换过程输入数据CDC 格式 1. RowKindI, order_id1001, amount100.00 2. RowKind-U, order_id1001, amount100.00 3. RowKindU, order_id1001, amount150.00 4. RowKind-D, order_id1001, amount150.00 输出数据Append-Only 格式FULL 格式 1. RowKindI, order_id1001, amount100.00, operation_typeINSERT 2. RowKindI, order_id1001, amount100.00, operation_typeUPDATE_BEFORE 3. RowKindI, order_id1001, amount150.00, operation_typeUPDATE_AFTER 4. RowKindI, order_id1001, amount150.00, operation_typeDELETE完整示例三用 FakeSource 构造测试数据无需真实数据库使用 FakeSource 直接构造带 RowKind 的测试数据快速验证插件的各种转换效果。这也是调试与集成测试中最常用的手段。env { parallelism 1 job.mode BATCH } source { FakeSource { plugin_output fake_cdc_data schema { fields { pk_id bigint name string score int } primaryKey { name pk_id columnNames [pk_id] } } rows [ { kind INSERT fields [1, A, 100] }, { kind INSERT fields [2, B, 100] }, { kind UPDATE_BEFORE fields [1, A, 100] }, { kind UPDATE_AFTER fields [1, A_updated, 95] }, { kind UPDATE_BEFORE fields [2, B, 100] }, { kind UPDATE_AFTER fields [2, B_updated, 98] }, { kind DELETE fields [1, A_updated, 95] } ] } } transform { RowKindExtractor { plugin_input fake_cdc_data plugin_output transformed_data custom_field_name change_type transform_type FULL } } sink { Console { plugin_input transformed_data } }预期输出I, pk_id1, nameA, score100, change_typeINSERT I, pk_id2, nameB, score100, change_typeINSERT I, pk_id1, nameA, score100, change_typeUPDATE_BEFORE I, pk_id1, nameA_updated, score95, change_typeUPDATE_AFTER I, pk_id2, nameB, score100, change_typeUPDATE_BEFORE I, pk_id2, nameB_updated, score98, change_typeUPDATE_AFTER I, pk_id1, nameA_updated, score95, change_typeDELETE注意示例中 FakeSource 通过rows[].kind字段直接声明每条数据的 RowKind取值与transform_type的 FULL 枚举一致primaryKey配置则是模拟带主键的 CDC 表结构便于后续在数据湖中按主键做回放或去重分析。使用注意事项与最佳实践综合文档与源码使用时请留意以下几点字段名冲突即报错custom_field_name一旦与输入表已有字段重名插件在初始化阶段就会抛出异常因此配置前应确认上游 Schema 中没有同名列SHORT 更省空间FULL 更可读同样的变更语义SHORT 每行仅多出 2 个字符在千万级以上的大表历史留存场景中差异明显需要人工审计或直接对接 BI 分析时建议 FULLAppend-Only 不等于幂等转换后的流里同一主键的 UPDATE_BEFORE / UPDATE_AFTER / DELETE 都变成了独立 INSERT 行下游如果需要最终一致的数据视图需要结合主键做增量合并例如利用 Iceberg 的 upsert 能力或按时间戳取最新行多表场景开箱可用插件支持multi_tables等通用多表选项多库多表 CDC 任务中每个表都会得到结构相同的新增字段建议组合使用RowKindExtractor 与 FilterRowKind按 RowKind 过滤数据行功能互补——前者负责保语义、转追加后者负责按类型裁剪数据流可根据业务需要编排在同一 transform 链中。总结RowKindExtractor 是 SeaTunnel 处理 CDC 数据与 Append-Only 下游对接时的关键转换插件一个参数控制新字段命名一个参数控制格式两层源码逻辑提取原始 RowKind 改写为 INSERT即可完成变更为追加的平滑降级。配合 Iceberg 等数据湖的追加写入它让完整的数据库变更历史得以以低成本、可分析的形式长期沉淀是 CDC 实时入湖、入仓链路中一个简单而实用的组件。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进