From e3ac61e745a89d3cfa3df74a79f6d4e67ad147f8 Mon Sep 17 00:00:00 2001 From: Esteban Zimanyi Date: Tue, 1 Sep 2026 00:21:01 +0200 Subject: [PATCH] State what the binding is MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The README documents a tree this repository does not have. It describes `UDTs`, `UDFs` and `utils` folders, a `MeosDatatype` base class, `PeriodUDT`, `UDTRegistrator`, `UDFRegistrator`, `PeriodUDFRegistrator`, an `org.mobiltydb.Examples` package and an AIS example — none of which exist. `src/main` holds one class, `MeosMemory`. It also states Spark 3.4.0 and Java 17 against the live Spark 3.5.1 and Java 21, prints `cmake -D MEOS=ON..` without `-DALL=ON`, and its Unit Test section says the tests are disabled in the build configuration while CI runs them on every push. It now states the binding as it is: a Spark SQL surface GENERATED from the MEOS catalog, with the chain that derives it, the one-command refresh, how to register and call the surface, what the suite proves, and the file list. The requirement that `libmeos` be built with `-DALL=ON` against the same commit as the jar is stated with its consequence, since JNR-FFI resolves lazily and a narrower library fails at the call rather than at load. Every name in the examples is read from `GeneratedSurfaceTest`, which executes them: `registerAll`, `numInstants`, `atTime`, `eIntersects`, `nearestApproachDistance`, `eDwithinPairs` with its `explode` shape, `tint_in`, `tint_out`, `temporal_num_instants`, `tnumber_twavg`. The session example uses the `local[1]` master the suite uses rather than a wider one nothing here has exercised. The ULB, OGC and logo matter is unchanged. --- README.md | 343 +++++++++++++++--------------------------------------- 1 file changed, 97 insertions(+), 246 deletions(-) diff --git a/README.md b/README.md index 0e2d118c..a0840863 100644 --- a/README.md +++ b/README.md @@ -4,7 +4,7 @@ An open-source large-scale geospatial trajectory data analytics platform based o ::: note 📝 -MobilitySpark explores the advantages of [MobilityDB](https://github.com/MobilityDB/MobilityDB) datatypes and functions in the Spark environment. However, it should be noted that the current implementation is not a full implementation and is still far from complete. +MobilitySpark brings [MobilityDB](https://github.com/MobilityDB/MobilityDB) datatypes and functions into the Spark environment. Its Spark SQL surface is *generated* from the MEOS catalog rather than hand-written, so it tracks MobilityDB master rather than a snapshot of it. ::: @@ -19,299 +19,150 @@ More information about MobilityDB, including publications, presentations, etc., ## Table of Contents -- [Features](https://github.com/MobilityDB/MobilitySpark#features) -- [Requirement](https://github.com/MobilityDB/MobilitySpark#requirements) -- [Installation](https://github.com/MobilityDB/MobilitySpark#build-the-project) -- [Usage](https://github.com/MobilityDB/MobilitySpark#usage) -- [Understanding SparkMeos](https://github.com/MobilityDB/MobilitySpark#understanding-sparkmeos) - - [The Code Structure](https://github.com/MobilityDB/MobilitySpark#the-code-structure) - - [UDTs](https://github.com/MobilityDB/MobilitySpark#udts) - - [UDFs](https://github.com/MobilityDB/MobilitySpark#udfs) -- [Future Work](https://github.com/MobilityDB/MobilitySpark#future-work) +- [How it is built](#how-it-is-built) +- [Requirements](#requirements) +- [Building](#building) +- [Using the binding](#using-the-binding) +- [Running the tests](#running-the-tests) +- [Project structure](#project-structure) -## Features -- ✅ Common user-defined-type (UDT) implementation, defined as MeosDatatype. - - Includes instantiation of some MobilityDB types such as PeriodUDT, PeriodSetUDT, TimestampSetUDT. -- 📝 Some functions defined as UDF’s. -- 🌟 Examples. +## How it is built -## Requirements - -- 🚀 MobilityDB installed with MEOS -- 🔧 JMEOS working version -- ⚡ Spark 3.4.0 -- 📝 Maven 4 -- ☕ Java 17 (recommended) - -## Usage +MobilitySpark holds no hand-written UDF registrations. The whole Spark SQL surface is emitted +at build time by a generator reading two inputs, both derived from one MobilityDB commit: -### Build the project - -First, make sure to have installed Maven in your machine. In case of MacOS, you can install via [Homebrew](https://formulae.brew.sh/formula/maven). The other step will be to install MobilityDB with MEOS in your machine, as the meos library files are required. To install MobilityDB with MEOS please follow the following instructions in your terminal: - -```bash -git clone https://github.com/MobilityDB/MobilityDB -mkdir MobilityDB/build -cd MobilityDB/build -**cmake -D MEOS=ON.. // This flag is important, to indicate the installer to build with the meos tools.** -make -sudo make install +``` +MobilityDB (master) + -> MEOS-API meos-idl.json, the catalog of the MEOS C surface + -> JMEOS org.jmeos:meos, the JVM FFI projection of that catalog + -> MobilitySpark the Spark UDF layer, generated from both ``` -Once done, we can proceed to build the project. Spark will be downloaded as a maven dependency automatically, and JMEOS is already packaged within the project. - -To build the project using Maven command-line tools, follow these steps: - -1. Open a terminal or command prompt. -2. Navigate to the directory where the project's `pom.xml` file is located. -3. Run the command `mvn clean install` to build the project and create the JAR file. -4. Once the build is complete, the JAR file can be found in the `target` directory. - -Note that you may need to have Maven installed on your system in order to run the command. +The generator is `tools/codegen_jvm.py --engine spark`, which **JMEOS owns**. This repository +stages it: `tools/codegen_jvm.py`, `tools/codegen_spark_udfs.py` and `tools/meos-idl.json` are +gitignored, and the refresh chain and CI copy them in. Nothing under `target/generated-sources` +is committed either. The one hand-written class in `src/main` is `MeosMemory`, which frees the +native pointers MEOS returns. -### Set up Intellij IDEA (Optional) +`GENERATION.md` is the full contract. -Additionally, if you are using Intellij IDEA you can use similar setup to run your spark project. -```java - - - - - - - -``` +## Requirements -### Running the examples +- Java 21 (the build sets `maven.compiler.source`/`target` to 21) +- Maven 3.8+ +- Python 3, for the generator +- `libmeos.so` built from MobilityDB master **with every family on**, and the + `org.jmeos:meos:1.0` jar built against the same commit -::: note 📝 +That last point is not optional bookkeeping. The generated surface names every function the +catalog carries, and JNR-FFI resolves symbols lazily, so a `libmeos` built without `-DALL=ON` +fails at the call rather than at load. `tools/refresh-from-master.sh` below produces a matching +pair; build MEOS by hand and the flag is `-DMEOS=ON -DALL=ON`. -Please note that the examples provided in MobilitySpark are for simple demonstration purposes only. They do not represent a full implementation of the MEOS functionality and should not be used as such. However, they do provide a good starting point for understanding how Spark interacts with MEOS. -::: +## Building -Once you have built the project run the command: +One command runs the whole chain — deriving `libmeos` and the catalog from MobilityDB master, +building the JMEOS jar against them, regenerating the UDF layer and running the suite: ```bash -java --add-exports java.base/sun.nio.ch=ALL-UNNAMED \ - --add-exports java.base/sun.security.action=ALL-UNNAMED \ - --add-opens java.base/java.nio=ALL-UNNAMED \ - --add-opens java.base/java.lang=ALL-UNNAMED \ - -cp target/classes/ org.mobiltydb.Examples. +tools/refresh-from-master.sh ``` -where , is the name of the example class to execute. - -It is recommended though to run the example directly in IntelliJ idea with the following VM configurations: +To refresh against a local MobilityDB branch instead of master: ```bash ---add-exports -java.base/sun.nio.ch=ALL-UNNAMED ---add-exports -java.base/sun.security.action=ALL-UNNAMED ---add-opens -java.base/java.nio=ALL-UNNAMED ---add-opens -java.base/java.lang=ALL-UNNAMED +tools/refresh-from-master.sh --mdb ~/src/MobilityDB ``` -### AIS Dataset +It runs the shared `MobilityDB/MEOS-API` `refresh-jvm-chain.sh` over this repository, cloning +MEOS-API the first time into `.meos-chain/`. This repository's own last leg is +`tools/refresh.conf`. -We also implemented AIS Dataset example. +With the catalog and the jar already in place, the ordinary Maven build regenerates and tests: -```java -// Read AIS Dataset - ais = ais.withColumn("point", callUDF("tGeogPointIn", col("latitude"), col("longitude"), col("t"))) - .withColumn("sog", callUDF("tFloatIn", col("sog"), col("t"))); - ais = ais.drop("latitude", "longitude"); - ais.show(); +```bash +mvn -B clean test ``` -In the example, we did aggregation using spark and custom UDF from SparkMeos to assemble the dataset. +`generate-sources` runs the generator; `build-helper` adds `target/generated-sources/spark` as a +source root. No skip-the-tests variant is offered, deliberately: the suite is what distinguishes a +regenerated surface from a merely well-formed one, and the tree-hygiene job fails a tree that +introduces a skip flag. -```java -// Assemble AIS Dataset - Dataset trajectories = ais.groupBy("mmsi") - .agg(callUDF("tGeogPointSeqIn", functions.collect_list(col("point"))).as("trajectory"), - callUDF("tFloatSeqIn", functions.collect_list(col("sog"))).as("sog")); - trajectories.show(); -``` -Furthermore, we did simple analytics. +## Using the binding -```java -Dataset originalCounts = ais.groupBy("mmsi") - .count() - .withColumnRenamed("count", "original #points"); +Register the generated surface on a `SparkSession`, then call the functions from Spark SQL: - Dataset instantsCounts = trajectories - .withColumn("SparkMEOS #points", callUDF("tGeogPointSeqNumInstant", trajectories.col("trajectory"))); - - Dataset startTimeStamp = trajectories - .withColumn("Start Timestamp", callUDF("tGeogPointSeqStartTimestamp", trajectories.col("trajectory"))); +```java +import org.mobilitydb.spark.generated.GeneratedSpatioTemporalUDFs; - originalCounts.join(instantsCounts, "mmsi").join(startTimeStamp, "mmsi"). - select("mmsi", "SparkMEOS #points", "original #points", "Start Timestamp").show(); +SparkSession spark = SparkSession.builder().appName("mobilityspark").master("local[1]").getOrCreate(); +GeneratedSpatioTemporalUDFs.registerAll(spark); ``` -## Understanding SparkMeos - -### The Code Structure - -The SparkMeos project is divided mainly into two folders: `UDTs` and `UDFs`, as well as a `utils` folder. The `UDTs` folder contains the user-defined types that have been implemented in the project, while the `UDFs` folder contains the user-defined functions. The `utils` folder contains utility classes that are used throughout the project. - -This structure allows for easy organization and management of the project's components, making it easier to maintain and further develop the project. +Temporal values travel as strings in MEOS hex-WKB, and geometries as hex-EWKB, so any Spark type +system carries them: -### UDTs +```sql +-- accessors under their canonical MobilityDB SQL names +SELECT numInstants(trip) FROM trips; +SELECT atTime(trip, '[2001-01-01, 2001-01-03]') FROM trips; -The `MeosDatatype` class in SparkMeos is the core class of the project. It has the signature `public abstract class MeosDatatype extends UserDefinedType`. All UDTs in SparkMeos should inherit from this class to implement their behavior. The class assumes that all Meos datatypes utilize a BinaryType datatype for standardization purposes. The serialization and deserialization methods are already implemented in the class, and the only methods that need to be redefined when creating a `MeosDatatype` are `userClass()`, specifying the JMEOS class that this datatype is linked to, and `fromString(String s)`, to create the object from a string. +-- spatial relationships and distance +SELECT eIntersects(trip, 'LINESTRING(1.5 0, 1.5 5)') FROM trips; +SELECT nearestApproachDistance(a.trip, b.trip) FROM trips a, trips b; -For example, let’s evaluate the implementation of the PeriodUDT: - -```java -@SQLUserDefinedType(udt = PeriodUDT.class) -public class PeriodUDT extends MeosDatatype { - /** - * Provides the Java class associated with this UDT. - * @return The Period class type. - */ - @Override - public Class userClass() { - return Period.class; - } - @Override - protected Period fromString(String s) throws SQLException{ - return new Period(s); - } -} +-- an N-by-N pair surface over two arrays of trips, consumed with explode: +-- pr.i and pr.j index back into the (0-based) input arrays +SELECT pr.i, pr.j +FROM (SELECT explode(eDwithinPairs(trips_a, trips_b, 1000.0)) AS pr FROM trip_arrays); ``` -The previous code block is written in Java and defines the implementation of the `PeriodUDT` user-defined type in the SparkMeos project. It includes the class signature, annotations, and inheritance of the `MeosDatatype` class, which is the core class of the project. The `PeriodUDT` class overrides two methods, `userClass()` and `fromString(String s)`, to specify the JMEOS class that this datatype is linked to and to create the object from a string, respectively. The `@SQLUserDefinedType` annotation is used to specify the UDT class that the `PeriodUDT` class is associated with. - -Now we can instantiate a Period datatype in Spark! - -```java -meos_initialize("UTC"); - - SparkSession spark = SparkSession - .builder() - .appName("Java Spark SQL basic example") - .config("spark.master", "local") - .getOrCreate(); - - UDTRegistrator.registerUDTs(spark); - //UDFRegistrator.registerUDFs(spark); - PeriodUDFRegistrator.registerAllUDFs(spark); - // Create some example Period objects - OffsetDateTime now = OffsetDateTime.now(); - Period period1 = new Period(now, now.plusHours(1)); - Period period2 = new Period(now.plusHours(1), now.plusHours(3)); - Period period3 = new Period(now.plusHours(2), now.plusHours(3)); - - List data = Arrays.asList( - RowFactory.create(period1), - RowFactory.create(period2), - RowFactory.create(period3) - ); - - StructType schema = new StructType() - .add("period", new PeriodUDT()); - - // Create a DataFrame with a single column of Periods - Dataset df = spark.createDataFrame(data, schema); - - // Register the DataFrame as a temporary view - df.createOrReplaceTempView("Periods"); - - // Use Spark SQL to query the view - Dataset result = spark.sql("SELECT * FROM Periods"); - - System.out.println("Example 1: Show all Periods."); - df.printSchema(); - // Show the result - result.show(false); - -meos_finalize(); -``` +The C-level entry points are registered too (`tint_in`, `tint_out`, `temporal_num_instants`, +`tnumber_twavg`, …), as is the portable bare-name operator dialect the catalog's `byOperator` map +defines. Because the names come from the catalog, a rename upstream arrives here by +regeneration rather than by editing this repository. -The previous code shows how to instantiate a `Period` datatype in Spark, create a `DataFrame` with a single column of `Periods`, register the `DataFrame` as a temporary view, and use Spark SQL to query the view. The resulting `DataFrame` is then printed to the console. The `meos_initialize` and `meos_finalize` functions are used to initialize and finalize the JMEOS middleware, respectively. However let me highlight some important things when using SparkMeos: +Free what you keep: pointers returned across the FFI boundary are raw native addresses the JVM +garbage collector does not track. `MeosMemory.free(Pointer...)` releases them, and a query that +leaks one per row will exhaust the native heap on a cross join. -- `meos_initialize()` and `meos_finalize()` functions are used to initialize and finalize the JMEOS middleware, respectively. These functions should always be called before and after using the JMEOS middleware. -- Registering the datatypes and UDFs using the respective registrators (`UDTRegistrator.registerUDTs()` and `UDFRegistrator.registerUDFs()`) is necessary for Spark to recognize and use the UDTs and UDFs in the project. -This will print something similar to: +## Running the tests ```bash -Example 1: Show all Periods. -root - |-- period: period (nullable = true) - -23/08/28 10:33:33 INFO CodeGenerator: Code generated in 224.520375 ms -23/08/28 10:33:34 INFO CodeGenerator: Code generated in 75.834833 ms -+------------------------------------------------+ -|period | -+------------------------------------------------+ -|[2023-08-28 10:33:31+02, 2023-08-28 11:33:31+02)| -|[2023-08-28 11:33:31+02, 2023-08-28 13:33:31+02)| -|[2023-08-28 12:33:31+02, 2023-08-28 13:33:31+02)| -+------------------------------------------------+ +mvn -B clean test ``` -Notice how the Spark schema now recognizes the period datatype. - -### UDFs +`GeneratedSurfaceTest` drives the generated surface from known hex-WKB literals and asserts that +it *binds and executes*, not merely that it compiles: scalar accessors and I/O round-trips, double +/ boolean / byte marshalling, the cbuffer and npoint families, the JSON-path surface, value-array +accessors, the N-by-N array UDFs, the canonical `@sqlfn` names with runtime argument-kind +dispatch, a folded out-parameter, and the H3 cell prefilter. The bare-name operators are read from +the catalog's own `byOperator` map rather than hard-coded, so a dialect rename updates the test by +itself. -Each UDT in SparkMeos needs to register its own UDFs in Spark because Spark does not automatically recognize the functions defined in the UDT class. By registering the UDFs, Spark is able to recognize and use them in SQL queries and data manipulation operations. Additionally, registering the UDFs allows for a more organized and modular approach to defining functions in the project, making it easier to maintain and modify as necessary. +MEOS keeps process-global state and cannot be re-initialised in a JVM that has finalised it, so +Surefire runs one JVM per test class (`forkCount=1`, `reuseForks=false`). Keep that configuration. -Given the extensive number of functions in MobilityDB, this POC only implements a few of these functions. -For example: +## Project structure -```java -spark.sql( - "SELECT periodExpand(" + - "stringToPeriod('[2023-08-07 14:10:49+02, 2023-08-07 15:10:49+02)'), " + - "stringToPeriod('[2019-09-08 02:00:01+02, 2019-09-10 02:00:01+02)')) " + - "as period" -).show(false); +``` +GENERATION.md the generation contract +pom.xml generator wiring, dependencies, Surefire policy +.github/workflows/maven.yml tree hygiene, then provision + generate + test +src/main/java/org/mobilitydb/spark/ + MeosMemory.java frees the native pointers MEOS returns +src/test/java/org/mobilitydb/spark/ + GeneratedSurfaceTest.java the generated surface binds and executes +tools/refresh-from-master.sh the whole chain in one command +tools/refresh.conf this repository's last leg of that chain ``` -The `spark.sql()` command shown queries Spark to execute a SQL statement. In this particular example, the SQL statement is selecting a UDF called `periodExpand()` that takes two arguments, which are two periods converted from strings using the `stringToPeriod()` function. The `periodExpand()` function returns a third period that has the lower bound of the earliest period and the upper bound of the latest period. The resulting period is then displayed as an output using the `show()` function. - -As well as with the UDT’s, each UDF should be registered, this happens when calling `UDFRegistrator.registerUDFs(spark);`. - -## Unit Test - -Currently, we have implemented small unittest for our project. However, the unittest only works for Intellij IDEA and will fail if we run them through maven. We are currently disabling the unittest in the build configuration. - -## Future Work - -Since this is only a POC to demonstrate that MobilityDB and MEOS can be ported into Spark, we want to highlight the next steps and future work to achieve a full implementation: - -- **Complete implementation and integration of MEOS data types into SparkMEOS.** - - This includes the core, basic, boxes, temporal, and time data types from JMEOS. - - So far, we have implemented the time data types as well as some extra few. -- **Generation of the UDFs with MobilityDB syntax.** - - The currently implementation wraps some of the functions either into a UDF1 or UDF2 class, but wrapping all the existing functions will be an almost impossible task to do. - - A proposal is to create a UDF generator script that, just like in JMEOS, creates all the UDFs based on the expected data and return type. - - Since some of the functions can return many different data types, future implementations should take into account the use of function overloading and polymorphism to fit into Spark SQL this behavior. This is achievable since all UDTs use MeosDatatype under the hood, and also each instance of T follows the hierarchy defined in JMEOS. -- **************Proper UDT implementation************** - - Right now, we have abstraction called MeosDataType for implementing the UDT, and able to implement MEOS DataType as UDT in Spark. In the future we want to extend the abstraction of Meos DataType by implementing the base type like Temporal, Sequence, etc as an UDT interface to allow polymorphism. -- **Creation of the JDBC Spark / Postgres Dialect with the UDTs.** - - Initial work in this has been explored (please see the jdbc-controller branch), but a more robust and satisfactory implementation is needed. - - We need to implement a Dialect that correctly maps each UDT to its MobilityDB equivalent. For example we need to map period to PeriodUDT and perform the necessary transformations. - - This happens under the JdbcDialect class from Spark. It is necessary to override some of this class’ methods such as getCatalystType for table reading and getJdbcType for table writing. - - The PostgresDriver should also be extended to handle and register first the UDT datatypes, or map the UDT datatypes to their corresponding PGObject (which is implemented in the JMEOS datatypes). -- **Implementing MobilityDB’s optimizations into Spark Catalyst.** - - This point will represent very extensive work but critical to exploit the benefits of Spark and distributed computing to it’s maximum. - - Catalyst is Spark’s optimizer. All queries go through it to generate an execution plan. - - It will be critical to implement distribution behavior to reduce random shuffling when executing operations in Spark. For instance, we can partition MobilityDB’s datatypes by 2 different approaches: by time or by geographic position. In the case of position, catalyst should implement the partition strategy to process data belonging to the same tile in the same cluster. +Everything else the build needs — the catalog, the generator, the generated sources and the jar — +is produced rather than committed.