You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Flink中仅查看及保留upsert-kafka表的UA(UPDATE_AFTER, +U)数据至Kafka主题

要实现只读取并存储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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.18 04:00:39