This is an automated email from the ASF dual-hosted git repository.
fhueske 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 f2370b9b3f4 [FLINK-40079][table] Reject PTF calls with sys-args if
they are disabled (#28675)
f2370b9b3f4 is described below
commit f2370b9b3f4366081b8533dbed9805a26590c1ce
Author: Fabian Hueske <[email protected]>
AuthorDate: Mon Jul 20 18:49:24 2026 +0200
[FLINK-40079][table] Reject PTF calls with sys-args if they are disabled
(#28675)
* [FLINK-40079][table] Reject PTF calls with sys-args if they are disabled
* Add a check in SqlValidator to reject PTF calls with system-args
(on_time, uid) in SQL querys if the function disabled them.
* Add a check in ResolveCallByArgumentsRule to reject system-args in
functions that disabled them from Table API.
Generated-By: Claude Opus 4.8 (1M context)
---
.../resolver/rules/ResolveCallByArgumentsRule.java | 3 ++
.../table/types/inference/SystemTypeInference.java | 28 +++++++++++
.../planner/calcite/FlinkCalciteSqlValidator.java | 30 +++++++++++
.../stream/sql/MLPredictTableFunctionTest.java | 30 ++++++++++-
.../plan/stream/sql/ProcessTableFunctionTest.java | 58 ++++++++++++++++++++++
.../plan/stream/sql/SnapshotTableFunctionTest.java | 9 ++--
.../plan/stream/sql/MLPredictTableFunctionTest.xml | 2 +-
7 files changed, 152 insertions(+), 8 deletions(-)
diff --git
a/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/expressions/resolver/rules/ResolveCallByArgumentsRule.java
b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/expressions/resolver/rules/ResolveCallByArgumentsRule.java
index 5e4a02cb9fb..7364536099a 100644
---
a/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/expressions/resolver/rules/ResolveCallByArgumentsRule.java
+++
b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/expressions/resolver/rules/ResolveCallByArgumentsRule.java
@@ -299,6 +299,9 @@ final class ResolveCallByArgumentsRule implements
ResolverRule {
functionName));
}
+ SystemTypeInference.checkNoSystemArguments(
+ inference.disableSystemArguments(), namedArgs.keySet(),
functionName);
+
fillInDefaultNamedArguments(declaredArgs, namedArgs);
fillInPtfSpecificNamedArguments(
functionName, definition, declaredArgs, namedArgs,
actualArgs);
diff --git
a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/SystemTypeInference.java
b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/SystemTypeInference.java
index 971831d6ac4..a1f9fbc5a9f 100644
---
a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/SystemTypeInference.java
+++
b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/inference/SystemTypeInference.java
@@ -46,6 +46,7 @@ import javax.annotation.Nullable;
import java.util.ArrayList;
import java.util.Arrays;
+import java.util.Collection;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
@@ -126,6 +127,33 @@ public class SystemTypeInference {
return !UID_FORMAT.test(uid);
}
+ /**
+ * Rejects the implicit system arguments ({@code on_time}, {@code uid})
for a function that
+ * disables them via {@link TypeInference#disableSystemArguments()}.
+ *
+ * <p>The system arguments are not part of such a function's signature.
Enforcing this from
+ * every translation path (SQL operand checking and Table API call
resolution) rejects them
+ * consistently, regardless of whether the function is processed by the
generic PTF rule or a
+ * dedicated optimizer rule (e.g. ML_PREDICT, LATERAL SNAPSHOT).
+ */
+ public static void checkNoSystemArguments(
+ boolean sysArgsDisabled,
+ Collection<String> suppliedArgumentNames,
+ String functionName) {
+ if (!sysArgsDisabled) {
+ return;
+ }
+ for (StaticArgument systemArg : PROCESS_TABLE_FUNCTION_SYSTEM_ARGS) {
+ if (suppliedArgumentNames.contains(systemArg.getName())) {
+ throw new ValidationException(
+ String.format(
+ "Invalid function call. The '%s' argument is
not supported "
+ + "because function '%s' does not use
system arguments.",
+ systemArg.getName(), functionName));
+ }
+ }
+ }
+
//
--------------------------------------------------------------------------------------------
private static void checkScalarArgsOnly(List<StaticArgument> defaultArgs) {
diff --git
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/calcite/FlinkCalciteSqlValidator.java
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/calcite/FlinkCalciteSqlValidator.java
index 651a86925c3..0eb925352f0 100644
---
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/calcite/FlinkCalciteSqlValidator.java
+++
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/calcite/FlinkCalciteSqlValidator.java
@@ -30,9 +30,11 @@ import org.apache.flink.table.data.TimestampData;
import org.apache.flink.table.functions.FunctionKind;
import org.apache.flink.table.planner.catalog.CatalogSchemaModel;
import org.apache.flink.table.planner.catalog.CatalogSchemaTable;
+import org.apache.flink.table.planner.functions.bridging.BridgingSqlFunction;
import org.apache.flink.table.planner.plan.FlinkCalciteCatalogReader;
import org.apache.flink.table.planner.plan.utils.FlinkRexUtil;
import org.apache.flink.table.planner.utils.ShortcutUtils;
+import org.apache.flink.table.types.inference.SystemTypeInference;
import org.apache.flink.table.types.logical.DecimalType;
import org.apache.calcite.plan.RelOptCluster;
@@ -89,6 +91,7 @@ import java.math.BigDecimal;
import java.time.ZoneId;
import java.util.ArrayList;
import java.util.Collections;
+import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Objects;
@@ -383,6 +386,7 @@ public final class FlinkCalciteSqlValidator extends
FlinkSqlParsingValidator {
final SqlBasicCall call = (SqlBasicCall) node;
checkNoNamedAndPositionalMixedArgs(call);
+ checkDisabledSystemArgs(call);
// Special case for MODEL
if (node instanceof SqlExplicitModelCall) {
@@ -456,6 +460,32 @@ public final class FlinkCalciteSqlValidator extends
FlinkSqlParsingValidator {
}
}
+ /**
+ * Rejects the implicit PTF system arguments (on_time, uid) for functions
that disable them.
+ *
+ * <p>This must happen before Calcite permutes named arguments, because
unknown named arguments
+ * are silently dropped during permutation and would otherwise be lost.
The actual rule and
+ * error message live in {@link
SystemTypeInference#checkNoSystemArguments} so that the Table
+ * API path (which resolves calls without this validator) enforces it
identically.
+ */
+ private static void checkDisabledSystemArgs(SqlBasicCall call) {
+ final SqlOperator operator = call.getOperator();
+ if (!(operator instanceof BridgingSqlFunction)
+ || !((BridgingSqlFunction)
operator).getTypeInference().disableSystemArguments()) {
+ return;
+ }
+ final Set<String> suppliedArgNames = new HashSet<>();
+ for (SqlNode operand : call.getOperandList()) {
+ if (operand != null && operand.getKind() ==
SqlKind.ARGUMENT_ASSIGNMENT) {
+ final SqlNode nameNode = ((SqlCall) operand).operand(1);
+ if (nameNode instanceof SqlIdentifier) {
+ suppliedArgNames.add(((SqlIdentifier)
nameNode).getSimple());
+ }
+ }
+ }
+ SystemTypeInference.checkNoSystemArguments(true, suppliedArgNames,
operator.getName());
+ }
+
@Override
public SqlNode maybeCast(SqlNode node, RelDataType currentType,
RelDataType desiredType) {
return super.maybeCast(node, currentType, desiredType);
diff --git
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/MLPredictTableFunctionTest.java
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/MLPredictTableFunctionTest.java
index b409e02cf9d..93a9c5b937d 100644
---
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/MLPredictTableFunctionTest.java
+++
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/MLPredictTableFunctionTest.java
@@ -30,6 +30,7 @@ import org.junit.jupiter.api.Test;
import java.util.Collections;
+import static org.apache.flink.core.testutils.FlinkAssertions.anyCauseMatches;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
/** Tests for model table value function in stream mode. */
@@ -76,8 +77,7 @@ public class MLPredictTableFunctionTest extends
MLPredictTableFunctionTestBase {
@Test
public void testInputTableIsInsertOnlyStream() {
String sql =
- "SELECT *\n"
- + "FROM TABLE(ML_PREDICT(TABLE MyTable, MODEL MyModel,
DESCRIPTOR(a, b)))";
+ "SELECT *\n" + "FROM ML_PREDICT(TABLE MyTable, MODEL MyModel,
DESCRIPTOR(a, b))";
util.verifyRelPlan(
sql,
JavaScalaConversionUtil.toScala(
@@ -103,4 +103,30 @@ public class MLPredictTableFunctionTest extends
MLPredictTableFunctionTestBase {
.hasMessageContaining(
"StreamPhysicalMLPredictTableFunction doesn't support
consuming update and delete changes which is produced by node
TableSourceScan(table=[[default_catalog, default_database, CdcTable]],
fields=[a, b])");
}
+
+ @Test
+ void testOnTimeArgumentNotAllowed() {
+ // ML_PREDICT disables the implicit system arguments. Supplying
`on_time` must be rejected
+ // at the SQL level even though ML_PREDICT is handled by a dedicated
optimizer rule.
+ String sql =
+ "SELECT * FROM ML_PREDICT(INPUT => TABLE MyTable, MODEL =>
MODEL MyModel, "
+ + "ARGS => DESCRIPTOR(a, b), on_time =>
DESCRIPTOR(rowtime))";
+ assertThatThrownBy(() -> util.verifyRelPlan(sql))
+ .satisfies(
+ anyCauseMatches(
+ "The 'on_time' argument is not supported
because function "
+ + "'ML_PREDICT' does not use system
arguments."));
+ }
+
+ @Test
+ void testUidArgumentNotAllowed() {
+ String sql =
+ "SELECT * FROM ML_PREDICT(INPUT => TABLE MyTable, MODEL =>
MODEL MyModel, "
+ + "ARGS => DESCRIPTOR(a, b), uid => 'my-uid')";
+ assertThatThrownBy(() -> util.verifyRelPlan(sql))
+ .satisfies(
+ anyCauseMatches(
+ "The 'uid' argument is not supported because
function "
+ + "'ML_PREDICT' does not use system
arguments."));
+ }
}
diff --git
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.java
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.java
index d34b3f3bc8d..7500bef9c98 100644
---
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.java
+++
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/ProcessTableFunctionTest.java
@@ -76,6 +76,7 @@ import static
org.apache.flink.table.annotation.ArgumentTrait.ROW_SEMANTIC_TABLE
import static
org.apache.flink.table.annotation.ArgumentTrait.SET_SEMANTIC_TABLE;
import static org.apache.flink.table.annotation.ArgumentTrait.SUPPORT_UPDATES;
import static org.apache.flink.table.api.Expressions.$;
+import static org.apache.flink.table.api.Expressions.lit;
import static org.apache.flink.table.api.Expressions.row;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
@@ -283,6 +284,63 @@ class ProcessTableFunctionTest extends TableTestBase {
"Disabling system arguments is not supported
for user-defined PTF."));
}
+ @Test
+ void testOnTimeArgRejectedForDisabledPtf() {
+ util.addTemporarySystemFunction("f", NoSystemArgsTableFunction.class);
+ assertThatThrownBy(
+ () ->
+ util.verifyRelPlan(
+ "SELECT * FROM f(r => TABLE
t_watermarked, i => 1, "
+ + "on_time =>
DESCRIPTOR(ts));"))
+ .satisfies(
+ anyCauseMatches(
+ "The 'on_time' argument is not supported
because function "
+ + "'f' does not use system
arguments."));
+ }
+
+ @Test
+ void testUidArgRejectedForDisabledPtf() {
+ util.addTemporarySystemFunction("f", NoSystemArgsScalarFunction.class);
+ assertThatThrownBy(() -> util.verifyRelPlan("SELECT * FROM f(i => 1,
uid => 'my-uid');"))
+ .satisfies(
+ anyCauseMatches(
+ "The 'uid' argument is not supported because
function "
+ + "'f' does not use system
arguments."));
+ }
+
+ @Test
+ void testSystemArgRejectedByNameBeforeTypeCheck() {
+ // System arguments are rejected by name, rather than a type mismatch.
+ util.addTemporarySystemFunction("f", NoSystemArgsTableFunction.class);
+ assertThatThrownBy(
+ () ->
+ util.verifyRelPlan(
+ "SELECT * FROM f(r => TABLE t, i => 1,
on_time => 1);"))
+ .satisfies(
+ anyCauseMatches(
+ "The 'on_time' argument is not supported
because function "
+ + "'f' does not use system
arguments."));
+ }
+
+ @Test
+ void testSystemArgRejectedForDisabledPtfViaTableApi() {
+ // The same enforcement applies to the Table API path, which resolves
calls via
+ // ResolveCallByArgumentsRule instead of the SQL validator.
+ util.addTemporarySystemFunction("f", NoSystemArgsTableFunction.class);
+ assertThatThrownBy(
+ () ->
+ util.tableEnv()
+ .fromCall(
+ "f",
+
util.tableEnv().from("t").asArgument("r"),
+ lit(1).asArgument("i"),
+
lit("my-uid").asArgument("uid")))
+ .satisfies(
+ anyCauseMatches(
+ "The 'uid' argument is not supported because
function "
+ + "'f' does not use system
arguments."));
+ }
+
@Test
void testUidPipelineSplitIntoTwoFunctions() {
util.addTemporarySystemFunction("f", SetSemanticTableFunction.class);
diff --git
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/SnapshotTableFunctionTest.java
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/SnapshotTableFunctionTest.java
index 6d254b220b8..eef09609aeb 100644
---
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/SnapshotTableFunctionTest.java
+++
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/SnapshotTableFunctionTest.java
@@ -24,7 +24,6 @@ import org.apache.flink.table.planner.utils.TableTestBase;
import org.apache.flink.table.planner.utils.TableTestUtil;
import org.junit.jupiter.api.BeforeEach;
-import org.junit.jupiter.api.Disabled;
import org.junit.jupiter.api.Test;
import static org.apache.flink.core.testutils.FlinkAssertions.anyCauseMatches;
@@ -121,9 +120,6 @@ public class SnapshotTableFunctionTest extends
TableTestBase {
}
@Test
- @Disabled(
- "SNAPSHOT sets disableSystemArguments(true), but that flag is
currently not enforced. "
- + "Re-enable once FLINK-40079 is fixed.")
void testSystemArgumentsNotAllowed() {
// SNAPSHOT disables the implicit system arguments (e.g. `on_time`).
Passing one in a
// LATERAL context must be rejected because the argument is not part
of the function
@@ -136,7 +132,10 @@ public class SnapshotTableFunctionTest extends
TableTestBase {
+ "input => TABLE Rates, "
+ "on_time =>
DESCRIPTOR(rate_time))) AS r "
+ "WHERE o.currency =
r.currency"))
- .satisfies(anyCauseMatches("on_time"));
+ .satisfies(
+ anyCauseMatches(
+ "The 'on_time' argument is not supported
because function "
+ + "'SNAPSHOT' does not use system
arguments."));
}
@Test
diff --git
a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/MLPredictTableFunctionTest.xml
b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/MLPredictTableFunctionTest.xml
index b1d51479529..778451f067e 100644
---
a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/MLPredictTableFunctionTest.xml
+++
b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/MLPredictTableFunctionTest.xml
@@ -869,7 +869,7 @@ Calc(select=[a, b, c, d, rowtime,
PROCTIME_MATERIALIZE(proctime) AS proctime, c0
<TestCase name="testInputTableIsInsertOnlyStream">
<Resource name="sql">
<![CDATA[SELECT *
-FROM TABLE(ML_PREDICT(TABLE MyTable, MODEL MyModel, DESCRIPTOR(a, b)))]]>
+FROM ML_PREDICT(TABLE MyTable, MODEL MyModel, DESCRIPTOR(a, b))]]>
</Resource>
<Resource name="ast">
<![CDATA[