如何在Flink中仅查看及保留upsert-kafka表的UA(UPDATE_AFTER, +U)数据至Kafka主题
解决Flink Upsert-Kafka仅保留UPDATE_AFTER数据的方案
要实现只读取并存储Upsert-Kafka主题中的UPDATE_AFTER(+U)数据,核心是利用Flink Upsert-Kafka连接器暴露的op元数据字段进行过滤,以下是具体实现方法:
1. 定义包含元数据字段的源表
在建源表时,必须显式声明op虚拟列,该字段会自动获取Kafka消息对应的操作类型,取值包括INSERT、UPDATE_BEFORE、UPDATE_AFTER、DELETE:
CREATE TABLE source_upsert_kafka ( id STRING PRIMARY KEY NOT ENFORCED, name STRING, age INT, -- 引入操作类型元数据字段 `op` STRING METADATA FROM 'value.op' VIRTUAL ) WITH ( 'connector' = 'upsert-kafka', 'topic' = 'your_source_topic', 'properties.bootstrap.servers' = 'kafka_host:9092', 'key.format' = 'json', 'value.format' = 'json' );
2. 过滤仅保留UPDATE_AFTER数据
在查询时通过WHERE子句筛选op = 'UPDATE_AFTER'的记录,即可只获取UA数据:
SELECT id, name, age FROM source_upsert_kafka WHERE `op` = 'UPDATE_AFTER';
3. 将过滤后的数据写入目标Kafka主题
根据目标主题的需求,选择普通Kafka连接器或Upsert-Kafka连接器作为sink:
方案A:写入普通Kafka主题
如果目标不需要upsert语义,直接用普通Kafka连接器:
CREATE TABLE target_kafka ( id STRING, name STRING, age INT ) WITH ( 'connector' = 'kafka', 'topic' = 'your_target_topic', 'properties.bootstrap.servers' = 'kafka_host:9092', 'key.format' = 'json', 'value.format' = 'json', 'sink.partitioner' = 'round-robin' ); -- 插入过滤后的UA数据 INSERT INTO target_kafka SELECT id, name, age FROM source_upsert_kafka WHERE `op` = 'UPDATE_AFTER';
方案B:写入Upsert-Kafka主题
如果目标需要保持upsert语义,继续用Upsert-Kafka连接器:
CREATE TABLE target_upsert_kafka ( id STRING PRIMARY KEY NOT ENFORCED, name STRING, age INT ) WITH ( 'connector' = 'upsert-kafka', 'topic' = 'your_target_upsert_topic', 'properties.bootstrap.servers' = 'kafka_host:9092', 'key.format' = 'json', 'value.format' = 'json' ); -- 插入过滤后的UA数据 INSERT INTO target_upsert_kafka SELECT id, name, age FROM source_upsert_kafka WHERE `op` = 'UPDATE_AFTER';
4. DataStream API方式(可选)
如果用DataStream API开发,可在流处理阶段直接过滤:
// 先将表转换为DataStream DataStream<RowData> sourceStream = tableEnv.toDataStream(sourceTable); // 过滤仅保留UPDATE_AFTER记录(假设op是最后一列) DataStream<RowData> filteredStream = sourceStream.filter(row -> { String op = row.getString(row.getArity() - 1); return "UPDATE_AFTER".equals(op); }); // 将过滤后的流转换回表,再写入目标Kafka Table filteredTable = tableEnv.fromDataStream(filteredStream); tableEnv.executeSql("INSERT INTO target_kafka SELECT * FROM " + filteredTable);
内容的提问来源于stack exchange,提问作者olivia
相关产品推荐
相关产品推荐

