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

如何在Flink SQL两张表的Union结果中实现高效去重?

问题分析与解决方案

你遇到的是Flink SQL优化器在Union多表后,无法将row_number()去重逻辑转换为高效Deduplicate算子的问题。核心原因是Union操作打破了优化器对单表流特性的识别,再加上proctime在多表场景下的节点独立性,导致优化器退而采用通用的GroupAggregate + Rank执行路径。

下面是几个可落地的解决方案:

方案1:先单表去重,再Union

既然单表去重能触发Deduplicate优化,那就先对Kafka和JDBC表分别做去重,再将结果Union。这样每个子查询都会生成高效的去重算子,整体性能最优。

修改后的SQL:

-- Kafka表单表去重
create view kafka_deduplicated as
select * from (
    select
        *,
        row_number() over (partition by recordId order by proctime asc) as row_num
    from `kafka_catalog`.`default_database`.`dummy`
) where row_num = 1;

-- JDBC表单表去重
create view jdbc_deduplicated as
select * from (
    select
        *,
        row_number() over (partition by recordId order by proctime asc) as row_num
    from `jdbc_catalog`.`default_database`.`dummy`
) where row_num = 1;

-- Union去重后的结果
create view dummy_full_history as
select * from kafka_deduplicated
union
select * from jdbc_deduplicated;

方案2:改用事件时间(recordTime)作为排序键

如果recordTime是全局递增且能唯一标识记录顺序的字段(比如数据生成时间),可以把排序字段从proctime换成recordTime。优化器对基于事件时间的去重逻辑识别度更高,因为事件时间是数据本身的属性,不会因不同表的处理节点差异产生混乱。

修改后的去重视图:

create view dummy_full_history as
select * from (
    select
        *,
        row_number() over (partition by recordId order by recordTime asc) as row_num
    from dummy_union
) where row_num = 1;

方案3:升级Flink版本

低版本Flink(如1.14及以下)对Union后的流去重优化支持有限,升级到1.15+版本后,优化器对多源流的去重逻辑识别能力有明显提升,大概率能自动将row_number()转换为Deduplicate算子。

Flink 1.13及以上支持SELECT DISTINCT ON语法,这是更简洁的去重写法,优化器对这种语法的识别优先级更高,更容易触发Deduplicate优化:

create view dummy_full_history as
select distinct on (recordId) *
from dummy_union
order by recordId, proctime asc;

额外注意点

  • 确保两个表的recordId是全局唯一的去重键,否则去重逻辑本身存在漏洞。
  • JDBC表是批处理源,Kafka是流处理源,Union后属于流批混合的流表,优化器处理会更谨慎,尽量给优化器提供明确的语义提示(比如先单表去重)。

内容的提问来源于stack exchange,提问作者Grzes

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 04:31:01