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
//
--------------------------------------------------------------------------------------------