- 数据集成
- ETL
- 大数据
- 批处理
- 流处理
- 变更数据捕获
【免费下载链接】seatunnel
SeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.
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 | +I | 0 | 插入操作 |
UPDATE_BEFORE | -U | 1 | 更新操作的前像(旧值),用于需要先撤回旧行的非幂等更新 |
UPDATE_AFTER | +U | 2 | 更新操作的后像(新值),也可单独表示幂等更新 |
DELETE | -D | 3 | 删除操作 |
这些取值在 changelog(变更日志)流中各有含义:UPDATE_BEFORE与UPDATE_AFTER通常成对出现,用于建模"先撤回旧值、再写入新值"的更新;而基于主键的幂等更新可以只发出UPDATE_AFTER。批处理模式下,FakeSource 等普通数据源产生的行类型固定为INSERT。
参数说明:include_kinds 与 exclude_kinds
FilterRowKind 只有两个核心参数,且二者互斥,只能配置其中一个:
| 参数名 | 类型 | 是否必须 | 默认值 | 说明 |
|---|---|---|---|---|
include_kinds | array | 二选一 | 无 | 要包含的行类型列表,仅保留列表中的行类型 |
exclude_kinds | array | 二选一 | 无 | 要排除的行类型列表,丢弃列表中的行类型 |
取值为RowKind枚举名,即INSERT、UPDATE_BEFORE、UPDATE_AFTER、DELETE,例如:
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"); }几个关键实现细节:
- 二选一分支:
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_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_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
- 大数据
- 批处理
- 流处理
- 变更数据捕获
【免费下载链接】seatunnel
SeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.
相关推荐
SeaTunnel FilterRowKind 转换插件详解:按 RowKind(INSERT/UPDATE/DELETE)精准过滤数据行
SeaTunnel FilterRowKind 转换插件详解:按 RowKind(INSERT/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),仅供参考