SeaTunnel FilterRowKind 转换插件实战:按行类型(RowKind)精准过滤数据
2026/9/20 4:10:38 网站建设 项目流程
  • 数据集成
  • ETL
  • 大数据
  • 批处理
  • 流处理
  • 变更数据捕获

【免费下载链接】seatunnel

SeaTunnel 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继承自FilterRowTransform(FilterRowTransform.java),后者复用了inputCatalogTable的 schema 与表标识(transformTableSchema()transformTableIdentifier()均为直接 copy),因此过滤前后数据集结构完全不变,这使它非常适合插入在任意转换链中做"前置裁剪",例如先过滤掉冗余的UPDATE_BEFORE行,再交给下游 Sink 或聚合类转换处理。

前置知识:RowKind 的四种行类型

SeaTunnel 的行类型定义在 RowKind.java,共有 4 种取值:

RowKind短标识字节值含义
INSERT+I0插入操作
UPDATE_BEFORE-U1更新操作的前像(旧值),用于需要先撤回旧行的非幂等更新
UPDATE_AFTER+U2更新操作的后像(新值),也可单独表示幂等更新
DELETE-D3删除操作

这些取值在 changelog(变更日志)流中各有含义:UPDATE_BEFOREUPDATE_AFTER通常成对出现,用于建模"先撤回旧值、再写入新值"的更新;而基于主键的幂等更新可以只发出UPDATE_AFTER。批处理模式下,FakeSource 等普通数据源产生的行类型固定为INSERT

参数说明:include_kinds 与 exclude_kinds

FilterRowKind 只有两个核心参数,且二者互斥,只能配置其中一个:

参数名类型是否必须默认值说明
include_kindsarray二选一要包含的行类型列表,仅保留列表中的行类型
exclude_kindsarray二选一要排除的行类型列表,丢弃列表中的行类型

取值为RowKind枚举名,即INSERTUPDATE_BEFOREUPDATE_AFTERDELETE,例如:

transform { FilterRowKind { include_kinds = ["INSERT", "UPDATE_AFTER"] } }

参数定义位于 FilterRowKinkTransformConfig.java,两个选项均声明为List<RowKind>类型且无默认值——这意味着"什么都不配"是不合法的。

互斥校验是如何生效的

工厂类 FilterRowKindTransformFactory.java 通过OptionRule声明了两层约束:

  • exclusive(INCLUDE_KINDS, EXCLUDE_KINDS):两个参数不能同时配置,只能二选一;
  • Conditions.notEmpty(...):配置的那个参数不允许为空数组

对应的单元测试 FilterRowKindTransformFactoryTest.java 覆盖了全部非法组合:两个都不配、两个都配、配置空数组,均会抛出OptionValidationException

源码级解析:过滤逻辑是如何执行的

核心过滤逻辑集中在 FilterRowKindTransform.java:

private 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"); }

几个关键实现细节:

  1. 二选一分支initConfiginclude_kinds是否为空为分界,把配置归一化为"仅 exclude"或"仅 include"两个互斥分支;
  2. 返回 null 即丢弃transformRow返回null表示该行被过滤掉,返回原行inputRow表示保留——这是FilterRowTransform体系通用的约定,因此该插件不产生任何新行,也不改变行内容
  3. 防御性兜底:即使绕过了配置校验(例如直接以代码构造插件实例),当两个集合都为空时也会抛出SeaTunnelRuntimeException,错误信息为 "Either excludeKinds or includeKinds must be configured",测试 testDirectConstructionWithEmptyKindsFailsOnTransform 专门验证了这一点;
  4. 多表支持:工厂创建的是 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_kinds

transform { FilterRowKind { plugin_input = "cdc_source" plugin_output = "only_insert" include_kinds = ["INSERT"] } }

示例三:CDC 变更流中剔除更新前像

在 MySQL CDC 同步场景中,UPDATE事件通常产生UPDATE_BEFORE+UPDATE_AFTER两行。如果下游目标表采用 upsert 语义,UPDATE_BEFORE是冗余的,可以这样裁剪:

transform { FilterRowKind { include_kinds = ["INSERT", "UPDATE_AFTER", "DELETE"] } }

这样既保留了完整的增删语义,又去掉了仅用于"撤回"的旧值行,减少写入量。

示例四:多表作业

多表场景下,FilterRowKind 会对每个匹配到的表独立应用相同过滤规则(可配合multi_tablestable_match_regex等通用参数使用),具体多表配置方式可参考 Transform 多表支持。

常见问题

1. 两个参数都配置会怎样?配置校验直接失败。OptionRule.exclusive禁止同时出现include_kindsexclude_kinds,作业启动阶段即抛出OptionValidationException,错误信息中会同时提到这两个参数。

2. 一个都不配置会怎样?两个参数均无默认值,作业校验阶段报错;若绕过校验直接构造插件,运行时也会抛出 "Either excludeKinds or includeKinds must be configured"。

3. 配置空数组可以吗?不可以。Conditions.notEmpty约束要求数组非空,空数组同样无法通过校验。

4. 大小写敏感吗?取值是 Java 枚举RowKind的枚举名,必须严格使用INSERTUPDATE_BEFOREUPDATE_AFTERDELETE这四种写法。

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_kindsexclude_kinds参数实现"白名单保留"或"黑名单剔除",且不改变表结构与行内容。无论是 CDC 变更流裁剪、批处理数据筛选,还是多表作业的统一过滤,都可以在 Transform 阶段一行配置完成。结合 实现源码 与 单元测试 阅读,可以更透彻地理解其校验规则与过滤语义。

  • 数据集成
  • ETL
  • 大数据
  • 批处理
  • 流处理
  • 变更数据捕获

【免费下载链接】seatunnel

SeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.

项目地址:https://gitcode.com/GitHub_Trending/se/seatunnel
点击查看免费下载

相关推荐

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询