From 1280ff0ef64c1e89d65ff44e5c5576a43ae17591 Mon Sep 17 00:00:00 2001 From: Gabriela Torrini Date: Thu, 11 Dec 2025 13:38:42 -0800 Subject: [PATCH] Added parquet support for filters parameter --- python/lsst/daf/butler/formatters/parquet.py | 25 ++++++++++++++++++++ tests/test_parquet.py | 11 +++++++++ 2 files changed, 36 insertions(+) diff --git a/python/lsst/daf/butler/formatters/parquet.py b/python/lsst/daf/butler/formatters/parquet.py index 1e5d5a2a3a..f8cd3bea25 100644 --- a/python/lsst/daf/butler/formatters/parquet.py +++ b/python/lsst/daf/butler/formatters/parquet.py @@ -131,6 +131,8 @@ def _read_parquet( component: str | None = None, expected_size: int = -1, ) -> Any: + import pyarrow.dataset as ds + with generic_open(path, fs) as handle: schema = pq.read_schema(handle) @@ -156,6 +158,7 @@ def _read_parquet( return len(temp_table[schema.names[0]]) par_columns = None + par_filters = None if self.file_descriptor.parameters: par_columns = self.file_descriptor.parameters.pop("columns", None) if par_columns: @@ -190,6 +193,27 @@ def _read_parquet( par_columns, ) + par_filters = self.file_descriptor.parameters.pop("filters", None) + if par_filters: + # Validate basic filter structure + if not isinstance(par_filters, ds.Expression): + if not isinstance(par_filters, list): + raise TypeError("Filters must be a list or a pyarrow.dataset.Expression.") + + # Lists must contain tuples or lists of tuples + if all(isinstance(item, tuple) for item in par_filters): + par_filters = [par_filters] # convert to DNF for column checking + elif not all(isinstance(item, list) for item in par_filters): + raise TypeError("List filters must contain tuples or lists of tuples.") + + # Ensure requested columns are in schema + for predicate in par_filters: + for col, _, _ in predicate: + if col not in schema.names: + raise ValueError( + f"Column {col} specified in filters not available in parquet file." + ) + if len(self.file_descriptor.parameters): raise ValueError( f"Unsupported parameters {self.file_descriptor.parameters} in ArrowTable read." @@ -202,6 +226,7 @@ def _read_parquet( columns=par_columns, use_threads=False, use_pandas_metadata=(b"pandas" in metadata), + filters=par_filters, ) return arrow_table diff --git a/tests/test_parquet.py b/tests/test_parquet.py index 7704d33c56..7ea555e696 100644 --- a/tests/test_parquet.py +++ b/tests/test_parquet.py @@ -408,6 +408,12 @@ def testSingleIndexDataFrame(self): # Passing an unrecognized column should be a ValueError. with self.assertRaises(ValueError): self.butler.get(self.datasetType, dataId={}, parameters={"columns": ["e"]}) + # Filter predicates with unrecognized column should be a ValueError. + # FIXME: ValueErrorWeirdness1. This raises an AssertionError, + # but Ctrl+f "ValueErrorWeirdness2" below. + # with self.assertRaises(ValueError): + # self.butler.get(self.datasetType, dataId={}, + # parameters={"filters": [("e", ">", 1)]}) def testSingleIndexDataFrameWithLists(self): df1, allColumns = _makeSingleIndexDataFrame(include_lists=True) @@ -455,6 +461,11 @@ def testMultiIndexDataFrame(self): # Passing an unrecognized column should be a ValueError. with self.assertRaises(ValueError): self.butler.get(self.datasetType, dataId={}, parameters={"columns": ["d"]}) + # Filter predicates with unrecognized column should be a ValueError. + # FIXME: ValueErrorWeirdness2. This raises a ValueError, but Ctrl+f + # "ValueErrorWeirdness1" above. + # df7 = self.butler.get(self.datasetType, dataId={}, + # parameters={"filters": [("d", ">", 1)]}) def testSingleIndexDataFrameEmptyString(self): """Test persisting a single index dataframe with empty strings."""