This is an automated email from the ASF dual-hosted git repository.
twalthr 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 0a2bad7b052 [FLINK-40271][table] Fix SQL serialization of Table API
group window properties
0a2bad7b052 is described below
commit 0a2bad7b052538adecda1f68ba8b558c3aa8f097
Author: Timo Theusner <[email protected]>
AuthorDate: Thu Aug 13 15:41:25 2026 +0200
[FLINK-40271][table] Fix SQL serialization of Table API group window
properties
This closes #28886.
---
.../operations/WindowAggregateQueryOperation.java | 158 +++++++++++++++++----
.../utils/OperationExpressionsUtils.java | 23 ++-
.../operations/utils/OperationTreeBuilder.java | 2 +-
.../operations/utils/SortOperationFactory.java | 2 +-
.../flink/table/test/program/TableTestProgram.java | 13 ++
.../functions/BuiltInFunctionDefinitions.java | 10 --
.../table/api/QueryOperationSqlSemanticTest.java | 1 +
.../api/QueryOperationSqlSerializationTest.java | 78 ++++++----
.../table/api/QueryOperationTestPrograms.java | 107 ++++++++++++++
9 files changed, 331 insertions(+), 63 deletions(-)
diff --git
a/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/WindowAggregateQueryOperation.java
b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/WindowAggregateQueryOperation.java
index aebced129b1..21cef9e7694 100644
---
a/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/WindowAggregateQueryOperation.java
+++
b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/WindowAggregateQueryOperation.java
@@ -21,11 +21,17 @@ package org.apache.flink.table.operations;
import org.apache.flink.annotation.Internal;
import org.apache.flink.table.api.TableException;
import org.apache.flink.table.catalog.ResolvedSchema;
+import org.apache.flink.table.expressions.CallExpression;
import org.apache.flink.table.expressions.FieldReferenceExpression;
import org.apache.flink.table.expressions.ResolvedExpression;
import org.apache.flink.table.expressions.SqlFactory;
import org.apache.flink.table.expressions.ValueLiteralExpression;
+import org.apache.flink.table.functions.BuiltInFunctionDefinitions;
+import org.apache.flink.table.functions.FunctionDefinition;
import org.apache.flink.table.operations.utils.OperationExpressionsUtils;
+import org.apache.flink.table.types.logical.LogicalType;
+import org.apache.flink.table.types.logical.utils.LogicalTypeChecks;
+import org.apache.flink.table.utils.EncodingUtils;
import org.apache.flink.util.StringUtils;
import javax.annotation.Nullable;
@@ -35,7 +41,6 @@ import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Optional;
-import java.util.function.Function;
import java.util.stream.Collectors;
import java.util.stream.Stream;
@@ -53,6 +58,14 @@ import static
org.apache.flink.util.Preconditions.checkNotNull;
public class WindowAggregateQueryOperation implements QueryOperation {
private static final String INPUT_ALIAS = "$$T_WIN_AGG";
+
+ // Output columns of a windowing TVF, which is what window properties
become in SQL.
+ private static final String WINDOW_START_COLUMN =
+ BuiltInFunctionDefinitions.WINDOW_START.getSqlName();
+ private static final String WINDOW_END_COLUMN =
+ BuiltInFunctionDefinitions.WINDOW_END.getSqlName();
+ private static final String WINDOW_TIME_COLUMN = "window_time";
+
private final List<ResolvedExpression> groupingExpressions;
private final List<ResolvedExpression> aggregateExpressions;
private final List<ResolvedExpression> windowPropertiesExpressions;
@@ -94,35 +107,112 @@ public class WindowAggregateQueryOperation implements
QueryOperation {
@Override
public String asSerializableString(SqlFactory sqlFactory) {
+ final List<WindowColumn> windowColumns = resolveWindowColumns();
return String.format(
"SELECT %s FROM TABLE(%s\n) %s GROUP BY %s",
- Stream.of(
- groupingExpressions.stream(),
- aggregateExpressions.stream(),
- windowPropertiesExpressions.stream())
- .flatMap(Function.identity())
- .map(
- expr ->
-
OperationExpressionsUtils.scopeReferencesWithAlias(
- INPUT_ALIAS, expr))
- .map(
- resolvedExpression ->
-
resolvedExpression.asSerializableString(sqlFactory))
- .collect(Collectors.joining(", ")),
+ serializeSelectList(windowColumns, sqlFactory),
OperationUtils.indent(
groupWindow.asSerializableString(
child.asSerializableString(sqlFactory),
sqlFactory)),
INPUT_ALIAS,
- Stream.concat(
- Stream.of("window_start", "window_end"),
- groupingExpressions.stream()
- .map(
- expr ->
-
OperationExpressionsUtils
-
.scopeReferencesWithAlias(
-
INPUT_ALIAS, expr))
- .map(expr ->
expr.asSerializableString(sqlFactory)))
- .collect(Collectors.joining(", ")));
+ serializeGroupBy(windowColumns, sqlFactory));
+ }
+
+ private List<WindowColumn> resolveWindowColumns() {
+ return windowPropertiesExpressions.stream()
+ .map(property -> new WindowColumn(aliasOf(property),
windowColumnOf(property)))
+ .collect(Collectors.toList());
+ }
+
+ private static String aliasOf(ResolvedExpression aliasedProperty) {
+ return OperationExpressionsUtils.extractName(aliasedProperty)
+ .orElseThrow(
+ () ->
+ new TableException(
+ "Expected a named alias over a window
property. Got: "
+ + aliasedProperty));
+ }
+
+ /** The windowing TVF output column that the given window property
denotes. */
+ private String windowColumnOf(ResolvedExpression aliasedProperty) {
+ final FunctionDefinition property = windowPropertyOf(aliasedProperty);
+ if (BuiltInFunctionDefinitions.WINDOW_START == property) {
+ return WINDOW_START_COLUMN;
+ }
+ if (BuiltInFunctionDefinitions.WINDOW_END == property) {
+ return WINDOW_END_COLUMN;
+ }
+ if (BuiltInFunctionDefinitions.ROWTIME == property) {
+ return WINDOW_TIME_COLUMN;
+ }
+ if (BuiltInFunctionDefinitions.PROCTIME == property) {
+ checkWindowIsProcessingTime();
+ return WINDOW_TIME_COLUMN;
+ }
+ throw new TableException("Unsupported window property: " + property);
+ }
+
+ private static FunctionDefinition windowPropertyOf(ResolvedExpression
aliasedProperty) {
+ final List<ResolvedExpression> children =
aliasedProperty.getResolvedChildren();
+ if (!children.isEmpty() && children.get(0) instanceof CallExpression) {
+ final FunctionDefinition property =
+ ((CallExpression) children.get(0)).getFunctionDefinition();
+ if
(OperationExpressionsUtils.WINDOW_PROPERTIES.contains(property)) {
+ return property;
+ }
+ }
+ throw new TableException(
+ "Expected an aliased window property call. Got: " +
aliasedProperty);
+ }
+
+ /**
+ * Rejects a processing-time property on an event-time group window.
+ *
+ * <p>The {@code window_time} column of a windowing TVF derives its time
attribute kind from the
+ * window, not from the requested property, so it cannot express a
processing-time attribute of
+ * an event-time window. Since a {@code rowtime} property on a
processing-time window is
+ * rejected during expression resolution (see {@code
WindowTimeIndictorInputTypeStrategy}), we
+ * reject the other case here rather than serializing it.
+ */
+ private void checkWindowIsProcessingTime() {
+ final LogicalType windowTimeAttribute =
+
groupWindow.getTimeAttribute().getOutputDataType().getLogicalType();
+ if (LogicalTypeChecks.isProctimeAttribute(windowTimeAttribute)) {
+ return;
+ }
+ throw new TableException(
+ String.format(
+ "The processing-time property of the event-time group
window '%s' "
+ + "cannot be expressed in windowing-TVF
syntax. The window_time "
+ + "column of a windowing TVF always has the
time attribute kind "
+ + "of the window itself, so it cannot
represent a wall-clock "
+ + "processing-time attribute. Define the group
window on a "
+ + "processing-time attribute instead.",
+ groupWindow.getAlias()));
+ }
+
+ private String serializeSelectList(List<WindowColumn> windowColumns,
SqlFactory sqlFactory) {
+ return Stream.concat(
+ Stream.concat(groupingExpressions.stream(),
aggregateExpressions.stream())
+ .map(expr -> scopedToInput(expr, sqlFactory)),
+ windowColumns.stream().map(WindowColumn::asSelectItem))
+ .collect(Collectors.joining(", "));
+ }
+
+ private String scopedToInput(ResolvedExpression expr, SqlFactory
sqlFactory) {
+ return OperationExpressionsUtils.scopeReferencesWithAlias(INPUT_ALIAS,
expr)
+ .asSerializableString(sqlFactory);
+ }
+
+ private String serializeGroupBy(List<WindowColumn> windowColumns,
SqlFactory sqlFactory) {
+ final Stream<String> groupedWindowColumns =
+ windowColumns.stream().anyMatch(WindowColumn::isWindowTime)
+ ? Stream.of(WINDOW_START_COLUMN, WINDOW_END_COLUMN,
WINDOW_TIME_COLUMN)
+ : Stream.of(WINDOW_START_COLUMN, WINDOW_END_COLUMN);
+ return Stream.concat(
+ groupedWindowColumns,
+ groupingExpressions.stream().map(expr ->
scopedToInput(expr, sqlFactory)))
+ .collect(Collectors.joining(", "));
}
public List<ResolvedExpression> getGroupingExpressions() {
@@ -151,6 +241,26 @@ public class WindowAggregateQueryOperation implements
QueryOperation {
return visitor.visit(this);
}
+ /** A windowing TVF output column, projected under the alias of a window
property. */
+ private static final class WindowColumn {
+
+ private final String alias;
+ private final String column;
+
+ private WindowColumn(String alias, String column) {
+ this.alias = alias;
+ this.column = column;
+ }
+
+ private boolean isWindowTime() {
+ return WINDOW_TIME_COLUMN.equals(column);
+ }
+
+ private String asSelectItem() {
+ return String.format("(%s) AS %s", column,
EncodingUtils.escapeIdentifier(alias));
+ }
+ }
+
/** Wrapper for resolved expressions of a {@link
org.apache.flink.table.api.GroupWindow}. */
@Internal
public static class ResolvedGroupWindow {
diff --git
a/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/utils/OperationExpressionsUtils.java
b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/utils/OperationExpressionsUtils.java
index 4f39f3ecb84..81b4daf380c 100644
---
a/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/utils/OperationExpressionsUtils.java
+++
b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/utils/OperationExpressionsUtils.java
@@ -34,11 +34,14 @@ import
org.apache.flink.table.functions.BuiltInFunctionDefinitions;
import org.apache.flink.table.functions.FunctionDefinition;
import org.apache.flink.table.operations.QueryOperation;
+import java.util.Arrays;
import java.util.Collections;
+import java.util.HashSet;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Optional;
+import java.util.Set;
import java.util.stream.Collectors;
import static
org.apache.flink.table.expressions.ApiExpressionUtils.isFunctionOfKind;
@@ -47,7 +50,6 @@ import static
org.apache.flink.table.expressions.ApiExpressionUtils.unresolvedRe
import static
org.apache.flink.table.expressions.ApiExpressionUtils.valueLiteral;
import static org.apache.flink.table.expressions.ExpressionUtils.extractValue;
import static org.apache.flink.table.functions.BuiltInFunctionDefinitions.AS;
-import static
org.apache.flink.table.functions.BuiltInFunctionDefinitions.WINDOW_PROPERTIES;
import static org.apache.flink.table.functions.FunctionKind.AGGREGATE;
/**
@@ -59,6 +61,25 @@ import static
org.apache.flink.table.functions.FunctionKind.AGGREGATE;
@Internal
public class OperationExpressionsUtils {
+ /**
+ * Functions that denote a property of the {@link
org.apache.flink.table.api.GroupWindow} they
+ * are called on.
+ */
+ public static final Set<FunctionDefinition> WINDOW_PROPERTIES =
+ new HashSet<>(
+ Arrays.asList(
+ BuiltInFunctionDefinitions.WINDOW_START,
+ BuiltInFunctionDefinitions.WINDOW_END,
+ BuiltInFunctionDefinitions.PROCTIME,
+ BuiltInFunctionDefinitions.ROWTIME));
+
+ /** Functions that declare the sort order of an expression in a {@code
orderBy}. */
+ public static final Set<FunctionDefinition> ORDERING =
+ new HashSet<>(
+ Arrays.asList(
+ BuiltInFunctionDefinitions.ORDER_ASC,
+ BuiltInFunctionDefinitions.ORDER_DESC));
+
//
--------------------------------------------------------------------------------------------
// Pre-expression resolution utils
//
--------------------------------------------------------------------------------------------
diff --git
a/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/utils/OperationTreeBuilder.java
b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/utils/OperationTreeBuilder.java
index 2cf8691f165..6eb66edfb28 100644
---
a/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/utils/OperationTreeBuilder.java
+++
b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/utils/OperationTreeBuilder.java
@@ -926,7 +926,7 @@ public final class OperationTreeBuilder {
@Override
public Void visit(UnresolvedCallExpression call) {
FunctionDefinition functionDefinition =
call.getFunctionDefinition();
- if
(BuiltInFunctionDefinitions.WINDOW_PROPERTIES.contains(functionDefinition)) {
+ if
(OperationExpressionsUtils.WINDOW_PROPERTIES.contains(functionDefinition)) {
throw new ValidationException(exceptionMessage);
}
call.getChildren().forEach(expr -> expr.accept(this));
diff --git
a/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/utils/SortOperationFactory.java
b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/utils/SortOperationFactory.java
index 5c9604ac164..9ab50345b55 100644
---
a/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/utils/SortOperationFactory.java
+++
b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/utils/SortOperationFactory.java
@@ -31,8 +31,8 @@ import java.util.Collections;
import java.util.List;
import java.util.stream.Collectors;
-import static
org.apache.flink.table.functions.BuiltInFunctionDefinitions.ORDERING;
import static
org.apache.flink.table.functions.BuiltInFunctionDefinitions.ORDER_ASC;
+import static
org.apache.flink.table.operations.utils.OperationExpressionsUtils.ORDERING;
/** Utility class for creating a valid {@link SortQueryOperation} operation. */
@Internal
diff --git
a/flink-table/flink-table-api-java/src/test/java/org/apache/flink/table/test/program/TableTestProgram.java
b/flink-table/flink-table-api-java/src/test/java/org/apache/flink/table/test/program/TableTestProgram.java
index aa1c11a7601..30d8aab4809 100644
---
a/flink-table/flink-table-api-java/src/test/java/org/apache/flink/table/test/program/TableTestProgram.java
+++
b/flink-table/flink-table-api-java/src/test/java/org/apache/flink/table/test/program/TableTestProgram.java
@@ -219,6 +219,19 @@ public class TableTestProgram {
.collect(Collectors.toList());
}
+ /**
+ * A helper method to avoid boilerplate code. It assumes that only a
single Table API statement
+ * is tested.
+ */
+ public TableApiTestStep getRunTableApiTestStep() {
+ final List<TestStep> tableApiSteps =
+ runSteps.stream()
+ .filter(s -> s.getKind() == TestKind.TABLE_API)
+ .collect(Collectors.toList());
+ Preconditions.checkArgument(tableApiSteps.size() == 1, "Single Table
API step expected.");
+ return (TableApiTestStep) tableApiSteps.get(0);
+ }
+
/**
* A helper method to avoid boilerplate code. It assumes that only a
single SQL statement is
* tested.
diff --git
a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/functions/BuiltInFunctionDefinitions.java
b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/functions/BuiltInFunctionDefinitions.java
index a69cbcbf0c6..707c6ea5bfe 100644
---
a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/functions/BuiltInFunctionDefinitions.java
+++
b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/functions/BuiltInFunctionDefinitions.java
@@ -58,11 +58,9 @@ import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
import java.util.EnumSet;
-import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Optional;
-import java.util.Set;
import java.util.stream.Collectors;
import static org.apache.flink.table.api.DataTypes.BIGINT;
@@ -3475,14 +3473,6 @@ public final class BuiltInFunctionDefinitions {
.outputTypeStrategy(TypeStrategies.MISSING)
.build();
- public static final Set<FunctionDefinition> WINDOW_PROPERTIES =
- new HashSet<>(Arrays.asList(WINDOW_START, WINDOW_END, PROCTIME,
ROWTIME));
-
- public static final Set<FunctionDefinition> TIME_ATTRIBUTES =
- new HashSet<>(Arrays.asList(PROCTIME, ROWTIME));
-
- public static final List<FunctionDefinition> ORDERING =
Arrays.asList(ORDER_ASC, ORDER_DESC);
-
/**
* True when {@code key} appears among the {@code op_mapping} keys. Each
map key may itself be a
* comma-separated list (e.g. {@code "INSERT,UPDATE_AFTER"}) - each part
is trimmed and compared
diff --git
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/api/QueryOperationSqlSemanticTest.java
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/api/QueryOperationSqlSemanticTest.java
index 26735396b58..2862284602d 100644
---
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/api/QueryOperationSqlSemanticTest.java
+++
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/api/QueryOperationSqlSemanticTest.java
@@ -49,6 +49,7 @@ public class QueryOperationSqlSemanticTest extends
SemanticTestBase {
QueryOperationTestPrograms.AGGREGATE_HAVING_QUERY_OPERATION,
QueryOperationTestPrograms.LIMIT_QUERY_OPERATION,
QueryOperationTestPrograms.WINDOW_AGGREGATE_QUERY_OPERATION,
+
QueryOperationTestPrograms.WINDOW_AGGREGATE_ROWTIME_QUERY_OPERATION,
QueryOperationTestPrograms.UNION_ALL_QUERY_OPERATION,
QueryOperationTestPrograms.LATERAL_JOIN_QUERY_OPERATION,
QueryOperationTestPrograms.GROUP_HOP_WINDOW_EVENT_TIME,
diff --git
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/api/QueryOperationSqlSerializationTest.java
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/api/QueryOperationSqlSerializationTest.java
index c513c277cd5..d3175a773bb 100644
---
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/api/QueryOperationSqlSerializationTest.java
+++
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/api/QueryOperationSqlSerializationTest.java
@@ -31,6 +31,7 @@ import org.apache.flink.table.test.program.TableTestProgram;
import org.apache.flink.table.test.program.TableTestProgramRunner;
import org.apache.flink.table.test.program.TestStep.TestKind;
+import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.TestInstance;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.MethodSource;
@@ -42,7 +43,11 @@ import java.util.HashMap;
import java.util.List;
import java.util.Map;
+import static org.apache.flink.table.api.Expressions.$;
+import static org.apache.flink.table.api.Expressions.lit;
import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatCode;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
/** Tests for serialization of {@link
org.apache.flink.table.operations.QueryOperation}. */
@TestInstance(TestInstance.Lifecycle.PER_CLASS)
@@ -64,6 +69,8 @@ public class QueryOperationSqlSerializationTest implements
TableTestProgramRunne
QueryOperationTestPrograms.AGGREGATE_HAVING_QUERY_OPERATION,
QueryOperationTestPrograms.LIMIT_QUERY_OPERATION,
QueryOperationTestPrograms.WINDOW_AGGREGATE_QUERY_OPERATION,
+
QueryOperationTestPrograms.WINDOW_AGGREGATE_ROWTIME_QUERY_OPERATION,
+
QueryOperationTestPrograms.WINDOW_AGGREGATE_PROCTIME_QUERY_OPERATION,
QueryOperationTestPrograms.UNION_ALL_QUERY_OPERATION,
QueryOperationTestPrograms.LATERAL_JOIN_QUERY_OPERATION,
QueryOperationTestPrograms.SQL_QUERY_OPERATION,
@@ -87,19 +94,9 @@ public class QueryOperationSqlSerializationTest implements
TableTestProgramRunne
void testSqlSerialization(TableTestProgram program) {
final TableEnvironment env = setupEnv(program);
- final TableApiTestStep tableApiStep =
- (TableApiTestStep)
- program.runSteps.stream()
- .filter(s -> s instanceof TableApiTestStep)
- .findFirst()
- .get();
-
- final SqlTestStep sqlStep =
- (SqlTestStep)
- program.runSteps.stream()
- .filter(s -> s instanceof SqlTestStep)
- .findFirst()
- .get();
+ final TableApiTestStep tableApiStep = program.getRunTableApiTestStep();
+ final SqlTestStep sqlStep = program.getRunSqlTestStep();
+
final Table table = tableApiStep.toTable(env);
assertThat(table.getQueryOperation().asSerializableString(new
InlineFunctionSqlFactory()))
.isEqualTo(sqlStep.sql);
@@ -110,19 +107,8 @@ public class QueryOperationSqlSerializationTest implements
TableTestProgramRunne
void testSqlAsJobNameForQueryOperation(TableTestProgram program) {
final TableEnvironmentImpl env = (TableEnvironmentImpl)
setupEnv(program);
- final TableApiTestStep tableApiStep =
- (TableApiTestStep)
- program.runSteps.stream()
- .filter(s -> s instanceof TableApiTestStep)
- .findFirst()
- .get();
-
- final SqlTestStep sqlStep =
- (SqlTestStep)
- program.runSteps.stream()
- .filter(s -> s instanceof SqlTestStep)
- .findFirst()
- .get();
+ final TableApiTestStep tableApiStep = program.getRunTableApiTestStep();
+ final SqlTestStep sqlStep = program.getRunSqlTestStep();
final Table table = tableApiStep.toTable(env);
@@ -138,6 +124,46 @@ public class QueryOperationSqlSerializationTest implements
TableTestProgramRunne
assertThat(streamGraph.getJobName()).isEqualTo(sqlStep.sql);
}
+ @Test
+ void testProctimeWindowGeneratedSqlPlans() {
+ final TableTestProgram program =
+
QueryOperationTestPrograms.WINDOW_AGGREGATE_PROCTIME_QUERY_OPERATION;
+ final TableEnvironment env = setupEnv(program);
+ final Table tableApiTable =
program.getRunTableApiTestStep().toTable(env);
+
+ final String generatedSql =
+ tableApiTable
+ .getQueryOperation()
+ .asSerializableString(new InlineFunctionSqlFactory());
+
+ final Table sqlTable = env.sqlQuery(generatedSql);
+
+ assertThat(sqlTable.getResolvedSchema().getColumnNames())
+ .isEqualTo(tableApiTable.getResolvedSchema().getColumnNames());
+ assertThatCode(sqlTable::explain).doesNotThrowAnyException();
+ }
+
+ @Test
+ void
testProctimePropertyOfEventTimeWindowCannotBeExpressedInWindowingTvfSyntax() {
+ final TableEnvironment env =
+
setupEnv(QueryOperationTestPrograms.WINDOW_AGGREGATE_ROWTIME_QUERY_OPERATION);
+
+ final Table table =
+ env.from("s")
+
.window(Tumble.over(lit(5).seconds()).on($("ts")).as("w"))
+ .groupBy($("w"), $("b"))
+ .select($("b"), $("w").proctime(), $("a").sum());
+
+ assertThatThrownBy(
+ () ->
+ table.getQueryOperation()
+ .asSerializableString(new
InlineFunctionSqlFactory()))
+ .isInstanceOf(TableException.class)
+ .hasMessageContaining(
+ "The processing-time property of the event-time group
window 'w' cannot be "
+ + "expressed in windowing-TVF syntax.");
+ }
+
private static TableEnvironment setupEnv(TableTestProgram program) {
final TableEnvironment env =
TableEnvironment.create(
diff --git
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/api/QueryOperationTestPrograms.java
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/api/QueryOperationTestPrograms.java
index 16ea3502ab8..0035c9154df 100644
---
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/api/QueryOperationTestPrograms.java
+++
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/api/QueryOperationTestPrograms.java
@@ -286,6 +286,113 @@ public class QueryOperationTestPrograms {
+ ") $$T_PROJECT")
.build();
+ static final TableTestProgram WINDOW_AGGREGATE_ROWTIME_QUERY_OPERATION =
+ TableTestProgram.of(
+ "window-aggregate-rowtime-query-operation",
+ "verifies sql serialization of the window rowtime
property")
+ .setupTableSource(
+ SourceTestStep.newBuilder("s")
+ .addSchema(
+ "a bigint",
+ "b string",
+ "ts TIMESTAMP_LTZ(3)",
+ "WATERMARK FOR ts AS ts - INTERVAL
'1' SECOND")
+ .producedValues(
+ Row.of(2L, "apple",
dayOfSeconds(0)),
+ Row.of(3L, "apple",
dayOfSeconds(4)),
+ Row.of(1L, "apple",
dayOfSeconds(7)))
+ .build())
+ .setupTableSink(
+ SinkTestStep.newBuilder("sink")
+ .addSchema(
+ "b string",
+ "w_start TIMESTAMP_LTZ(3)",
+ "w_end TIMESTAMP_LTZ(3)",
+ "w_rowtime TIMESTAMP_LTZ(3)",
+ "a_sum bigint")
+ .consumedValues(
+ Row.of(
+ "apple",
+ dayOfSeconds(0),
+ dayOfSeconds(5),
+
dayOfSeconds(5).minusMillis(1),
+ 5L),
+ Row.of(
+ "apple",
+ dayOfSeconds(5),
+ dayOfSeconds(10),
+
dayOfSeconds(10).minusMillis(1),
+ 1L))
+ .build())
+ .runTableApi(
+ t ->
+ t.from("s")
+ .window(
+
Tumble.over(lit(5).seconds())
+ .on($("ts"))
+ .as("w"))
+ .groupBy($("w"), $("b"))
+ .select(
+ $("b"),
+ $("w").start(),
+ $("w").end(),
+ $("w").rowtime(),
+ $("a").sum()),
+ "sink")
+ .runSql(
+ "SELECT `$$T_PROJECT`.`b`, `$$T_PROJECT`.`EXPR$0`,
`$$T_PROJECT`.`EXPR$1`, "
+ + "`$$T_PROJECT`.`EXPR$2`,
`$$T_PROJECT`.`EXPR$3` FROM (\n"
+ + " SELECT `$$T_WIN_AGG`.`b`,
(SUM(`$$T_WIN_AGG`.`a`)) AS `EXPR$3`, "
+ + "(window_start) AS `EXPR$0`,
(window_end) AS `EXPR$1`, "
+ + "(window_time) AS `EXPR$2` FROM TABLE(\n"
+ + " TUMBLE((\n"
+ + " SELECT `$$T_SOURCE`.`a`,
`$$T_SOURCE`.`b`, "
+ + "`$$T_SOURCE`.`ts` FROM
`default_catalog`.`default_database`.`s` $$T_SOURCE\n"
+ + " ), DESCRIPTOR(`ts`), INTERVAL
'0 00:00:05.000' DAY(2) TO SECOND(3))\n"
+ + " ) $$T_WIN_AGG GROUP BY
window_start, window_end, window_time, "
+ + "`$$T_WIN_AGG`.`b`\n"
+ + ") $$T_PROJECT")
+ .build();
+
+ /**
+ * Can not be tested with {@link
+ *
org.apache.flink.table.planner.plan.nodes.exec.testutils.SemanticTestBase} as a
base class,
+ * because a processing-time window only emits from a processing-time
timer, and a bounded test
+ * source terminates the job before the timer fires, so the query produces
no rows to assert on.
+ */
+ static final TableTestProgram WINDOW_AGGREGATE_PROCTIME_QUERY_OPERATION =
+ TableTestProgram.of(
+ "window-aggregate-proctime-query-operation",
+ "verifies sql serialization of the window proctime
property")
+ .setupTableSource(
+ SourceTestStep.newBuilder("s")
+ .addSchema("a bigint", "b string",
"proctime AS PROCTIME()")
+ .producedValues(Row.of(2L, "apple"),
Row.of(3L, "apple"))
+ .build())
+ .runTableApi(
+ t ->
+ t.from("s")
+ .window(
+
Tumble.over(lit(5).seconds())
+ .on($("proctime"))
+ .as("w"))
+ .groupBy($("w"), $("b"))
+ .select($("b"), $("w").proctime(),
$("a").sum()),
+ "sink")
+ .runSql(
+ "SELECT `$$T_PROJECT`.`b`, `$$T_PROJECT`.`EXPR$0`,
"
+ + "`$$T_PROJECT`.`EXPR$1` FROM (\n"
+ + " SELECT `$$T_WIN_AGG`.`b`,
(SUM(`$$T_WIN_AGG`.`a`)) AS `EXPR$1`, "
+ + "(window_time) AS `EXPR$0` FROM TABLE(\n"
+ + " TUMBLE((\n"
+ + " SELECT `$$T_SOURCE`.`a`,
`$$T_SOURCE`.`b`, "
+ + "`$$T_SOURCE`.`proctime` FROM
`default_catalog`.`default_database`.`s` $$T_SOURCE\n"
+ + " ), DESCRIPTOR(`proctime`),
INTERVAL '0 00:00:05.000' DAY(2) TO SECOND(3))\n"
+ + " ) $$T_WIN_AGG GROUP BY
window_start, window_end, window_time, "
+ + "`$$T_WIN_AGG`.`b`\n"
+ + ") $$T_PROJECT")
+ .build();
+
private static Instant dayOfSeconds(int second) {
return LocalDateTime.of(2024, 1, 1, 0, 0,
second).atZone(ZoneId.of("UTC")).toInstant();
}