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

Flink 1.17.1中StreamPhysicalOverAggregate报错的解决方法咨询

问题原因

你遇到的错误是因为Flink优化器将ROW_NUMBER() OVER (...) WHERE row_num=1的逻辑转换为Deduplicate算子,但该算子输出的是包含更新/删除的变更流,而StreamPhysicalOverAggregate算子不支持消费这类变更流,本质是Flink对这种去重场景的执行计划优化与流处理的变更语义不兼容。

解决方案

针对Flink 1.17.1,提供两种可行的修改方案:

方案一:改用Flink原生去重语法(推荐)

Flink支持DISTINCT ON语法直接实现按主键取最新数据的逻辑,避免生成冲突的算子计划:

CREATE OR REPLACE VIEW identity_entitlement_source_keyed AS (
    SELECT DISTINCT ON (id) *
    FROM identity_entitlement_poc
    ORDER BY id, eventTime DESC NULLS LAST
);

方案二:用聚合函数替代窗口函数

通过MAX...KEEP DENSE_RANK的聚合语法获取每个id的最新记录,这种写法会生成聚合算子而非OverAggregate,兼容变更流处理:

CREATE OR REPLACE VIEW identity_entitlement_source_keyed AS (
    SELECT
        id,
        MAX(entitlementAssignments) KEEP (DENSE_RANK LAST ORDER BY eventTime) AS entitlementAssignments
    FROM identity_entitlement_poc
    GROUP BY id
);

额外注意事项

如果业务需要保留eventTime字段,可以调整写法干扰优化器,避免生成Deduplicate节点:

CREATE OR REPLACE VIEW identity_entitlement_source_keyed AS (
    SELECT
        id,
        entitlementAssignments,
        eventTime
    FROM (
        SELECT
            id,
            entitlementAssignments,
            eventTime,
            ROW_NUMBER() OVER (PARTITION BY id ORDER BY eventTime DESC NULLS LAST) AS row_num,
            COUNT(*) OVER (PARTITION BY id) AS cnt
        FROM identity_entitlement_poc
    ) t
    WHERE row_num = 1
);

这种写法性能略低于前两种方案,仅在需要保留全量字段时使用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 10:16:26