[
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.
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));
}
h3. Additional Context
The same bug has already been fixed internally in the {{release-11}} branch. A
test case covering the interleaved partition-column scenario would prevent
regressions.
was:
## Description
### 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
```java
for (int i = 0; i < partitionKeyCount; i++) {
RecordData.FieldGetter fieldGetter =
fieldGetters.get(columnCount - partitionKeyCount - 1 + i);
...
}
```
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.
### 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.
### 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.
### 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.
### Suggested Fix
Locate each partition column via
`schema.getColumnNames().indexOf(partitionKey)`:
```java
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));
}
```
### Additional Context
The same bug has already been fixed internally in the `release-11` branch. A
test case covering the interleaved partition-column scenario would prevent
regressions.
---
如果你需要,我也可以顺便生成对应的中文版,或者把这段内容直接写入一个 `.md` 文件方便你复制粘贴。
可以,下面是可直接用于提交到 Apache Flink CDC 的英文 issue 描述。
---
## Title
`[Bug][pipeline-connector][maxcompute] Partition values extracted from wrong
column index when partition columns are not contiguous at the end of schema`
## Description
### 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
```java
for (int i = 0; i < partitionKeyCount; i++) {
RecordData.FieldGetter fieldGetter =
fieldGetters.get(columnCount - partitionKeyCount - 1 + i);
...
}
```
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.
### 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.
### 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.
### 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.
### Suggested Fix
Locate each partition column via
`schema.getColumnNames().indexOf(partitionKey)`:
```java
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));
}
```
### Additional Context
The same bug has already been fixed internally in the `release-11` branch. A
test case covering the interleaved partition-column scenario would prevent
regressions.
> [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.
> 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));
> }
> h3. Additional Context
> The same bug has already been fixed internally in the {{release-11}} branch.
> A test case covering the interleaved partition-column scenario would prevent
> regressions.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)