twalthr commented on code in PR #28784: URL: https://github.com/apache/flink/pull/28784#discussion_r3622740970
##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionTestPrograms.java:
##########
@@ -1961,4 +1961,129 @@ 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_SET_SEMANTIC_TABLE_FROM_SESSION_VIEW_WITH_MULTI_PARTITION_BY =
+ TableTestProgram.of(
+
"process-set-from-session-view-with-multi-partition-by",
+ "set semantic table partitioned by
multiple columns, sourced from a view wrapping a SESSION window aggregate")
+ .setupTemporarySystemFunction("f",
SetSemanticTableFunction.class)
+ .setupTableSource(
+ SourceTestStep.newBuilder("t")
+ .addSchema(
+ "suite_name STRING",
+ "test_name STRING",
+ "ts TIMESTAMP_LTZ(3)",
+ "WATERMARK FOR ts AS ts -
INTERVAL '0.001' SECOND")
+ .producedValues(
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(0)),
+ Row.of(
+ "suiteB",
+ "test2",
+
Instant.ofEpochMilli(1)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(2)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(3)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(4)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(5)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(6)))
+ .build())
+ .setupSql(
+ "CREATE VIEW v AS "
+ + "SELECT suite_name, test_name,
COUNT(*) AS c "
+ + "FROM SESSION(TABLE t PARTITION
BY (suite_name, test_name), DESCRIPTOR(ts), INTERVAL '0.002' SECOND) "
+ + "GROUP BY suite_name, test_name,
window_start, window_end")
+ .setupTableSink(
+ SinkTestStep.newBuilder("sink")
+ .addSchema(
+ "`suite_name` STRING",
+ "`test_name` STRING",
+ "`out` STRING")
+ .consumedValues(
+ "+I[suiteB, test2,
{+I[suiteB, test2, 1], 1}]",
+ "+I[suiteA, test1,
{+I[suiteA, test1, 6], 1}]")
+ .build())
+ .runSql(
+ "INSERT INTO sink SELECT * FROM f(r =>
TABLE v PARTITION BY (suite_name, test_name), i => 1)")
+ .build();
+
+ public static final TableTestProgram
+
PROCESS_SET_SEMANTIC_TABLE_FROM_SESSION_VIEW_WITH_MULTI_PARTITION_BY_AND_ORDER_BY
=
+ TableTestProgram.of(
+
"process-order-by-multi-partition-key-and-order-by",
+ "set semantic table partitioned and
ordered by multiple columns, sourced from a view wrapping a SESSION window
aggregate")
+ .setupTemporarySystemFunction("f",
SetSemanticTableFunction.class)
+ .setupTableSource(
+ SourceTestStep.newBuilder("t")
+ .addSchema(
+ "suite_name STRING",
+ "test_name STRING",
+ "ts TIMESTAMP_LTZ(3)",
+ "WATERMARK FOR ts AS ts -
INTERVAL '0.001' SECOND")
+ .producedValues(
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(0)),
+ Row.of(
+ "suiteB",
+ "test2",
+
Instant.ofEpochMilli(1)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(2)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(3)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(4)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(5)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(6)))
+ .build())
+ .setupSql(
+ "CREATE VIEW v AS "
+ + "SELECT suite_name, test_name,
window_time, COUNT(*) AS c "
+ + "FROM TABLE(SESSION(TABLE t
PARTITION BY (suite_name, test_name), DESCRIPTOR(ts), INTERVAL '0.002' SECOND))
"
Review Comment:
```suggestion
+ "FROM SESSION(TABLE t
PARTITION BY (suite_name, test_name), DESCRIPTOR(ts), INTERVAL '0.002' SECOND) "
```
##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionTestPrograms.java:
##########
@@ -1961,4 +1961,129 @@ 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_SET_SEMANTIC_TABLE_FROM_SESSION_VIEW_WITH_MULTI_PARTITION_BY =
+ TableTestProgram.of(
+
"process-set-from-session-view-with-multi-partition-by",
+ "set semantic table partitioned by
multiple columns, sourced from a view wrapping a SESSION window aggregate")
+ .setupTemporarySystemFunction("f",
SetSemanticTableFunction.class)
+ .setupTableSource(
+ SourceTestStep.newBuilder("t")
+ .addSchema(
+ "suite_name STRING",
+ "test_name STRING",
+ "ts TIMESTAMP_LTZ(3)",
+ "WATERMARK FOR ts AS ts -
INTERVAL '0.001' SECOND")
+ .producedValues(
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(0)),
+ Row.of(
+ "suiteB",
+ "test2",
+
Instant.ofEpochMilli(1)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(2)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(3)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(4)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(5)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(6)))
+ .build())
+ .setupSql(
+ "CREATE VIEW v AS "
+ + "SELECT suite_name, test_name,
COUNT(*) AS c "
+ + "FROM SESSION(TABLE t PARTITION
BY (suite_name, test_name), DESCRIPTOR(ts), INTERVAL '0.002' SECOND) "
+ + "GROUP BY suite_name, test_name,
window_start, window_end")
+ .setupTableSink(
+ SinkTestStep.newBuilder("sink")
+ .addSchema(
+ "`suite_name` STRING",
+ "`test_name` STRING",
+ "`out` STRING")
+ .consumedValues(
+ "+I[suiteB, test2,
{+I[suiteB, test2, 1], 1}]",
+ "+I[suiteA, test1,
{+I[suiteA, test1, 6], 1}]")
+ .build())
+ .runSql(
+ "INSERT INTO sink SELECT * FROM f(r =>
TABLE v PARTITION BY (suite_name, test_name), i => 1)")
+ .build();
+
+ public static final TableTestProgram
+
PROCESS_SET_SEMANTIC_TABLE_FROM_SESSION_VIEW_WITH_MULTI_PARTITION_BY_AND_ORDER_BY
=
+ TableTestProgram.of(
+
"process-order-by-multi-partition-key-and-order-by",
+ "set semantic table partitioned and
ordered by multiple columns, sourced from a view wrapping a SESSION window
aggregate")
+ .setupTemporarySystemFunction("f",
SetSemanticTableFunction.class)
+ .setupTableSource(
+ SourceTestStep.newBuilder("t")
+ .addSchema(
+ "suite_name STRING",
+ "test_name STRING",
+ "ts TIMESTAMP_LTZ(3)",
+ "WATERMARK FOR ts AS ts -
INTERVAL '0.001' SECOND")
+ .producedValues(
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(0)),
+ Row.of(
+ "suiteB",
+ "test2",
+
Instant.ofEpochMilli(1)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(2)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(3)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(4)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(5)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(6)))
+ .build())
+ .setupSql(
+ "CREATE VIEW v AS "
+ + "SELECT suite_name, test_name,
window_time, COUNT(*) AS c "
+ + "FROM TABLE(SESSION(TABLE t
PARTITION BY (suite_name, test_name), DESCRIPTOR(ts), INTERVAL '0.002' SECOND))
"
+ + "GROUP BY suite_name, test_name,
window_start, window_end, window_time")
+ .setupTableSink(
+ SinkTestStep.newBuilder("sink")
+ .addSchema(
+ "`suite_name` STRING",
+ "`test_name` STRING",
+ "`out` STRING")
+ .consumedValues(
+ "+I[suiteB, test2,
{+I[suiteB, test2, 1970-01-01T00:00:00.002Z, 1], 1}]",
+ "+I[suiteA, test1,
{+I[suiteA, test1, 1970-01-01T00:00:00.007Z, 6], 1}]")
Review Comment:
can we modify the output such that `c DESC` has a meaning?
##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionTestPrograms.java:
##########
@@ -1961,4 +1961,129 @@ 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_SET_SEMANTIC_TABLE_FROM_SESSION_VIEW_WITH_MULTI_PARTITION_BY =
Review Comment:
```suggestion
PROCESS_MULTI_PARTITION_BY =
```
##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionTestPrograms.java:
##########
@@ -1961,4 +1961,129 @@ 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_SET_SEMANTIC_TABLE_FROM_SESSION_VIEW_WITH_MULTI_PARTITION_BY =
+ TableTestProgram.of(
+
"process-set-from-session-view-with-multi-partition-by",
+ "set semantic table partitioned by
multiple columns, sourced from a view wrapping a SESSION window aggregate")
+ .setupTemporarySystemFunction("f",
SetSemanticTableFunction.class)
+ .setupTableSource(
+ SourceTestStep.newBuilder("t")
+ .addSchema(
+ "suite_name STRING",
+ "test_name STRING",
+ "ts TIMESTAMP_LTZ(3)",
+ "WATERMARK FOR ts AS ts -
INTERVAL '0.001' SECOND")
+ .producedValues(
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(0)),
+ Row.of(
+ "suiteB",
+ "test2",
+
Instant.ofEpochMilli(1)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(2)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(3)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(4)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(5)),
+ Row.of(
+ "suiteA",
+ "test1",
+
Instant.ofEpochMilli(6)))
+ .build())
+ .setupSql(
+ "CREATE VIEW v AS "
+ + "SELECT suite_name, test_name,
COUNT(*) AS c "
+ + "FROM SESSION(TABLE t PARTITION
BY (suite_name, test_name), DESCRIPTOR(ts), INTERVAL '0.002' SECOND) "
+ + "GROUP BY suite_name, test_name,
window_start, window_end")
+ .setupTableSink(
+ SinkTestStep.newBuilder("sink")
+ .addSchema(
+ "`suite_name` STRING",
+ "`test_name` STRING",
+ "`out` STRING")
+ .consumedValues(
+ "+I[suiteB, test2,
{+I[suiteB, test2, 1], 1}]",
+ "+I[suiteA, test1,
{+I[suiteA, test1, 6], 1}]")
+ .build())
+ .runSql(
+ "INSERT INTO sink SELECT * FROM f(r =>
TABLE v PARTITION BY (suite_name, test_name), i => 1)")
+ .build();
+
+ public static final TableTestProgram
+
PROCESS_SET_SEMANTIC_TABLE_FROM_SESSION_VIEW_WITH_MULTI_PARTITION_BY_AND_ORDER_BY
=
Review Comment:
```suggestion
PROCESS_MULTI_PARTITION_BY_AND_ORDER_BY =
```
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
