Skip to content

[FLINK-40528][table-planner] Don't evaluate non-key expressions for partial delete StreamExecCalc - #29106

Open
twalthr wants to merge 1 commit into
apache:masterfrom
twalthr:FLINK-40528
Open

[FLINK-40528][table-planner] Don't evaluate non-key expressions for partial delete StreamExecCalc#29106
twalthr wants to merge 1 commit into
apache:masterfrom
twalthr:FLINK-40528

Conversation

@twalthr

@twalthr twalthr commented Sep 4, 2026

Copy link
Copy Markdown
Contributor

What is the purpose of the change

FLINK-37475 (#26306) lets the planner drop
ChangelogNormalize when a source only produces partial "delete-by-key" tombstones (key columns
set, all other columns null/absent) and the sink accepts the
same, avoiding the cost of materializing full row state for upsert-style (e.g. upsert-Kafka)
pipelines.

The planner's FlinkChangelogModeInferenceProgram decided whether a Calc could keep forwarding
DELETE_BY_KEY purely by checking whether its filter condition only touches key columns — it
never checked whether the projection does. As a result, a Calc with no filter (or a
key-safe one) was always allowed to forward DELETE_BY_KEY unchanged, even though its projection
could evaluate expressions over non-key columns that are legitimately null/absent in a
delete-by-key row. This could throw at runtime (e.g. constructing ROW(id, arr) where arr is a
NOT NULL ARRAY column that is actually null for the tombstone) or silently produce wrong
output (e.g. value + 2 evaluated on a null value).

This PR fixes the runtime side for the common case (StreamExecCalc, a plain SQL projection):
when a -D row enters a Calc that is forwarding DELETE_BY_KEY, only the output columns that
are a trivial pass-through of one of the Calc's own upsert keys are evaluated; every other
column becomes a typed NULL instead of evaluating its (potentially unsafe) expression.

Python calc and (regular) async calc share the exact same planner-level treatment as plain
Calc (they extend the same StreamPhysicalCalcBase), but their exec nodes
(CommonExecPythonCalc/CommonExecAsyncCalc) invoke an external/remote function per row and do
not get this runtime protection. For these, the planner is instead made conservative: it now
never forwards DELETE_BY_KEY through any Calc variant other than the plain one, forcing a
ChangelogNormalize to stay upstream so a remote/external call is never made with a partial,
possibly-null row.

Brief change log

  • StreamExecCalc/CommonExecCalc/BatchExecCalc gain an optional int[] partialDeleteKeys
    field (only ever set on StreamExecCalc, persisted in the JSON plan only when non-null),
    naming the output column indices that are safe to evaluate on a delete-by-key row.
  • StreamPhysicalCalc computes partialDeleteKeys from the Calc's own resolved
    DeleteKindTrait and FlinkRelMetadataQuery.getUpsertKeys, taking the union of all
    candidate upsert-key sets (a Calc's output can have more than one, e.g. a duplicated/aliased
    key column) so a real key column is never dropped in favor of an unrelated one.
  • CalcCodeGenerator generates a runtime branch: for a -D row, only the key-derived output
    columns are evaluated (via GenerateUtils.generateNullLiteral. Each branch is generated inside its own pushed local-ref scope
    so a cached sub-expression (e.g. the BinaryRowWriter code backing a ROW(...) constructor)
    is never hoisted into the unconditional/bottom scope and run for every row regardless of
    branch.
  • FlinkChangelogModeInferenceProgram: only the plain StreamPhysicalCalc keeps the existing
    "forward DELETE_BY_KEY if the filter is key-safe" check; every other StreamPhysicalCalcBase
    subclass (Python calc, async calc, and any future subclass) now always requires FULL_DELETE.
  • Added/extended TableTestPrograms in DeletesByKeyPrograms/DeletesByKeySemanticTests
    (non-key row-constructor expression, a key-safe filter combined with it, a duplicated/cast
    key column, and an async calc regression) and CalcTestPrograms/CalcRestoreTest (a
    compiled-plan restore test asserting partialDeleteKeys round-trips through the plan JSON).

Verifying this change

This change added tests and can be verified as follows:

  • DeletesByKeySemanticTests: covers
    plain delete-by-key passthrough, a full-delete fallback, a non-key projection expression, the
    same combined with a key-safe filter, a NOT NULL ARRAY wrapped in a ROW(...) constructor,
    a duplicated/cast key column producing multiple upsert-key candidates, and an async calc
    (verifying the planner now forces a full, materialized delete instead of letting the
    delete-by-key tombstone reach the remote function).
  • CalcRestoreTest#calc-partial-delete-with-expression-and-filter: a new compiled-plan restore
    test where the delete-by-key row is only produced after restoring from the compiled plan
    JSON, verifying partialDeleteKeys round-trips correctly through plan serialization.
  • Full regression pass across CalcITCase, CalcRestoreTest, CalcBatchRestoreTest,
    AsyncCalcRestoreTest, PythonCalcJsonPlanTest, PythonAsyncCalcJsonPlanTest,
    ChangelogNormalizeOptimizationTest, FlinkCalcMergeRuleTest,
    PushCalcPastChangelogNormalizeRuleTest, and CalcMergeTest — confirms no existing golden
    plan or behavior changed for any Calc node that isn't in a delete-by-key pipeline.
  • Manually verified each new/changed check is meaningful by temporarily reverting the
    corresponding fix and confirming the associated test fails (a crash or wrong/partial output
    without the fix, the expected result with it).

Does this pull request potentially affect one of the following parts:

  • Dependencies (does it add or upgrade a dependency): no
  • The public API, i.e., is any changed class annotated with @Public(Evolving): no (all
    changed classes are internal planner/exec-node classes)
  • The serializers: no
  • The runtime per-record code paths (performance sensitive): yes — CalcCodeGenerator is on
    the hot path for every Calc; the change adds one Option match per generated operator with
    no extra generated code when partialDeleteKeys is unset (the overwhelming majority of
    Calc nodes), and only branches at runtime for Calcs that are already in a delete-by-key
    pipeline.
  • Anything that affects deployment or recovery: JobManager (and its components), Checkpointing,
    Kubernetes/Yarn, ZooKeeper: no
  • The S3 file system connector: no

Documentation

  • Does this pull request introduce a new feature? no
  • If yes, how is the feature documented? not applicable

Was generative AI tooling used to co-author this PR?
  • Yes (please specify the tool below)

Generated-by: Claude Sonnet 5

@flinkbot

flinkbot commented Sep 4, 2026

Copy link
Copy Markdown
Collaborator

CI report:

Bot commands The @flinkbot bot supports the following commands:
  • @flinkbot run azure re-run the last Azure build

},
"abilities" : [ {
"type" : "FilterPushDown",
"predicates" : [ ]

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

probably should be handled separately

} ],
"type" : "BOOLEAN NOT NULL"
},
"partialDeleteKeys" : [ 0 ],

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

will it be in plan in case of empty?

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants