Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 26 additions & 0 deletions src/test/java/org/mobilitydb/spark/GeneratedSqlSurfaceTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -280,6 +280,32 @@ void temporalAggregatesAnswerWhatMobilityDBAnswers() {
+ "tintFromText(s) AS v FROM VALUES ('" + String.join("'), ('", tints) + "') AS t(s))"));
}

@Test
void aggregatesOfEveryStateAnswerWhatMobilityDBAnswers() {
// What MobilityDB answers for the same aggregate over the same values with the time zone
// set to UTC, each value in a partition of its own, so every partial state crosses the
// shuffle: the skip list of tAvg and tCentroid as the bytes taggstate_serialize writes,
// the set and span set of the unions as setstate_serialize and spansetstate_serialize
// write them, and the box of extent as its WKB.
assertEquals("{[1@2001-01-01 00:00:00+00, 2@2001-01-02 00:00:00+00), "
+ "[2.5@2001-01-02 00:00:00+00, 3.5@2001-01-03 00:00:00+00], "
+ "(4@2001-01-03 00:00:00+00, 5@2001-01-04 00:00:00+00]}",
aggregate("tAvg", "tfloatFromText", "[1@2001-01-01 00:00:00+00, 3@2001-01-03 00:00:00+00]",
"[3@2001-01-02 00:00:00+00, 5@2001-01-04 00:00:00+00]"));
assertEquals("{[POINT(1 0)@2001-01-01 00:00:00+00, POINT(3 2)@2001-01-03 00:00:00+00]}",
aggregate("tCentroid", "tgeompointFromText",
"[Point(0 0)@2001-01-01 00:00:00+00, Point(2 2)@2001-01-03 00:00:00+00]",
"[Point(2 0)@2001-01-01 00:00:00+00, Point(4 2)@2001-01-03 00:00:00+00]"));
assertEquals("{1, 2, 3, 5}", aggregate("setUnionAgg", "intsetFromText", "{1, 3}", "{2, 3, 5}"));
assertEquals("{[1, 5), [7, 8)}",
aggregate("spanUnionAgg", "intspanFromText", "[1, 3)", "[2, 5)", "[7, 8)"));
assertEquals("STBOX X((0,1),(2,4))",
aggregate("extent", "stboxFromText", "STBOX X((1,1),(2,2))", "STBOX X((0,3),(1,4))"));
assertEquals("TBOXFLOAT XT([1, 5],[2001-01-01 00:00:00+00, 2001-01-04 00:00:00+00])",
aggregate("extent", "tfloatFromText", "[1@2001-01-01 00:00:00+00, 3@2001-01-03 00:00:00+00]",
"[3@2001-01-02 00:00:00+00, 5@2001-01-04 00:00:00+00]"));
}

@Test
void peakCountOverMergedPieces() {
// The pieces of one trip are merged first, so the count counts trips, not pieces: two
Expand Down
7 changes: 7 additions & 0 deletions src/test/java/org/mobilitydb/spark/GeneratedSurfaceTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -505,6 +505,13 @@ void temporal_aggregates_answer_what_mobilitydb_answers() {
aggregate("tAndAgg", "tbool_in", "tbool_out", tbools));
assertEquals("{[t@2001-01-01 00:00:00+00, t@2001-01-04 00:00:00+00]}",
aggregate("tOrAgg", "tbool_in", "tbool_out", tbools));
// The skip list of tAvg carries the sum and the count of each instant, which the bytes
// taggstate_serialize writes carry across the shuffle.
assertEquals("{[1@2001-01-01 00:00:00+00, 2@2001-01-02 00:00:00+00), "
+ "[2.5@2001-01-02 00:00:00+00, 3.5@2001-01-03 00:00:00+00], "
+ "(4@2001-01-03 00:00:00+00, 5@2001-01-04 00:00:00+00]}",
aggregate("tAvg", "tfloat_in", "asText", "[1@2001-01-01, 3@2001-01-03]",
"[3@2001-01-02, 5@2001-01-04]"));
// A temporal integer and a temporal float in two partitions are values of two
// types, which a typed aggregate refuses as MobilityDB's typed aggregates do.
Exception e = assertThrows(Exception.class, () -> scalar("SELECT tSum(v) FROM (SELECT "
Expand Down
Loading