[ 
https://issues.apache.org/jira/browse/FLINK-40301?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Qilong Wang updated FLINK-40301:
--------------------------------
    Description: 13233  (was: 
`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:

```sql
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;
```

Before the fix, the generated plan contains two adjacent Calc nodes on the 
output side of `ChangelogNormalize`:

```text
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(...)
```

`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:

```text
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(...)
```

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.)

> 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
>             Fix For: 2.4.0
>
>
> 13233



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to