Skip to content

Run the aggregates of the catalog from the roles it states on both typed SQL surfaces - #167

Merged
estebanzimanyi merged 1 commit into
MobilityDB:mainfrom
estebanzimanyi:codegen/catalog-aggregates
Oct 10, 2026
Merged

estebanzimanyi merged 1 commit into
MobilityDB:mainfrom
estebanzimanyi:codegen/catalog-aggregates

Conversation

@estebanzimanyi

Copy link
Copy Markdown
Member

An aggregate of the catalog's aggregates section runs from the public MEOS function of each role
it states: the transition folds a value into the state, the combine joins two partial states, the
final answers the result, and a partial state crosses between workers as the bytes the serialize
function writes and the deserialize function reads, or, for a state of a SQL type (the box of
extent), as the form that type's codec writes. _catalog_aggregates of tools/codegen_jvm.py
chooses, for both engines, the aggregates taking one value with a combine and a state either
internal with a serialize and a deserialize function or of a SQL type the surface holds, and
MeosAggregate holds their roles: a combine answers in the place of its first state, returning
it or freeing it for a larger one, except that a skip-list combine answers its second state when
the first is empty, and the state it does not answer is released by releasing what its final
answers, the final consuming its state.

The typed Flink surface registers each as an AggregateFunction whose accumulator, MeosAggState,
holds the bytes of the state and no native pointer; the typed Spark surface registers each
through MeosAggregates, a partition keeping its state between two values and writing it when it
leaves the executor. Both register extent, tCount, tSum, tAvg, tCentroid, tdensity, tnpoints,
mergeAgg, tMinAgg, tMaxAgg, tAndAgg, tOrAgg, setUnionAgg, spanUnionAgg and spansetUnionAgg, a
name a scalar of the surface carries going to the scalar as main's generator names the
aggregates, setUnion, spanUnion and spansetUnion being the unions of a value and a set. The
Spark UDF arm takes its temporal aggregates from the same section, a partial state crossing as
the bytes of its serialize function, so the pairing of temporal_tagg_finalfn with
temporal_to_taggstate leaves both arms and temporal_to_taggstate joins the gaps ledger.

Witness: over MobilityDB bdc4e41234 main's generator registers no aggregate on the typed Flink
surface, and on the typed Spark surface and the Spark UDF arm 7 whose final function is
temporal_tagg_finalfn, a partial state crossing as the temporal value that function answers,
which carries neither the sums and counts of tAvg nor the coordinates of tCentroid; MobilitySpark
fails on "Cannot resolve routine tAvg" and MobilityFlink's test does not compile, the surface
holding no class TAvg.

Measured over MobilityDB bdc4e41234 with the catalog of MEOS-API d259ddbfea: both typed surfaces
register 15 aggregates of 141 overloads besides the 1163 SQL functions and 8931 overloads main's
generator registers; the 67 overloads of more than one argument (wCount, appendInstant), the 40
without a combine (appendInstant, appendSequence) and the 12 over a value of no surface type
(setUnion over an integer) are counted as left out. The Spark UDF arm registers 11 temporal
aggregates against 7, adding tAvg, tCentroid, tdensity and tnpoints, and the gaps ledger loses
the 11 role functions it reaches and gains temporal_to_taggstate. MobilitySpark passes its 38
tests and MobilityFlink its GeneratedSqlSurfaceTest with assertions that tCount, tAvg,
tCentroid, setUnionAgg, spanUnionAgg and extent answer what MobilityDB answers over values in
partitions of their own, in a streaming job and in a two-phase batch job, and that tAvg joins
two partial states one of which crosses the accumulator's serializer.

Why: an aggregate SQL states answers on both engines as PostgreSQL answers it, from the roles
MEOS carries, with no hand rule pairing two functions.

…ped SQL surfaces

An aggregate of the catalog's aggregates section runs from the public MEOS function of each role
it states: the transition folds a value into the state, the combine joins two partial states, the
final answers the result, and a partial state crosses between workers as the bytes the serialize
function writes and the deserialize function reads, or, for a state of a SQL type (the box of
extent), as the form that type's codec writes. _catalog_aggregates of tools/codegen_jvm.py
chooses, for both engines, the aggregates taking one value with a combine and a state either
internal with a serialize and a deserialize function or of a SQL type the surface holds, and
MeosAggregate holds their roles: a combine answers in the place of its first state, returning
it or freeing it for a larger one, except that a skip-list combine answers its second state when
the first is empty, and the state it does not answer is released by releasing what its final
answers, the final consuming its state.

The typed Flink surface registers each as an AggregateFunction whose accumulator, MeosAggState,
holds the bytes of the state and no native pointer; the typed Spark surface registers each
through MeosAggregates, a partition keeping its state between two values and writing it when it
leaves the executor. Both register extent, tCount, tSum, tAvg, tCentroid, tdensity, tnpoints,
mergeAgg, tMinAgg, tMaxAgg, tAndAgg, tOrAgg, setUnionAgg, spanUnionAgg and spansetUnionAgg, a
name a scalar of the surface carries going to the scalar as main's generator names the
aggregates, setUnion, spanUnion and spansetUnion being the unions of a value and a set. The
Spark UDF arm takes its temporal aggregates from the same section, a partial state crossing as
the bytes of its serialize function, so the pairing of temporal_tagg_finalfn with
temporal_to_taggstate leaves both arms and temporal_to_taggstate joins the gaps ledger.

Witness: over MobilityDB bdc4e41234 main's generator registers no aggregate on the typed Flink
surface, and on the typed Spark surface and the Spark UDF arm 7 whose final function is
temporal_tagg_finalfn, a partial state crossing as the temporal value that function answers,
which carries neither the sums and counts of tAvg nor the coordinates of tCentroid; MobilitySpark
fails on "Cannot resolve routine `tAvg`" and MobilityFlink's test does not compile, the surface
holding no class TAvg.

Measured over MobilityDB bdc4e41234 with the catalog of MEOS-API d259ddbfea: both typed surfaces
register 15 aggregates of 141 overloads besides the 1163 SQL functions and 8931 overloads main's
generator registers; the 67 overloads of more than one argument (wCount, appendInstant), the 40
without a combine (appendInstant, appendSequence) and the 12 over a value of no surface type
(setUnion over an integer) are counted as left out. The Spark UDF arm registers 11 temporal
aggregates against 7, adding tAvg, tCentroid, tdensity and tnpoints, and the gaps ledger loses
the 11 role functions it reaches and gains temporal_to_taggstate. MobilitySpark passes its 38
tests and MobilityFlink its GeneratedSqlSurfaceTest with assertions that tCount, tAvg,
tCentroid, setUnionAgg, spanUnionAgg and extent answer what MobilityDB answers over values in
partitions of their own, in a streaming job and in a two-phase batch job, and that tAvg joins
two partial states one of which crosses the accumulator's serializer.

Why: an aggregate SQL states answers on both engines as PostgreSQL answers it, from the roles
MEOS carries, with no hand rule pairing two functions.
@estebanzimanyi
estebanzimanyi merged commit 66d285b into MobilityDB:main Oct 10, 2026
2 checks passed
@estebanzimanyi
estebanzimanyi deleted the codegen/catalog-aggregates branch October 10, 2026 11:13
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant