diff --git a/CHANGELOG.md b/CHANGELOG.md index 117796a..a1535b8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -19,6 +19,10 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0. ### Changed +- GeoParquet files get at least ten row groups: a file under 250,000 rows is written in smaller + groups instead of 25,000 rows each, so every group covers a small area and a reader can + skip it. A group has at least 2,048 rows, so a file under 20,480 rows gets fewer groups. + Larger files are unchanged. - **Breaking:** `vec merge` keeps all properties by default, `--include` restricts them to the core properties plus the given ones and `--exclude` removes any property except for the geometry. It reports required properties that are not included, also collection-only diff --git a/tests/test_convert_duckdb.py b/tests/test_convert_duckdb.py index 91cc5b8..1b6cd46 100644 --- a/tests/test_convert_duckdb.py +++ b/tests/test_convert_duckdb.py @@ -9,6 +9,7 @@ import shapely from loguru import logger +from vecorel_cli.conversion import duckdb as duckdb_module from vecorel_cli.conversion.base import BaseConverter from vecorel_cli.conversion.duckdb import DuckDBBaseConverter, _constants_table from vecorel_cli.validate import ValidateData @@ -886,3 +887,31 @@ def test_merge_parquet_keeps_a_crs_duckdb_would_drop(tmp_folder): crs = geo["columns"][geo["primary_column"]]["crs"] assert crs is not None, "the merged file lost the CRS of its parts" assert "3857" in json.dumps(crs) + + +@pytest.mark.parametrize("spatial_order", [True, False]) +def test_duckdb_converter_gives_a_small_file_enough_row_groups( + tmp_folder, monkeypatch, spatial_order +): + from vecorel_cli.encoding.geoparquet import GeoParquet + + # DuckDB builds row groups from 2048-row vectors, so the file has to be large enough + n = 25_000 + monkeypatch.setattr(GeoParquet, "row_group_size", 20_000) + if not spatial_order: + # as for a CRS without an area of use + monkeypatch.setattr(duckdb_module, "hilbert_reference_bounds", lambda *args: None) + gdf = gpd.GeoDataFrame( + { + "id": [str(i) for i in range(n)], + "name": ["a"] * n, + "geometry": [shapely.box(i * 1e-5, 0, i * 1e-5 + 1e-5, 1e-5) for i in range(n)], + }, + crs="EPSG:4326", + ) + src = tmp_folder / "many.parquet" + gdf.to_parquet(src) + dest = tmp_folder / "many-out.parquet" + Converter().convert(dest, input_files={str(src): src.name}) + + assert pq.ParquetFile(dest).metadata.num_row_groups >= GeoParquet.min_row_groups diff --git a/tests/test_encoding_geoparquet.py b/tests/test_encoding_geoparquet.py index 567c8a6..60a3e4c 100644 --- a/tests/test_encoding_geoparquet.py +++ b/tests/test_encoding_geoparquet.py @@ -252,3 +252,28 @@ def test_read_hydrates_array_and_object_constants(tmp_parquet_file): data = GeoParquet(tmp_parquet_file).read(hydrate=True) assert list(data["tags"]) == [["a", "b"]] * 3 assert list(data["attrs"]) == [{"k": "v"}] * 3 + + +def test_a_small_file_still_gets_enough_row_groups(tmp_parquet_file, monkeypatch): + # with too few groups each one spans most of the file, so a reader can skip none + from geopandas import GeoDataFrame + from shapely.geometry import box + + monkeypatch.setattr(GeoParquet, "row_group_size", 100) + monkeypatch.setattr(GeoParquet, "min_row_group_rows", 10) + gdf = GeoDataFrame( + {"id": [str(i) for i in range(500)], "geometry": [box(i, 0, i + 1, 1) for i in range(500)]}, + crs="EPSG:4326", + ) + gp = GeoParquet(tmp_parquet_file) + gp.set_collection( + { + "schemas": {"c": ["https://vecorel.org/specification/v0.1.0/schema.yaml"]}, + "collection": "c", + } + ) + gp.write(gdf) + + assert pq.ParquetFile(tmp_parquet_file).metadata.num_row_groups == GeoParquet.min_row_groups + assert GeoParquet.row_group_size_for(5_000) == 100 # a large file keeps the cap + assert GeoParquet.row_group_size_for(20) == 10 # and a tiny one the floor diff --git a/vecorel_cli/conversion/duckdb.py b/vecorel_cli/conversion/duckdb.py index be741b6..0393115 100644 --- a/vecorel_cli/conversion/duckdb.py +++ b/vecorel_cli/conversion/duckdb.py @@ -561,6 +561,8 @@ def in_collections(cids): # GeoDataFrame-based codepath with pq.ParquetFile(output_file) as pf: meta = pf.schema_arrow.metadata or {} + final_group_size = GeoParquet.row_group_size_for(pf.metadata.num_rows) + keys_path, is_sorted = None, True if b"geo" in meta: geo = json.loads(meta[b"geo"]) primary = geo["primary_column"] @@ -570,20 +572,22 @@ def in_collections(cids): self.warning("CRS declares no area of use; skipping spatial ordering") else: keys_path, is_sorted = self._write_hilbert_keys(output_file, primary, bounds) - try: - if not is_sorted: - self._sort_output( - con, - output_file, - keys_path, - compression, - collection_json, - row_group_size, - ) - self.info("Sorted output into Hilbert order") - finally: - if os.path.exists(keys_path): - os.unlink(keys_path) + try: + # a small file is rewritten for its row groups even when already sorted + if not is_sorted or final_group_size != row_group_size: + self._rewrite_output( + con, + output_file, + keys_path, + compression, + collection_json, + final_group_size, + ) + if not is_sorted: + self.info("Sorted output into Hilbert order") + finally: + if keys_path and os.path.exists(keys_path): + os.unlink(keys_path) gp = GeoParquet(output_file) gp.set_collection(collection) @@ -838,26 +842,27 @@ def _write_hilbert_keys(self, output_file, primary, bounds): raise return keys_path, is_sorted - # Rewrites the file in the order given by the Hilbert keys. - # DuckDB sorts externally (spilling to disk if needed), so this works for - # datasets that don't fit into memory. - def _sort_output( + # Rewrites the file in the given row group size, and in the order given by the + # Hilbert keys if there are any. DuckDB sorts externally (spilling to disk if + # needed), so this works for datasets that don't fit into memory. + def _rewrite_output( self, con, output_file, keys_path, compression, collection_json, row_group_size ): directory = os.path.dirname(output_file) or "." tmp_path = None + query = f"SELECT d.* FROM read_parquet({_sql_path(output_file)}) d" + if keys_path: + query += f""" + POSITIONAL JOIN read_parquet({_sql_path(keys_path)}) k + ORDER BY k.hilbert, k.ordinal + """ try: with NamedTemporaryFile("wb", delete=False, dir=directory, suffix=".parquet") as tmp: tmp_path = tmp.name con.execute( f""" - COPY ( - SELECT d.* - FROM read_parquet({_sql_path(output_file)}) d - POSITIONAL JOIN read_parquet({_sql_path(keys_path)}) k - ORDER BY k.hilbert, k.ordinal - ) TO ? ( + COPY ({query}) TO ? ( FORMAT parquet, ROW_GROUP_SIZE {row_group_size}, compression ?, diff --git a/vecorel_cli/encoding/geoparquet.py b/vecorel_cli/encoding/geoparquet.py index d359f49..ed7c768 100644 --- a/vecorel_cli/encoding/geoparquet.py +++ b/vecorel_cli/encoding/geoparquet.py @@ -35,6 +35,19 @@ class GeoParquet(BaseEncoding): ext = [".parquet", ".geoparquet"] media_type = "application/vnd.apache.parquet" row_group_size = 25000 + # with fewer groups each one spans most of a small file, so a reader can skip none + min_row_groups = 10 + # DuckDB writes row groups in whole 2048-row vectors, so both write paths use that unit + min_row_group_rows = 2048 + + @classmethod + def row_group_size_for(cls, num_rows: int) -> int: + """Rows per group: at most row_group_size, but small enough for min_row_groups groups.""" + unit = cls.min_row_group_rows + if num_rows >= cls.min_row_groups * cls.row_group_size: + return cls.row_group_size + wanted = num_rows // cls.min_row_groups // unit * unit + return max(unit, min(cls.row_group_size, wanted)) def __init__(self, file: Union[Path, URL, str]): super().__init__(file) @@ -268,7 +281,7 @@ def write( coerce_timestamps="ms", compression=compression, schema_version=geoparquet_version, - row_group_size=self.row_group_size, + row_group_size=self.row_group_size_for(len(data)), write_covering_bbox=bool(geoparquet_version != "1.0.0"), compression_level=compression_level, )