raminqaf commented on code in PR #29411:
URL: https://github.com/apache/flink/pull/29411#discussion_r4208618890
##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/VariantSemanticTest.java:
##########
@@ -277,25 +277,106 @@ public class VariantSemanticTest extends
SemanticTestBase {
.build();
static final TableTestProgram VARIANT_AS_AGG_KEY =
- TableTestProgram.of("variant-as-agg-key", "validates variant as
agg key")
+ TableTestProgram.of(
+ "variant-as-agg-key", "validates that a variant
cannot be an agg key")
.setupTableSource(
SourceTestStep.newBuilder("t")
.addSchema("k VARIANT", "v INTEGER")
- .producedValues(
- Row.of(BUILDER.of(1), 1),
- Row.of(BUILDER.of(2), 2),
- Row.of(BUILDER.of(1), 2))
+ .producedValues(Row.of(BUILDER.of(1), 1))
.build())
+ .runFailingSql(
+ "SELECT k, SUM(v) AS total FROM t GROUP BY k",
+ ValidationException.class,
+ "Column 'k' of type VARIANT cannot be used as a
grouping key, "
+ + "because the type has no equality. "
+ + "Cast the value to a comparable type
first.")
+ .build();
+
+ static final SourceTestStep OBJECTS_WITH_DIFFERENT_KEYS_SOURCE =
+ SourceTestStep.newBuilder("t")
+ .addSchema("s STRING")
+ .producedValues(
+ Row.of("{\"a\":1,\"b\":2}"),
+ Row.of("{\"b\":2,\"a\":1}"),
+ Row.of("{\"a\":1,\"c\":3}"))
+ .build();
+
+ static final TableTestProgram VARIANT_FIELD_CAST_AS_AGG_KEY =
+ TableTestProgram.of(
+ "variant-field-cast-as-agg-key",
+ "validates that casting a variant field before
grouping forms one group")
+ .setupTableSource(OBJECTS_WITH_DIFFERENT_KEYS_SOURCE)
.setupTableSink(
SinkTestStep.newBuilder("sink_t")
- .addSchema("k VARIANT", "total INTEGER")
+ .addSchema("a INT", "cnt BIGINT")
.consumedValues(
- Row.of(1, 1),
- Row.of(2, 2),
- Row.ofKind(RowKind.UPDATE_BEFORE,
1, 1),
- Row.ofKind(RowKind.UPDATE_AFTER,
1, 3))
+ Row.of(1, 1L),
+ Row.ofKind(RowKind.UPDATE_BEFORE,
1, 1L),
+ Row.ofKind(RowKind.UPDATE_AFTER,
1, 2L),
+ Row.ofKind(RowKind.UPDATE_BEFORE,
1, 2L),
+ Row.ofKind(RowKind.UPDATE_AFTER,
1, 3L))
+ .build())
+ .runSql(
+ "INSERT INTO sink_t SELECT a, COUNT(*) FROM "
+ + "(SELECT CAST(PARSE_JSON(s)['a'] AS INT)
AS a FROM t) GROUP BY a")
+ .build();
+
+ static final TableTestProgram VARIANT_FIELD_AS_AGG_KEY =
+ TableTestProgram.of(
+ "variant-field-as-agg-key",
+ "validates that a variant field cannot be an agg
key")
+ .setupTableSource(OBJECTS_WITH_DIFFERENT_KEYS_SOURCE)
+ .runFailingSql(
+ "SELECT a, COUNT(*) FROM (SELECT
PARSE_JSON(s)['a'] AS a FROM t) GROUP BY a",
+ ValidationException.class,
+ "Column 'a' of type VARIANT cannot be used as a
grouping key")
+ .build();
+
+ static final TableTestProgram VARIANT_OBJECT_AS_DISTINCT_KEY =
+ TableTestProgram.of(
+ "variant-object-as-distinct-key",
+ "validates that a variant cannot be a DISTINCT
key")
+ .setupTableSource(OBJECTS_WITH_DIFFERENT_KEYS_SOURCE)
+ .runFailingSql(
+ "SELECT COUNT(DISTINCT PARSE_JSON(s)) FROM t",
+ ValidationException.class,
+ "of type VARIANT cannot be used as an argument of
a DISTINCT aggregate")
+ .runFailingSql(
+ "SELECT DISTINCT PARSE_JSON(s) FROM t",
+ ValidationException.class,
+ "of type VARIANT cannot be used as a grouping key")
+ .build();
+
+ static final TableTestProgram VARIANT_FIELD_EQUALS_VARIANT =
+ TableTestProgram.of(
+ "variant-field-equals-variant",
+ "validates that variants cannot be compared with
=")
+ .setupTableSource(OBJECTS_WITH_DIFFERENT_KEYS_SOURCE)
+ .runFailingSql(
+ "SELECT s FROM t WHERE PARSE_JSON(s)['a'] =
PARSE_JSON('1')",
+ ValidationException.class,
+ "An expression of type VARIANT cannot be used in a
comparison with '='")
+ .build();
+
+ static final TableTestProgram VARIANT_FIELD_AS_JOIN_KEY =
+ TableTestProgram.of(
+ "variant-field-as-join-key",
+ "validates that a variant field cannot be a join
key")
+ .setupTableSource(
+ SourceTestStep.newBuilder("l")
+ .addSchema("s STRING")
+
.producedValues(Row.of("{\"id\":7,\"x\":1}"))
.build())
- .runSql("INSERT INTO sink_t SELECT k, SUM(v) AS total FROM
t GROUP BY k")
+ .setupTableSource(
+ SourceTestStep.newBuilder("r")
+ .addSchema("s STRING")
+
.producedValues(Row.of("{\"y\":2,\"id\":7}"))
+ .build())
+ .runFailingSql(
+ "SELECT PARSE_JSON(l.s)['x'], PARSE_JSON(r.s)['y']
"
+ + "FROM l JOIN r ON PARSE_JSON(l.s)['id']
= PARSE_JSON(r.s)['id']",
Review Comment:
added
--
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]