From 46e0e844dfecd45dedfe4bdac3619ff18bbd01b7 Mon Sep 17 00:00:00 2001 From: Farid Date: Mon, 7 Sep 2026 21:15:48 +0200 Subject: [PATCH] Wrap nullable struct fields in an Avro union in to_avro() _parse_struct() wraps nullable array and map fields in a [type, null] Avro union via _is_nullable(), but nullable struct fields (e.g. the 'candidate' field of a ZTF alert) were passed through as plain records. Spark's own to_avro() serializer always writes a union discriminator byte for nullable fields regardless of the published schema, so any strict Avro reader following the schema this function produces fails to decode nullable struct fields with an out-of-range/index error. Downstream, this has been worked around ad hoc by patching the generated schema JSON after the fact (see astrolabsoftware/ztf.fink-portal.org's spark_ztf_inference_feed.py); fixing it here removes the need for that per-caller patch. --- fink_utils/spark/schema_converter.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/fink_utils/spark/schema_converter.py b/fink_utils/spark/schema_converter.py index 53dc942..6dfffe1 100644 --- a/fink_utils/spark/schema_converter.py +++ b/fink_utils/spark/schema_converter.py @@ -181,7 +181,8 @@ def _parse_struct( elif "type" in field["type"]: subData = field["type"] if subData["type"] == "struct": - outField["type"] = _parse_struct(subData, field["name"]) + avro_type = _parse_struct(subData, field["name"]) + outField["type"] = _is_nullable(field, avro_type) elif subData["type"] == "array": avro_type = _parse_array(subData, field["name"]) outField["type"] = _is_nullable(field, avro_type)