更新时间:2026-06-29 GMT+08:00
FlinkSQL支持RowKind过滤能力
使用场景
上游算子往下游发送数据,需要根据RowKind过滤符合要求的数据时,可以使用该特性。
使用限制
- 仅支持FlinkSQL Hint方式使用该特性。
- 使用该特性后作业Sink Connector需要支持Changelog消息,如:hudi、upsert-kafka、print等。
- 本章节仅适用于MRS 3.6.0-LTS及之后版本。
使用方法
在FlinkSQL中添加 /*+ ROWKIND('INSERT','UPDATE_AFTER','UPDATE_BEFORE','DELETE') */ ,括号中指定需要保留的RowKind类型,对应的数据就可以发送到下游。根据业务需求选择需要保留的RowKind类型,未列出的RowKind类型数据将被过滤。
RowKind类型:
- INSERT:插入数据,表示新增的记录。
- UPDATE_AFTER:更新后数据,表示更新后的新记录。
- UPDATE_BEFORE:更新前数据,表示更新前的旧记录。
- DELETE:删除数据,表示被删除的记录。
SQL示例:
create table datagen1(pid int, uid int) with('connector' = 'datagen');
create table datagen2(pid int, uid int) with('connector' = 'datagen');
create table printsink(pid int, uid int) with('connector' = 'print');
insert into
printsink
SELECT
/*+ ROWKIND('INSERT','UPDATE_AFTER') */
datagen1.pid,
datagen2.uid
FROM
datagen1
left join datagen2 on datagen1.pid = datagen2.pid; 父主题: Flink企业级能力增强