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
相关产品推荐
相关产品推荐

