From 2e8a29ac4413d1b8bdd57eca7c0f5c90f7730a3b Mon Sep 17 00:00:00 2001 From: "David H. Irving" Date: Fri, 26 Jun 2026 11:59:15 -0700 Subject: [PATCH 1/4] Add conversion of multiblock to parquet --- .../lsst/pipe/base/quantum_graph/_convert.py | 75 +++++++++++++++++++ 1 file changed, 75 insertions(+) create mode 100644 python/lsst/pipe/base/quantum_graph/_convert.py diff --git a/python/lsst/pipe/base/quantum_graph/_convert.py b/python/lsst/pipe/base/quantum_graph/_convert.py new file mode 100644 index 000000000..766008063 --- /dev/null +++ b/python/lsst/pipe/base/quantum_graph/_convert.py @@ -0,0 +1,75 @@ +import tempfile +import zipfile +from collections.abc import Iterator +from contextlib import contextmanager +from itertools import batched + +import pyarrow as pa +import zstandard +from duckdb import DuckDBPyConnection, connect + +from ._multiblock import MultiblockReader + + +def convert_multiblock_to_parquet(input_zip_path: str, output_zip_path: str) -> None: + with zipfile.ZipFile(input_zip_path, "r") as zf: + if (cdict_path := zipfile.Path(zf, "compression_dict")).exists(): + cdict = zstandard.ZstdCompressionDict(cdict_path.read_bytes()) + decompressor = zstandard.ZstdDecompressor(cdict) + with _initialize_duckdb_connection() as db: + for filename in ("quanta", "datasets", "metadata", "logs"): + reader = _get_record_batch_reader(filename, zf, decompressor) + db.execute( + """ + COPY (SELECT (data::JSON)::VARIANT as data FROM reader) + TO $output_filename (FORMAT 'parquet', COMPRESSION 'zstd') + """, + {"output_filename": f"{filename}.parquet"}, + ) + + +@contextmanager +def _initialize_duckdb_connection(memory_limit: str = "8GB") -> Iterator[DuckDBPyConnection]: + """Set up an in-memory DuckDB database, applying a memory limit and using + the system tempdir for its working directory. + """ + with ( + tempfile.TemporaryDirectory(suffix=".duckdb") as tmpdir, + connect( + config={ + # DuckDB will use up to 80% of system memory if you don't + # explicitly set it. + "memory_limit": memory_limit, + # DuckDB will use the current working directory if you don't + # set it, which is frequently a small quota-limited volume on + # USDF. + "temp_directory": tmpdir, + # For paranoia -- some of the scratch volumes at USDF are + # petabyte-sized. + "max_temp_directory_size": "128GB", + } + ) as conn, + ): + yield conn + + +def _fetch_json_bytes( + filename: str, zip: zipfile.ZipFile, decompressor: zstandard.ZstdDecompressor +) -> Iterator[bytes]: + for data in MultiblockReader.read_all_bytes_in_zip(zip, filename, int_size=8, page_size=8 * 1024 * 1024): + yield decompressor.decompress(data) + + +def _to_batches(schema: pa.Schema, it: Iterator[bytes]) -> Iterator[pa.RecordBatch]: + for batch in batched(it, 1000): + array = pa.array(batch) + yield pa.RecordBatch.from_arrays([array], schema=schema) + + +def _get_record_batch_reader( + filename: str, zip: zipfile.ZipFile, decompressor: zstandard.ZstdDecompressor +) -> pa.RecordBatchReader: + schema = pa.schema([("data", pa.string())]) + return pa.RecordBatchReader.from_batches( + schema, _to_batches(schema, _fetch_json_bytes(filename, zip, decompressor)) + ) From 4ef1fa10c1fbb7b21bddfe2d91f3a4b47dfe0168 Mon Sep 17 00:00:00 2001 From: "David H. Irving" Date: Fri, 26 Jun 2026 12:38:35 -0700 Subject: [PATCH 2/4] Augment output with UUID column and sort by it --- .../lsst/pipe/base/quantum_graph/_convert.py | 19 +++++++++++++++---- 1 file changed, 15 insertions(+), 4 deletions(-) diff --git a/python/lsst/pipe/base/quantum_graph/_convert.py b/python/lsst/pipe/base/quantum_graph/_convert.py index 766008063..35965ec63 100644 --- a/python/lsst/pipe/base/quantum_graph/_convert.py +++ b/python/lsst/pipe/base/quantum_graph/_convert.py @@ -17,11 +17,22 @@ def convert_multiblock_to_parquet(input_zip_path: str, output_zip_path: str) -> cdict = zstandard.ZstdCompressionDict(cdict_path.read_bytes()) decompressor = zstandard.ZstdDecompressor(cdict) with _initialize_duckdb_connection() as db: - for filename in ("quanta", "datasets", "metadata", "logs"): + for filename, id_field in ( + ("quanta", "quantum_id"), + ("datasets", "dataset_id"), + ("metadata", None), + ("logs", None), + ): reader = _get_record_batch_reader(filename, zf, decompressor) + id_column = f"data.{id_field}::UUID AS id," if id_field is not None else "" + order_by = "ORDER BY id" if id_field is not None else "" db.execute( - """ - COPY (SELECT (data::JSON)::VARIANT as data FROM reader) + f""" + COPY ( + SELECT {id_column} data FROM + (SELECT (raw_data::JSON)::VARIANT as data FROM reader) + {order_by} + ) TO $output_filename (FORMAT 'parquet', COMPRESSION 'zstd') """, {"output_filename": f"{filename}.parquet"}, @@ -69,7 +80,7 @@ def _to_batches(schema: pa.Schema, it: Iterator[bytes]) -> Iterator[pa.RecordBat def _get_record_batch_reader( filename: str, zip: zipfile.ZipFile, decompressor: zstandard.ZstdDecompressor ) -> pa.RecordBatchReader: - schema = pa.schema([("data", pa.string())]) + schema = pa.schema([("raw_data", pa.string())]) return pa.RecordBatchReader.from_batches( schema, _to_batches(schema, _fetch_json_bytes(filename, zip, decompressor)) ) From ffb4243a305507c0925f0ab6d054a4be0342427b Mon Sep 17 00:00:00 2001 From: "David H. Irving" Date: Fri, 26 Jun 2026 14:06:15 -0700 Subject: [PATCH 3/4] add runner script for conversion --- python/convert_multiblock.py | 13 +++++++++++++ 1 file changed, 13 insertions(+) create mode 100644 python/convert_multiblock.py diff --git a/python/convert_multiblock.py b/python/convert_multiblock.py new file mode 100644 index 000000000..0dcf1b658 --- /dev/null +++ b/python/convert_multiblock.py @@ -0,0 +1,13 @@ +import click + +from lsst.pipe.base.quantum_graph._convert import convert_multiblock_to_parquet + + +@click.command +@click.argument("path") +def run(path: str) -> None: + convert_multiblock_to_parquet(path, "output.zip") + + +if __name__ == "__main__": + run() From dad4ab2f2d0d16b0349d27c651ef50b12c2d9e8b Mon Sep 17 00:00:00 2001 From: "David H. Irving" Date: Fri, 26 Jun 2026 15:50:34 -0700 Subject: [PATCH 4/4] prevent OOM with larger files --- .../lsst/pipe/base/quantum_graph/_convert.py | 42 +++++++++---------- 1 file changed, 20 insertions(+), 22 deletions(-) diff --git a/python/lsst/pipe/base/quantum_graph/_convert.py b/python/lsst/pipe/base/quantum_graph/_convert.py index 35965ec63..0bc49aa6e 100644 --- a/python/lsst/pipe/base/quantum_graph/_convert.py +++ b/python/lsst/pipe/base/quantum_graph/_convert.py @@ -6,7 +6,7 @@ import pyarrow as pa import zstandard -from duckdb import DuckDBPyConnection, connect +from duckdb import ColumnExpression, DuckDBPyConnection, connect from ._multiblock import MultiblockReader @@ -16,31 +16,29 @@ def convert_multiblock_to_parquet(input_zip_path: str, output_zip_path: str) -> if (cdict_path := zipfile.Path(zf, "compression_dict")).exists(): cdict = zstandard.ZstdCompressionDict(cdict_path.read_bytes()) decompressor = zstandard.ZstdDecompressor(cdict) - with _initialize_duckdb_connection() as db: - for filename, id_field in ( - ("quanta", "quantum_id"), - ("datasets", "dataset_id"), - ("metadata", None), - ("logs", None), + with initialize_duckdb_connection() as db: + for filename, id_field, variant_encode, row_group_size in ( + ("quanta", "quantum_id", True, 122880), + ("datasets", "dataset_id", True, 122880), + # These don't really work as Parquet -- the individual values + # are too large, so we run out of memory easily when writing. + # ("metadata", None, False, 2048), + # ("logs", None, False, 2048), ): reader = _get_record_batch_reader(filename, zf, decompressor) - id_column = f"data.{id_field}::UUID AS id," if id_field is not None else "" - order_by = "ORDER BY id" if id_field is not None else "" - db.execute( - f""" - COPY ( - SELECT {id_column} data FROM - (SELECT (raw_data::JSON)::VARIANT as data FROM reader) - {order_by} - ) - TO $output_filename (FORMAT 'parquet', COMPRESSION 'zstd') - """, - {"output_filename": f"{filename}.parquet"}, - ) + sql = db.from_arrow(reader) + if variant_encode: + sql = sql.select("(data::JSON)::VARIANT as data") + if id_field is not None: + sql = sql.select( + ColumnExpression(f"data.{id_field}").cast("UUID").alias("id"), + "data", + ).order("id") + sql.to_parquet(f"{filename}.parquet", compression="zstd", row_group_size=row_group_size) @contextmanager -def _initialize_duckdb_connection(memory_limit: str = "8GB") -> Iterator[DuckDBPyConnection]: +def initialize_duckdb_connection(memory_limit: str = "16GB") -> Iterator[DuckDBPyConnection]: """Set up an in-memory DuckDB database, applying a memory limit and using the system tempdir for its working directory. """ @@ -80,7 +78,7 @@ def _to_batches(schema: pa.Schema, it: Iterator[bytes]) -> Iterator[pa.RecordBat def _get_record_batch_reader( filename: str, zip: zipfile.ZipFile, decompressor: zstandard.ZstdDecompressor ) -> pa.RecordBatchReader: - schema = pa.schema([("raw_data", pa.string())]) + schema = pa.schema([("data", pa.string())]) return pa.RecordBatchReader.from_batches( schema, _to_batches(schema, _fetch_json_bytes(filename, zip, decompressor)) )