This is an automated email from the ASF dual-hosted git repository.
snuyanzin pushed a commit to branch release-2.3
in repository https://gitbox.apache.org/repos/asf/flink.git
The following commit(s) were added to refs/heads/release-2.3 by this push:
new 1cbbfa589c2 [FLINK-40391][table] PTF with non default convertion-class
of table arg might fail
1cbbfa589c2 is described below
commit 1cbbfa589c2e4577430dd01d7905234d7e8a377a
Author: Sergey Nuyanzin <[email protected]>
AuthorDate: Fri Aug 14 23:08:43 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 e2b5e76167a..b113671cbe1 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
@@ -105,6 +105,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 887eab68f45..171226f708c 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;
@@ -1941,4 +1942,18 @@ public class ProcessTableFunctionTestPrograms {
// Also in constructed types: ROW (table input) vs.
STRUCTURED (expected).
.runSql("INSERT INTO sink SELECT * FROM f(p => TABLE v, b
=> 42)")
.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 8a75cf35416..c8f41318731 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;
@@ -1155,6 +1156,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
//
--------------------------------------------------------------------------------------------