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

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


The following commit(s) were added to refs/heads/master by this push:
     new 3c930cb2083 [FLINK-40391][table] PTF with non default convertion-class 
of table arg might fail
3c930cb2083 is described below

commit 3c930cb20830d836bea7371ca2d76160b8c4169b
Author: Sergey Nuyanzin <[email protected]>
AuthorDate: Fri Aug 14 18:42:56 2026 +0200

    [FLINK-40391][table] PTF with non default convertion-class of table arg 
might fail
---
 .../flink/table/types/inference/TypeInferenceUtil.java    |  5 ++++-
 .../exec/stream/ProcessTableFunctionSemanticTests.java    |  3 ++-
 .../exec/stream/ProcessTableFunctionTestPrograms.java     | 15 +++++++++++++++
 .../nodes/exec/stream/ProcessTableFunctionTestUtils.java  |  8 ++++++++
 4 files changed, 29 insertions(+), 2 deletions(-)

diff --git 
a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/TypeInferenceUtil.java
 
b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/TypeInferenceUtil.java
index d1d7c15d743..d7bc1669fb9 100644
--- 
a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/TypeInferenceUtil.java
+++ 
b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/TypeInferenceUtil.java
@@ -505,7 +505,10 @@ public final class TypeInferenceUtil {
                                             final DataType actualType =
                                                     
castCallContext.getArgumentDataTypes().get(pos);
                                             if (expectedType == null) {
-                                                return actualType;
+                                                return expectedArg
+                                                        .getConversionClass()
+                                                        
.map(actualType::bridgedTo)
+                                                        .orElse(actualType);
                                             }
                                             if (!supportsImplicitCast(
                                                     
actualType.getLogicalType(),
diff --git 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionSemanticTests.java
 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionSemanticTests.java
index 2357235b118..4b6a105cb2e 100644
--- 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionSemanticTests.java
+++ 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionSemanticTests.java
@@ -97,6 +97,7 @@ public class ProcessTableFunctionSemanticTests extends 
SemanticTestBase {
                 ProcessTableFunctionTestPrograms.PROCESS_ORDER_BY,
                 ProcessTableFunctionTestPrograms.PROCESS_MULTI_INPUT_ORDER_BY,
                 ProcessTableFunctionTestPrograms.PROCESS_ORDER_BY_TABLE_API,
-                ProcessTableFunctionTestPrograms.PROCESS_IMPLICIT_CASTS);
+                ProcessTableFunctionTestPrograms.PROCESS_IMPLICIT_CASTS,
+                
ProcessTableFunctionTestPrograms.PROCESS_ROW_DATA_CONVERSION_TABLE);
     }
 }
diff --git 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionTestPrograms.java
 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionTestPrograms.java
index 327c3167fe5..5a9b260c220 100644
--- 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionTestPrograms.java
+++ 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionTestPrograms.java
@@ -52,6 +52,7 @@ import 
org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctio
 import 
org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctionTestUtils.PojoStateTimeFunction;
 import 
org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctionTestUtils.PojoWithDefaultStateFunction;
 import 
org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctionTestUtils.RequiredTimeFunction;
+import 
org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctionTestUtils.RowDataRowSemanticTableFunction;
 import 
org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctionTestUtils.RowSemanticTableFunction;
 import 
org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctionTestUtils.RowSemanticTablePassThroughFunction;
 import 
org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctionTestUtils.ScalarArgsFunction;
@@ -2043,4 +2044,18 @@ public class ProcessTableFunctionTestPrograms {
                             "INSERT INTO sink SELECT * FROM f("
                                     + "r => TABLE v PARTITION BY (suite_name, 
test_name) ORDER BY (window_time ASC, c DESC), i => 1)")
                     .build();
+
+    public static final TableTestProgram PROCESS_ROW_DATA_CONVERSION_TABLE =
+            TableTestProgram.of(
+                            "process-row-data-conversion",
+                            "table argument with a non-default RowData 
conversion class")
+                    .setupTemporarySystemFunction("f", 
RowDataRowSemanticTableFunction.class)
+                    .setupSql(BASIC_VALUES)
+                    .setupTableSink(
+                            SinkTestStep.newBuilder("sink")
+                                    .addSchema(BASE_SINK_SCHEMA)
+                                    .consumedValues("+I[{Hello Bob!}]", 
"+I[{Hello Alice!}]")
+                                    .build())
+                    .runSql("INSERT INTO sink SELECT * FROM f(input => TABLE 
t)")
+                    .build();
 }
diff --git 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionTestUtils.java
 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionTestUtils.java
index 89fc127a413..57d185fd78d 100644
--- 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionTestUtils.java
+++ 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionTestUtils.java
@@ -29,6 +29,7 @@ import org.apache.flink.table.api.dataview.ListView;
 import org.apache.flink.table.api.dataview.MapView;
 import org.apache.flink.table.catalog.DataTypeFactory;
 import org.apache.flink.table.connector.ChangelogMode;
+import org.apache.flink.table.data.RowData;
 import org.apache.flink.table.functions.ChangelogFunction;
 import org.apache.flink.table.functions.ProcessTableFunction;
 import org.apache.flink.table.functions.ScalarFunction;
@@ -1188,6 +1189,13 @@ public class ProcessTableFunctionTestUtils {
         }
     }
 
+    /** Testing function with non default conversion class. */
+    public static class RowDataRowSemanticTableFunction extends 
AppendProcessTableFunctionBase {
+        public void eval(@ArgumentHint(ROW_SEMANTIC_TABLE) RowData input) {
+            collectObjects("Hello " + input.getString(0) + "!");
+        }
+    }
+
     // 
--------------------------------------------------------------------------------------------
     // Helpers
     // 
--------------------------------------------------------------------------------------------

Reply via email to