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 @@ -109,4 +109,4 @@ jobs:
uses: MobilityDB/MEOS-API/.github/actions/check-test-outcome@master
with:
log: ${{ runner.temp }}/build.log
min-tests: "19"
min-tests: "28"
19 changes: 19 additions & 0 deletions GENERATION.md
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,25 @@ MobilityDB @ master
CI stages the master-derived catalog to `tools/meos-idl.json` before the build; the catalog
is gitignored, not committed.

## The Flink SQL surface

`codegen_jvm.py --engine flink-sql` emits the MobilityDB SQL surface for Flink SQL from the
same catalog and jar: one Flink function per SQL name the catalog states (a signature's
`sqlName`, else the function's `@sqlfn`), each overload an `eval` method, and one RAW type per
MEOS value type the signatures use. A MEOS value crosses Flink as the serialized form its
catalog codec writes, hex WKB where the catalog states a WKB decoder and an `asHexWKB`
encoder and text otherwise, never as a native pointer, so Flink can copy, checkpoint and group
it. A SQL array argument stands for the C array the catalog pairs with a count
(`shape.inputArrays`), and the eval passes the length of the Flink array as that count. A
signature whose types the surface cannot carry (`Datum`, aggregate state, an array element with
no codec) is skipped and counted in the generator's report.

`org.mobilitydb.flink.sql.MobilityFlinkSql.registerAll(tEnv)` registers every function as a
catalog function under its MobilityDB SQL name. A Flink built-in of the same name (`lower`,
`round`, `abs`, …) resolves before a catalog function, and a name Flink's parser reserves
(`overlaps`, `contains`, `union`, …) is written quoted, `` `overlaps`(a, b) ``. Maven
`generate-sources` runs the engine into `target/generated-sql`, beside the facades.

## Generate-then-retire — the green-CI version is the probe

Hand-written facades/glue are replaced by the generated forwarders **family by family,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,7 @@
*
* <p><b>Typical usage</b> — per-vehicle-pair "did they come within
* 100m of each other in the last 5 minutes?" via
* {@code MeosOpsTGeo.edwithin_tgeo_tgeo} (tier = {@code cross-stream}):
* {@code MeosOpsFreeGeo.edwithin_tgeo_tgeo} (tier = {@code cross-stream}):
*
* <pre>{@code
* KeyedStream<VehiclePosition, Integer> a = streamA.keyBy(VehiclePosition::regionId);
Expand All @@ -64,7 +64,7 @@
* (left, right, ctx) -> {
* Pointer leftT = left.toTGeoPointer();
* Pointer rightT = right.toTGeoPointer();
* if (MeosOpsTGeo.edwithin_tgeo_tgeo(leftT, rightT, 100.0) != 0) {
* if (MeosOpsFreeGeo.edwithin_tgeo_tgeo(leftT, rightT, 100.0) != 0) {
* return new MeetingEvent(left.id(), right.id(), ctx.getLeftTimestamp());
* }
* return null; // no output for non-matches
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -44,14 +44,14 @@
* independently.
*
* <p><b>Typical usage</b>: scalar-predicate filter against the
* generated {@code MeosOpsTBox.overlaps_tbox_tbox} (tier =
* generated {@code MeosOpsFreeCore.overlaps_tbox_tbox} (tier =
* {@code stateless}):
*
* <pre>{@code
* DataStream<TbiePair> in = ...;
* DataStream<TbiePair> overlapping = in.filter(
* new MeosStatelessFilter<>(
* pair -> MeosOpsTBox.overlaps_tbox_tbox(pair.a, pair.b)));
* pair -> MeosOpsFreeCore.overlaps_tbox_tbox(pair.a, pair.b)));
* }</pre>
*
* <p>For int-coded predicates (JMEOS returns {@code int} for some MEOS
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -52,14 +52,14 @@
*
* <p><b>Typical usage</b>: register a stateless MEOS predicate / arithmetic
* call as a per-event map step in a DataStream pipeline. Example with
* the generated {@code MeosOpsTBox.overlaps_tbox_tbox} (tier =
* the generated {@code MeosOpsFreeCore.overlaps_tbox_tbox} (tier =
* {@code stateless}, per the codegen manifest):
*
* <pre>{@code
* DataStream<TbiePair> in = ...; // (tboxA, tboxB)
* DataStream<Boolean> overlap = in.map(
* new MeosStatelessMap<>(
* pair -> MeosOpsTBox.overlaps_tbox_tbox(pair.a, pair.b)));
* pair -> MeosOpsFreeCore.overlaps_tbox_tbox(pair.a, pair.b)));
* }</pre>
*
* <p><b>Tier coverage</b>: as of the codegen state on the parent PR,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,8 +38,8 @@
*
* <p>The {@code windowed} tier is "output cardinality changes; needs a
* window". The canonical examples are
* {@code temporal_length(tgeo)} (one length per trajectory window),
* {@code temporal_twavg(tnumber)} (one time-weighted average per
* {@code tpoint_length(tpoint)} (one length per trajectory window),
* {@code tnumber_twavg(tnumber)} (one time-weighted average per
* window), and the per-class {@code _trajectory} / {@code _time} /
* {@code _timespan} accessors that reduce a full sequence to a single
* derived value.
Expand All @@ -61,7 +61,7 @@
* per-window.
*
* <p><b>Typical usage</b> — per-vehicle per-tumbling-window
* trajectory length via {@code MeosOpsTemporal.temporal_length} (tier
* trajectory length via {@code MeosOpsTPoint.tpoint_length} (tier
* = {@code windowed}):
*
* <pre>{@code
Expand All @@ -73,7 +73,7 @@
* .process(new MeosWindowedAggregate<Integer, VehiclePoint, VehicleLength, TimeWindow>(
* (window, events, ctx) -> {
* Pointer trajectory = buildTrajectoryFromPoints(events); // adopter helper
* double length = MeosOpsTemporal.temporal_length(trajectory);
* double length = MeosOpsTPoint.tpoint_length(trajectory);
* return new VehicleLength(ctx.getCurrentKey(), window.getStart(), length);
* }));
* }</pre>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,11 +38,11 @@ The pattern is the same across all four tiers:
```java
// 1. Pick the generated MeosOps method
// (Javadoc tier marker tells you which wiring to use)
boolean overlap = MeosOpsTBox.overlaps_tbox_tbox(boxA, boxB); // tier = stateless
boolean overlap = MeosOpsFreeCore.overlaps_tbox_tbox(boxA, boxB); // tier = stateless

// 2. Wrap with the matching wiring
MeosStatelessFilter<TboxPair> filter = MeosStatelessFilter.fromIntPredicate(
pair -> MeosOpsTBox.overlaps_tbox_tbox(pair.a, pair.b));
pair -> MeosOpsFreeCore.overlaps_tbox_tbox(pair.a, pair.b));

// 3. Apply to the DataStream
DataStream<TboxPair> overlapping = stream.filter(filter);
Expand All @@ -61,9 +61,9 @@ through a 3-stage DataStream pipeline using two of the generated
facades wired through `MeosStatelessMap` + `MeosStatelessFilter`:

1. Parse a stream of TBox WKT strings via
`MeosOpsFreeCore.tbox_in` (io-meta, no state).
`MeosOpsTBox.tbox_in` (io-meta, no state).
2. Filter to those overlapping a fixed query box via
`MeosOpsTBox.overlaps_tbox_tbox` (stateless predicate).
`MeosOpsFreeCore.overlaps_tbox_tbox` (stateless predicate).
3. Serialize each survivor to hex-WKB via
`MeosOpsTBox.tbox_as_hexwkb` (io-meta, no state).

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@
* <li>Parses each into a JMEOS {@code Pointer} via
* {@code MeosOpsTBox.tbox_in} (tier = {@code io-meta}).</li>
* <li>Filters to those that overlap with a fixed query TBox via
* {@code MeosOpsTBox.overlaps_tbox_tbox} wrapped as a
* {@code MeosOpsFreeCore.overlaps_tbox_tbox} wrapped as a
* {@link MeosStatelessFilter} (tier = {@code stateless}).</li>
* <li>Maps each surviving TBox to its serialized WKB hex via
* {@code MeosOpsTBox.tbox_as_hexwkb} wrapped as a
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -59,9 +59,9 @@ static void finalizeMeos() {

@Test
void cbuffer() {
Pointer cb = MeosOpsFreeCbuffer.cbuffer_make(MeosOpsFreeGeo.geom_in("POINT(1 1)", 0), 0.5);
Pointer cb = MeosOpsCbuffer.cbuffer_make(MeosOpsGeometry.geom_in("POINT(1 1)", 0), 0.5);
assertNotNull(cb);
assertEquals(0.5, MeosOpsFreeCbuffer.cbuffer_radius(cb), 1e-9);
assertNotNull(MeosOpsFreeCbuffer.cbuffer_out(cb, 6));
assertEquals(0.5, MeosOpsCbuffer.cbuffer_radius(cb), 1e-9);
assertNotNull(MeosOpsCbuffer.cbuffer_out(cb, 6));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -86,8 +86,8 @@ void geoStbox() {

@Test
void geoGeometry() {
Pointer geom = MeosOpsFreeGeo.geom_in("POINT(1 1)", 0);
Pointer geom = MeosOpsGeometry.geom_in("POINT(1 1)", 0);
assertNotNull(geom);
assertTrue(MeosOpsFreeGeo.geo_as_text(geom, 6).toUpperCase().contains("POINT"));
assertTrue(MeosOpsGeo.geo_as_text(geom, 6).toUpperCase().contains("POINT"));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -58,9 +58,9 @@ static void finalizeMeos() {

@Test
void npoint() {
Pointer np = MeosOpsFreeNpoint.npoint_make(1, 0.5);
Pointer np = MeosOpsNpoint.npoint_make(1, 0.5);
assertNotNull(np);
assertEquals(1, MeosOpsFreeNpoint.npoint_route(np));
assertEquals(0.5, MeosOpsFreeNpoint.npoint_position(np), 1e-9);
assertEquals(1, MeosOpsNpoint.npoint_route(np));
assertEquals(0.5, MeosOpsNpoint.npoint_position(np), 1e-9);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -59,9 +59,9 @@ static void finalizeMeos() {

@Test
void pose() {
Pointer pose = MeosOpsFreePose.pose_in("Pose(Point(1 1), 0.5)");
Pointer pose = MeosOpsPose.pose_in("Pose(Point(1 1), 0.5)");
assertNotNull(pose);
assertNotNull(MeosOpsFreePose.pose_out(pose, 6));
assertEquals(0.5, MeosOpsFreePose.pose_yaw(pose), 1e-9);
assertNotNull(MeosOpsPose.pose_out(pose, 6));
assertEquals(0.5, MeosOpsPose.pose_yaw(pose), 1e-9);
}
}
37 changes: 35 additions & 2 deletions binding/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -12,8 +12,9 @@

<artifactId>mobility-flink-binding</artifactId>

<!-- The Flink MEOS binding: the org.mobilitydb.meos.MeosOps* facades generated from the
JMEOS catalog at build time, plus the generic Flink DataStream wiring layer over them.
<!-- The Flink MEOS binding: the org.mobilitydb.meos.MeosOps* facades and the
org.mobilitydb.flink.sql SQL surface, both generated from the JMEOS catalog at build
time, plus the generic Flink DataStream wiring layer over the facades.
A pure projection of the MEOS surface — no application or benchmark code. -->

<dependencies>
Expand All @@ -32,6 +33,17 @@
<artifactId>flink-clients</artifactId>
<version>${flink.version}</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-api-java</artifactId>
<version>${flink.version}</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-planner_2.12</artifactId>
<version>${flink.version}</version>
<scope>test</scope>
</dependency>

<dependency>
<groupId>org.jmeos</groupId>
Expand Down Expand Up @@ -101,6 +113,26 @@
</arguments>
</configuration>
</execution>
<!-- The org.mobilitydb.flink.sql SQL surface, from the same catalog and jar. -->
<execution>
<id>generate-sql-surface</id>
<phase>generate-sources</phase>
<goals><goal>exec</goal></goals>
<configuration>
<executable>python3</executable>
<arguments>
<argument>${project.basedir}/../tools/codegen_jvm.py</argument>
<argument>--engine</argument>
<argument>flink-sql</argument>
<argument>--catalog</argument>
<argument>${project.basedir}/../tools/meos-idl.json</argument>
<argument>--jar</argument>
<argument>${settings.localRepository}/org/jmeos/meos/${jmeos.version}/meos-${jmeos.version}.jar</argument>
<argument>--out</argument>
<argument>${project.build.directory}/generated-sql</argument>
</arguments>
</configuration>
</execution>
</executions>
</plugin>
<plugin>
Expand All @@ -115,6 +147,7 @@
<configuration>
<sources>
<source>${project.build.directory}/generated-facades/src/main/java</source>
<source>${project.build.directory}/generated-sql/src/main/java</source>
</sources>
</configuration>
</execution>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,7 @@
*
* <p><b>Typical usage</b> — per-vehicle-pair "did they come within
* 100m of each other in the last 5 minutes?" via
* {@code MeosOpsTGeo.edwithin_tgeo_tgeo} (tier = {@code cross-stream}):
* {@code MeosOpsFreeGeo.edwithin_tgeo_tgeo} (tier = {@code cross-stream}):
*
* <pre>{@code
* KeyedStream<VehiclePosition, Integer> a = streamA.keyBy(VehiclePosition::regionId);
Expand All @@ -64,7 +64,7 @@
* (left, right, ctx) -> {
* Pointer leftT = left.toTGeoPointer();
* Pointer rightT = right.toTGeoPointer();
* if (MeosOpsTGeo.edwithin_tgeo_tgeo(leftT, rightT, 100.0) != 0) {
* if (MeosOpsFreeGeo.edwithin_tgeo_tgeo(leftT, rightT, 100.0) != 0) {
* return new MeetingEvent(left.id(), right.id(), ctx.getLeftTimestamp());
* }
* return null; // no output for non-matches
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -44,14 +44,14 @@
* independently.
*
* <p><b>Typical usage</b>: scalar-predicate filter against the
* generated {@code MeosOpsTBox.overlaps_tbox_tbox} (tier =
* generated {@code MeosOpsFreeCore.overlaps_tbox_tbox} (tier =
* {@code stateless}):
*
* <pre>{@code
* DataStream<TbiePair> in = ...;
* DataStream<TbiePair> overlapping = in.filter(
* new MeosStatelessFilter<>(
* pair -> MeosOpsTBox.overlaps_tbox_tbox(pair.a, pair.b)));
* pair -> MeosOpsFreeCore.overlaps_tbox_tbox(pair.a, pair.b)));
* }</pre>
*
* <p>For int-coded predicates (JMEOS returns {@code int} for some MEOS
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -52,14 +52,14 @@
*
* <p><b>Typical usage</b>: register a stateless MEOS predicate / arithmetic
* call as a per-event map step in a DataStream pipeline. Example with
* the generated {@code MeosOpsTBox.overlaps_tbox_tbox} (tier =
* the generated {@code MeosOpsFreeCore.overlaps_tbox_tbox} (tier =
* {@code stateless}, per the codegen manifest):
*
* <pre>{@code
* DataStream<TbiePair> in = ...; // (tboxA, tboxB)
* DataStream<Boolean> overlap = in.map(
* new MeosStatelessMap<>(
* pair -> MeosOpsTBox.overlaps_tbox_tbox(pair.a, pair.b)));
* pair -> MeosOpsFreeCore.overlaps_tbox_tbox(pair.a, pair.b)));
* }</pre>
*
* <p><b>Tier coverage</b>: as of the codegen state on the parent PR,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,8 +38,8 @@
*
* <p>The {@code windowed} tier is "output cardinality changes; needs a
* window". The canonical examples are
* {@code temporal_length(tgeo)} (one length per trajectory window),
* {@code temporal_twavg(tnumber)} (one time-weighted average per
* {@code tpoint_length(tpoint)} (one length per trajectory window),
* {@code tnumber_twavg(tnumber)} (one time-weighted average per
* window), and the per-class {@code _trajectory} / {@code _time} /
* {@code _timespan} accessors that reduce a full sequence to a single
* derived value.
Expand All @@ -61,7 +61,7 @@
* per-window.
*
* <p><b>Typical usage</b> — per-vehicle per-tumbling-window
* trajectory length via {@code MeosOpsTemporal.temporal_length} (tier
* trajectory length via {@code MeosOpsTPoint.tpoint_length} (tier
* = {@code windowed}):
*
* <pre>{@code
Expand All @@ -73,7 +73,7 @@
* .process(new MeosWindowedAggregate<Integer, VehiclePoint, VehicleLength, TimeWindow>(
* (window, events, ctx) -> {
* Pointer trajectory = buildTrajectoryFromPoints(events); // adopter helper
* double length = MeosOpsTemporal.temporal_length(trajectory);
* double length = MeosOpsTPoint.tpoint_length(trajectory);
* return new VehicleLength(ctx.getCurrentKey(), window.getStart(), length);
* }));
* }</pre>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,11 +38,11 @@ The pattern is the same across all four tiers:
```java
// 1. Pick the generated MeosOps method
// (Javadoc tier marker tells you which wiring to use)
boolean overlap = MeosOpsTBox.overlaps_tbox_tbox(boxA, boxB); // tier = stateless
boolean overlap = MeosOpsFreeCore.overlaps_tbox_tbox(boxA, boxB); // tier = stateless

// 2. Wrap with the matching wiring
MeosStatelessFilter<TboxPair> filter = MeosStatelessFilter.fromIntPredicate(
pair -> MeosOpsTBox.overlaps_tbox_tbox(pair.a, pair.b));
pair -> MeosOpsFreeCore.overlaps_tbox_tbox(pair.a, pair.b));

// 3. Apply to the DataStream
DataStream<TboxPair> overlapping = stream.filter(filter);
Expand All @@ -61,9 +61,9 @@ through a 3-stage DataStream pipeline using two of the generated
facades wired through `MeosStatelessMap` + `MeosStatelessFilter`:

1. Parse a stream of TBox WKT strings via
`MeosOpsFreeCore.tbox_in` (io-meta, no state).
`MeosOpsTBox.tbox_in` (io-meta, no state).
2. Filter to those overlapping a fixed query box via
`MeosOpsTBox.overlaps_tbox_tbox` (stateless predicate).
`MeosOpsFreeCore.overlaps_tbox_tbox` (stateless predicate).
3. Serialize each survivor to hex-WKB via
`MeosOpsTBox.tbox_as_hexwkb` (io-meta, no state).

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@
* <li>Parses each into a JMEOS {@code Pointer} via
* {@code MeosOpsTBox.tbox_in} (tier = {@code io-meta}).</li>
* <li>Filters to those that overlap with a fixed query TBox via
* {@code MeosOpsTBox.overlaps_tbox_tbox} wrapped as a
* {@code MeosOpsFreeCore.overlaps_tbox_tbox} wrapped as a
* {@link MeosStatelessFilter} (tier = {@code stateless}).</li>
* <li>Maps each surviving TBox to its serialized WKB hex via
* {@code MeosOpsTBox.tbox_as_hexwkb} wrapped as a
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -59,9 +59,9 @@ static void finalizeMeos() {

@Test
void cbuffer() {
Pointer cb = MeosOpsFreeCbuffer.cbuffer_make(MeosOpsFreeGeo.geom_in("POINT(1 1)", 0), 0.5);
Pointer cb = MeosOpsCbuffer.cbuffer_make(MeosOpsGeometry.geom_in("POINT(1 1)", 0), 0.5);
assertNotNull(cb);
assertEquals(0.5, MeosOpsFreeCbuffer.cbuffer_radius(cb), 1e-9);
assertNotNull(MeosOpsFreeCbuffer.cbuffer_out(cb, 6));
assertEquals(0.5, MeosOpsCbuffer.cbuffer_radius(cb), 1e-9);
assertNotNull(MeosOpsCbuffer.cbuffer_out(cb, 6));
}
}
Loading
Loading