SeaTunnel FilterRowKind 转换插件实战:按行类型(RowKind)精准过滤数据

SeaTunnel FilterRowKind 转换插件实战:按行类型(RowKind)精准过滤数据 数据集成ETL大数据批处理流处理变更数据捕获【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/GitHub_Trending/se/seatunnel点击查看免费下载FilterRowKind 是 SeaTunnel 内置的 V2 转换插件用于按数据行的类型RowKind过滤数据典型场景是在 CDC、变更数据捕获或 upsert 数据流中只保留 INSERT、DELETE 等指定类型的行。读完本文你将掌握 RowKind 四种取值的确切含义、include_kinds/exclude_kinds的配置规则与互斥约束并能基于源码理解其过滤实现原理写出可直接运行的作业配置。插件定位FilterRowKind 能做什么在 SeaTunnel 的 Transform 体系中FilterRowKind 属于行过滤类转换插件它不修改任何字段、不改变表结构只根据每一行携带的行类型标记决定保留还是丢弃该行。从源码结构看FilterRowKindTransform继承自FilterRowTransformFilterRowTransform.java后者复用了inputCatalogTable的 schema 与表标识transformTableSchema()与transformTableIdentifier()均为直接 copy因此过滤前后数据集结构完全不变这使它非常适合插入在任意转换链中做前置裁剪例如先过滤掉冗余的UPDATE_BEFORE行再交给下游 Sink 或聚合类转换处理。前置知识RowKind 的四种行类型SeaTunnel 的行类型定义在 RowKind.java共有 4 种取值RowKind短标识字节值含义INSERTI0插入操作UPDATE_BEFORE-U1更新操作的前像旧值用于需要先撤回旧行的非幂等更新UPDATE_AFTERU2更新操作的后像新值也可单独表示幂等更新DELETE-D3删除操作这些取值在 changelog变更日志流中各有含义UPDATE_BEFORE与UPDATE_AFTER通常成对出现用于建模先撤回旧值、再写入新值的更新而基于主键的幂等更新可以只发出UPDATE_AFTER。批处理模式下FakeSource 等普通数据源产生的行类型固定为INSERT。参数说明include_kinds 与 exclude_kindsFilterRowKind 只有两个核心参数且二者互斥只能配置其中一个参数名类型是否必须默认值说明include_kindsarray二选一无要包含的行类型列表仅保留列表中的行类型exclude_kindsarray二选一无要排除的行类型列表丢弃列表中的行类型取值为RowKind枚举名即INSERT、UPDATE_BEFORE、UPDATE_AFTER、DELETE例如transform { FilterRowKind { include_kinds [INSERT, UPDATE_AFTER] } }参数定义位于 FilterRowKinkTransformConfig.java两个选项均声明为ListRowKind类型且无默认值——这意味着什么都不配是不合法的。互斥校验是如何生效的工厂类 FilterRowKindTransformFactory.java 通过OptionRule声明了两层约束exclusive(INCLUDE_KINDS, EXCLUDE_KINDS)两个参数不能同时配置只能二选一Conditions.notEmpty(...)配置的那个参数不允许为空数组。对应的单元测试 FilterRowKindTransformFactoryTest.java 覆盖了全部非法组合两个都不配、两个都配、配置空数组均会抛出OptionValidationException。源码级解析过滤逻辑是如何执行的核心过滤逻辑集中在 FilterRowKindTransform.javaprivate void initConfig(ReadonlyConfig config) { if (config.get(FilterRowKinkTransformConfig.INCLUDE_KINDS) null) { excludeKinds new HashSet(config.get(FilterRowKinkTransformConfig.EXCLUDE_KINDS)); } else { includeKinds new HashSet(config.get(FilterRowKinkTransformConfig.INCLUDE_KINDS)); } } Override protected SeaTunnelRow transformRow(SeaTunnelRow inputRow) { if (!this.excludeKinds.isEmpty()) { return this.excludeKinds.contains(inputRow.getRowKind()) ? null : inputRow; } if (!this.includeKinds.isEmpty()) { return this.includeKinds.contains(inputRow.getRowKind()) ? inputRow : null; } throw new SeaTunnelRuntimeException( CommonErrorCodeDeprecated.UNSUPPORTED_OPERATION, Transform config error! Either excludeKinds or includeKinds must be configured); }几个关键实现细节二选一分支initConfig以include_kinds是否为空为分界把配置归一化为仅 exclude或仅 include两个互斥分支返回 null 即丢弃transformRow返回null表示该行被过滤掉返回原行inputRow表示保留——这是FilterRowTransform体系通用的约定因此该插件不产生任何新行也不改变行内容防御性兜底即使绕过了配置校验例如直接以代码构造插件实例当两个集合都为空时也会抛出SeaTunnelRuntimeException错误信息为 Either excludeKinds or includeKinds must be configured测试 testDirectConstructionWithEmptyKindsFailsOnTransform 专门验证了这一点多表支持工厂创建的是 FieldRowKindMultiCatalogTransform它基于AbstractMultiCatalogMapTransform对每个输入 CatalogTable 分别构建一个FilterRowKindTransform因此该插件同样可用于多表multi-table作业。使用示例示例一排除 INSERT原文档示例FakeSource 生成的数据行类型固定为INSERT。如果使用 FilterRowKind 并排除INSERT那么下游 Sink 将收不到任何行env { job.mode BATCH } source { FakeSource { plugin_output fake row.num 100 schema { fields { id int name string age int } } } } transform { FilterRowKind { plugin_input fake plugin_output fake1 exclude_kinds [INSERT] } } sink { Console { plugin_input fake1 } }该示例演示了插件最直接的用法100 行INSERT数据全部被过滤Console Sink 无输出。示例二仅保留指定行类型include批处理 CDC 混合场景中若只想保留新增数据使用include_kindstransform { FilterRowKind { plugin_input cdc_source plugin_output only_insert include_kinds [INSERT] } }示例三CDC 变更流中剔除更新前像在 MySQL CDC 同步场景中UPDATE事件通常产生UPDATE_BEFOREUPDATE_AFTER两行。如果下游目标表采用 upsert 语义UPDATE_BEFORE是冗余的可以这样裁剪transform { FilterRowKind { include_kinds [INSERT, UPDATE_AFTER, DELETE] } }这样既保留了完整的增删语义又去掉了仅用于撤回的旧值行减少写入量。示例四多表作业多表场景下FilterRowKind 会对每个匹配到的表独立应用相同过滤规则可配合multi_tables、table_match_regex等通用参数使用具体多表配置方式可参考 Transform 多表支持。常见问题1. 两个参数都配置会怎样配置校验直接失败。OptionRule.exclusive禁止同时出现include_kinds与exclude_kinds作业启动阶段即抛出OptionValidationException错误信息中会同时提到这两个参数。2. 一个都不配置会怎样两个参数均无默认值作业校验阶段报错若绕过校验直接构造插件运行时也会抛出 Either excludeKinds or includeKinds must be configured。3. 配置空数组可以吗不可以。Conditions.notEmpty约束要求数组非空空数组同样无法通过校验。4. 大小写敏感吗取值是 Java 枚举RowKind的枚举名必须严格使用INSERT、UPDATE_BEFORE、UPDATE_AFTER、DELETE这四种写法。5. 过滤后字段结构会变吗不会。FilterRowTransform对 schema 与表标识只做复制不做任何修改下游无需调整字段映射。与其他插件的配合RowKindExtractor如果需要把行类型作为一列写入目标而非直接过滤可参考 rowkind-extractor它能将RowKind提取为普通字段与 FilterRowKind 形成提取与过滤的分工Transform 通用参数plugin_input/plugin_output等编排参数旧名source_table_name/result_table_name已废弃的详细说明见 Transform 通用参数Transform 体系总览更多转换插件与整体设计见 Transforms 目录 与 Transform 插件体系。小结FilterRowKind 是 SeaTunnel 中一个轻量但非常实用的行级过滤插件它基于RowKind的四种取值INSERT/UPDATE_BEFORE/UPDATE_AFTER/DELETE通过互斥的include_kinds或exclude_kinds参数实现白名单保留或黑名单剔除且不改变表结构与行内容。无论是 CDC 变更流裁剪、批处理数据筛选还是多表作业的统一过滤都可以在 Transform 阶段一行配置完成。结合 实现源码 与 单元测试 阅读可以更透彻地理解其校验规则与过滤语义。赞分享数据集成ETL大数据批处理流处理变更数据捕获【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/GitHub_Trending/se/seatunnel点击查看免费下载相关推荐SeaTunnel FilterRowKind 转换插件详解按 RowKindINSERT/UPDATE/DELETE精准过滤数据行SeaTunnel FilterRowKind 转换插件详解按 RowKindINSERT/UPDATE/DELETE精准过滤数据行 本文围绕 SeaTu数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Metadata 转换插件指南把库名、表名、RowKind 等行级元数据提取为普通字段SeaTunnel Metadata 转换插件指南把库名、表名、RowKind 等行级元数据提取为普通字段 元数据提取Metadata是 SeaTunne数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Calcite Transform 插件实战用标准 SQL 逐行转换数据与向量运算SeaTunnel Calcite Transform 插件实战用标准 SQL 逐行转换数据与向量运算 本篇技术指南围绕 SeaTunnel 的 Calcit数据集成ETL大数据批处理流处理变更数据捕获创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考