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
2 changes: 1 addition & 1 deletion .github/workflows/maven.yml
Original file line number Diff line number Diff line change
Expand Up @@ -123,4 +123,4 @@ jobs:
uses: MobilityDB/MEOS-API/.github/actions/check-test-outcome@master
with:
log: ${{ runner.temp }}/build.log
min-tests: "13"
min-tests: "41"
18 changes: 18 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -104,6 +104,24 @@ own (`tint`, `tfloat`, `tgeompoint`, `floatspan`, ...), carried as the form its
writes, and each MobilityDB SQL name is registered once, its overload and result type chosen from
the argument types while Spark plans the call, as PostgreSQL resolves them:

Naming `MobilitySparkExtensions` in `spark.sql.extensions` registers the surface on every session
Spark builds: the session the builder answers, each `newSession()`, each Spark Connect session and
each Thrift Server session, with no further call:

```java
SparkSession spark = SparkSession.builder()
.config("spark.sql.extensions", "org.mobilitydb.spark.catalyst.MobilitySparkExtensions")
.getOrCreate();
```

```
spark-submit --conf spark.sql.extensions=org.mobilitydb.spark.catalyst.MobilitySparkExtensions ...
```

The same extension adds the optimizer rules of `org.mobilitydb.spark.catalyst`. A session built
without it registers the surface on itself through `MobilitySparkSql.registerAll`, which answers
alike on a session the extension already serves:

```java
import org.mobilitydb.spark.sql.MobilitySparkSql;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,21 +29,31 @@
import org.apache.spark.sql.SparkSessionExtensions;
import org.apache.spark.sql.catalyst.plans.logical.LogicalPlan;
import org.apache.spark.sql.catalyst.rules.Rule;
import org.mobilitydb.spark.sql.MobilitySparkSql;

import scala.Function1;
import scala.runtime.AbstractFunction1;
import scala.runtime.BoxedUnit;

/**
* The Catalyst rules MobilitySpark adds to a Spark session, named in `spark.sql.extensions`:
* The typed SQL surface and the Catalyst rules MobilitySpark adds to a Spark session, named in
* `spark.sql.extensions`:
*
* <pre>
* SparkSession.builder().config("spark.sql.extensions",
* "org.mobilitydb.spark.catalyst.MobilitySparkExtensions")
* </pre>
*
* Spark applies the named class to the session's extensions, as a Scala function of one argument,
* which this class provides through scala.runtime.AbstractFunction1.
* Spark applies the named class to the extensions of every session it builds, the default one,
* each newSession(), each Spark Connect session and each Thrift Server session, as a Scala
* function of one argument, which this class provides through scala.runtime.AbstractFunction1.
*
* It registers {@link MobilitySparkSql} on each of those sessions through a check rule, which
* Spark builds once per session when it builds the session's analyzer, before the session
* resolves its first function, as Apache Sedona's SedonaSqlExtensions registers its functions
* through SedonaContext.create. The rule itself checks nothing. A session thus holds the whole
* typed surface without calling {@link MobilitySparkSql#registerAll}, and a call to it on such a
* session answers as the surface already answers.
*
* It injects {@link OrderConjunctsByCost}, which moves a MobilitySpark function behind the
* comparisons beside it in one conjunction. Spark holds no cost for a user-defined function, so
Expand All @@ -56,6 +66,19 @@ public final class MobilitySparkExtensions

@Override
public BoxedUnit apply(SparkSessionExtensions extensions) {
extensions.injectCheckRule(
new AbstractFunction1<SparkSession, Function1<LogicalPlan, BoxedUnit>>() {
@Override
public Function1<LogicalPlan, BoxedUnit> apply(SparkSession session) {
MobilitySparkSql.registerAll(session);
return new AbstractFunction1<LogicalPlan, BoxedUnit>() {
@Override
public BoxedUnit apply(LogicalPlan plan) {
return BoxedUnit.UNIT;
}
};
}
});
extensions.injectOptimizerRule(new AbstractFunction1<SparkSession, Rule<LogicalPlan>>() {
@Override
public Rule<LogicalPlan> apply(SparkSession session) {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,88 @@
/*****************************************************************************
*
* This MobilityDB code is provided under The PostgreSQL License.
* Copyright (c) 2020-2026, Université libre de Bruxelles and MobilityDB
* contributors
*
* Permission to use, copy, modify, and distribute this software and its
* documentation for any purpose, without fee, and without a written
* agreement is hereby retained provided that the above copyright notice and
* this paragraph and the following two paragraphs appear in all copies.
*
* IN NO EVENT SHALL UNIVERSITE LIBRE DE BRUXELLES BE LIABLE TO ANY PARTY FOR
* DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES
* INCLUDING LOST PROFITS, ARISING OUT OF THE USE OF THIS SOFTWARE AND ITS
* DOCUMENTATION, EVEN IF UNIVERSITE LIBRE DE BRUXELLES HAS BEEN ADVISED OF
* THE POSSIBILITY OF SUCH DAMAGE.
*
* UNIVERSITE LIBRE DE BRUXELLES SPECIFICALLY DISCLAIMS ANY WARRANTIES,
* INCLUDING, BUT NOT LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY
* AND FITNESS FOR A PARTICULAR PURPOSE. THE SOFTWARE PROVIDED HEREUNDER IS
* ON AN "AS IS" BASIS, AND UNIVERSITE LIBRE DE BRUXELLES HAS NO OBLIGATIONS
* TO PROVIDE MAINTENANCE, SUPPORT, UPDATES, ENHANCEMENTS, OR MODIFICATIONS.
*
*****************************************************************************/

package org.mobilitydb.spark.catalyst;

import static org.junit.jupiter.api.Assertions.assertEquals;

import org.apache.spark.sql.SparkSession;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
import org.mobilitydb.spark.sql.MobilitySparkSql;

/**
* The typed SQL surface on every session of an application naming MobilitySparkExtensions in
* `spark.sql.extensions`: the session the builder answers and each session newSession() answers,
* as Spark Connect and the Thrift Server build one per client, resolve a MobilitySpark function and
* type without a call to MobilitySparkSql.registerAll, and a call to it on such a session answers
* alike.
*/
public class MobilitySparkExtensionsTest {

private static final String INSTANT = "tfloat(CAST(1.5 AS DOUBLE), TIMESTAMP '2026-03-01 00:00:00')";

private static SparkSession spark;

@BeforeAll
static void session() {
spark = SparkSession.builder().appName("session-surface").master("local[1]")
.config("spark.ui.enabled", "false")
.config("spark.sql.session.timeZone", "UTC")
.config("spark.sql.extensions", MobilitySparkExtensions.class.getName())
.getOrCreate();
}

@AfterAll
static void stop() {
if (spark != null) {
spark.stop();
}
}

private static void assertSurface(SparkSession session) {
assertEquals("1.5@2026-03-01 00:00:00+00",
session.sql("SELECT asText(" + INSTANT + ")").collectAsList().get(0).getString(0));
assertEquals("tfloat",
session.sql("SELECT " + INSTANT).schema().fields()[0].dataType().simpleString());
}

@Test
void theBuiltSessionHoldsTheSurface() {
assertSurface(spark);
}

@Test
void aNewSessionHoldsTheSurface() {
assertSurface(spark.newSession());
}

@Test
void registerAllOnTheSessionAnswersAlike() {
SparkSession session = spark.newSession();
MobilitySparkSql.registerAll(session);
assertSurface(session);
}
}
Loading