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

Yanquan Lv updated FLINK-40739:
-------------------------------
    Description: 
h2. Title

{{[Bug][pipeline-connector][maxcompute] Partition values extracted from wrong 
column index when partition columns are not contiguous at the end of schema}}
h2. Description
h3. Problem

In {{{}SessionManageOperator#extractPartition{}}}, the MaxCompute pipeline sink 
computes the index of a partition column by assuming that all partition columns 
are contiguous and located at the end of the schema:

[https://github.com/apache/flink-cdc/blob/master/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-maxcompute/src/main/java/org/apache/flink/cdc/connectors/maxcompute/coordinator/SessionManageOperator.java#L268]
{code:java}
for (int i = 0; i < partitionKeyCount; i++) {
    RecordData.FieldGetter fieldGetter =
            fieldGetters.get(columnCount - partitionKeyCount - 1 + i);
    ...
} {code}
The {{-1}} offset is incorrect. When partition columns are *not* contiguously 
placed at the end of the column list, the sink reads values from the wrong 
field getters and writes data into the wrong MaxCompute partitions.
h3. Example

Consider a table whose columns are ordered as:
||Index||Column||Is Partition Key||
|0|id| |
|1|pt|yes|
|2|name| |
|3|dt|yes|

Here {{columnCount = 4}} and {{{}partitionKeyCount = 2{}}}.

With the current code:
 * For {{{}i = 0{}}}: index = {{4 - 2 - 1 + 0 = 1}} → reads {{pt}} (correct by 
coincidence)
 * For {{{}i = 1{}}}: index = {{4 - 2 - 1 + 1 = 2}} → reads {{name}} (wrong! 
should read {{dt}} at index 3)

As a result, the partition spec will contain {{{}pt=<correct>, dt=<value of 
name>{}}}, which silently corrupts data.
h3. Impact
 * Data is written into incorrect partitions when the source schema has 
partition columns interleaved with regular columns.
 * This can happen after schema evolution (e.g., {{{}ADD COLUMN{}}}) or when 
the user explicitly defines columns in a non-trailing order.

h3. Expected Behavior

The sink should resolve each partition column by its name instead of by a fixed 
trailing offset, so that it works regardless of column ordering.

  was:
h2. Title

{{[Bug][pipeline-connector][maxcompute] Partition values extracted from wrong 
column index when partition columns are not contiguous at the end of schema}}
h2. Description
h3. Problem

In {{{}SessionManageOperator#extractPartition{}}}, the MaxCompute pipeline sink 
computes the index of a partition column by assuming that all partition columns 
are contiguous and located at the end of the schema:

[https://github.com/apache/flink-cdc/blob/master/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-maxcompute/src/main/java/org/apache/flink/cdc/connectors/maxcompute/coordinator/SessionManageOperator.java#L268]
{code:java}
for (int i = 0; i < partitionKeyCount; i++) {
    RecordData.FieldGetter fieldGetter =
            fieldGetters.get(columnCount - partitionKeyCount - 1 + i);
    ...
} {code}
The {{-1}} offset is incorrect. When partition columns are *not* contiguously 
placed at the end of the column list, the sink reads values from the wrong 
field getters and writes data into the wrong MaxCompute partitions.
h3. Example

Consider a table whose columns are ordered as:
||Index||Column||Is Partition Key||
|0|id| |
|1|pt|yes|
|2|name| |
|3|dt|yes|

Here {{columnCount = 4}} and {{{}partitionKeyCount = 2{}}}.

With the current code:
 * For {{{}i = 0{}}}: index = {{4 - 2 - 1 + 0 = 1}} → reads {{pt}} (correct by 
coincidence)
 * For {{{}i = 1{}}}: index = {{4 - 2 - 1 + 1 = 2}} → reads {{name}} (wrong! 
should read {{dt}} at index 3)

As a result, the partition spec will contain {{{}pt=<correct>, dt=<value of 
name>{}}}, which silently corrupts data.
h3. Impact
 * Data is written into incorrect partitions when the source schema has 
partition columns interleaved with regular columns.
 * This can happen after schema evolution (e.g., {{{}ADD COLUMN{}}}) or when 
the user explicitly defines columns in a non-trailing order.

h3. Expected Behavior

The sink should resolve each partition column by its name instead of by a fixed 
trailing offset, so that it works regardless of column ordering.
h3. Suggested Fix

Locate each partition column via 
{{{}schema.getColumnNames().indexOf(partitionKey){}}}:
for (String partitionKey : schema.partitionKeys()) {
int partitionColumnIndex = schema.getColumnNames().indexOf(partitionKey);
if (partitionColumnIndex < 0 || partitionColumnIndex >= columnCount)

{ throw new IllegalStateException( String.format("Unable to find partition 
column \"%s\" in schema %s", partitionKey, schema)); }

RecordData.FieldGetter fieldGetter = fieldGetters.get(partitionColumnIndex);
Object value = fieldGetter.getFieldOrNull(recordData);
partitionSpec.set(partitionKey, Objects.toString(value));
}


> [Bug][pipeline-connector][maxcompute] Partition values extracted from wrong 
> column index when partition columns are not contiguous at the end of schema
> -------------------------------------------------------------------------------------------------------------------------------------------------------
>
>                 Key: FLINK-40739
>                 URL: https://issues.apache.org/jira/browse/FLINK-40739
>             Project: Flink
>          Issue Type: Bug
>          Components: Flink CDC
>    Affects Versions: cdc-3.2.0, cdc-3.3.0, cdc-3.2.1, cdc-3.4.0, cdc-3.5.0
>            Reporter: Yanquan Lv
>            Priority: Minor
>
> h2. Title
> {{[Bug][pipeline-connector][maxcompute] Partition values extracted from wrong 
> column index when partition columns are not contiguous at the end of schema}}
> h2. Description
> h3. Problem
> In {{{}SessionManageOperator#extractPartition{}}}, the MaxCompute pipeline 
> sink computes the index of a partition column by assuming that all partition 
> columns are contiguous and located at the end of the schema:
> [https://github.com/apache/flink-cdc/blob/master/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-maxcompute/src/main/java/org/apache/flink/cdc/connectors/maxcompute/coordinator/SessionManageOperator.java#L268]
> {code:java}
> for (int i = 0; i < partitionKeyCount; i++) {
>     RecordData.FieldGetter fieldGetter =
>             fieldGetters.get(columnCount - partitionKeyCount - 1 + i);
>     ...
> } {code}
> The {{-1}} offset is incorrect. When partition columns are *not* contiguously 
> placed at the end of the column list, the sink reads values from the wrong 
> field getters and writes data into the wrong MaxCompute partitions.
> h3. Example
> Consider a table whose columns are ordered as:
> ||Index||Column||Is Partition Key||
> |0|id| |
> |1|pt|yes|
> |2|name| |
> |3|dt|yes|
> Here {{columnCount = 4}} and {{{}partitionKeyCount = 2{}}}.
> With the current code:
>  * For {{{}i = 0{}}}: index = {{4 - 2 - 1 + 0 = 1}} → reads {{pt}} (correct 
> by coincidence)
>  * For {{{}i = 1{}}}: index = {{4 - 2 - 1 + 1 = 2}} → reads {{name}} (wrong! 
> should read {{dt}} at index 3)
> As a result, the partition spec will contain {{{}pt=<correct>, dt=<value of 
> name>{}}}, which silently corrupts data.
> h3. Impact
>  * Data is written into incorrect partitions when the source schema has 
> partition columns interleaved with regular columns.
>  * This can happen after schema evolution (e.g., {{{}ADD COLUMN{}}}) or when 
> the user explicitly defines columns in a non-trailing order.
> h3. Expected Behavior
> The sink should resolve each partition column by its name instead of by a 
> fixed trailing offset, so that it works regardless of column ordering.



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

Reply via email to