This is an automated email from the ASF dual-hosted git repository.

lvyanquan pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/flink-cdc.git


The following commit(s) were added to refs/heads/master by this push:
     new 4f239e71f8 [hotfix][runtime] Respect UDF precedence during type 
validation (#4505)
4f239e71f8 is described below

commit 4f239e71f877278e17df0d988326c4e8f3fdcc76
Author: haruki <[email protected]>
AuthorDate: Thu Aug 20 11:51:50 2026 +0800

    [hotfix][runtime] Respect UDF precedence during type validation (#4505)
    
    Co-authored-by: 春栖 <[email protected]>
---
 .../flink/cdc/runtime/parser/TransformParser.java  | 80 ++++++++++++++++++++--
 .../transform/PostTransformOperatorTest.java       |  2 +-
 .../cdc/runtime/parser/TransformParserTest.java    | 22 ++++++
 3 files changed, 99 insertions(+), 5 deletions(-)

diff --git 
a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/TransformParser.java
 
b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/TransformParser.java
index 1020a39705..02691968c2 100644
--- 
a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/TransformParser.java
+++ 
b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/TransformParser.java
@@ -55,8 +55,10 @@ import org.apache.calcite.sql.SqlIdentifier;
 import org.apache.calcite.sql.SqlKind;
 import org.apache.calcite.sql.SqlNode;
 import org.apache.calcite.sql.SqlNodeList;
+import org.apache.calcite.sql.SqlOperator;
 import org.apache.calcite.sql.SqlOperatorTable;
 import org.apache.calcite.sql.SqlSelect;
+import org.apache.calcite.sql.SqlSyntax;
 import org.apache.calcite.sql.parser.SqlParseException;
 import org.apache.calcite.sql.parser.SqlParser;
 import org.apache.calcite.sql.parser.SqlParserPos;
@@ -66,9 +68,11 @@ import org.apache.calcite.sql.type.SqlReturnTypeInference;
 import org.apache.calcite.sql.type.SqlTypeFactoryImpl;
 import org.apache.calcite.sql.type.SqlTypeName;
 import org.apache.calcite.sql.util.SqlOperatorTables;
+import org.apache.calcite.sql.util.SqlShuttle;
 import org.apache.calcite.sql.validate.SqlConformance;
 import org.apache.calcite.sql.validate.SqlConformanceEnum;
 import org.apache.calcite.sql.validate.SqlDelegatingConformance;
+import org.apache.calcite.sql.validate.SqlNameMatcher;
 import org.apache.calcite.sql.validate.SqlValidator;
 import org.apache.calcite.sql.validate.SqlValidatorUtil;
 import org.apache.calcite.sql2rel.SqlToRelConverter;
@@ -214,18 +218,22 @@ public class TransformParser {
         TransformSqlOperatorTable transformSqlOperatorTable = 
TransformSqlOperatorTable.instance();
         SqlOperatorTable udfOperatorTable = SqlOperatorTables.of(udfFunctions);
         SqlOperatorTable aiFunctionOperatorTable = 
AiFunctionSqlOperatorTable.create();
+        // Calcite looks up function candidates again when deriving a call's 
type. Rebind
+        // same-name calls and expose only UDF candidates to keep validation 
consistent with
+        // UDF-first code generation.
         SqlValidator validator =
                 SqlValidatorUtil.newValidator(
-                        SqlOperatorTables.chain(
-                                transformSqlOperatorTable,
+                        new UdfFirstSqlOperatorTable(
                                 udfOperatorTable,
-                                aiFunctionOperatorTable),
+                                SqlOperatorTables.chain(
+                                        transformSqlOperatorTable, 
aiFunctionOperatorTable)),
                         calciteCatalogReader,
                         factory,
                         SqlValidator.Config.DEFAULT
                                 .withIdentifierExpansion(true)
                                 .withConformance(SqlConformanceEnum.MYSQL_5));
-        SqlNode validateSqlNode = validator.validate(sqlNode);
+        SqlNode validateSqlNode =
+                validator.validate(resolveUserDefinedFunctions(sqlNode, 
udfFunctions));
         SqlToRelConverter sqlToRelConverter =
                 new SqlToRelConverter(
                         null,
@@ -240,6 +248,70 @@ public class TransformParser {
         return relRoot.rel;
     }
 
+    private static SqlNode resolveUserDefinedFunctions(
+            SqlNode sqlNode, List<SqlFunction> udfFunctions) {
+        return sqlNode.accept(
+                new SqlShuttle() {
+                    @Override
+                    public SqlNode visit(SqlCall call) {
+                        SqlNode visited = super.visit(call);
+                        if (visited instanceof SqlBasicCall) {
+                            SqlBasicCall basicCall = (SqlBasicCall) visited;
+                            if (basicCall.getOperator().getSyntax().family != 
SqlSyntax.FUNCTION) {
+                                return visited;
+                            }
+                            udfFunctions.stream()
+                                    .filter(
+                                            udf ->
+                                                    udf.getName()
+                                                            .equalsIgnoreCase(
+                                                                    basicCall
+                                                                            
.getOperator()
+                                                                            
.getName()))
+                                    .findFirst()
+                                    .ifPresent(basicCall::setOperator);
+                        }
+                        return visited;
+                    }
+                });
+    }
+
+    private static final class UdfFirstSqlOperatorTable implements 
SqlOperatorTable {
+        private final SqlOperatorTable udfOperatorTable;
+        private final SqlOperatorTable fallbackOperatorTable;
+
+        private UdfFirstSqlOperatorTable(
+                SqlOperatorTable udfOperatorTable, SqlOperatorTable 
fallbackOperatorTable) {
+            this.udfOperatorTable = udfOperatorTable;
+            this.fallbackOperatorTable = fallbackOperatorTable;
+        }
+
+        @Override
+        public void lookupOperatorOverloads(
+                SqlIdentifier opName,
+                @Nullable SqlFunctionCategory category,
+                SqlSyntax syntax,
+                List<SqlOperator> operatorList,
+                SqlNameMatcher nameMatcher) {
+            List<SqlOperator> udfOperators = new ArrayList<>();
+            udfOperatorTable.lookupOperatorOverloads(
+                    opName, category, syntax, udfOperators, nameMatcher);
+            if (udfOperators.isEmpty()) {
+                fallbackOperatorTable.lookupOperatorOverloads(
+                        opName, category, syntax, operatorList, nameMatcher);
+            } else {
+                operatorList.addAll(udfOperators);
+            }
+        }
+
+        @Override
+        public List<SqlOperator> getOperatorList() {
+            List<SqlOperator> operators = new 
ArrayList<>(udfOperatorTable.getOperatorList());
+            operators.addAll(fallbackOperatorTable.getOperatorList());
+            return operators;
+        }
+    }
+
     public static SqlSelect parseSelect(String statement) {
         SqlNode sqlNode;
         try {
diff --git 
a/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/operators/transform/PostTransformOperatorTest.java
 
b/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/operators/transform/PostTransformOperatorTest.java
index 8bfa4a03ed..2e02b05a3b 100644
--- 
a/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/operators/transform/PostTransformOperatorTest.java
+++ 
b/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/operators/transform/PostTransformOperatorTest.java
@@ -488,7 +488,7 @@ class PostTransformOperatorTest {
                 PostTransformOperator.newBuilder()
                         .addTransform(
                                 CUSTOMERS_TABLEID.identifier(),
-                                "*, CAST(IFNULL(1, 0) AS VARCHAR) AS 
udf_ifnull, "
+                                "*, IFNULL(1, 0) AS udf_ifnull, "
                                         + "TRY_CAST(col1) AS udf_try_cast, "
                                         + "NULLIF('%s', col1) AS udf_nullif",
                                 null)
diff --git 
a/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/parser/TransformParserTest.java
 
b/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/parser/TransformParserTest.java
index 50fcc1e455..06db8f9f47 100644
--- 
a/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/parser/TransformParserTest.java
+++ 
b/flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/parser/TransformParserTest.java
@@ -1230,6 +1230,28 @@ class TransformParserTest {
                 udfDescriptors);
     }
 
+    @Test
+    void testUdfDoesNotReplaceStructuralOperator() {
+        List<UserDefinedFunctionDescriptor> udfDescriptors =
+                Collections.singletonList(
+                        new UserDefinedFunctionDescriptor(
+                                "as",
+                                
"org.apache.flink.cdc.udf.examples.java.AddOneFunctionClass"));
+
+        List<ProjectionColumn> projectionColumns =
+                TransformParser.generateProjectionColumns(
+                        "id AS alias",
+                        DUMMY_COLUMNS,
+                        udfDescriptors,
+                        new SupportedMetadataColumn[0]);
+
+        Assertions.assertThat(projectionColumns).hasSize(1);
+        ProjectionColumn projectionColumn = projectionColumns.get(0);
+        
Assertions.assertThat(projectionColumn.getColumnName()).isEqualTo("alias");
+        
Assertions.assertThat(projectionColumn.getDataType()).isEqualTo(DataTypes.INT());
+        
Assertions.assertThat(projectionColumn.getScriptExpression()).isEqualTo("$0");
+    }
+
     @Test
     public void testTranslateUdfFilterToJaninoExpressionWithColumnNameMap() {
         List<Column> columns =

Reply via email to