[
https://issues.apache.org/jira/browse/FLINK-40301?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
ASF GitHub Bot updated FLINK-40301:
-----------------------------------
Labels: pull-request-available (was: )
> Adjacent Calc nodes created by watermark transposition prevent projection
> pushdown through ChangelogNormalize
> -------------------------------------------------------------------------------------------------------------
>
> Key: FLINK-40301
> URL: https://issues.apache.org/jira/browse/FLINK-40301
> Project: Flink
> Issue Type: Bug
> Components: Table SQL / Planner
> Affects Versions: 2.4.0
> Reporter: Qilong Wang
> Priority: Major
> Labels: pull-request-available
> Fix For: 2.4.0
>
>
> `WatermarkAssignerChangelogNormalizeTransposeRule` may create adjacent
> `StreamPhysicalCalc` nodes on the output side of `ChangelogNormalize`.
> Because these newly created Calc nodes are not merged before
> `PushCalcPastChangelogNormalizeRule` runs, fields that are not required by
> `ChangelogNormalize` cannot be projected out from its input. This
> unnecessarily increases the shuffle row width and the state size of
> `ChangelogNormalize`.
> For example:
> {code:java}
> CREATE TABLE t2 (
> k VARBINARY,
> ingestion_time TIMESTAMP(3) METADATA FROM 'ts',
> a VARCHAR NOT NULL,
> f BOOLEAN NOT NULL,
> WATERMARK FOR ingestion_time AS ingestion_time,
> PRIMARY KEY (a) NOT ENFORCED
> ) WITH (
> 'connector' = 'values',
> 'readable-metadata' = 'ts:TIMESTAMP(3)',
> 'changelog-mode' = 'I,UA,D'
> );
> SELECT f, COUNT
> FROM (
> SELECT ingestion_time, f, a
> FROM t2
> )
> GROUP BY f; {code}
> Before the fix, the generated plan contains two adjacent Calc nodes on the
> output side of `ChangelogNormalize`:
> {code:java}
> GroupAggregate(groupBy=[f])
> +- Exchange(distribution=[hash[f]])
> +- Calc(select=[f])
> +- Calc(select=[f, ingestion_time])
> +- ChangelogNormalize(key=[a])
> +- Exchange(distribution=[hash[a]])
> +- WatermarkAssigner(rowtime=[ingestion_time])
> +- TableSourceScan(...) {code}
> `PushCalcPastChangelogNormalizeRule` only sees the immediate Calc with fields
> `[f, ingestion_time]`. After adding the unique key `a` required by
> `ChangelogNormalize`, all input fields are still referenced, so no projection
> can be pushed to the input side of `ChangelogNormalize`.
> The expected plan is:
> {code:java}
> GroupAggregate(groupBy=[f])
> +- Exchange(distribution=[hash[f]])
> +- Calc(select=[f])
> +- ChangelogNormalize(key=[a])
> +- Exchange(distribution=[hash[a]])
> +- Calc(select=[f, a])
> +- WatermarkAssigner(rowtime=[ingestion_time])
> +- TableSourceScan(...) {code}
> Merging the adjacent Calc nodes before running
> `PushCalcPastChangelogNormalizeRule` allows the rule to push `[f, a]` to the
> input side of `ChangelogNormalize`.
> The `ingestion_time` field remains available to `WatermarkAssigner`, but it
> no longer enters the shuffle or the state of `ChangelogNormalize`.
> `ChangelogNormalize` is still required in this example because the downstream
> aggregation requires `UPDATE_BEFORE` messages.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)