Repository navigation
Run the aggregates of the catalog from the roles it states on both typed SQL surfaces - #167
Merged
estebanzimanyi merged 1 commit intoOct 10, 2026
Conversation
…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.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
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 surfaceholding 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.