From 75b596dcdc7bf49292c3cc0826eac49e2d1cf12d Mon Sep 17 00:00:00 2001 From: Timo Walther Date: Tue, 15 Sep 2026 15:21:04 +0200 Subject: [PATCH] [FLINK-40528][table-planner] Address additional feedback for partial-delete StreamExecCalc --- .../physical/stream/StreamPhysicalCalc.scala | 10 +++++++--- .../FlinkChangelogModeInferenceProgram.scala | 15 ++++++++++++++- .../nodes/exec/stream/DeletesByKeyPrograms.java | 2 +- 3 files changed, 22 insertions(+), 5 deletions(-) diff --git a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/nodes/physical/stream/StreamPhysicalCalc.scala b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/nodes/physical/stream/StreamPhysicalCalc.scala index a13938bf78618d..79060d396af298 100644 --- a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/nodes/physical/stream/StreamPhysicalCalc.scala +++ b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/nodes/physical/stream/StreamPhysicalCalc.scala @@ -87,9 +87,13 @@ class StreamPhysicalCalc( } // Every column in every candidate is, by construction of - // FlinkRelMdUniqueKeys.getProjectUniqueKeys, guaranteed to be a trivial - // pass-through of an input field - never a risky expression to evaluate. - val keyIndices = outputUpsertKeys.flatMap(bitSet => bitSet.map(_.intValue())).toSet.toArray + // FlinkRelMdUniqueKeys.getProjectUniqueKeys, guaranteed to be a pass-through + // of an input field or an injective expression. + val keyIndices = outputUpsertKeys + .flatMap(bitSet => bitSet.map(_.intValue())) + .toSet + .toArray + .sorted if (keyIndices.nonEmpty) { keyIndices } else { diff --git a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/optimize/program/FlinkChangelogModeInferenceProgram.scala b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/optimize/program/FlinkChangelogModeInferenceProgram.scala index 300691a6982a69..e72d00e1aab01f 100644 --- a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/optimize/program/FlinkChangelogModeInferenceProgram.scala +++ b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/optimize/program/FlinkChangelogModeInferenceProgram.scala @@ -1436,7 +1436,7 @@ class FlinkChangelogModeInferenceProgram extends FlinkOptimizeProgram[StreamOpti case calc: StreamPhysicalCalc => if ( requiredTrait == DeleteKindTrait.DELETE_BY_KEY && - isNonUpsertKeyCondition(calc) + (isNonUpsertKeyCondition(calc) || !hasOutputUpsertKey(calc)) ) { None } else { @@ -1697,6 +1697,19 @@ class FlinkChangelogModeInferenceProgram extends FlinkOptimizeProgram[StreamOpti } } + /** + * Whether this calc's own output still has an upsert key after its projection. A DELETE_BY_KEY + * tombstone passed through this calc must still carry a key in its output, otherwise nothing + * downstream would know what to delete. + * + * This method is an extra safety net, in case downstream consumers don't require an upsert key. + */ + private def hasOutputUpsertKey(calc: StreamPhysicalCalcBase): Boolean = { + val fmq = FlinkRelMetadataQuery.reuseOrCreate(calc.getCluster.getMetadataQuery) + val upsertKeys = fmq.getUpsertKeys(calc) + upsertKeys != null && upsertKeys.exists(!_.isEmpty) + } + private def isNonUpsertKeyCondition(calc: StreamPhysicalCalcBase): Boolean = { val program = calc.getProgram if (program.getCondition == null) { diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/DeletesByKeyPrograms.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/DeletesByKeyPrograms.java index c230b2fbfb17f2..d9b9b334437e4e 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/DeletesByKeyPrograms.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/DeletesByKeyPrograms.java @@ -300,7 +300,7 @@ public final class DeletesByKeyPrograms { TableTestProgram.of( "delete-by-key-delete-by-key-with-expression", "NOT NULL constraints have no effect. The row constructor expression" - + "is not evaluated for partial deletion.") + + " is not evaluated for partial deletion.") .setupTableSource( SourceTestStep.newBuilder("source_t") .addSchema(