From fadc2ad63dd7f641306477a553c212658b13304f Mon Sep 17 00:00:00 2001 From: Esteban Zimanyi Date: Sat, 10 Oct 2026 15:55:25 +0200 Subject: [PATCH] Register the SQL surface on every session Spark builds MobilitySparkExtensions, named in spark.sql.extensions, registers MobilitySparkSql on every session Spark builds: the one SparkSession.builder() answers, each newSession(), each Spark Connect session and each Thrift Server session. It injects a check rule whose builder calls MobilitySparkSql.registerAll on the SparkSession it is given, and Spark runs that builder once per SparkSession when it builds its analyzer, before any function of it resolves. The rule itself checks nothing. Apache Sedona 1.8.0's SedonaSqlExtensions registers its functions the same way, a check rule whose builder calls SedonaContext.create(spark). A call to MobilitySparkSql.registerAll on a session the extension serves answers alike. The README states the extension as the way to hold the surface. MobilitySparkExtensionsTest checks that a SparkSession built with only the extension named, and one newSession() answers from it, resolve asText(tfloat(1.5, TIMESTAMP '2026-03-01 00:00:00')) to 1.5@2026-03-01 00:00:00+00 and tfloat(...) to the type tfloat, and that a call to registerAll on such a SparkSession answers alike. The floor of the suite is its 41 tests. Witness: with the extension injecting only its optimizer rules, the two session tests fail on "[UNRESOLVED_ROUTINE] Cannot resolve routine `asText`". Measured over MobilityDB 1f41894b35, the catalog of MEOS-API ff5f838173 and the generator of JMEOS 66d285b742: the 41 tests pass with no warning. Over Spark Connect, a server of the Spark 4.0.1 image with Sedona 1.8.0 and Delta 4.0.0 naming the three extensions answers the asText query with 1.5@2026-03-01 00:00:00+00 on two client sessions, beside ST_AsText(ST_Point(1, 2)). Why: Spark Connect and the Thrift Server hand every client a new session, which no caller of registerAll reaches, so a MobilitySpark function resolves there only when the extension registers it. --- .github/workflows/maven.yml | 2 +- README.md | 18 ++++ .../catalyst/MobilitySparkExtensions.java | 29 +++++- .../catalyst/MobilitySparkExtensionsTest.java | 88 +++++++++++++++++++ 4 files changed, 133 insertions(+), 4 deletions(-) create mode 100644 src/test/java/org/mobilitydb/spark/catalyst/MobilitySparkExtensionsTest.java diff --git a/.github/workflows/maven.yml b/.github/workflows/maven.yml index e9585d0d..e6d0a5c7 100644 --- a/.github/workflows/maven.yml +++ b/.github/workflows/maven.yml @@ -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" diff --git a/README.md b/README.md index e570fd4a..00e51e9c 100644 --- a/README.md +++ b/README.md @@ -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; diff --git a/src/main/java/org/mobilitydb/spark/catalyst/MobilitySparkExtensions.java b/src/main/java/org/mobilitydb/spark/catalyst/MobilitySparkExtensions.java index 14fbb4e7..98293cac 100644 --- a/src/main/java/org/mobilitydb/spark/catalyst/MobilitySparkExtensions.java +++ b/src/main/java/org/mobilitydb/spark/catalyst/MobilitySparkExtensions.java @@ -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`: * *
  *   SparkSession.builder().config("spark.sql.extensions",
  *       "org.mobilitydb.spark.catalyst.MobilitySparkExtensions")
  * 
* - * 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 @@ -56,6 +66,19 @@ public final class MobilitySparkExtensions @Override public BoxedUnit apply(SparkSessionExtensions extensions) { + extensions.injectCheckRule( + new AbstractFunction1>() { + @Override + public Function1 apply(SparkSession session) { + MobilitySparkSql.registerAll(session); + return new AbstractFunction1() { + @Override + public BoxedUnit apply(LogicalPlan plan) { + return BoxedUnit.UNIT; + } + }; + } + }); extensions.injectOptimizerRule(new AbstractFunction1>() { @Override public Rule apply(SparkSession session) { diff --git a/src/test/java/org/mobilitydb/spark/catalyst/MobilitySparkExtensionsTest.java b/src/test/java/org/mobilitydb/spark/catalyst/MobilitySparkExtensionsTest.java new file mode 100644 index 00000000..037748ef --- /dev/null +++ b/src/test/java/org/mobilitydb/spark/catalyst/MobilitySparkExtensionsTest.java @@ -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); + } +}