This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new a4c78825f2 [cdc] Fix computed column args upper-cased in
case-insensitive mode (#9717)
a4c78825f2 is described below
commit a4c78825f25193b7e1fc8fa061e115d76d57d033
Author: Li Guo <[email protected]>
AuthorDate: Thu Sep 10 22:52:32 2026 -0700
[cdc] Fix computed column args upper-cased in case-insensitive mode (#9717)
---
.../paimon/flink/action/cdc/ComputedColumn.java | 33 +++++++++-
.../flink/action/cdc/ComputedColumnUtils.java | 44 ++++++++-----
.../action/cdc/format/AbstractRecordParser.java | 3 +-
.../cdc/mongodb/strategy/MongoVersionStrategy.java | 3 +-
.../flink/action/cdc/mysql/MySqlRecordParser.java | 4 +-
.../action/cdc/postgres/PostgresRecordParser.java | 3 +-
.../flink/action/cdc/ComputedColumnUtilsTest.java | 74 ++++++++++++++++++++++
.../cdc/kafka/KafkaCanalSyncTableActionITCase.java | 20 ++++--
8 files changed, 153 insertions(+), 31 deletions(-)
diff --git
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/ComputedColumn.java
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/ComputedColumn.java
index 5e9041a120..2e4fc4e39d 100644
---
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/ComputedColumn.java
+++
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/ComputedColumn.java
@@ -23,6 +23,7 @@ import org.apache.paimon.types.DataType;
import javax.annotation.Nullable;
import java.io.Serializable;
+import java.util.Map;
/**
* A Computed column's value is computed from input columns. Only expression
with at most two inputs
@@ -34,10 +35,12 @@ public class ComputedColumn implements Serializable {
private final String columnName;
private final Expression expression;
+ private final boolean caseSensitive;
- public ComputedColumn(String columnName, Expression expression) {
+ public ComputedColumn(String columnName, Expression expression, boolean
caseSensitive) {
this.columnName = columnName;
this.expression = expression;
+ this.caseSensitive = caseSensitive;
}
public String columnName() {
@@ -61,4 +64,32 @@ public class ComputedColumn implements Serializable {
}
return expression.eval(input);
}
+
+ /**
+ * Evaluates this column against a source record. The referenced field is
matched by exact name,
+ * or ignoring case when the catalog is case-insensitive (record keys keep
the case of the
+ * source system).
+ */
+ @Nullable
+ public String evalFromRecord(Map<String, String> rowData) {
+ return eval(referencedValue(rowData));
+ }
+
+ @Nullable
+ private String referencedValue(Map<String, String> rowData) {
+ String reference = fieldReference();
+ if (reference == null) {
+ return null;
+ }
+ String value = rowData.get(reference);
+ if (caseSensitive || rowData.containsKey(reference)) {
+ return value;
+ }
+ for (Map.Entry<String, String> entry : rowData.entrySet()) {
+ if (reference.equalsIgnoreCase(entry.getKey())) {
+ return entry.getValue();
+ }
+ }
+ return null;
+ }
}
diff --git
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/ComputedColumnUtils.java
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/ComputedColumnUtils.java
index 7bd9dd561d..c85f33b6aa 100644
---
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/ComputedColumnUtils.java
+++
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/ComputedColumnUtils.java
@@ -26,11 +26,14 @@ import org.apache.paimon.utils.Preconditions;
import org.apache.flink.api.java.tuple.Tuple2;
import java.util.ArrayList;
+import java.util.HashMap;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;
+import static org.apache.paimon.utils.StringUtils.toLowerCaseIfNeed;
+
/** Utility methods for {@link ComputedColumn}, such as build. */
public class ComputedColumnUtils {
@@ -39,13 +42,21 @@ public class ComputedColumnUtils {
return buildComputedColumns(computedColumnArgs, physicFields, true);
}
- /** The caseSensitive only affects check. We don't change field names at
building phase. */
+ /**
+ * Column names and expression arguments are kept as given. When the
catalog is
+ * case-insensitive, names are only case-converted where they are matched:
in the dependency
+ * sort, in the type lookup of referenced fields and in the lookup of the
referenced field in
+ * source records.
+ */
public static List<ComputedColumn> buildComputedColumns(
List<String> computedColumnArgs, List<DataField> physicFields,
boolean caseSensitive) {
Map<String, DataType> typeMapping =
physicFields.stream()
.collect(
- Collectors.toMap(DataField::name,
DataField::type, (v1, v2) -> v2));
+ Collectors.toMap(
+ field ->
toLowerCaseIfNeed(field.name(), caseSensitive),
+ DataField::type,
+ (v1, v2) -> v2));
// sort computed column args by dependencies
LinkedHashMap<String, Tuple2<String, String[]>> sortedArgs =
@@ -58,11 +69,12 @@ public class ComputedColumnUtils {
String[] args = columnArg.getValue().f1;
Expression expr = Expression.create(typeMapping, caseSensitive,
exprName, args);
- ComputedColumn cmpColumn = new ComputedColumn(columnName, expr);
- computedColumns.add(new ComputedColumn(columnName, expr));
+ ComputedColumn cmpColumn = new ComputedColumn(columnName, expr,
caseSensitive);
+ computedColumns.add(cmpColumn);
- // remember the column type for later reference by other computed
columns
- typeMapping.put(columnName, cmpColumn.columnType());
+ // remember the column type for later reference by other computed
columns, in the form
+ // ReferencedField looks it up
+ typeMapping.put(toLowerCaseIfNeed(columnName, caseSensitive),
cmpColumn.columnType());
}
return computedColumns;
@@ -70,14 +82,12 @@ public class ComputedColumnUtils {
private static LinkedHashMap<String, Tuple2<String, String[]>>
sortComputedColumnArgs(
List<String> computedColumnArgs, boolean caseSensitive) {
- List<String> argList =
- computedColumnArgs.stream()
- .map(x -> caseSensitive ? x : x.toUpperCase())
- .collect(Collectors.toList());
-
+ // Only the keys used for the dependency sort are case-converted.
Names and arguments stay
+ // as typed: literals such as the date_format pattern must reach the
expression unchanged.
LinkedHashMap<String, Tuple2<String, String[]>> eqMap = new
LinkedHashMap<>();
LinkedHashMap<String, String> refMap = new LinkedHashMap<>();
- for (String arg : argList) {
+ Map<String, String> originalNames = new HashMap<>();
+ for (String arg : computedColumnArgs) {
String[] kv = arg.split("=");
if (kv.length != 2) {
throw new IllegalArgumentException(
@@ -85,6 +95,7 @@ public class ComputedColumnUtils {
"Invalid computed column argument: %s. Please
use format 'column-name=expr-name(args, ...)'.",
arg));
}
+ String columnName = kv[0].trim();
String expression = kv[1].trim();
// parse expression
int left = expression.indexOf('(');
@@ -97,9 +108,11 @@ public class ComputedColumnUtils {
String exprName = expression.substring(0, left);
String[] args = expression.substring(left + 1, right).split(",");
+ String sortName = toLowerCaseIfNeed(columnName, caseSensitive);
+ eqMap.put(columnName, Tuple2.of(exprName, args));
// args[0] may be empty string, eg. "cal_col=now()"
- eqMap.put(kv[0].trim(), Tuple2.of(exprName, args));
- refMap.put(kv[0].trim(), args[0].trim());
+ refMap.put(sortName, toLowerCaseIfNeed(args[0].trim(),
caseSensitive));
+ originalNames.put(sortName, columnName);
}
List<String> sortedKeys = DfsSort.sortKeys(refMap);
@@ -107,7 +120,8 @@ public class ComputedColumnUtils {
LinkedHashMap<String, Tuple2<String, String[]>> sortedMap =
new LinkedHashMap<>(refMap.size());
for (String key : sortedKeys) {
- sortedMap.put(key, eqMap.get(key));
+ String columnName = originalNames.get(key);
+ sortedMap.put(columnName, eqMap.get(columnName));
}
return sortedMap;
}
diff --git
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/format/AbstractRecordParser.java
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/format/AbstractRecordParser.java
index 5c092d881e..d5ea8f9d17 100644
---
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/format/AbstractRecordParser.java
+++
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/format/AbstractRecordParser.java
@@ -118,8 +118,7 @@ public abstract class AbstractRecordParser
Map<String, String> rowData, CdcSchema.Builder schemaBuilder) {
computedColumns.forEach(
computedColumn -> {
- String result =
-
computedColumn.eval(rowData.get(computedColumn.fieldReference()));
+ String result = computedColumn.evalFromRecord(rowData);
rowData.put(computedColumn.columnName(), result);
schemaBuilder.column(computedColumn.columnName(),
computedColumn.columnType());
diff --git
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/mongodb/strategy/MongoVersionStrategy.java
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/mongodb/strategy/MongoVersionStrategy.java
index fcabcaffbb..0a72f71160 100644
---
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/mongodb/strategy/MongoVersionStrategy.java
+++
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/mongodb/strategy/MongoVersionStrategy.java
@@ -160,8 +160,7 @@ public interface MongoVersionStrategy {
computedColumns.forEach(
computedColumn -> {
String columnName = computedColumn.columnName();
- String fieldReference = computedColumn.fieldReference();
- String computedValue =
computedColumn.eval(parsedRow.get(fieldReference));
+ String computedValue =
computedColumn.evalFromRecord(parsedRow);
resultMap.put(columnName, computedValue);
schemaBuilder.column(columnName,
computedColumn.columnType());
diff --git
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/mysql/MySqlRecordParser.java
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/mysql/MySqlRecordParser.java
index 6c8f2ae324..9c4caee5ef 100644
---
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/mysql/MySqlRecordParser.java
+++
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/mysql/MySqlRecordParser.java
@@ -279,9 +279,7 @@ public class MySqlRecordParser implements
FlatMapFunction<CdcSourceRecord, RichC
// generate values of computed columns
for (ComputedColumn computedColumn : computedColumns) {
- String refName = computedColumn.fieldReference();
-
- resultMap.put(computedColumn.columnName(),
computedColumn.eval(resultMap.get(refName)));
+ resultMap.put(computedColumn.columnName(),
computedColumn.evalFromRecord(resultMap));
// remember the computed column data type for later reference by
other computed columns
schemaBuilder.column(computedColumn.columnName(),
computedColumn.columnType());
diff --git
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/postgres/PostgresRecordParser.java
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/postgres/PostgresRecordParser.java
index 003ae7b529..dc89d9408e 100644
---
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/postgres/PostgresRecordParser.java
+++
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/postgres/PostgresRecordParser.java
@@ -352,8 +352,7 @@ public class PostgresRecordParser
// generate values of computed columns
for (ComputedColumn computedColumn : computedColumns) {
- String refName = computedColumn.fieldReference();
- resultMap.put(computedColumn.columnName(),
computedColumn.eval(resultMap.get(refName)));
+ resultMap.put(computedColumn.columnName(),
computedColumn.evalFromRecord(resultMap));
}
for (CdcMetadataConverter metadataConverter : metadataConverters) {
diff --git
a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/ComputedColumnUtilsTest.java
b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/ComputedColumnUtilsTest.java
index d04f91053b..a740d41afe 100644
---
a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/ComputedColumnUtilsTest.java
+++
b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/ComputedColumnUtilsTest.java
@@ -25,12 +25,15 @@ import org.junit.jupiter.api.Test;
import java.util.Arrays;
import java.util.Collections;
+import java.util.HashMap;
import java.util.List;
+import java.util.Map;
import java.util.stream.Collectors;
import static
org.apache.paimon.flink.action.cdc.ComputedColumnUtils.buildComputedColumns;
import static org.junit.Assert.assertThrows;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
/** Test for ComputedColumnUtils. */
public class ComputedColumnUtilsTest {
@@ -53,6 +56,77 @@ public class ComputedColumnUtilsTest {
columns.stream().map(ComputedColumn::columnName).collect(Collectors.toList()));
}
+ @Test
+ public void testCaseInsensitiveKeepsArgumentsAsGiven() {
+ List<DataField> physicalFields =
+ Arrays.asList(
+ new DataField(0, "id", DataTypes.INT()),
+ new DataField(1, "create_time",
DataTypes.TIMESTAMP(3)));
+ List<ComputedColumn> columns =
+ buildComputedColumns(
+ Arrays.asList(
+ "dt=date_format(CREATE_TIME,yyyy-MM-dd)",
+ "tag=cast(hello, STRING)"),
+ physicalFields,
+ false);
+
+ assertEquals(
+ Arrays.asList("dt", "tag"),
+
columns.stream().map(ComputedColumn::columnName).collect(Collectors.toList()));
+ // the pattern must not be case-converted: YYYY-MM-DD means week year
and day of year
+ assertEquals("2024-12-30", columns.get(0).eval("2024-12-30
10:00:00.000"));
+ assertEquals("hello", columns.get(1).eval(null));
+ }
+
+ @Test
+ public void testEvalFromRecordCaseInsensitive() {
+ List<DataField> physicalFields =
+ Arrays.asList(new DataField(0, "create_time",
DataTypes.TIMESTAMP(3)));
+ Map<String, String> record = new HashMap<>();
+ // the source system keeps upper case names, the user wrote the
reference in lower case
+ record.put("CREATE_TIME", "2024-12-30 10:00:00.000");
+
+ ComputedColumn caseInsensitive =
+ buildComputedColumns(
+
Collections.singletonList("dt=date_format(create_time,yyyy-MM-dd)"),
+ physicalFields,
+ false)
+ .get(0);
+ assertEquals("2024-12-30", caseInsensitive.evalFromRecord(record));
+
+ ComputedColumn caseSensitive =
+ buildComputedColumns(
+
Collections.singletonList("dt=date_format(create_time,yyyy-MM-dd)"),
+ physicalFields,
+ true)
+ .get(0);
+ assertNull(caseSensitive.evalFromRecord(record));
+ record.put("create_time", "2024-12-31 10:00:00.000");
+ assertEquals("2024-12-31", caseSensitive.evalFromRecord(record));
+ // an exact match wins over a match that ignores case
+ assertEquals("2024-12-31", caseInsensitive.evalFromRecord(record));
+ // a present key with a null value is a null input, not a reason to
look at other columns
+ record.put("create_time", null);
+ assertNull(caseInsensitive.evalFromRecord(record));
+ }
+
+ @Test
+ public void testCaseInsensitiveReferenceBetweenComputedColumns() {
+ List<DataField> physicalFields = Arrays.asList(new DataField(0,
"_date", DataTypes.DATE()));
+ List<ComputedColumn> columns =
+ buildComputedColumns(
+ Arrays.asList("_YEAR_STR=substring(_year, 0, 2)",
"_year=year(_DATE)"),
+ physicalFields,
+ false);
+
+ assertEquals(
+ Arrays.asList("_year", "_YEAR_STR"),
+
columns.stream().map(ComputedColumn::columnName).collect(Collectors.toList()));
+ assertEquals(DataTypes.INT(), columns.get(0).columnType());
+ assertEquals(DataTypes.STRING(), columns.get(1).columnType());
+ assertEquals("20",
columns.get(1).eval(columns.get(0).eval("2023-03-23")));
+ }
+
@Test
public void testCycleReference() {
List<String> calColArgs =
diff --git
a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/kafka/KafkaCanalSyncTableActionITCase.java
b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/kafka/KafkaCanalSyncTableActionITCase.java
index df48d9bd4e..984c0c102b 100644
---
a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/kafka/KafkaCanalSyncTableActionITCase.java
+++
b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/kafka/KafkaCanalSyncTableActionITCase.java
@@ -1109,9 +1109,12 @@ public class KafkaCanalSyncTableActionITCase extends
KafkaSyncTableActionITCase
createFileStoreTable(
RowType.of(
new DataType[] {
- DataTypes.INT().notNull(), DataTypes.DATE(),
DataTypes.INT(),
+ DataTypes.INT().notNull(),
+ DataTypes.DATE(),
+ DataTypes.INT(),
+ DataTypes.STRING()
},
- new String[] {"_id", "_date", "_year"}),
+ new String[] {"_id", "_date", "_year",
"_date_str"}),
Collections.emptyList(),
Collections.singletonList("_id"),
Collections.emptyList(),
@@ -1129,7 +1132,9 @@ public class KafkaCanalSyncTableActionITCase extends
KafkaSyncTableActionITCase
.withCatalogConfig(
Collections.singletonMap(
CatalogOptions.CASE_SENSITIVE.key(),
"false"))
- .withComputedColumnArgs("_YEAR=year(_DATE)")
+ // the pattern literal must survive the
case-insensitive handling
+ .withComputedColumnArgs(
+ "_YEAR=year(_DATE)",
"_DATE_STR=date_format(_DATE,yyyy-MM-dd)")
.build();
runActionWithDefaultEnv(action);
@@ -1140,11 +1145,14 @@ public class KafkaCanalSyncTableActionITCase extends
KafkaSyncTableActionITCase
RowType rowType =
RowType.of(
new DataType[] {
- DataTypes.INT().notNull(), DataTypes.DATE(),
DataTypes.INT()
+ DataTypes.INT().notNull(),
+ DataTypes.DATE(),
+ DataTypes.INT(),
+ DataTypes.STRING()
},
- new String[] {"_id", "_date", "_year"});
+ new String[] {"_id", "_date", "_year", "_date_str"});
waitForResult(
- Arrays.asList("+I[1, 19439, 2023]", "+I[2, NULL, NULL]"),
+ Arrays.asList("+I[1, 19439, 2023, 2023-03-23]", "+I[2, NULL,
NULL, NULL]"),
getFileStoreTable(tableName),
rowType,
Collections.singletonList("_id"));