From cb65c8fbcde994c3d7f539f45398c37b7d64144b Mon Sep 17 00:00:00 2001 From: Esteban Zimanyi Date: Fri, 9 Oct 2026 19:10:33 +0200 Subject: [PATCH] Carry an array of fixed-size values as the contiguous C array MEOS reads A SQL array of a value whose C type is a fixed-size catalog struct reaches MEOS as the contiguous array a constructor reads, spanset_make(Span *spans, int count) and set(npoint[]) on npointset_make(const Npoint *values, int count) among them, once the catalog states the array in shape.inputArrays. _array_arg of tools/codegen_jvm.py decodes each element as it decodes the pointers of an array of values and MeosSqlRuntime.Inputs.structs copies each into the C array, each of the size #struct_layout of tools/codegen_spark_udfs.py gives the struct, the layout SqlModel reads once for both engines. The typed Flink and Spark surfaces register spanset(intspan[]), spanset(bigintspan[]), spanset(floatspan[]), spanset(datespan[]), spanset(tstzspan[]), set(npoint[]) and set(cbuffer[]). The Spark UDF arm refuses such a parameter under its C name, since a UDF decoding one value would hand MEOS the count the caller gives: spanset_make, stboxarr_round, npointset_make and cbufferset_make have no C-named UDF and the gaps ledger lists them as array: *. Witness: over MobilityDB 2a1d52a7ee both typed surfaces refuse the seven signatures as arity:sql, and the C-named Spark UDF spanset_make(spans, count) decodes one span and passes count to spanset_make, so a count above 1 reads past the value. Measured over MobilityDB 2a1d52a7ee with the catalog of MEOS-API catalog/struct-value-input-arrays: both typed surfaces register 1154 SQL functions and 8754 overloads, against 8747 with main's generator, the seven passing Span, Npoint and Cbuffer elements of 24, 16 and 32 bytes, the sizes gcc gives the structs of the installed headers; the Spark UDF surface registers the same functions as main's but those four, and the gaps step reports none new with the ledger at 764 functions. MobilitySpark passes its 35 tests and MobilityFlink its 21 binding and 12 benchmark tests with assertions that spanset over two int spans answers {[1, 3), [5, 7)} and set over three npoints, one repeated, and over two cbuffers answers 2 values; with main's generator MobilitySpark fails on "function spanset(array) does not exist". Over the catalog of MEOS-API master e84319b both typed surfaces register the 8747 overloads main's generator does and the Spark UDF surface the same functions, the generated sources differing by Inputs.structs alone. Why: a SQL array of values answers on every engine as PostgreSQL answers it, and a binding hands MEOS the array it reads, never one value standing for several. --- tools/codegen_jvm.py | 45 ++++++++++++++++++++++++++++--------- tools/codegen_spark_udfs.py | 8 +++++++ tools/spark-udf-gaps.txt | 6 +++-- 3 files changed, 47 insertions(+), 12 deletions(-) diff --git a/tools/codegen_jvm.py b/tools/codegen_jvm.py index fbba24cf8..0b478ea63 100644 --- a/tools/codegen_jvm.py +++ b/tools/codegen_jvm.py @@ -629,6 +629,13 @@ def __init__(self, cat, jmeos, pkg=SQL_PKG, engine='flink'): self.jmeos = jmeos self.fns = cat['functions'] self.by_name = {f['name']: f for f in self.fns} + # The catalog struct layouts, under #struct_layout of codegen_spark_udfs.py, the rule the + # Spark arm sizes them by: the rows of a set-returning signature and a contiguous array + # of structs (#_array_arg) read them, on both engines. + spark = _spark_module() + spark.STRUCTS.update({s['name']: s for s in cat.get('structs') or []}) + self.layout = lambda name: (spark.struct_layout(name) # noqa: E731 + if name in spark.STRUCTS else None) self.enc = cat.get('typeEncodings', {}) self.enums = {e['name'] for e in cat.get('enums', [])} # The value of each macro and enum member a wrapper binds by name. @@ -1025,6 +1032,15 @@ def _array_arg(m, elem, a, p, jt, name, temps): cls = f'{m.pkg}.types.{m.value_class[elem]}' temps.append((t, f'{name}, {cls}::decode', 'values')) return f'{cls}[]', t + # A contiguous array of structs (spanset_make reads Span *spans): each element decodes as + # above and its bytes are copied into the C array, each the size #struct_layout of + # codegen_spark_udfs.py gives the catalog struct. + lay = m.layout(_base(el['canonical'])) if c.count('*') == 1 else None + if elem in m.value_class and lay is not None \ + and m.sql_cbase.get(elem) == _base(el['canonical']): + cls = f'{m.pkg}.types.{m.value_class[elem]}' + temps.append((t, f'{name}, {cls}::decode, {lay[0]}', 'structs')) + return f'{cls}[]', t hit = SQL_ARRAY_SCALAR.get((elem, _base(el['c']))) \ or SQL_ARRAY_SCALAR.get((elem, _base(el['canonical']))) if hit and c.count('*') == 1: @@ -1175,6 +1191,7 @@ def _emit_eval(ov, defaults=None): L.append(' try {') for t, e, kind in temps: v = {'value': f'_in.value({e})', 'values': f'_in.values({e})', + 'structs': f'_in.structs({e})', 'buffer': f'_in.hold({e})'}[kind] L.append(f' Pointer {t} = {v};') L += [f' {s}' for s in body] @@ -1860,6 +1877,22 @@ def _value_class_src(m, sql): return b; } + /** The values decoded one by one by dec into the contiguous C array of structs of + * the given size MEOS reads, as {@link #values} decodes them into an array of + * pointers: each decoded value is kept for release and its bytes are copied in. */ + public Pointer structs(V[] vs, + java.util.function.Function dec, int size) { + Pointer b = hold(buffer(vs.length, size)); + for (int i = 0; i < vs.length; i++) { + Pointer p = value(dec.apply(element(vs, i))); + if (p == null) { + throw new IllegalArgumentException("an array element does not decode"); + } + p.transferTo(0, b, (long) i * size, size); + } + return b; + } + boolean holds(Pointer r) { for (int i = 0; i < n; i++) { if (owned[i] != null && owned[i].address() == r.address()) { @@ -2074,11 +2107,7 @@ def run_flink_sql(args): for sql in m.value_class: (root / 'types' / f'{m.value_class[sql]}.java').write_text(_value_class_src(m, sql)) - # The catalog struct layouts, under #struct_layout of codegen_spark_udfs.py, the rule the - # Spark arm sizes them by. - spark = _spark_module() - spark.STRUCTS.update({s['name']: s for s in cat.get('structs') or []}) - layout = lambda name: spark.struct_layout(name) if name in spark.STRUCTS else None # noqa: E731 + layout = m.layout names = defaultdict(list) # SQL name -> eval methods seen = defaultdict(set) @@ -2778,11 +2807,7 @@ def run_spark_sql(args): for sql in m.value_class: (root / 'types' / f'{m.value_class[sql]}.java').write_text(_spark_value_class_src(m, sql)) - # The catalog struct layouts, under #struct_layout of codegen_spark_udfs.py, as the flink-sql - # engine reads them for a set-returning signature's rows. - spark = _spark_module() - spark.STRUCTS.update({s['name']: s for s in cat.get('structs') or []}) - layout = lambda name: spark.struct_layout(name) if name in spark.STRUCTS else None # noqa: E731 + layout = m.layout names = defaultdict(list) # SQL name -> (eval lines, Spark arg types, ret, classes) seen = defaultdict(set) diff --git a/tools/codegen_spark_udfs.py b/tools/codegen_spark_udfs.py index a9ae710f9..55a2e126f 100644 --- a/tools/codegen_spark_udfs.py +++ b/tools/codegen_spark_udfs.py @@ -281,6 +281,14 @@ def supported(f): if r is None: b = base(f["returnType"]["canonical"]) return ("internal" if b in INTERNAL or b == "__INTERNAL__" else "ret:"+norm(f["returnType"]["canonical"])) + # An array the catalog names in shape.inputArrays is no single value, though a contiguous + # array of structs (spanset_make reads Span *spans) has the C type of one: a UDF decoding + # one value would hand MEOS `count` elements to read past it. The typed SQL surfaces carry + # such an array (#_array_arg of codegen_jvm.py); here it is refused. + arrays = {a["param"] for a in (f.get("shape") or {}).get("inputArrays") or ()} + for p in in_params: + if p["name"] in arrays and (arg_kind(p["canonical"]) or ("",))[0] == "ptr": + return "array:" + norm(p["canonical"]) for p in in_params: if arg_kind(p["canonical"]) is None: b = base(p["canonical"]) diff --git a/tools/spark-udf-gaps.txt b/tools/spark-udf-gaps.txt index 95e180c29..d4954e16c 100644 --- a/tools/spark-udf-gaps.txt +++ b/tools/spark-udf-gaps.txt @@ -52,7 +52,7 @@ cbuffer_as_ewkb array-or-out-param:size_out; unsupported-return:uint8_t * cbuffer_hash ret:uint32_t cbufferarr_round unsupported-return:struct Cbuffer ** cbufferarr_to_geom internal -cbufferset_make internal +cbufferset_make array:Cbuffer * cbufferset_value_n internal cmp_date_timestamp arg:Timestamp cmp_timestamp_date arg:Timestamp @@ -247,7 +247,7 @@ ne_timestamp_timestamptz arg:Timestamp ne_timestamptz_timestamp arg:Timestamp npoint_as_ewkb array-or-out-param:size_out; unsupported-return:uint8_t * npoint_hash ret:uint32_t -npointset_make internal +npointset_make array:Npoint * npointset_value_n internal nsegment_as_ewkb array-or-out-param:size_out; unsupported-return:uint8_t * overabove_tpcbox_tpcbox arg:TPCBox * @@ -428,6 +428,7 @@ setstate_deserialize array-or-out-param:bytes setstate_serialize array-or-out-param:size_out; unsupported-return:uint8_t * span_hash ret:uint32_t spanset_hash ret:uint32_t +spanset_make array:Span * spanset_spanarr internal spanset_spans array-or-out-param:count spanset_split_each_n_spans array-or-out-param:count @@ -465,6 +466,7 @@ stbox_tmax arg:TimestampTz * stbox_tmin arg:TimestampTz * stbox_to_box3d internal stbox_to_gbox internal +stboxarr_round array:STBox * taggstate_deserialize array-or-out-param:bytes; no-encoder:SkipList taggstate_serialize no-decoder:SkipList; array-or-out-param:size_out; unsupported-return:uint8_t * tbigint_time_boxes array-or-out-param:count