diff --git a/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/flink/FlinkPipelineTransformITCase.java b/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/flink/FlinkPipelineTransformITCase.java
index 9a774642e2f..23c60fdb4e1 100644
--- a/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/flink/FlinkPipelineTransformITCase.java
+++ b/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/flink/FlinkPipelineTransformITCase.java
@@ -2656,9 +2656,9 @@ void testTransformErrorMessage() {
.cause()
.isExactlyInstanceOf(FlinkRuntimeException.class)
.hasMessage(
- "Failed to compile expression TransformExpressionKey{originalExpression='id1 > 0', compiledExpression='greaterThan($0, 0)', argumentNames=[__time_zone__, __epoch_time__], argumentClasses=[class java.lang.String, class java.lang.Long], returnClass=class java.lang.Boolean, columnNameMap={id1=$0}}")
+ "Failed to compile expression TransformExpressionKey{originalExpression='id1 > 0', compiledScript='return greaterThan($0, 0);', argumentNames=[__time_zone__, __epoch_time__], argumentClasses=[class java.lang.String, class java.lang.Long], returnClass=class java.lang.Boolean, columnNameMap={id1=$0}}")
.cause()
- .hasMessageContaining("Compiled expression: greaterThan($0, 0)")
+ .hasMessageContaining("Compiled script: return greaterThan($0, 0);")
.hasMessageContaining("Column name map: {$0 -> id1}")
.rootCause()
.isExactlyInstanceOf(CompileException.class)
@@ -2724,7 +2724,7 @@ void testTransformErrorMessage() {
.hasMessageContaining(
"Failed to evaluate filtering expression for table `default_namespace.default_schema.mytable1`.\n"
+ "\tOriginal expression: name + 1 > 0\n"
- + "\tCompiled expression: greaterThan($0 + 1, 0)\n"
+ + "\tCompiled script: return greaterThan($0 + 1, 0);\n"
+ "\tColumn name map: {$0 -> name}")
.rootCause()
.isExactlyInstanceOf(RuntimeException.class)
diff --git a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/ProjectionColumn.java b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/ProjectionColumn.java
index fdbaf267851..df655794267 100644
--- a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/ProjectionColumn.java
+++ b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/ProjectionColumn.java
@@ -17,10 +17,12 @@
package org.apache.flink.cdc.runtime.operators.transform;
+import org.apache.flink.cdc.common.converter.JavaClassConverter;
import org.apache.flink.cdc.common.schema.Column;
import org.apache.flink.cdc.common.types.DataType;
import org.apache.flink.cdc.common.utils.StringUtils;
import org.apache.flink.cdc.runtime.operators.transform.exceptions.TransformException;
+import org.apache.flink.cdc.runtime.parser.GeneratedExpression;
import java.io.Serializable;
import java.util.ArrayList;
@@ -39,8 +41,7 @@
*
* - column: column information parsed from projection.
*
- expression: a string for column expression split from the user-defined projection.
- *
- scriptExpression: a string for column script expression compiled from the column
- * expression.
+ *
- generatedExpression: statement-level code generated from the column expression.
*
- originalColumnNames: a list for recording the name of all columns used by the column
* expression.
*
@@ -49,19 +50,19 @@ public class ProjectionColumn implements Serializable {
private static final long serialVersionUID = 1L;
private final Column column;
private final String expression;
- private final String scriptExpression;
+ private final GeneratedExpression generatedExpression;
private final List originalColumnNames;
private final Map columnNameMap;
public ProjectionColumn(
Column column,
String expression,
- String scriptExpression,
+ GeneratedExpression generatedExpression,
List originalColumnNames,
Map columnNameMap) {
this.column = column;
this.expression = expression;
- this.scriptExpression = scriptExpression;
+ this.generatedExpression = generatedExpression;
this.originalColumnNames = originalColumnNames;
this.columnNameMap = columnNameMap;
}
@@ -70,7 +71,7 @@ public ProjectionColumn copy() {
return new ProjectionColumn(
column.copy(column.getName()),
expression,
- scriptExpression,
+ generatedExpression,
new ArrayList<>(originalColumnNames),
new HashMap<>(columnNameMap));
}
@@ -92,7 +93,15 @@ public String getExpression() {
}
public String getScriptExpression() {
- return scriptExpression;
+ return getCompiledScript();
+ }
+
+ public GeneratedExpression getGeneratedExpression() {
+ return generatedExpression;
+ }
+
+ public String getCompiledScript() {
+ return generatedExpression.asScript();
}
public List getOriginalColumnNames() {
@@ -108,7 +117,7 @@ public String getColumnNameMapAsString() {
}
public boolean isValidTransformedProjectionColumn() {
- return !StringUtils.isNullOrWhitespaceOnly(scriptExpression);
+ return !StringUtils.isNullOrWhitespaceOnly(generatedExpression.getResultTerm());
}
/**
@@ -120,7 +129,12 @@ public static ProjectionColumn ofForwarded(Column column, String mappedColumnNam
String name = column.getName();
Map columnNameMap = Collections.singletonMap(name, mappedColumnName);
return new ProjectionColumn(
- column, name, mappedColumnName, Collections.singletonList(name), columnNameMap);
+ column,
+ name,
+ GeneratedExpression.fromExpression(
+ mappedColumnName, JavaClassConverter.toJavaClass(column.getType())),
+ Collections.singletonList(name),
+ columnNameMap);
}
/**
@@ -136,7 +150,8 @@ public static ProjectionColumn ofAliased(
return new ProjectionColumn(
column.copy(newName),
originalName,
- mappedColumnName,
+ GeneratedExpression.fromExpression(
+ mappedColumnName, JavaClassConverter.toJavaClass(column.getType())),
Collections.singletonList(originalName),
columnNameMap);
}
@@ -150,13 +165,13 @@ public static ProjectionColumn ofCalculated(
String columnName,
DataType dataType,
String expression,
- String scriptExpression,
+ GeneratedExpression generatedExpression,
List originalColumnNames,
Map columnNameMap) {
return new ProjectionColumn(
Column.physicalColumn(columnName, dataType),
expression,
- scriptExpression,
+ generatedExpression,
originalColumnNames,
columnNameMap);
}
@@ -169,8 +184,8 @@ public String toString() {
+ ", expression='"
+ expression
+ '\''
- + ", scriptExpression='"
- + scriptExpression
+ + ", compiledScript='"
+ + getCompiledScript().replace("\n", "\\n")
+ '\''
+ ", originalColumnNames="
+ originalColumnNames
diff --git a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/ProjectionColumnProcessor.java b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/ProjectionColumnProcessor.java
index db5c37ca258..cbf9af7161b 100644
--- a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/ProjectionColumnProcessor.java
+++ b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/ProjectionColumnProcessor.java
@@ -22,8 +22,6 @@
import org.apache.flink.cdc.common.source.SupportedMetadataColumn;
import org.apache.flink.cdc.runtime.parser.JaninoCompiler;
-import org.codehaus.janino.ExpressionEvaluator;
-
import java.lang.reflect.InvocationTargetException;
import java.util.ArrayList;
import java.util.LinkedHashSet;
@@ -45,7 +43,7 @@ public class ProjectionColumnProcessor {
private final TransformExpressionKey transformExpressionKey;
private final Map supportedMetadataColumns;
private final List