From d6cfa194bd1ddf471ebb18c5b6dd628ce5dda23d Mon Sep 17 00:00:00 2001 From: Ivor Bosloper Date: Sat, 26 Sep 2026 11:36:16 +0200 Subject: [PATCH 1/6] merge_parquet: hydrate the constants the parts disagree on Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 4 +++ tests/test_convert_duckdb.py | 43 ++++++++++++++++++++++++++++++++ vecorel_cli/conversion/duckdb.py | 43 ++++++++++++++++++++++++++------ 3 files changed, 82 insertions(+), 8 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 776a999..1514909 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,10 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0. ## [Unreleased] +- Fix: `merge_parquet` no longer drops a property that each part kept as a constant in + its collection but on which the parts disagree; it becomes a column again, as in + `vec merge`. + ## [v0.3.1] - 2026-09-24 - Fix: ZIP archives compressed with Deflate64 are extracted through `inflate64` (a diff --git a/tests/test_convert_duckdb.py b/tests/test_convert_duckdb.py index ba182cf..35fe8f5 100644 --- a/tests/test_convert_duckdb.py +++ b/tests/test_convert_duckdb.py @@ -621,6 +621,49 @@ def test_merge_parquet_unions_parts_that_differ(tmp_folder): assert sorted(x for x in table.column("id").to_pylist()) == ["0", "1"] +def test_merge_parquet_hydrates_constants_the_parts_disagree_on(tmp_folder): + """Each part keeps its constants in its collection; one the parts disagree on + goes back into the rows, as `vec merge` does, and a shared one stays put.""" + parts = [] + for index, region in enumerate(("north", "south")): + gdf = gpd.GeoDataFrame( + {"id": [f"{index}"], "name": ["a"], "geometry": [shapely.box(index, 0, index + 1, 1)]}, + crs="EPSG:4326", + ) + src = tmp_folder / f"h_src_{index}.parquet" + gdf.to_parquet(src) + Part = type( + "Part", + (DuckDBBaseConverter,), + { + **CONFIG, + "column_additions": {"region": region, "country": "NL"}, + "missing_schemas": { + "properties": { + "name": {"type": "string"}, + "region": {"type": "string"}, + "country": {"type": "string"}, + } + }, + }, + ) + part = tmp_folder / f"h_part_{index}.parquet" + Part().convert(part, input_files={str(src): src.name}) + assert "region" in json.loads(pq.read_schema(part).metadata[b"collection"]) + parts.append(part) + + dest = tmp_folder / "hydrated.parquet" + Converter().merge_parquet(parts, dest) + + table = pq.read_table(dest) + rows = dict(zip(table.column("id").to_pylist(), table.column("region").to_pylist())) + assert rows == {"0": "north", "1": "south"} + collection = json.loads(table.schema.metadata[b"collection"]) + assert "region" not in collection + assert collection["country"] == "NL" and "country" not in table.schema.names + assert ValidateData().validate(dest, num=100, schema_map={}).errors == [] + + def test_merge_parquet_checks_ids_that_convert_generated(tmp_folder, capsys): """Each part numbered its own rows from zero, so the merge has to check what convert() is allowed to take for granted.""" diff --git a/vecorel_cli/conversion/duckdb.py b/vecorel_cli/conversion/duckdb.py index 80d63ee..99c9b8a 100644 --- a/vecorel_cli/conversion/duckdb.py +++ b/vecorel_cli/conversion/duckdb.py @@ -594,17 +594,44 @@ def merge_parquet(self, paths: list, output_file, collection=None, **kwargs) -> source_crs = self._common_crs(con, paths) - # union_by_name, because a converter drops a column a source file does not have, - # so two parts of one dataset can legitimately differ; the targets then come from - # the union rather than from whichever part happens to be first - sources = "[" + ",".join(_sql_path(path) for path in paths) + "]" - source_query = f"SELECT * FROM read_parquet({sources}, union_by_name=true)" - targets = [row[0] for row in con.execute(f"DESCRIBE {source_query}").fetchall()] - + collections = [GeoParquet(path).get_collection() for path in paths] if collection is None: from ..vecorel.ops import merge_collections - collection = merge_collections([GeoParquet(path).get_collection() for path in paths]) + collection = merge_collections(collections) + + # A part keeps its constants in its collection; one the merged collection does not + # carry, because the parts disagree on it, goes back into the rows, as `vec merge` does + hydrate = sorted( + { + key + for part in collections + for key in part.keys() + if key not in collection and key not in part.get_collection_only_properties() + } + - {"schemas"} + ) + + # union by name, because a converter drops a column a source file does not have, + # so two parts of one dataset can legitimately differ; the targets then come from + # the union rather than from whichever part happens to be first + if hydrate: + selects = [] + for path, part in zip(paths, collections): + names = pq.read_schema(path).names + literals = [ + f'{_sql_literal(part.get(key))} AS "{key}"' + for key in hydrate + if key not in names + ] + selects.append( + f"SELECT {', '.join(['*', *literals])} FROM read_parquet({_sql_path(path)})" + ) + source_query = " UNION ALL BY NAME ".join(selects) + else: + sources = "[" + ",".join(_sql_path(path) for path in paths) + "]" + source_query = f"SELECT * FROM read_parquet({sources}, union_by_name=true)" + targets = [row[0] for row in con.execute(f"DESCRIBE {source_query}").fetchall()] return self.write_query( con, From fe54e9c57b894fa1617fa35bad5339cf4fb39960 Mon Sep 17 00:00:00 2001 From: Matthias Mohr Date: Mon, 28 Sep 2026 13:53:17 +0200 Subject: [PATCH 2/6] Fix merge gaps, merge with DuckDB in vec merge (#56) * Fix merge gaps, merge with DuckDB in vec merge * Merge: hydrate only what the parts disagree on (review #1) * Merge: keep rows as they are with DuckDB, like in memory (review #2) * Report constants that don't fit their schema type (review #3) * Merge: reject two versions of one schema in a collection (review #4) * DuckDB: key ids by collection only with several collections (review #5) * Merge: apply excludes to the loaded data in memory (review #6) * DuckDB: quote required property names with _sql_name (review #7) * Share the CRS comparison of merge and DuckDB (review #8) * Merge: warn when includes drop a required collection-only property (review #9) * Merge: accept a collection that can't be determined (review #10) * Add strict mode to vec merge, minor improvements and bug fixes (#60) Co-authored-by: Matthias Mohr Co-authored-by: Ivor --- CHANGELOG.md | 254 ++++++++--- README.md | 38 +- tests/test_convert_duckdb.py | 81 +++- tests/test_encoding_geoparquet.py | 61 +++ tests/test_merge.py | 612 +++++++++++++++++++++++++++ tests/test_ops.py | 124 +++++- tests/test_validate.py | 65 +++ vecorel_cli/conversion/duckdb.py | 270 ++++++++---- vecorel_cli/create_geoparquet.py | 4 +- vecorel_cli/encoding/base.py | 23 +- vecorel_cli/encoding/geoparquet.py | 37 ++ vecorel_cli/merge.py | 131 +++++- vecorel_cli/parquet/types.py | 42 +- vecorel_cli/validation/base.py | 2 +- vecorel_cli/validation/geoparquet.py | 37 +- vecorel_cli/vecorel/ops.py | 220 +++++++++- vecorel_cli/vecorel/schemas.py | 23 +- vecorel_cli/vecorel/util.py | 19 + 18 files changed, 1873 insertions(+), 170 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 1514909..61a8766 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,29 +7,81 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0. ## [Unreleased] -- Fix: `merge_parquet` no longer drops a property that each part kept as a constant in - its collection but on which the parts disagree; it becomes a column again, as in - `vec merge`. +### Added + +- `vec merge` merges local GeoParquet files that are all in the target CRS with DuckDB, + without loading them into memory; `--engine` chooses the engine explicitly. +- `vec merge --crs first` keeps the CRS of the first dataset, the default is still EPSG:4326. +- `merge_parquet` accepts `properties` to restrict the merged properties. +- `vec merge --no-strict` warns instead of failing if the merged dataset would be invalid + and writes it anyway; `merge_parquet` has a `strict` parameter (default: `True`). + See the README for the differences between the modes. + +### Changed + +- **Breaking:** `vec merge` keeps all properties by default, `--include` restricts them to + the core properties plus the given ones and `--exclude` removes any property. It reports + required properties that are not included, also collection-only ones. +- **Breaking:** `vec merge` is strict by default: it fails if the merged dataset would be + invalid, e.g. because an id repeats within a collection. +- The message for missing required values names the missing properties with counts, + instead of showing the SQL condition. + +### Fixed + +- `merge_parquet` no longer drops a property that each part kept as a constant in + its collection but on which the parts disagree; it becomes a column again, with the + data type of its schema, as in `vec merge`. +- Merging no longer overwrites the schemas of a collection that occurs in multiple + datasets, it unites them and rejects two versions of the same schema in one collection. +- Merging warns about collection-only properties that differ between the datasets + and have to be removed. +- Merging fills in missing collection values of datasets that only list a single + collection in `schemas`. +- `vec merge` handles constants correctly: arrays and objects are no longer spread over + the rows or dropped, constants that all datasets share stay in the collection, and + constants get the data type of their schema, e.g. binary values are decoded instead + of written as base64 text. +- Merging reports a constant that doesn't fit the data type of its schema, with the + property and the file, instead of failing with an Arrow error. +- `merge_parquet` checks and numbers repeating ids per collection instead of across + all collections, and checks the required properties per collection. +- `merge_parquet` recomputes the bbox, which was empty for the rows of GeoParquet 1.0 parts. +- Columns without any value are no longer moved to the collection, which wrote NaN + (invalid JSON) or failed for nullable integers. +- `get_pyarrow_type` no longer modifies the schema of objects with `patternProperties`. +- Validation checks that all features in GeoParquet files have a collection that is + listed in `schemas`, and the required properties per collection. +- Validation checks the schemas of all collections, not only the first one. +- Validation no longer fails with a `KeyError` for GeoParquet files with multiple + collections, but without a collection column. ## [v0.3.1] - 2026-09-24 -- Fix: ZIP archives compressed with Deflate64 are extracted through `inflate64` (a +### Fixed + +- ZIP archives compressed with Deflate64 are extracted through `inflate64` (a dependency of py7zr, now declared). They required `zipfile-deflate64`, which was not declared and has no wheels for Python 3.11+. ## [v0.3.0] - 2026-09-24 -- Fix: a download that ends before its `Content-Length` is reached now fails instead - of caching the partial file, which a later run would otherwise reuse as if it were - complete (#46). -- Conversions no longer drop rows with missing required values (`max_dropped_share` - is removed): any row without a value for a required property fails the conversion, - so each dataset handles such rows explicitly (#33). -- Conversions report how many geometry parts the geometry repair removed; - previously they vanished silently (#33). +### Added + - `download_files()` now handles multi-volume 7z archives: URIs ending in `.7z.001`, `.7z.002`, ... that share a name are downloaded together and extracted as one 7z stream, with the target paths read from whichever part carries them (fiboa/cli#312). +- `id_columns` (with `id_separator`, default `-`): composes `id` by joining the + named columns, after the column migrations ran. Integer-typed float columns are cast + losslessly, so an id part does not render as `4.0`. +- Conversions report how many geometry parts the geometry repair removed; + previously they vanished silently (#33). + +### Changed + +- Conversions no longer drop rows with missing required values: any row without a value + for a required property fails the conversion, so each dataset handles such rows + explicitly (#33). - Converters no longer split multi-part geometries into one row per polygon (#47). Features keep the geometry modeling of the source, so ids, row counts and attribute values (such as an area) stay 1:1 with it. Rows without a polygonal geometry @@ -39,17 +91,26 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0. ids, or rows repeating across merged parts — are numbered with a `~` suffix (`id~1`, `id~2`, ...). - Converters number the rows (over all source files) as `id` when no column is mapped - to `id`, or when the mapped column is missing from the data (#47). The `index_as_id` - flag and the `"id": "id"` mapping it required are removed; previously the per-file - numbering also repeated ids after a multi-file read. -- Add `id_columns` (with `id_separator`, default `-`): composes `id` by joining the - named columns, after the column migrations ran. Integer-typed float columns are cast - losslessly, so an id part does not render as `4.0`. + to `id`, or when the mapped column is missing from the data (#47). Previously the + per-file numbering also repeated ids after a multi-file read. - `vec improve --explode-geometries` numbers the ids of the parts it creates (`~`), so the result keeps unique ids. +### Removed + +- `max_dropped_share`, as rows with missing required values are no longer dropped (#33). +- The `index_as_id` flag and the `"id": "id"` mapping it required (#47). + +### Fixed + +- A download that ends before its `Content-Length` is reached now fails instead + of caching the partial file, which a later run would otherwise reuse as if it were + complete (#46). + ## [v0.2.20] - 2026-09-17 +### Changed + - `BaseConverter.convert()` chooses a default variant when none is given, before `get_urls()` runs: the latest year when the variants are years, otherwise the first declared one (`default_variant()`). A converter that overrides `get_urls()` no longer has @@ -57,30 +118,44 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0. ## [v0.2.19] - 2026-09-17 -- Add `DuckDBBaseConverter.merge_parquet()`, which combines Vecorel GeoParquet files +### Added + +- `DuckDBBaseConverter.merge_parquet()`, which combines Vecorel GeoParquet files into one, checked and sorted over the whole set. `convert()` now ends in the same `write_query()`, so there is one route from a query to a packaged file (#36). -- Fix: a constant pinned to the feature level was passed to DuckDB as a bound parameter, - which a `COPY` binds before its subquery's, so the output went to a file named after the - constant. -- Add `BaseConverter.dehydrate` (default `True`). Set it to `False` when a conversion +- `BaseConverter.dehydrate` (default `True`). Set it to `False` when a conversion writes one part of a dataset: constants would otherwise be judged over the part and a property that varies between parts is lost (#35). -- Fix: an integer column with a null value read back as float64, so validation - rejected every value in it (#37). Integers are now read into pandas' nullable - dtypes, which also keeps int64 values that float64 cannot represent exactly. + +### Changed + - Require `aiohttp>=3.13.5`, which allows downloads from servers that send duplicate headers, such as Zenodo (#41). -- Fix: a failed download left an empty file in the cache that later runs treated as + +### Fixed + +- A constant pinned to the feature level was passed to DuckDB as a bound parameter, + which a `COPY` binds before its subquery's, so the output went to a file named after the + constant. +- An integer column with a null value read back as float64, so validation + rejected every value in it (#37). Integers are now read into pandas' nullable + dtypes, which also keeps int64 values that float64 cannot represent exactly. +- A failed download left an empty file in the cache that later runs treated as cached, so they never retried and failed with an unrelated error (#43). Downloads now stream to a `.part` file that is renamed only after a clean close. ## [v0.2.18] - 2026-09-14 -- Add an experimental `DuckDBBaseConverter` to convert large Parquet-based datasets +### Added + +- An experimental `DuckDBBaseConverter` to convert large Parquet-based datasets without loading them into memory. Its output matches the default converter (geometry handling, Hilbert order, data types, metadata, file packaging). `duckdb` is a new dependency. +- Converters warn when no column is mapped to `id` and when the id column is not unique. + +### Changed + - Converters record the collection id in the collection metadata and no longer add a constant `collection` column. Previously files without constant columns were written without any collection id. @@ -88,7 +163,6 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0. missing geometries). Missing required values are dropped only up to the new `max_dropped_share` (default 1%), above it the conversion fails. - Converters fail when both `sources` and `variants` are declared. -- Converters warn when no column is mapped to `id` and when the id column is not unique. - Converters load all schemas upfront with retries, so a temporary network issue no longer kills a long conversion at the very end. - Send `User-Agent: vecorel-cli` on HTTP downloads instead of fsspec's default. Servers that @@ -97,6 +171,8 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0. ## [v0.2.17] - 2026-09-03 +### Fixed + - Read GeoJSON as UTF-8, which the format mandates, instead of following the platform locale. The locale default is cp1252 on Windows, where `Grünland` silently became `Grünland` — in converters, and in `describe`, `merge` and `validate`. Writing now states UTF-8 explicitly too, @@ -106,7 +182,12 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0. ## [v0.2.16] - 2026-08-29 +### Changed + - Sort converter output by Hilbert distance against the CRS's total bounds instead of by raw WKB byte order. + +### Fixed + - Exit with a non-zero exit code when a command reports a failure (e.g. `vec validate-schema` on an invalid schema) [#27](https://github.com/vecorel/cli/issues/27) - Keep collection-only properties in the collection metadata when merging collections, e.g. in `create-geoparquet` and `merge`. @@ -115,100 +196,175 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0. ## [v0.2.15] - 2026-02-16 -- Add a `get_default_collection` to the Registry +### Added + +- A `get_default_collection` to the Registry ## [v0.2.14] - 2026-02-13 +### Added + - Enable verbose mode via env `VECOREL_VERBOSE` set to `1` -- Converters: - - Load the most specific class for converters - - Move block size check so that it applies for all downloads - - Avoid error with license set to None - - Made glob recursive, so it can be used for multiple directories + +### Changed + +- Converters: Load the most specific class for converters +- Converters: Made glob recursive, so it can be used for multiple directories + +### Fixed + +- Converters: Move block size check so that it applies for all downloads +- Converters: Avoid error with license set to None ## [v0.2.13] - 2026-02-13 +### Added + +- Option to set compression level, zstd defaults to 15 +- Support for Python 3.14 + +### Changed + - Change default compression to zstd -- Add option to set compression level, zstd defaults to 15 -- Add support for Python 3.14, remove support for Python 3.10 - Replace flatdict from pypi with a local version to avoid pkg_resource install issues - Updated dependencies (especially catering for future pandas versions) +### Removed + +- Support for Python 3.10 + ## [v0.2.12] - 2025-12-08 -- Change default temporal property to datetime +### Added + - Enable Converter.columns list and tuple types -- Add BaseConverter get_columns hook to customize columns after reading the file +- BaseConverter get_columns hook to customize columns after reading the file + +### Changed + +- Change default temporal property to datetime - Update STAC processing extension ## [v0.2.11] - 2025-10-09 +### Fixed + - XML/HTML-like tags (with < and > characters) in logs are properly escaped for loguru ## [v0.2.10] - 2025-10-09 -- XML tags in logs are properly escaped for loguru +### Changed + - Set the default temporal_property for STAC collection creation to `determination:datetime` instead of `determination_datetime` +### Fixed + +- XML tags in logs are properly escaped for loguru + ## [v0.2.9] - 2025-10-09 +### Added + - Converters: `column_filters` allows to inverse the mask + +### Fixed + - Fix use of license and provider in converter list - Various small bug fixes and type hint fixes ## [v0.2.8] - 2025-09-13 +### Fixed + - Fix issue with schema requests due to changes in the "firewall" by ReadTheDocs that sits in front of the PROJJSON schema ## [v0.2.7] - 2025-08-29 +### Changed + - `create-stac-collection`: - Don't set empty strings / only provide properties that have value - Detect the collection id more robustly ## [v0.2.6] - 2025-08-29 +### Changed + - Move ValidateData.required_schemas to Registry.required_extensions and adapted ValidateData accordingly ## [v0.2.5] - 2025-08-29 -- Fix deprecation warning for `re.sub` -- Add `unrar` dependency -- `create-stac-collection`: - - Set temporal property parameter from none to the actual configured default - - Support for GeoJSON input +### Added + +- `unrar` dependency +- `create-stac-collection`: Support for GeoJSON input - Allow to set a list of required schemas for validation +### Fixed + +- Fix deprecation warning for `re.sub` +- `create-stac-collection`: Set temporal property parameter from none to the actual configured default + ## [v0.2.4] - 2025-08-27 -- Encode numpy datatypes correctly when exporting to JSON +### Changed + - Code refactoring +### Fixed + +- Encode numpy datatypes correctly when exporting to JSON + ## [v0.2.3] - 2025-08-26 +### Changed + - Better support for merging schemas ## [v0.2.2] - 2025-08-25 +### Added + +- Return value to `ConvertData.check_datasets` + +### Changed + - Make the whole library easier to rebrand and reuse -- Separate CLI creation from `__init__.py` files to avoid import race coditions -- Add return value to `ConvertData.check_datasets` + +### Fixed + +- Separate CLI creation from `__init__.py` files to avoid import race conditions - Fix geopandas `datetime64` data type conversion issue ## [v0.2.1] - 2025-08-25 +### Changed + - Updated to use the Geometry Metrics Extension + +### Fixed + - Fixed various hardcoded vecorel instances in rename-extension - Fixed registry to be overridable by other CLI tools ## [v0.2.0] - 2025-08-15 +### Added + +- Internal `py-package` parameter to the `convert` command + +### Changed + - Migrate from vecorel.github.io to vecorel.org -- Add internal `py-package` parameter to the `convert` command + +### Fixed + - Bugfixes ## [v0.1.0] - 2025-08-15 +### Added + - First release based on vecorel CLI 0.1.0 [Unreleased]: diff --git a/README.md b/README.md index 3487c87..9688241 100644 --- a/README.md +++ b/README.md @@ -121,7 +121,43 @@ Check `vec describe --help` for more details. Merges multiple Vecorel datasets to a combined Vecorel dataset: -- `vec merge ec_ee.parquet ec_lv.parquet -o merged.parquet -e https://vecorel.org/hcat-extension/v0.1.0/schema.yaml -i ec:hcat_name -i ec:hcat_code -i ec:translated_name` +- `vec merge ec_ee.parquet ec_lv.parquet -o merged.parquet` +- Only the core properties and some additional properties: `vec merge ec_ee.parquet ec_lv.parquet -o merged.parquet -i ec:hcat_name -i ec:hcat_code` +- All properties except for some: `vec merge ec_ee.parquet ec_lv.parquet -o merged.parquet -e ec:translated_name` + +The merged dataset is in EPSG:4326 by default. Use `--crs` to choose another CRS, +or `--crs first` to keep the CRS of the first dataset. + +Local GeoParquet files that are all in the target CRS are merged with DuckDB, so they don't need to fit into memory. +All other datasets (e.g. GeoJSON or datasets that need to be reprojected) are merged in memory. +Use `--engine` to choose the engine explicitly. + +`-i` and `-e` apply to all properties, including collection-level metadata. +Constants that differ between the datasets are moved from the collection metadata to the features. + +#### Strict and non-strict mode + +By default, `vec merge` is strict: problems that make the merged dataset invalid are errors +and no dataset is written. +With `--no-strict`, `vec merge` is fail-safe: these problems are reported as warnings and +the dataset is written anyway. Check it with `vec validate` afterwards. + +| Problem | Strict (default) | `--no-strict` | +| ------- | ---------------- | ------------- | +| A required property has no value for some features | Error | Warning, the property is written as nullable | +| An id repeats within a collection | Error | Warning | +| The collection of a feature can't be determined | Error | Warning, the collection is left empty | +| `-i` or `-e` removes a required property | Error | Warning | +| A required collection-only property differs between the datasets | Error | Warning, the property is removed | +| An optional collection-only property differs between the datasets | Warning, the property is removed | Warning, the property is removed | +| A feature has an empty or missing geometry | Error | Warning, the feature is kept | +| A collection-level value doesn't fit the data type of its schema | Error | Warning, the value is left empty | +| A collection implements multiple versions of a schema, e.g. of an extension | Error | Error | +| The schemas of the datasets conflict | Error | Error | + +The strict mode only checks what a merge can break or check with little effort. +It doesn't validate the values against the schemas, e.g. patterns or value ranges, +so a merged dataset is only valid if the values of the source datasets are valid. Check `vec merge --help` for more details. diff --git a/tests/test_convert_duckdb.py b/tests/test_convert_duckdb.py index 35fe8f5..91cc5b8 100644 --- a/tests/test_convert_duckdb.py +++ b/tests/test_convert_duckdb.py @@ -1,4 +1,5 @@ import json +import re import sys import geopandas as gpd @@ -9,7 +10,7 @@ from loguru import logger from vecorel_cli.conversion.base import BaseConverter -from vecorel_cli.conversion.duckdb import DuckDBBaseConverter +from vecorel_cli.conversion.duckdb import DuckDBBaseConverter, _constants_table from vecorel_cli.validate import ValidateData from vecorel_cli.vecorel.hilbert import hilbert_keys_for_table @@ -664,6 +665,84 @@ def test_merge_parquet_hydrates_constants_the_parts_disagree_on(tmp_folder): assert ValidateData().validate(dest, num=100, schema_map={}).errors == [] +def _collection_part(folder, cid, ids, area=None): + gdf = gpd.GeoDataFrame( + { + "id": ids, + "name": ["a"] * len(ids), + "geometry": [shapely.box(i, 0, i + 1, 1) for i in range(len(ids))], + }, + crs="EPSG:4326", + ) + src = folder / f"c_src_{cid}_{len(ids)}.parquet" + gdf.to_parquet(src) + config = {**CONFIG, "id": cid} + if area is not None: + config["column_additions"] = {"area": area} + config["missing_schemas"] = { + "properties": {"name": {"type": "string"}, "area": {"type": "double"}} + } + part = folder / f"c_part_{cid}_{len(ids)}.parquet" + type("Part", (DuckDBBaseConverter,), config)().convert(part, input_files={str(src): src.name}) + return part + + +def test_merge_parquet_numbers_ids_per_collection(tmp_folder): + """Ids only have to be unique per collection, so only repeats within one get numbered.""" + parts = [ + _collection_part(tmp_folder, "a", ["1", "2"]), + _collection_part(tmp_folder, "b", ["1", "2"]), + _collection_part(tmp_folder, "b", ["1"]), + ] + dest = tmp_folder / "per_collection.parquet" + Converter().merge_parquet(parts, dest) + + table = pq.read_table(dest, columns=["collection", "id"]) + assert sorted(zip(table.column("collection").to_pylist(), table.column("id").to_pylist())) == [ + ("a", "1"), + ("a", "2"), + ("b", "1~1"), + ("b", "1~2"), + ("b", "2"), + ] + assert ValidateData().validate(dest, num=100, schema_map={}).errors == [] + + +def test_merge_parquet_hydrates_nan_as_null(tmp_folder): + """A NaN in the collection comes from a column without values, so it's a null.""" + parts = [ + _collection_part(tmp_folder, "a", ["1"], area=float("nan")), + _collection_part(tmp_folder, "b", ["2"], area=1.5), + ] + dest = tmp_folder / "nan.parquet" + Converter().merge_parquet(parts, dest) + + table = pq.read_table(dest, columns=["id", "area"]) + assert dict(zip(table.column("id").to_pylist(), table.column("area").to_pylist())) == { + "1": None, + "2": 1.5, + } + + +@pytest.mark.parametrize( + "value, dtype", [(300, "uint8"), ("not a date", "date-time")], ids=["overflow", "date-time"] +) +def test_constants_that_do_not_fit_their_type_are_reported(value, dtype, capsys): + log = Converter() + logger.remove() + logger.add(sys.stdout, format="{message}", level="DEBUG", colorize=False) + table = _constants_table({"x": value}, {"x": {"type": dtype}}, log=log, source="part.parquet") + + # the value would fail the writer, so it's left empty + assert table.column("x").to_pylist() == [None] + message = f"x: Value {value!r} doesn't fit data type {dtype}" + out = capsys.readouterr().out + assert message in out and "(in part.parquet)" in out + + with pytest.raises(ValueError, match=re.escape(message)): + _constants_table({"x": value}, {"x": {"type": dtype}}, log=log, strict=True) + + def test_merge_parquet_checks_ids_that_convert_generated(tmp_folder, capsys): """Each part numbered its own rows from zero, so the merge has to check what convert() is allowed to take for granted.""" diff --git a/tests/test_encoding_geoparquet.py b/tests/test_encoding_geoparquet.py index e8cfa78..0f3a1bb 100644 --- a/tests/test_encoding_geoparquet.py +++ b/tests/test_encoding_geoparquet.py @@ -148,3 +148,64 @@ def test_read_keeps_integers_that_have_a_null(tmp_parquet_file): from vecorel_cli.validate import ValidateData assert ValidateData().validate(tmp_parquet_file).errors == [] + + +def test_write_keeps_all_null_columns_in_the_data(tmp_parquet_file): + # NaN and NA aren't valid JSON, so they can't move to the collection + from geopandas import GeoDataFrame + from shapely.geometry import box + + gdf = GeoDataFrame( + { + "id": ["1", "2"], + "code": pd.array([None, None], dtype="Int64"), + "area": [np.nan, np.nan], + "name": [None, None], + "geometry": [box(0, 0, 1, 1), box(1, 0, 2, 1)], + }, + 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) + + collection = GeoParquet(tmp_parquet_file).get_collection() + assert not {"code", "area", "name"} & set(collection) + assert {"code", "area", "name"} <= set(pq.read_schema(tmp_parquet_file).names) + + +def test_get_pyarrow_type_keeps_the_schema(): + from vecorel_cli.parquet.types import get_pyarrow_type + + schema = {"type": "object", "patternProperties": {".*": {"type": "string"}}} + assert get_pyarrow_type(schema) == pa.map_(pa.string(), pa.string()) + assert get_pyarrow_type(schema) == pa.map_(pa.string(), pa.string()) + + +def test_read_hydrates_array_and_object_constants(tmp_parquet_file): + from geopandas import GeoDataFrame + from shapely.geometry import box + + gdf = GeoDataFrame( + {"id": ["1", "2", "3"], "geometry": [box(0, 0, 1, 1)] * 3}, + crs="EPSG:4326", + ) + gp = GeoParquet(tmp_parquet_file) + gp.set_collection( + { + "schemas": {"c": ["https://vecorel.org/specification/v0.1.0/schema.yaml"]}, + "collection": "c", + "tags": ["a", "b"], + "attrs": {"k": "v"}, + } + ) + gp.write(gdf, dehydrate=False) + + data = GeoParquet(tmp_parquet_file).read(hydrate=True) + assert list(data["tags"]) == [["a", "b"]] * 3 + assert list(data["attrs"]) == [{"k": "v"}] * 3 diff --git a/tests/test_merge.py b/tests/test_merge.py index 479ae36..7f9b79d 100644 --- a/tests/test_merge.py +++ b/tests/test_merge.py @@ -1,9 +1,26 @@ +import json +import re +import sys from pathlib import Path +import geopandas as gpd +import pandas as pd +import pyarrow as pa +import pyarrow.parquet as pq import pytest +import shapely +from loguru import logger +from vecorel_cli.cli.logger import LoggerMixin +from vecorel_cli.encoding.geojson import GeoJSON from vecorel_cli.encoding.geoparquet import GeoParquet from vecorel_cli.merge import MergeDatasets +from vecorel_cli.validate import ValidateData +from vecorel_cli.vecorel.schemas import Schemas + +CORE = Schemas.get_core_uri() +ADMIN = "https://vecorel.org/administrative-division-extension/v0.1.0/schema.yaml" +ENGINES = ["duckdb", "geopandas"] def test_merge(tmp_parquet_file: Path): @@ -47,6 +64,28 @@ def test_merge(tmp_parquet_file: Path): ] +def test_merge_excludes_without_reading_twice(tmp_parquet_file: Path, monkeypatch): + reads = [] + read = GeoJSON.read + + def spy(self, num=None, **kwargs): + if num is None: + reads.append(self.uri) + return read(self, num=num, **kwargs) + + monkeypatch.setattr(GeoJSON, "read", spy) + MergeDatasets().merge( + source=["tests/data-files/inspire.parquet", "tests/data-files/admin.json"], + target=tmp_parquet_file, + excludes=["foo"], + ) + + assert len(reads) == 1 + columns = GeoParquet(tmp_parquet_file).get_properties() + assert "foo" not in columns + assert "admin:country_code" in columns + + def test_merge_invalid_file(tmp_folder): out = tmp_folder / "output.parquet" merge = MergeDatasets() @@ -54,3 +93,576 @@ def test_merge_invalid_file(tmp_folder): merge.merge("invalid.parquet", out) with pytest.raises(FileNotFoundError): merge.merge(["invalid.parquet"], out) + + +def _part( + folder, + name, + cid, + n, + ids=None, + columns=None, + collection=None, + schemas=None, + custom=None, + with_collection=True, + geoparquet_version=None, + crs="EPSG:4326", +): + """A part with n rows; only the given collection-level values are in its collection.""" + gdf = gpd.GeoDataFrame( + { + "id": ids or [f"{name}{i}" for i in range(n)], + **(columns or {}), + "geometry": [shapely.box(i, 0, i + 1, 1) for i in range(n)], + }, + crs="EPSG:4326", + ).to_crs(crs) + path = folder / f"{name}.parquet" + meta = {"schemas": {cid: schemas or [CORE]}, **(collection or {})} + if with_collection: + meta["collection"] = cid + if custom: + meta["schemas:custom"] = custom + gp = GeoParquet(path) + gp.set_collection(meta) + gp.write(gdf, dehydrate=False, geoparquet_version=geoparquet_version) + return path + + +def _with_nullable_column(path, name, values, pa_type=pa.string()): + """Rewrites a column as nullable with the given values, as other tools could write it.""" + table = pq.read_table(path) + metadata = table.schema.metadata + index = table.schema.get_field_index(name) + if index >= 0: + table = table.remove_column(index) + table = table.append_column(pa.field(name, pa_type), pa.array(values, pa_type)) + pq.write_table(table.replace_schema_metadata(metadata), path) + return path + + +def _merge(folder, parts, engine="auto", **kwargs): + out = folder / f"merged-{engine}.parquet" + MergeDatasets().merge(source=[str(p) for p in parts], target=out, engine=engine, **kwargs) + return out + + +def _read(path): + table = pq.read_table(path) + collection = json.loads(table.schema.metadata[b"collection"]) + rows = table.drop_columns(["geometry"]).to_pylist() + rows.sort(key=lambda row: (row.get("collection") or "", row["id"])) + return rows, collection + + +def _errors(path): + return ValidateData().validate(path, num=100, schema_map={}).errors + + +@pytest.fixture +def log(): + # capsys only captures during the test call, so collect the messages in a list, + # after the first LoggerMixin has replaced the sinks + LoggerMixin() + messages = [] + logger.remove() + logger.add(messages.append, format="{message}", level="DEBUG", colorize=False) + + def read(): + out = "".join(messages) + messages.clear() + return out + + yield read + logger.remove() + logger.add(sys.stdout, format="{message}", level="DEBUG", colorize=False) + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_hydrates_array_and_object_constants(tmp_folder, engine): + custom = { + "properties": { + "tags": {"type": "array", "items": {"type": "string"}}, + "attrs": {"type": "object", "properties": {"k": {"type": "string"}}}, + } + } + a = _part( + tmp_folder, "a", "a", 2, collection={"tags": ["p", "q"], "attrs": {"k": "v"}}, custom=custom + ) + b = _part( + tmp_folder, "b", "b", 3, collection={"tags": ["r", "s"], "attrs": {"k": "w"}}, custom=custom + ) + out = _merge(tmp_folder, [a, b], engine) + + rows, _ = _read(out) + assert [(r["collection"], r["tags"], r["attrs"]) for r in rows] == [ + ("a", ["p", "q"], {"k": "v"}), + ("a", ["p", "q"], {"k": "v"}), + ("b", ["r", "s"], {"k": "w"}), + ("b", ["r", "s"], {"k": "w"}), + ("b", ["r", "s"], {"k": "w"}), + ] + assert _errors(out) == [] + + +@pytest.mark.parametrize("engine", ENGINES) +@pytest.mark.parametrize("with_schema", [False, True]) +def test_merge_keeps_shared_array_constants_in_the_collection(tmp_folder, engine, with_schema): + custom = {"properties": {"tags": {"type": "array", "items": {"type": "string"}}}} + custom = custom if with_schema else None + a = _part(tmp_folder, "a", "a", 2, collection={"tags": ["x", "y"]}, custom=custom) + b = _part(tmp_folder, "b", "b", 2, collection={"tags": ["x", "y"]}, custom=custom) + out = _merge(tmp_folder, [a, b], engine) + + rows, collection = _read(out) + assert collection["tags"] == ["x", "y"] + assert all("tags" not in row for row in rows) + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_hydrates_date_time_constants(tmp_folder, engine): + custom = {"properties": {"dt": {"type": "date-time"}}} + a = _part(tmp_folder, "a", "a", 2, collection={"dt": "2020-01-01T00:00:00Z"}, custom=custom) + dts = pd.to_datetime(["2021-01-01T00:00:00Z", "2021-06-01T00:00:00Z"]) + b = _part(tmp_folder, "b", "b", 2, columns={"dt": dts}, custom=custom) + out = _merge(tmp_folder, [a, b], engine) + + assert str(pq.read_schema(out).field("dt").type) == "timestamp[ms, tz=UTC]" + rows, _ = _read(out) + assert [r["dt"].isoformat() for r in rows] == [ + "2020-01-01T00:00:00+00:00", + "2020-01-01T00:00:00+00:00", + "2021-01-01T00:00:00+00:00", + "2021-06-01T00:00:00+00:00", + ] + assert _errors(out) == [] + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_keeps_ids_that_repeat_in_other_collections(tmp_folder, engine, log): + a = _part(tmp_folder, "a", "a", 2, ids=["1", "2"]) + b = _part(tmp_folder, "b", "b", 2, ids=["1", "2"]) + c = _part(tmp_folder, "c", "b", 1, ids=["1"]) + out = _merge(tmp_folder, [a, b, c], engine, strict=False) + + rows, _ = _read(out) + assert [(r["collection"], r["id"]) for r in rows] == [ + ("a", "1"), + ("a", "2"), + ("b", "1"), + ("b", "1"), + ("b", "2"), + ] + assert "repeat an id within their collection" in log() + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_unites_the_schemas_of_a_collection(tmp_folder, engine): + a = _part(tmp_folder, "a", "same", 2, columns={"admin:country_code": ["DE", "DE"]}) + b = _part( + tmp_folder, + "b", + "same", + 2, + columns={"admin:country_code": ["FR", "FR"]}, + schemas=[CORE, ADMIN], + ) + out = _merge(tmp_folder, [b, a], engine) + + _, collection = _read(out) + assert sorted(collection["schemas"]["same"]) == sorted([CORE, ADMIN]) + assert _errors(out) == [] + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_warns_about_collection_only_properties_that_differ(tmp_folder, engine, log): + custom = { + "properties": {"producer": {"type": "string"}, "source": {"type": "string"}}, + "collection": {"producer": True}, + } + a = _part( + tmp_folder, "a", "a", 2, collection={"producer": "Alice", "source": "x"}, custom=custom + ) + b = _part(tmp_folder, "b", "b", 2, collection={"producer": "Bob", "source": "x"}, custom=custom) + out = _merge(tmp_folder, [a, b], engine) + + rows, collection = _read(out) + assert "producer" not in collection + assert collection["source"] == "x" + assert all("producer" not in row for row in rows) + assert ( + "Collection-only properties differ between the datasets and are removed: producer" in log() + ) + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_keeps_all_properties_by_default(tmp_folder, engine): + custom = {"properties": {"name": {"type": "string"}}} + a = _part( + tmp_folder, + "a", + "a", + 2, + columns={"admin:country_code": ["DE", "DE"], "name": ["x", "y"]}, + schemas=[CORE, ADMIN], + custom=custom, + ) + b = _part(tmp_folder, "b", "b", 2, columns={"name": ["z", "w"]}, custom=custom) + out = _merge(tmp_folder, [a, b], engine) + + rows, _ = _read(out) + assert [(r["admin:country_code"], r["name"]) for r in rows] == [ + ("DE", "x"), + ("DE", "y"), + (None, "z"), + (None, "w"), + ] + # admin:country_code is only required in collection a + assert _errors(out) == [] + + out = _merge(tmp_folder, [a, b], engine, excludes=["name"]) + assert "name" not in pq.read_schema(out).names + assert "admin:country_code" in pq.read_schema(out).names + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_warns_when_includes_drop_required_properties(tmp_folder, engine, log): + a = _part( + tmp_folder, "a", "a", 2, columns={"admin:country_code": ["DE", "DE"]}, schemas=[CORE, ADMIN] + ) + b = _part( + tmp_folder, "b", "b", 2, columns={"admin:country_code": ["FR", "FR"]}, schemas=[CORE, ADMIN] + ) + _merge(tmp_folder, [a, b], engine, includes=["foo"], strict=False) + assert ( + "Required properties are not included, the merged file will be invalid: admin:country_code" + in log() + ) + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_fills_a_missing_collection(tmp_folder, engine): + a = _part(tmp_folder, "a", "a", 2, with_collection=False) + b = _part(tmp_folder, "b", "b", 2) + out = _merge(tmp_folder, [a, b], engine) + + rows, _ = _read(out) + assert [r["collection"] for r in rows] == ["a", "a", "b", "b"] + assert _errors(out) == [] + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_fills_missing_collection_values(tmp_folder, engine): + a = _with_nullable_column(_part(tmp_folder, "a", "a", 2), "collection", ["a", None]) + b = _part(tmp_folder, "b", "b", 2) + out = _merge(tmp_folder, [a, b], engine) + + rows, _ = _read(out) + assert [r["collection"] for r in rows] == ["a", "a", "b", "b"] + assert _errors(out) == [] + + +@pytest.mark.parametrize("engine", ENGINES) +@pytest.mark.parametrize("values", [None, ["a", None]], ids=["no-column", "null-values"]) +def test_merge_accepts_a_collection_it_cannot_determine(tmp_folder, engine, values): + schemas = {"a": [CORE], "x": [CORE]} + a = _part(tmp_folder, "a", "a", 2, collection={"schemas": schemas}, with_collection=False) + if values: + a = _with_nullable_column(a, "collection", values) + b = _part(tmp_folder, "b", "b", 2) + out = _merge(tmp_folder, [a, b], engine, strict=False) + + rows, _ = _read(out) + assert [r["collection"] for r in rows][-2:] == ["b", "b"] + assert len(rows) == 4 + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_excludes_the_collection(tmp_folder, engine, log): + a = _part(tmp_folder, "a", "a", 2) + b = _part(tmp_folder, "b", "b", 2) + out = _merge(tmp_folder, [a, b], engine, excludes=["collection"], strict=False) + + assert "collection" not in pq.read_schema(out).names + assert "the merged file will be invalid: collection" in log() + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_keeps_rows_without_a_required_value(tmp_folder, engine): + a = _part( + tmp_folder, "a", "a", 2, columns={"admin:country_code": ["DE", "DE"]}, schemas=[CORE, ADMIN] + ) + a = _with_nullable_column(a, "admin:country_code", ["DE", None]) + b = _part( + tmp_folder, "b", "b", 2, columns={"admin:country_code": ["FR", "FR"]}, schemas=[CORE, ADMIN] + ) + out = _merge(tmp_folder, [a, b], engine, strict=False) + + rows, _ = _read(out) + assert [r["admin:country_code"] for r in rows] == ["DE", None, "FR", "FR"] + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_keeps_empty_geometries(tmp_folder, engine): + a = _part(tmp_folder, "a", "a", 2) + b = _part(tmp_folder, "b", "b", 2) + gp = GeoParquet(b) + gdf = gp.read() + gdf.loc[0, "geometry"] = shapely.Polygon() + gp.write(gdf, dehydrate=False) + out = _merge(tmp_folder, [a, b], engine, strict=False) + + assert pq.read_metadata(out).num_rows == 4 + + +def test_merge_ignores_null_ids_for_duplicates(tmp_folder, log): + a = _with_nullable_column(_part(tmp_folder, "a", "a", 2), "id", [None, None]) + b = _part(tmp_folder, "b", "b", 2) + _merge(tmp_folder, [a, b], "geopandas", strict=False) + assert "repeat an id" not in log() + + +def test_merge_recomputes_the_bbox_of_geoparquet_1_0_parts(tmp_folder): + a = _part(tmp_folder, "a", "a", 2, geoparquet_version="1.0.0") + b = _part(tmp_folder, "b", "b", 2) + out = _merge(tmp_folder, [a, b], "duckdb") + + bboxes = pq.read_table(out, columns=["bbox"]).column("bbox").to_pylist() + assert len(bboxes) == 4 and None not in bboxes + + +def test_merge_engine_selection(tmp_folder, log): + a = _part(tmp_folder, "a", "a", 2) + b = _part(tmp_folder, "b", "b", 2) + c = _part(tmp_folder, "c", "c", 2, crs="EPSG:3857") + admin = "tests/data-files/admin.json" + + _merge(tmp_folder, [a, b]) + assert "Merging with DuckDB" in log() + _merge(tmp_folder, [a, b], crs="EPSG:4326") + assert "Merging with DuckDB" in log() + _merge(tmp_folder, [a, admin]) + assert "Merging in memory, as DuckDB only merges local GeoParquet files" in log() + _merge(tmp_folder, [a, c]) + assert "Merging in memory, as the datasets must be reprojected" in log() + _merge(tmp_folder, [a, b], crs="EPSG:3857") + assert "Merging in memory, as the datasets must be reprojected" in log() + with pytest.raises(ValueError, match="Can't merge with DuckDB"): + _merge(tmp_folder, [a, admin], "duckdb") + + +def test_merge_crs(tmp_folder, log): + a = _part(tmp_folder, "a", "a", 2) + c = _part(tmp_folder, "c", "c", 2, crs="EPSG:3857") + d = _part(tmp_folder, "d", "d", 2, crs="EPSG:3857") + + def crs_of(path): + return gpd.read_parquet(path).crs.to_epsg() + + # EPSG:4326 by default + assert crs_of(_merge(tmp_folder, [c, d])) == 4326 + assert "Merging in memory, as the datasets must be reprojected" in log() + + # the CRS of the first dataset + assert crs_of(_merge(tmp_folder, [c, d], crs="first")) == 3857 + assert "Merging with DuckDB" in log() + assert crs_of(_merge(tmp_folder, [c, a], crs="first")) == 3857 + assert "Merging in memory, as the datasets must be reprojected" in log() + + # a specific CRS + assert crs_of(_merge(tmp_folder, [a, d], crs="EPSG:3857")) == 3857 + assert "Merging in memory, as the datasets must be reprojected" in log() + assert crs_of(_merge(tmp_folder, [c, d], crs="EPSG:3857")) == 3857 + assert "Merging with DuckDB" in log() + + +def _with_collection(path, collection): + """Replaces the collection metadata of a file.""" + table = pq.read_table(path) + metadata = dict(table.schema.metadata) + metadata[b"collection"] = json.dumps(collection).encode("utf-8") + pq.write_table(table.replace_schema_metadata(metadata), path) + return path + + +def _check_modes(folder, parts, engine, log, message, **kwargs): + """With --no-strict the merge warns and writes the file, by default it fails, both with the message.""" + out = _merge(folder, parts, engine, strict=False, **kwargs) + assert message in log() + assert out.exists() + out.unlink() + with pytest.raises(ValueError, match=re.escape(message)): + _merge(folder, parts, engine, **kwargs) + return out + + +@pytest.mark.parametrize("engine", ENGINES) +@pytest.mark.parametrize("cids", [("a", "a"), ("a", "b")]) +def test_merge_modes_null_in_required_property(tmp_folder, engine, log, cids): + a = _part( + tmp_folder, + "a", + cids[0], + 2, + columns={"admin:country_code": ["DE", "DE"]}, + schemas=[CORE, ADMIN], + ) + a = _with_nullable_column(a, "admin:country_code", ["DE", None]) + b = _part( + tmp_folder, + "b", + cids[1], + 2, + columns={"admin:country_code": ["FR", "FR"]}, + schemas=[CORE, ADMIN], + ) + _check_modes( + tmp_folder, + [a, b], + engine, + log, + "Rows have no value for a required property: admin:country_code", + ) + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_modes_repeating_ids(tmp_folder, engine, log): + a = _part(tmp_folder, "a", "a", 2, ids=["1", "2"]) + b = _part(tmp_folder, "b", "a", 2, ids=["1", "3"]) + out = _check_modes(tmp_folder, [a, b], engine, log, "repeat an id within their collection") + assert out.name.startswith("merged") + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_modes_unknown_collection(tmp_folder, engine, log): + a = _with_nullable_column(_part(tmp_folder, "a", "a", 2), "collection", [None, None]) + a = _with_collection(a, {"schemas": {"a": [CORE], "x": [CORE]}}) + b = _part(tmp_folder, "b", "b", 2) + _check_modes(tmp_folder, [a, b], engine, log, "Can't determine the collection of the features") + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_modes_includes_remove_required_properties(tmp_folder, engine, log): + # -i also applies to collection-only properties + custom = { + "properties": {"producer": {"type": "string"}}, + "collection": {"producer": True}, + "required": ["producer"], + } + a = _part(tmp_folder, "a", "a", 2, collection={"producer": "X"}, custom=custom) + b = _part(tmp_folder, "b", "b", 2, collection={"producer": "X"}, custom=custom) + out = _check_modes( + tmp_folder, + [a, b], + engine, + log, + "Required properties are not included, the merged file will be invalid: producer", + includes=["foo"], + ) + assert out.name.startswith("merged") + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_modes_required_collection_only_property_differs(tmp_folder, engine, log): + custom = { + "properties": {"producer": {"type": "string"}}, + "collection": {"producer": True}, + "required": ["producer"], + } + a = _part(tmp_folder, "a", "a", 2, collection={"producer": "X"}, custom=custom) + b = _part(tmp_folder, "b", "b", 2, collection={"producer": "Y"}, custom=custom) + _check_modes( + tmp_folder, + [a, b], + engine, + log, + "Collection-only properties differ between the datasets and are removed: producer", + ) + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_modes_empty_geometries(tmp_folder, engine, log): + gdf = gpd.GeoDataFrame( + {"id": ["e0", "e1"], "geometry": [shapely.box(0, 0, 1, 1), shapely.Polygon()]}, + crs="EPSG:4326", + ) + a = tmp_folder / "empty.parquet" + gp = GeoParquet(a) + gp.set_collection({"schemas": {"a": [CORE]}, "collection": "a"}) + gp.write(gdf, dehydrate=False) + b = _part(tmp_folder, "b", "b", 2) + + _check_modes(tmp_folder, [a, b], engine, log, "1 of 4 rows have an empty or missing geometry") + # --no-strict keeps the rows + rows, _ = _read(_merge(tmp_folder, [a, b], engine, strict=False)) + assert len(rows) == 4 + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_modes_constant_that_does_not_fit_its_schema(tmp_folder, engine, log): + custom = {"properties": {"n": {"type": "uint8"}}} + a = _part(tmp_folder, "a", "a", 2, collection={"n": 300}, custom=custom) + b = _part(tmp_folder, "b", "b", 2, collection={"n": 1}, custom=custom) + out = _check_modes(tmp_folder, [a, b], engine, log, "n: Value 300 doesn't fit data type uint8") + + # --no-strict writes a null instead + rows, _ = _read(_merge(tmp_folder, [a, b], engine, strict=False)) + assert [(r["collection"], r["n"]) for r in rows] == [ + ("a", None), + ("a", None), + ("b", 1), + ("b", 1), + ] + assert out.name.startswith("merged") + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_hydrates_binary_constants(tmp_folder, engine): + custom = {"properties": {"blob": {"type": "binary"}}} + a = _part(tmp_folder, "a", "a", 2, collection={"blob": "aGVsbG8="}, custom=custom) + b = _part(tmp_folder, "b", "b", 2, collection={"blob": "d29ybGQ="}, custom=custom) + + rows, _ = _read(_merge(tmp_folder, [a, b], engine)) + assert [(r["collection"], r["blob"]) for r in rows] == [ + ("a", b"hello"), + ("a", b"hello"), + ("b", b"world"), + ("b", b"world"), + ] + + +def test_merge_modes_all_geometries_missing(tmp_folder, log): + paths = [] + for cid in ("a", "b"): + path = tmp_folder / f"{cid}.json" + features = [ + {"type": "Feature", "id": f"{cid}{i}", "geometry": None, "properties": {}} + for i in range(2) + ] + collection = {"schemas": {cid: [CORE]}, "collection": cid} + path.write_text( + json.dumps({"type": "FeatureCollection", **collection, "features": features}) + ) + paths.append(path) + + _check_modes( + tmp_folder, paths, "geopandas", log, "4 of 4 rows have an empty or missing geometry" + ) + + +def test_merge_writes_required_nested_properties_with_nulls_as_nullable(tmp_folder, log): + struct = pa.struct([("k", pa.string())]) + custom = { + "required": ["attrs"], + "properties": {"attrs": {"type": "object", "properties": {"k": {"type": "string"}}}}, + } + a = _part(tmp_folder, "a", "a", 2, columns={"attrs": [{"k": "v"}, {"k": "w"}]}, custom=custom) + a = _with_nullable_column(a, "attrs", [{"k": "v"}, None], struct) + b = _part(tmp_folder, "b", "a", 2, columns={"attrs": [{"k": "x"}, {"k": "y"}]}, custom=custom) + + out = _merge(tmp_folder, [a, b], "duckdb", strict=False) + assert "Rows have no value for a required property: attrs" in log() + assert pq.read_schema(out).field("attrs").nullable + rows, _ = _read(out) + assert [r["attrs"] for r in rows] == [{"k": "v"}, None, {"k": "x"}, {"k": "y"}] diff --git a/tests/test_ops.py b/tests/test_ops.py index 44aaaed..580588c 100644 --- a/tests/test_ops.py +++ b/tests/test_ops.py @@ -1,5 +1,7 @@ +import pytest + from vecorel_cli.vecorel.collection import Collection -from vecorel_cli.vecorel.ops import merge_collections +from vecorel_cli.vecorel.ops import merge_collections, warn_missing_required from vecorel_cli.vecorel.schemas import Schemas, VecorelSchema @@ -92,3 +94,123 @@ def test_merge_collections_keeps_collection_properties(): merged = merge_collections([collection1], properties=["only_in_first"]) assert "source_name" not in merged assert merged.get("only_in_first") == "x" + + +def test_merge_collections_unites_schemas_of_a_collection(): + core = "https://vecorel.org/specification/v0.1.0/schema.yaml" + ext = "https://example.com/ext.yaml" + merged = merge_collections( + [Collection({"schemas": {"c1": [core, ext]}}), Collection({"schemas": {"c1": [core]}})] + ) + assert merged.get_schemas() == Schemas({"c1": [core, ext]}) + + other_core = "https://vecorel.org/specification/v0.2.0/schema.yaml" + with pytest.raises(ValueError, match="conflicting core schemas"): + merge_collections( + [Collection({"schemas": {"c1": [core]}}), Collection({"schemas": {"c1": [other_core]}})] + ) + + fiboa = "https://fiboa.org/specification/v{}/schema.yaml" + with pytest.raises(ValueError, match="conflicting versions of a schema"): + merge_collections( + [ + Collection({"schemas": {"c1": [core, fiboa.format("0.2.0")]}}), + Collection({"schemas": {"c1": [core, fiboa.format("0.3.0")]}}), + ] + ) + # other collections may use another version + merge_collections( + [ + Collection({"schemas": {"c1": [core, fiboa.format("0.2.0")]}}), + Collection({"schemas": {"c2": [core, fiboa.format("0.3.0")]}}), + ] + ) + + +def test_merge_collections_warns_about_collection_only_properties(): + class Log: + warnings = [] + + def warning(self, message): + self.warnings.append(message) + + core = "https://vecorel.org/specification/v0.1.0/schema.yaml" + custom = {"properties": {"producer": {"type": "string"}}, "collection": {"producer": True}} + collections = [ + Collection({"schemas": {"c1": [core]}, "schemas:custom": custom, "producer": "A", "x": 1}), + Collection({"schemas": {"c2": [core]}, "schemas:custom": custom, "producer": "B", "x": 2}), + ] + log = Log() + merged = merge_collections(collections, log=log) + assert "producer" not in merged + # x is not collection-only, so it goes back into the features, which the caller handles + assert log.warnings == [ + "Collection-only properties differ between the datasets and are removed: producer" + ] + + +def test_warn_missing_required_collection_only_properties(tmp_path): + class Log: + warnings = [] + + def warning(self, message): + self.warnings.append(message) + + core = "https://vecorel.org/specification/v0.1.0/schema.yaml" + ext = "https://example.com/producer/v0.1.0/schema.yaml" + path = tmp_path / "schema.yaml" + path.write_text( + "$schema: https://vecorel.org/sdl/v0.2.0/schema.json\n" + "required: [producer]\n" + "collection: {producer: true}\n" + "properties: {producer: {type: string}}\n" + ) + schema_map = {ext: path} + properties = ["id", "geometry", "collection"] + collections = [ + Collection({"schemas": {"c1": [core, ext]}, "producer": "A"}), + Collection({"schemas": {"c2": [core, ext]}, "producer": "A"}), + ] + merged = merge_collections(collections, properties=properties) + assert "producer" not in merged + + log = Log() + warn_missing_required(collections, properties, schema_map, log) + message = "Required properties are not included, the merged file will be invalid: producer" + assert log.warnings == [message] + + with pytest.raises(ValueError, match=message): + warn_missing_required(collections, properties, schema_map, strict=True) + + +def test_merge_collections_uses_the_schema_map(tmp_path): + class Log: + warnings = [] + + def warning(self, message): + self.warnings.append(message) + + # the extension is only available through the schema map + core = "https://vecorel.org/specification/v0.1.0/schema.yaml" + ext = "https://example.com/producer/v0.1.0/schema.yaml" + path = tmp_path / "schema.yaml" + path.write_text( + "$schema: https://vecorel.org/sdl/v0.2.0/schema.json\n" + "required: [producer]\n" + "collection: {producer: true}\n" + "properties: {producer: {type: string}}\n" + ) + schema_map = {ext: path} + collections = [ + Collection({"schemas": {"c1": [core, ext]}, "producer": "A"}), + Collection({"schemas": {"c2": [core, ext]}, "producer": "B"}), + ] + message = "Collection-only properties differ between the datasets and are removed: producer" + + log = Log() + merged = merge_collections(collections, log=log, schema_map=schema_map) + assert "producer" not in merged + assert log.warnings == [message] + + with pytest.raises(ValueError, match=message): + merge_collections(collections, strict=True, schema_map=schema_map) diff --git a/tests/test_validate.py b/tests/test_validate.py index e0c65d1..6daaf19 100644 --- a/tests/test_validate.py +++ b/tests/test_validate.py @@ -108,3 +108,68 @@ def test_validate(test): assert error == expect assert not result.is_valid() + + +def _write_collections(path, collections, columns, schemas): + import geopandas as gpd + import shapely + + from vecorel_cli.encoding.geoparquet import GeoParquet + + n = len(collections) + gdf = gpd.GeoDataFrame( + { + "id": [str(i) for i in range(n)], + "collection": collections, + **columns, + "geometry": [shapely.box(i, 0, i + 1, 1) for i in range(n)], + }, + crs="EPSG:4326", + ) + gp = GeoParquet(path) + gp.set_collection({"schemas": schemas}) + gp.write(gdf, dehydrate=False) + + +def test_validate_rejects_features_without_a_known_collection(tmp_parquet_file): + core = "https://vecorel.org/specification/v0.1.0/schema.yaml" + _write_collections(tmp_parquet_file, ["a", None, "x"], {}, {"a": [core], "b": [core]}) + + errors = [str(e) for e in ValidateData().validate(tmp_parquet_file).errors] + assert errors == [ + "collection: 1 rows have no collection", + "collection: Not found in schemas: x", + ] + + +def test_validate_reports_a_missing_collection_column(tmp_parquet_file): + import geopandas as gpd + import shapely + + from vecorel_cli.encoding.geoparquet import GeoParquet + + core = "https://vecorel.org/specification/v0.1.0/schema.yaml" + gdf = gpd.GeoDataFrame( + {"id": ["1", "2"], "geometry": [shapely.box(0, 0, 1, 1)] * 2}, crs="EPSG:4326" + ) + gp = GeoParquet(tmp_parquet_file) + gp.set_collection({"schemas": {"a": [core], "b": [core]}}) + gp.write(gdf, dehydrate=False) + + errors = [str(e) for e in ValidateData().validate(tmp_parquet_file).errors] + assert errors == ["collection: Required field is missing"] + + +def test_validate_checks_required_properties_per_collection(tmp_parquet_file): + core = "https://vecorel.org/specification/v0.1.0/schema.yaml" + admin = "https://vecorel.org/administrative-division-extension/v0.1.0/schema.yaml" + schemas = {"a": [core, admin], "b": [core]} + columns = {"admin:country_code": ["DE", None, None]} + _write_collections(tmp_parquet_file, ["a", "a", "b"], columns, schemas) + + errors = [str(e) for e in ValidateData().validate(tmp_parquet_file).errors] + assert errors == ["admin:country_code: Required field has no value for collection 'a'"] + + columns = {"admin:country_code": ["DE", "DE", None]} + _write_collections(tmp_parquet_file, ["a", "a", "b"], columns, schemas) + assert ValidateData().validate(tmp_parquet_file).errors == [] diff --git a/vecorel_cli/conversion/duckdb.py b/vecorel_cli/conversion/duckdb.py index 99c9b8a..9cca45b 100644 --- a/vecorel_cli/conversion/duckdb.py +++ b/vecorel_cli/conversion/duckdb.py @@ -9,9 +9,13 @@ import pyarrow as pa import pyarrow.parquet as pq +from ..cli.logger import LoggerMixin from ..encoding.geojson import VecorelJSONEncoder from ..encoding.geoparquet import GeoParquet +from ..parquet.types import constant_array from ..vecorel.hilbert import hilbert_keys_for_table, hilbert_reference_bounds +from ..vecorel.ops import get_collection_id, merge_collections, report, warn_missing_required +from ..vecorel.util import find_differing_crs from .base import BaseConverter @@ -22,6 +26,11 @@ def _sql_path(path) -> str: return f"'{escaped}'" +def _sql_name(name) -> str: + escaped = str(name).replace('"', '""') + return f'"{escaped}"' + + # A COPY statement binds its own parameters before those of its subquery, so a value # placed in the SELECT cannot be a bound parameter: it would be read as the output path. def _sql_literal(value) -> str: @@ -45,17 +54,23 @@ def _sql_literal(value) -> str: return f"'{escaped}'" -def _normalize_crs(crs): - """The comparable pyproj CRS for a GeoParquet crs value, - which defaults to OGC:CRS84 when missing.""" - from pyproj import CRS - - return CRS.from_user_input(crs if crs is not None else "OGC:CRS84") - - -def _equal_crs(a, b) -> bool: - # GeoParquet coordinates are always x, y regardless of the CRS axis order - return a.equals(b, ignore_axis_order=True) +def _constants_table( + constants: dict, + props: dict, + log: Optional[LoggerMixin] = None, + source=None, + strict: bool = False, +) -> pa.Table: + """A single row with the given values, typed by the property schemas where available.""" + arrays = [] + for key, value in constants.items(): + try: + arrays.append(constant_array(value, props.get(key))) + except ValueError as e: + # the value would fail the writer, so it's left empty + report(f"{key}: {e} (in {source})", log, strict) + arrays.append(constant_array(None, props.get(key))) + return pa.table(arrays, names=list(constants.keys())) # This converter is experimental, use with caution. @@ -281,9 +296,8 @@ def _common_crs(self, con, sources: list): They are combined without reprojection, and DuckDB before 1.5 drops the CRS from the metadata, so it has to be read from the files themselves. """ - source_crs = None - reference = None - for i, source in enumerate(sources): + crs_values = [] + for source in sources: crs = None row = con.execute( "SELECT value FROM parquet_kv_metadata(?) WHERE key = 'geo'", [source] @@ -292,15 +306,14 @@ def _common_crs(self, con, sources: list): geo = json.loads(bytes(row[0])) primary = geo.get("primary_column", "") crs = geo.get("columns", {}).get(primary, {}).get("crs") - if i == 0: - source_crs = crs - reference = _normalize_crs(crs) - elif not _equal_crs(_normalize_crs(crs), reference): - raise ValueError( - f"The sources use different coordinate reference systems: {source} " - "differs from the first source. Reproject the sources to a common CRS." - ) - return source_crs + crs_values.append(crs) + differing = find_differing_crs(crs_values) + if differing is not None: + raise ValueError( + f"The sources use different coordinate reference systems: {sources[differing]} " + "differs from the first source. Reproject the sources to a common CRS." + ) + return crs_values[0] if crs_values else None def write_query( self, @@ -317,13 +330,19 @@ def write_query( compression_level: Optional[int] = None, geoparquet_version: Optional[str] = None, original_geometries: bool = False, + suffix_duplicate_ids: bool = True, + strict: bool = True, + merging: bool = False, ) -> str: """Write the rows a SELECT returns as a Vecorel GeoParquet file: drop rows that no required property or geometry survives, report an id that is not unique, normalize the geometries, sort into Hilbert order and package. `targets` names the properties the query returns, which is what the checks - run over. + run over. Ids only need to be unique per collection. + Without `strict`, a missing required value is only reported and empty + geometries are kept, as the in-memory merge does. + When `merging` strictly, empty geometries are an error instead of being dropped. """ compression = compression or "zstd" if compression == "zstd" and compression_level is None: @@ -346,23 +365,51 @@ def write_query( # Same required-value check, empty-geometry drop and id uniqueness # check as in the GeoDataFrame-based codepath, in one scan - schemas = collection.merge_schemas({}) collection_only = set(collection.get_collection_only_properties()) - required = [ - r - for r in schemas.get("required", []) - if r != "geometry" and r not in collection_only and r in selected_targets - ] + schema_groups = collection.get_schemas() + per_collection = "collection" in selected_targets and len(schema_groups) > 1 + + # The collections that require a property, None for all rows + required_by = {} + + def null_condition(schema, skip=("geometry",), cid=None): + required = [ + r + for r in schema.get("required", []) + if r not in skip and r not in collection_only and r in selected_targets + ] + for r in required: + required_by.setdefault(r, []).append(cid) + return " OR ".join(f"{_sql_name(target)} IS NULL" for target in required) or None + + if per_collection: + # Each collection only requires what its own schemas require + custom_schemas = collection.get_custom_schemas() + conditions = ['"collection" IS NULL'] + required_by["collection"] = [None] + for cid, group in schema_groups.items(): + schema = group.merge_schemas(custom_schemas=custom_schemas) + cond = null_condition(schema, skip=("geometry", "collection"), cid=cid) + if cond: + conditions.append(f'("collection" = {_sql_literal(cid)} AND ({cond}))') + null_cond = " OR ".join(conditions) + else: + null_cond = null_condition(collection.merge_schemas({})) stats = ["count(*)"] - null_cond = None - if required: - null_cond = " OR ".join(f'"{target}" IS NULL' for target in required) + if null_cond: stats.append(f"count(*) FILTER (WHERE {null_cond})") + # per property, so that the report can name what is missing + for name, cids in required_by.items(): + condition = f"{_sql_name(name)} IS NULL" + if None not in cids: + condition += f' AND "collection" IN ({", ".join(map(_sql_literal, cids))})' + stats.append(f"count(*) FILTER (WHERE {condition})") # row numbers are unique by construction check_ids = "id" in selected_targets and not ids_are_generated if check_ids: + id_key = 'struct_pack(c := "collection", i := "id")' if per_collection else '"id"' stats.append('count("id")') - stats.append('count(DISTINCT "id")') + stats.append(f'count(DISTINCT {id_key}) FILTER (WHERE "id" IS NOT NULL)') blank_cond = None repair_cond = None if "geometry" in selected_targets: @@ -383,30 +430,47 @@ def write_query( ) total = values.pop(0) invalid = values.pop(0) if null_cond else 0 + nulls = {name: values.pop(0) for name in required_by} if null_cond else {} if check_ids: non_null = values.pop(0) distinct = values.pop(0) - if distinct < non_null: + if distinct < non_null and suffix_duplicate_ids: self.warning( f"{type(self).__name__}: 'id' is not unique — {non_null - distinct:,} " f"of {non_null:,} rows repeat an id, so it cannot be `id`. Map a column " "that identifies a feature, or build one from the source's key columns. " "The repeating ids get a ~ suffix in the output." ) + elif distinct < non_null: + report( + f"{non_null - distinct:,} of {non_null:,} rows repeat an id within their collection", + self, + strict, + ) blanks = values.pop(0) if blank_cond else 0 repairs = values.pop(0) if repair_cond else 0 - if invalid: + missing = ", ".join( + f"{name} ({count:,} of {total:,} rows)" + for name, count in sorted(nulls.items()) + if count + ) + if invalid and (merging or not strict): + report(f"Rows have no value for a required property: {missing}", self, strict) + elif invalid: # A null in a required property is an error, whatever the count: # the writer rejects nulls in the non-nullable required fields # anyway, and silently dropping rows would make that data-quality # decision for the user (vecorel/cli#33) raise ValueError( - f"{invalid} of {total} rows have no value for a required property " - f"({null_cond}). Handle them in the converter: fix the mapping, " - "fill the values in column_migrations, or exclude the rows with " - "a column_filters entry." + f"Rows have no value for a required property: {missing}. " + "Handle them in the converter: fix the mapping, fill the values in " + "column_migrations, or exclude the rows with a column_filters entry." + ) + if blanks and merging: + report( + f"{blanks:,} of {total:,} rows have an empty or missing geometry", self, strict ) - if blanks: + elif blanks and strict: self.warning(f"Dropping {blanks} of {total} rows with an empty or missing geometry") source_query = f"SELECT * FROM ({source_query}) WHERE NOT ({blank_cond})" if repairs: @@ -467,9 +531,9 @@ def write_query( [output_file, compression, collection_json], ) - if "id" in selected_targets: + if "id" in selected_targets and suffix_duplicate_ids: self._suffix_duplicate_ids_in_file( - con, output_file, compression, collection_json, row_group_size + con, output_file, compression, collection_json, row_group_size, per_collection ) # Sort against the same CRS-derived Hilbert grid as the @@ -507,20 +571,22 @@ def write_query( compression_level=compression_level, geoparquet_version=geoparquet_version, crs=source_crs, + strict=strict, ) return output_file def _suffix_duplicate_ids_in_file( - self, con, output_file, compression, collection_json, row_group_size + self, con, output_file, compression, collection_json, row_group_size, per_collection ): - """Number the ids that appear on several rows (id~1, id~2, ...), like the + """Number the ids that appear on several rows of a collection (id~1, id~2, ...), like the GeoDataFrame-based codepath, when the source does not provide unique ids (the pre-write check only warns). Runs on the written file, so the id column keeps its type when nothing repeats; the rewrite may reorder rows, which the Hilbert sort afterwards puts right. A numbered id can collide with one the source already carries (x~1), so this repeats until nothing repeats.""" + key = '"collection", "id"' if per_collection else '"id"' reported = False while True: total, duplicated = con.execute( @@ -528,7 +594,7 @@ def _suffix_duplicate_ids_in_file( SELECT coalesce(sum(n), 0), coalesce(sum(n) FILTER (WHERE n > 1), 0) FROM ( SELECT count(*) AS n FROM read_parquet({_sql_path(output_file)}) - WHERE "id" IS NOT NULL GROUP BY "id" + WHERE "id" IS NOT NULL GROUP BY {key} ) """ ).fetchone() @@ -551,9 +617,9 @@ def _suffix_duplicate_ids_in_file( f""" COPY ( SELECT * EXCLUDE (file_row_number) REPLACE ( - CASE WHEN count(*) OVER (PARTITION BY "id") > 1 + CASE WHEN count(*) OVER (PARTITION BY {key}) > 1 THEN CAST("id" AS VARCHAR) || '~' || CAST( - row_number() OVER (PARTITION BY "id" ORDER BY file_row_number) + row_number() OVER (PARTITION BY {key} ORDER BY file_row_number) AS VARCHAR) ELSE CAST("id" AS VARCHAR) END AS "id") @@ -575,18 +641,31 @@ def _suffix_duplicate_ids_in_file( os.unlink(tmp_path) raise - def merge_parquet(self, paths: list, output_file, collection=None, **kwargs) -> str: + def merge_parquet( + self, + paths: list, + output_file, + collection=None, + properties=None, + strict: bool = True, + **kwargs, + ) -> str: """Combine Vecorel GeoParquet files into one, checked and sorted over the whole set rather than per file. The parts are written by a converter, so they need no column mapping and their geometries are already valid polygons; pass `original_geometries=False` - to run the geometry step anyway. + to run the geometry step anyway. `properties` restricts the properties that + are merged. The bbox is always recomputed, as not all parts may have one. + Problems that make the merged file invalid are errors if `strict`, + otherwise warnings. """ if not paths: raise ValueError("No paths to merge") paths = [str(path) for path in paths] kwargs.setdefault("original_geometries", True) + if properties is not None: + properties = set(properties) | {"geometry"} con = duckdb.connect() con.install_extension("spatial") @@ -596,41 +675,70 @@ def merge_parquet(self, paths: list, output_file, collection=None, **kwargs) -> collections = [GeoParquet(path).get_collection() for path in paths] if collection is None: - from ..vecorel.ops import merge_collections - - collection = merge_collections(collections) - - # A part keeps its constants in its collection; one the merged collection does not - # carry, because the parts disagree on it, goes back into the rows, as `vec merge` does - hydrate = sorted( - { - key - for part in collections - for key in part.keys() - if key not in collection and key not in part.get_collection_only_properties() - } - - {"schemas"} - ) + collection = merge_collections( + collections, properties=properties, log=self, strict=strict + ) + if properties is not None: + warn_missing_required(collections, properties, {}, self, strict) + props = collection.merge_schemas({}).get("properties", {}) # union by name, because a converter drops a column a source file does not have, # so two parts of one dataset can legitimately differ; the targets then come from # the union rather than from whichever part happens to be first - if hydrate: - selects = [] - for path, part in zip(paths, collections): - names = pq.read_schema(path).names - literals = [ - f'{_sql_literal(part.get(key))} AS "{key}"' - for key in hydrate - if key not in names - ] - selects.append( - f"SELECT {', '.join(['*', *literals])} FROM read_parquet({_sql_path(path)})" + selects = [] + for i, (path, part) in enumerate(zip(paths, collections)): + names = pq.read_schema(path).names + collection_only = set(part.get_collection_only_properties()) + # A part keeps its constants in its collection; one the merged collection does not + # carry, because the parts disagree on it, goes back into the rows, as `vec merge` does + constants = { + key: value + for key, value in part.items() + if key not in collection + and key not in collection_only + and key != "schemas" + and key not in names + and (properties is None or key in properties) + } + keep_collection = properties is None or "collection" in properties + fill = get_collection_id(part) if keep_collection else None + if fill is not None and "collection" not in names and "collection" not in collection: + # Features of multiple collections must state their collection + constants["collection"] = fill + elif fill is None and keep_collection: + # the statistics tell whether the column has nulls without a scan + if "collection" in names: + with pq.ParquetFile(path) as pq_file: + unknown = bool(GeoParquet._columns_with_nulls(pq_file, {"collection"})) + else: + unknown = "collection" not in collection + if unknown: + report( + f"Can't determine the collection of the features in {path}", self, strict + ) + + columns = [] + for name in names: + if name == "bbox" or (properties is not None and name not in properties): + continue + if name == "collection" and fill is not None: + columns.append(f'coalesce("collection", {_sql_literal(fill)}) AS "collection"') + else: + columns.append(_sql_name(name)) + query = f"SELECT {', '.join(columns)}" + source = f"read_parquet({_sql_path(path)})" + if constants: + # A registered table rather than SQL literals, so arrays, objects and + # temporal values keep their types + table = f"constants_{i}" + con.register( + table, + _constants_table(constants, props, log=self, source=path, strict=strict), ) - source_query = " UNION ALL BY NAME ".join(selects) - else: - sources = "[" + ",".join(_sql_path(path) for path in paths) + "]" - source_query = f"SELECT * FROM read_parquet({sources}, union_by_name=true)" + query += f", {table}.*" + source += f" CROSS JOIN {table}" + selects.append(f"{query} FROM {source}") + source_query = " UNION ALL BY NAME ".join(selects) targets = [row[0] for row in con.execute(f"DESCRIBE {source_query}").fetchall()] return self.write_query( @@ -640,6 +748,8 @@ def merge_parquet(self, paths: list, output_file, collection=None, **kwargs) -> collection, targets=targets, source_crs=source_crs, + strict=strict, + merging=True, **kwargs, ) diff --git a/vecorel_cli/create_geoparquet.py b/vecorel_cli/create_geoparquet.py index 7aefd76..fedf455 100644 --- a/vecorel_cli/create_geoparquet.py +++ b/vecorel_cli/create_geoparquet.py @@ -59,7 +59,9 @@ def create( # Read source data encodings = [create_encoding(s) for s in source] # Merge encodings into a single GeoDataFrame - geodata, collection = merge(encodings, properties=properties, schema_map=schema_map) + geodata, collection = merge( + encodings, properties=properties, schema_map=schema_map, log=self + ) # Write to target target_encoding = GeoParquet(target) diff --git a/vecorel_cli/encoding/base.py b/vecorel_cli/encoding/base.py index c6dbc85..84ee33a 100644 --- a/vecorel_cli/encoding/base.py +++ b/vecorel_cli/encoding/base.py @@ -1,6 +1,7 @@ from pathlib import Path from typing import Optional, Union +import pandas as pd from fsspec import AbstractFileSystem from geopandas import GeoDataFrame from yarl import URL @@ -111,20 +112,32 @@ def read( raise NotImplementedError("Not supported by encoding") def hydrate_from_collection( - self, data: GeoDataFrame, schema_map: SchemaMapping = {} + self, + data: GeoDataFrame, + schema_map: SchemaMapping = {}, + keys: Optional[list[str]] = None, ) -> GeoDataFrame: """ Merge the collection metadata into the GeoDataFrame. + + If `keys` is specified, only those keys are merged. """ collection = self.get_collection() collection_only = collection.get_collection_only_properties(schema_map=schema_map) - keys = list(collection.keys()) + if keys is None: + keys = list(collection.keys()) for key in keys: + if key not in collection: + continue value = collection[key] if key in collection_only: continue if key not in data.columns: - data[key] = value + if isinstance(value, (list, dict)): + # pandas would spread a list over the rows and align a dict to the index + data[key] = pd.Series([value] * len(data), index=data.index, dtype=object) + else: + data[key] = value collection.pop(key, None) return data @@ -155,6 +168,10 @@ def dehydrate_to_collection( if properties and key not in properties: continue + # All-null columns can't be expressed in JSON reliably (NaN, NA), keep them as they are + if data[key].isna().all(): + continue + # todo: This only works for scalar values, how to handle dicts, etc? try: if data[key].nunique(dropna=False) == 1: diff --git a/vecorel_cli/encoding/geoparquet.py b/vecorel_cli/encoding/geoparquet.py index db45033..d359f49 100644 --- a/vecorel_cli/encoding/geoparquet.py +++ b/vecorel_cli/encoding/geoparquet.py @@ -158,6 +158,8 @@ def write( compression: Optional[str] = "zstd", compression_level: Optional[int] = None, # default level for compression geoparquet_version: Optional[str] = None, + # False writes required columns that contain nulls as nullable instead of failing + strict: bool = True, **kwargs, # capture unknown arguments ) -> bool: if compression == "zstd" and compression_level is None: @@ -192,6 +194,8 @@ def write( pq_fields = [] for column in properties: required = column in required_props and not has_multiple_collections + if required and not strict and data[column].isna().any(): + required = False schema = props.get(column, {}) dtype = schema.get("type") @@ -278,6 +282,8 @@ def postprocess( compression_level: Optional[int] = None, geoparquet_version: Optional[str] = None, crs=None, # the CRS to record in the GeoParquet metadata, e.g. from the source file + # False keeps required columns that contain nulls nullable instead of failing + strict: bool = True, **kwargs, # capture unknown arguments ) -> bool: """ @@ -324,6 +330,7 @@ def postprocess( geoparquet_version, crs=crs, compression_changed=compression != existing_compression, + strict=strict, ) if tmp_path is None: return False @@ -338,6 +345,33 @@ def postprocess( self.pq_schema = None return True + @staticmethod + def _columns_with_nulls(pq_file: pq.ParquetFile, names: set[str]) -> set[str]: + """ + The columns that contain nulls according to the statistics, or that have none. + Nested columns only have statistics for their leaves, so they are read instead. + """ + metadata = pq_file.metadata + result = set() + flat = set() + for rg in range(metadata.num_row_groups): + row_group = metadata.row_group(rg) + for i in range(row_group.num_columns): + column = row_group.column(i) + name = column.path_in_schema + if name not in names: + continue + flat.add(name) + stats = column.statistics + if stats is None or not stats.has_null_count or stats.null_count > 0: + result.add(name) + for name in (names - flat) & set(pq_file.schema_arrow.names): + for rg in range(metadata.num_row_groups): + if pq_file.read_row_group(rg, columns=[name]).column(0).null_count > 0: + result.add(name) + break + return result + # Rewrites the Parquet file to a temp file and returns its path, # or returns None if the file needs no changes def _rewrite( @@ -349,6 +383,7 @@ def _rewrite( geoparquet_version: str, crs=None, compression_changed: bool = False, + strict: bool = True, ) -> Optional[str]: existing_schema = pq_file.schema_arrow col_names = existing_schema.names @@ -368,6 +403,8 @@ def _rewrite( if "id" in col_names: required_columns.add("id") required_columns |= {r for r in schemas.get("required", []) if r in col_names} + if not strict: + required_columns -= self._columns_with_nulls(pq_file, required_columns) if "bbox" in col_names: bbox_type = existing_schema.field("bbox").type diff --git a/vecorel_cli/merge.py b/vecorel_cli/merge.py index c86e365..09d61ff 100644 --- a/vecorel_cli/merge.py +++ b/vecorel_cli/merge.py @@ -1,18 +1,20 @@ from pathlib import Path -from typing import Union +from typing import Optional, Union import click from yarl import URL from .basecommand import BaseCommand, runnable from .cli.options import ( - CRS, VECOREL_FILES_ARG, VECOREL_TARGET, ) from .encoding.auto import create_encoding +from .encoding.base import BaseEncoding +from .encoding.geoparquet import GeoParquet from .registry import Registry from .vecorel.ops import merge as merge_ +from .vecorel.util import find_differing_crs class MergeDatasets(BaseCommand): @@ -22,27 +24,40 @@ class MergeDatasets(BaseCommand): Merges multiple {Registry.project} datasets to a combined {Registry.project} dataset. This simply appends the datasets to each other. - It does not check for duplicates or other constraints. - It adds a collection column to discriminate the source of the rows. + Ids that repeat within a collection are reported, but not changed. + Each feature keeps its collection, which is stored in a column if the + datasets have multiple collections. + + Local GeoParquet files that are all in the target CRS are merged with DuckDB, + which doesn't need to fit the data into memory. All other datasets are merged in memory. + + By default, problems that make the merged dataset invalid are errors. + With --no-strict, they are reported as warnings and the dataset is written anyway. """ default_crs = "EPSG:4326" + first_crs = "first" + engines = ["auto", "duckdb", "geopandas"] @staticmethod def get_cli_args(): return { "source": VECOREL_FILES_ARG, "target": VECOREL_TARGET(), - "crs": CRS(MergeDatasets.default_crs), + "crs": click.option( + "--crs", + type=click.STRING, + help=f"GeoParquet only: Coordinate Reference System (CRS) to use for the file. Use '{MergeDatasets.first_crs}' for the CRS of the first dataset.", + show_default=True, + default=MergeDatasets.default_crs, + ), "include": click.option( "--include", "-i", "includes", type=click.STRING, multiple=True, - help="Properties to include in addition to the core properties.", - show_default=True, - default=Registry.core_properties, + help="Properties to include in addition to the core properties. Includes all properties if not given.", ), "exclude": click.option( "--exclude", @@ -50,7 +65,20 @@ def get_cli_args(): "excludes", type=click.STRING, multiple=True, - help="Core properties to exclude.", + help="Properties to exclude.", + ), + "engine": click.option( + "--engine", + type=click.Choice(MergeDatasets.engines), + help="The engine to merge with, 'auto' uses DuckDB if possible.", + show_default=True, + default="auto", + ), + "strict": click.option( + "--strict/--no-strict", + default=True, + show_default=True, + help="Fail if the merged dataset would be invalid. With --no-strict, warn and write the dataset anyway.", ), } @@ -62,25 +90,96 @@ def merge( crs=None, includes=[], excludes=[], + engine="auto", + strict=True, ): if not isinstance(source, list): raise ValueError("Source must be a list.") if len(source) == 0: raise ValueError("No source files provided") + if engine not in self.engines: + raise ValueError(f"Engine must be one of {', '.join(self.engines)}") encodings = [create_encoding(s) for s in source] + for encoding in encodings: + if not encoding.exists(): + raise FileNotFoundError(f"File not found: {encoding.uri}") if isinstance(target, str): target = Path(target) + target = create_encoding(target) if not crs: crs = self.default_crs + elif crs == self.first_crs: + crs = None - properties = Registry.core_properties.copy() - properties.extend(includes) - properties = list(set(properties) - set(excludes)) + properties = None + if includes: + properties = list(set(Registry.core_properties) | set(includes)) - gdf, collection = merge_(encodings, crs=crs, properties=properties) + blocker = self.get_duckdb_blocker(encodings, target, crs) + if engine == "duckdb" and blocker: + raise ValueError(f"Can't merge with DuckDB: {blocker}") - target = create_encoding(target) - target.set_collection(collection) - target.write(gdf, properties=properties) + if engine == "duckdb" or (engine == "auto" and blocker is None): + from .conversion.duckdb import DuckDBBaseConverter + + if excludes: + if properties is None: + properties = self.get_available_properties(encodings) + properties = list(set(properties) - set(excludes)) + + self.info("Merging with DuckDB") + DuckDBBaseConverter().merge_parquet( + [e.uri for e in encodings], + target.uri, + properties=properties, + suffix_duplicate_ids=False, + strict=strict, + ) + else: + if engine == "auto": + self.info(f"Merging in memory, as {blocker}") + gdf, collection = merge_( + encodings, + crs=crs, + properties=properties, + log=self, + excludes=excludes, + strict=strict, + ) + target.set_collection(collection) + target.write(gdf, properties=properties, strict=strict) return target + + @staticmethod + def get_available_properties(encodings: list[BaseEncoding]) -> list[str]: + properties = set() + for encoding in encodings: + columns = encoding.get_properties() + if columns is None: + columns = encoding.read().columns + properties |= set(columns) + properties |= set(encoding.get_collection().keys()) + return list(properties) + + @staticmethod + def get_duckdb_blocker( + encodings: list[BaseEncoding], target: BaseEncoding, crs=None + ) -> Optional[str]: + """ + The reason why the datasets can't be merged with DuckDB, None if they can. + A `crs` of None stands for the CRS of the first dataset. + """ + for encoding in [*encodings, target]: + if not isinstance(encoding, GeoParquet) or not isinstance(encoding.uri, Path): + return "DuckDB only merges local GeoParquet files" + + crs_values = [] + for encoding in encodings: + geo = encoding.get_geoparquet_metadata() or {} + column = geo.get("columns", {}).get(geo.get("primary_column"), {}) + crs_values.append(column.get("crs")) + if find_differing_crs(crs_values, reference=crs or None) is not None: + return "the datasets must be reprojected to a common CRS" + + return None diff --git a/vecorel_cli/parquet/types.py b/vecorel_cli/parquet/types.py index d7ba6ab..99dbf43 100644 --- a/vecorel_cli/parquet/types.py +++ b/vecorel_cli/parquet/types.py @@ -1,4 +1,7 @@ +import base64 +import binascii import datetime +from typing import Optional import numpy as np import pandas as pd @@ -166,7 +169,7 @@ def get_pyarrow_type(schema): elif len(pattern_properties) > 0: if len(pattern_properties) > 1: raise Exception("Multiple pattern properties are not supported") - _, subschema = pattern_properties.popitem() + subschema = next(iter(pattern_properties.values())) values = get_pyarrow_type(subschema) return pa.map_(pa.string(), values) else: @@ -307,3 +310,40 @@ def get_pyarrow_type_for_geopandas(dtype): } SUPPORTED_PROTOCOLS = ["http", "https", "s3", "gs"] + + +def from_json_value(value, dtype=None): + """ + A collection-level value as a Python value for the data type. + Collection metadata is JSON, so temporal and binary values are encoded as strings. + """ + if value is None or (isinstance(value, float) and np.isnan(value)): + return None + if isinstance(value, str): + if dtype == "date-time": + ts = pd.Timestamp(value) + ts = ts.tz_localize("UTC") if ts.tzinfo is None else ts.tz_convert("UTC") + return ts.to_pydatetime() + if dtype == "date": + return datetime.date.fromisoformat(value[:10]) + if dtype == "binary": + return base64.b64decode(value, validate=True) + return value + + +def constant_array(value, schema: Optional[dict] = None, length: int = 1) -> pa.Array: + """ + An array that repeats a collection-level value, typed by its schema if given. + Raises a ValueError if the value doesn't fit the schema. + """ + dtype = (schema or {}).get("type") + try: + pa_type = get_pyarrow_type(schema) if dtype else None + except Exception: + pa_type = None + if pa_type is None: + return pa.array([from_json_value(value)] * length) + try: + return pa.array([from_json_value(value, dtype)] * length, type=pa_type) + except (pa.ArrowException, ValueError, TypeError, OverflowError, binascii.Error) as e: + raise ValueError(f"Value {value!r} doesn't fit data type {dtype}: {e}") from e diff --git a/vecorel_cli/validation/base.py b/vecorel_cli/validation/base.py index c5e5643..320ee50 100644 --- a/vecorel_cli/validation/base.py +++ b/vecorel_cli/validation/base.py @@ -87,7 +87,7 @@ def validate_schemas(self, groups: Schemas): # Check whether all required schemas are present, not needed in vecorel but other specs if len(self.required_schemas) == 0: - return + continue for pattern in self.required_schemas: found = False diff --git a/vecorel_cli/validation/geoparquet.py b/vecorel_cli/validation/geoparquet.py index e76d262..99de5c7 100644 --- a/vecorel_cli/validation/geoparquet.py +++ b/vecorel_cli/validation/geoparquet.py @@ -51,6 +51,9 @@ def _validate(self, num: Optional[int] = None, schema_map: SchemaMapping = {}): if key not in columns and not collection_props.get(key, False): self.error(f"{key}: Required field is missing") + if validate_data and "collection" in columns: + self.validate_collection_column(data, schemas, collection, schema_map) + # Validate whether the Parquet schema complies with the property schemas geo = self.encoding.get_geoparquet_metadata() property_schemas = schema.get("properties", {}) @@ -102,9 +105,10 @@ def _validate(self, num: Optional[int] = None, schema_map: SchemaMapping = {}): # Validate data of the column issues = [] - if validate_data and not has_multiple_collections: + # Without a collection column the rows can't be grouped, which is reported above + if validate_data and (not has_multiple_collections or "collection" not in columns): issues = validate_column(data[key], prop_schema) - elif validate_data and has_multiple_collections: + elif validate_data: # Validate data for each collection separately for cid, cschema in schemas.items(): vecorel_schema = cschema.merge_schemas( @@ -132,6 +136,35 @@ def _validate(self, num: Optional[int] = None, schema_map: SchemaMapping = {}): f"Data was not fully validated, only the first {num} rows were checked", ) + def validate_collection_column(self, data, schemas, collection, schema_map: SchemaMapping): + values = data["collection"] + missing = int(values.isna().sum()) + if missing: + self.error(f"collection: {missing} rows have no collection") + + unknown = set(values.dropna().unique()) - set(schemas.keys()) + if unknown: + self.error(f"collection: Not found in schemas: {', '.join(sorted(map(str, unknown)))}") + + if len(schemas) <= 1: + return + # Nullability can't express what each collection requires, so check the values + custom_schemas = collection.get_custom_schemas() + for cid, cschema in schemas.items(): + vecorel_schema = cschema.merge_schemas( + schema_map=schema_map, custom_schemas=custom_schemas, validator=self + ) + collection_props = vecorel_schema.get("collection", {}) + rows = data[values == cid] + for key in vecorel_schema.get("required", []): + if collection_props.get(key, False): + continue + if key not in rows.columns: + if len(rows) > 0: + self.error(f"{key}: Required field is missing for collection '{cid}'") + elif rows[key].isna().any(): + self.error(f"{key}: Required field has no value for collection '{cid}'") + def validate_geometry_column(self, key, prop_schema, geo): columns = geo.get("columns", {}) if key not in columns: diff --git a/vecorel_cli/vecorel/ops.py b/vecorel_cli/vecorel/ops.py index 94b5778..1ec4249 100644 --- a/vecorel_cli/vecorel/ops.py +++ b/vecorel_cli/vecorel/ops.py @@ -1,24 +1,93 @@ +from typing import Optional + import pandas as pd from geopandas import GeoDataFrame +from ..cli.logger import LoggerMixin from ..encoding.base import BaseEncoding +from ..parquet.types import NULLABLE_INTEGERS, constant_array from ..vecorel.collection import Collection from ..vecorel.schemas import Schemas, VecorelSchema from ..vecorel.typing import SchemaMapping +def report(message: str, log: Optional[LoggerMixin] = None, strict: bool = False): + """ + A problem that makes the merged file invalid: + an error in strict mode, otherwise a warning. + """ + if strict: + raise ValueError(message) + if log: + log.warning(message) + + def merge( encodings: list[BaseEncoding], crs=None, properties=None, schema_map: SchemaMapping = {}, + log: Optional[LoggerMixin] = None, + excludes: Optional[list[str]] = None, + strict: bool = False, ) -> tuple[GeoDataFrame, Collection]: + frames = [item.read(properties=properties, schema_map=schema_map) for item in encodings] + collections = [item.get_collection() for item in encodings] + if excludes: + if properties is None: + properties = set() + for gdf, collection in zip(frames, collections): + properties |= set(gdf.columns) | set(collection.keys()) + properties = list(set(properties) - set(excludes)) + frames = [gdf.drop(columns=[c for c in excludes if c in gdf.columns]) for gdf in frames] + merged_collection = merge_collections( + collections, properties=properties, log=log, strict=strict, schema_map=schema_map + ) + if properties is not None: + warn_missing_required(collections, properties, schema_map, log, strict) + props = merged_collection.merge_schemas(schema_map=schema_map).get("properties", {}) + data = [] - collections = [] + for item, gdf, collection in zip(encodings, frames, collections): + # Only what the merged collection doesn't carry goes back into the rows + keys = [ + key + for key in collection + if key not in merged_collection and (properties is None or key in properties) + ] + # Constants with a schema get their data type, as in DuckDB, e.g. binary is decoded; + # a value that doesn't fit the data type would fail the writer, so it's left empty + collection_only = set(collection.get_collection_only_properties(schema_map=schema_map)) + typed = {} + for key in keys: + if key in gdf.columns or key in collection_only or key not in collection: + continue + schema = props.get(key) + try: + array = constant_array(collection[key], schema, len(gdf)) + except ValueError as e: + report(f"{key}: {e} (in {item.uri})", log, strict) + array = constant_array(None, schema, len(gdf)) + if (schema or {}).get("type"): + typed[key] = array + gdf = item.hydrate_from_collection(gdf, schema_map=schema_map, keys=keys) + for key, array in typed.items(): + series = array.to_pandas(types_mapper=NULLABLE_INTEGERS.get) + series.index = gdf.index + gdf[key] = series - for item in encodings: - # Load the dataset - gdf = item.read(hydrate=True, properties=properties, schema_map=schema_map) + keep_collection = properties is None or "collection" in properties + cid = get_collection_id(collection) if keep_collection else None + if cid is not None and "collection" in gdf.columns: + gdf["collection"] = gdf["collection"].fillna(cid) + elif cid is not None: + gdf["collection"] = cid + elif keep_collection and ( + gdf["collection"].isna().any() + if "collection" in gdf.columns + else "collection" not in merged_collection + ): + report(f"Can't determine the collection of the features in {item.uri}", log, strict) if not crs: # If no CRS is given, use the first CRS that is available as the base CRS @@ -27,22 +96,128 @@ def merge( # Change the CRS if necessary gdf.to_crs(crs=crs, inplace=True) - # Add data to lists data.append(gdf) - collections.append(item.get_collection()) # Concatenate all GeoDataFrames to a single GeoDataFrame merged = GeoDataFrame(pd.concat(data, ignore_index=True)) - # Remove empty columns - merged.dropna(axis=1, how="all", inplace=True) + # Remove empty columns, except for the geometry, which is required + geometry = merged.geometry.name + merged = merged.drop( + columns=[c for c in merged.columns if c != geometry and merged[c].isna().all()] + ) + + if "id" in merged.columns: + with_id = merged[merged["id"].notna()] + key = [c for c in ("collection", "id") if c in merged.columns] + duplicates = int(with_id.duplicated(subset=key).sum()) + if duplicates: + report(f"{duplicates} rows repeat an id within their collection", log, strict) + + check_merged_data(merged, merged_collection, schema_map, log, strict, properties) + + return merged, merged_collection + + +def check_merged_data( + merged: GeoDataFrame, + collection: Collection, + schema_map: SchemaMapping = {}, + log: Optional[LoggerMixin] = None, + strict: bool = False, + properties=None, +): + """ + Reports empty geometries and missing required values, per collection. + Properties that are not selected are reported by warn_missing_required. + """ + geometries = merged.geometry + empty = int((geometries.isna() | geometries.is_empty).sum()) + if empty: + report(f"{empty} of {len(merged)} rows have an empty or missing geometry", log, strict) + + groups = collection.get_schemas() + multiple = len(groups) > 1 + collection_only = set(collection.get_collection_only_properties(schema_map=schema_map)) + custom_schemas = collection.get_custom_schemas() + missing = set() + for cid, group in groups.items(): + if multiple: + if "collection" not in merged.columns: + continue + rows = merged[merged["collection"] == cid] + else: + rows = merged + if len(rows) == 0: + continue + schema = group.merge_schemas(schema_map=schema_map, custom_schemas=custom_schemas) + for key in schema.get("required", []): + if ( + key == "geometry" + or key in collection_only + or key in collection + or (properties is not None and key not in properties) + ): + continue + if key not in rows.columns or rows[key].isna().any(): + missing.add(key) + if multiple and (properties is None or "collection" in properties): + if "collection" not in merged.columns or merged["collection"].isna().any(): + missing.add("collection") + if missing: + report( + f"Rows have no value for a required property: {', '.join(sorted(missing))}", + log, + strict, + ) + - # Merge all collections - collection = merge_collections(collections, properties=properties) +def get_collection_id(collection: Collection) -> Optional[str]: + """ + The collection id of a dataset that doesn't state it for all features: + the `collection` value in the collection metadata or the only collection in `schemas`. + None if neither determines it. + """ + cid = collection.get("collection") + if isinstance(cid, str) and len(cid) > 0: + return cid + schemas = collection.get_schemas() + if len(schemas) == 1: + return next(iter(schemas.keys())) + return None - return merged, collection + +def warn_missing_required( + collections: list[Collection], + properties: list[str], + schema_map: SchemaMapping, + log: Optional[LoggerMixin] = None, + strict: bool = False, +): + """ + Reports required properties that the selected properties don't include. + Checks the source collections, as the selection also drops the custom schemas + that require a property. + """ + missing = set() + for collection in collections: + schema = collection.merge_schemas(schema_map=schema_map) + missing |= set(schema.get("required", [])) - set(properties) + missing -= {"geometry", "schemas"} + if missing: + report( + f"Required properties are not included, the merged file will be invalid: {', '.join(sorted(missing))}", + log, + strict, + ) -def merge_collections(collections: list[Collection], properties=None) -> Collection: +def merge_collections( + collections: list[Collection], + properties=None, + log: Optional[LoggerMixin] = None, + strict: bool = False, + schema_map: SchemaMapping = {}, +) -> Collection: schemas = Schemas() custom_schemas = VecorelSchema() handled_separately = ("schemas", "schemas:custom") @@ -69,6 +244,27 @@ def merge_collections(collections: list[Collection], properties=None) -> Collect other_props = {k: v for k, v in other_props.items() if k in properties} custom_schemas = custom_schemas.pick(properties) + if log or strict: + # The other properties go back into the features, but collection-only properties can't + dropped = set() + required = set() + for c in collections: + keys = { + k + for k in c.keys() + if k not in other_props + and k not in handled_separately + and (properties is None or k in properties) + } + if keys: + dropped |= keys & set(c.get_collection_only_properties(schema_map=schema_map)) + required |= set(c.merge_schemas(schema_map=schema_map).get("required", [])) + message = "Collection-only properties differ between the datasets and are removed: " + if dropped & required: + report(message + ", ".join(sorted(dropped & required)), log, strict) + if dropped - required and log: + log.warning(message + ", ".join(sorted(dropped - required))) + collection = Collection({"schemas": schemas, **other_props}) collection.set_custom_schemas(custom_schemas) diff --git a/vecorel_cli/vecorel/schemas.py b/vecorel_cli/vecorel/schemas.py index c4d6a40..f3b740b 100644 --- a/vecorel_cli/vecorel/schemas.py +++ b/vecorel_cli/vecorel/schemas.py @@ -312,6 +312,7 @@ def merge_schemas( class Schemas(dict): spec_pattern = r"https://vecorel.org/specification/v([^/]+)/schema.yaml" spec_schema = "https://vecorel.org/specification/v{version}/schema.yaml" + version_pattern = r"/v\d+\.\d+\.\d+[^/]*/" @staticmethod def get_core_uri(version: Optional[str] = None) -> str: @@ -330,8 +331,26 @@ def __init__(self, schemas: Union[RawSchemas, "Schemas"] = {}): def get_all(self) -> list[CollectionSchemas]: return list(self.values()) - def add_all(self, schemas: "Schemas"): - self.update(schemas) + def add_all(self, schemas: Union[RawSchemas, "Schemas"]): + """Add the schemas of other collections, uniting the lists of a collection that exists already.""" + for collection, schema_list in schemas.items(): + merged = CollectionSchemas( + set(self.get(collection, set())) | set(schema_list), collection + ) + cores = {s for s in merged if re.match(Schemas.spec_pattern, s)} + if len(cores) > 1: + raise ValueError( + f"Collection '{collection}' has conflicting core schemas: {', '.join(sorted(cores))}" + ) + versions = {} + for uri in merged: + versions.setdefault(re.sub(Schemas.version_pattern, "/", uri), set()).add(uri) + for uris in versions.values(): + if len(uris) > 1: + raise ValueError( + f"Collection '{collection}' has conflicting versions of a schema: {', '.join(sorted(uris))}" + ) + self[collection] = merged def is_empty(self) -> bool: return len(self) == 0 diff --git a/vecorel_cli/vecorel/util.py b/vecorel_cli/vecorel/util.py index 1a9d30d..4e0812d 100644 --- a/vecorel_cli/vecorel/util.py +++ b/vecorel_cli/vecorel/util.py @@ -16,6 +16,25 @@ file_cache = {} +def find_differing_crs(crs_values: list, reference=None) -> Optional[int]: + """The index of the first GeoParquet crs value that differs from the reference + (by default the first value), None if they all agree. A missing crs is OGC:CRS84.""" + from pyproj import CRS + + def normalize(crs): + return CRS.from_user_input(crs if crs is not None else "OGC:CRS84") + + reference = normalize(reference) if reference is not None else None + for i, crs in enumerate(crs_values): + crs = normalize(crs) + if reference is None: + reference = crs + # GeoParquet coordinates are always x, y regardless of the CRS axis order + elif not crs.equals(reference, ignore_axis_order=True): + return i + return None + + def suffix_duplicate_ids(gdf): """Number the ids that appear on several rows (id~1, id~2, ...), which turns the id column into strings. Null ids are left alone. Returns the frame and From f7117c65b16473119a2dc95733167cbd3dc3cd02 Mon Sep 17 00:00:00 2001 From: Matthias Mohr Date: Mon, 28 Sep 2026 14:21:53 +0200 Subject: [PATCH 3/6] Error early when vecorel or extension versions diverge in a merge --- CHANGELOG.md | 9 +++-- README.md | 5 ++- tests/test_merge.py | 44 +++++++++++++++++++++++ tests/test_ops.py | 31 ++++++++++++++++- tests/test_validate.py | 15 ++++++++ vecorel_cli/conversion/duckdb.py | 52 ++++++++++++++++++++-------- vecorel_cli/encoding/geojson.py | 4 +-- vecorel_cli/validation/geoparquet.py | 12 +++++-- vecorel_cli/vecorel/ops.py | 23 ++++++++++++ 9 files changed, 171 insertions(+), 24 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 61a8766..92551ed 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -24,6 +24,8 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0. required properties that are not included, also collection-only ones. - **Breaking:** `vec merge` is strict by default: it fails if the merged dataset would be invalid, e.g. because an id repeats within a collection. +- **Breaking:** Merging fails before reading any data if the datasets use different + versions of the Vecorel specification or of an extension. - The message for missing required values names the missing properties with counts, instead of showing the SQL condition. @@ -33,7 +35,7 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0. its collection but on which the parts disagree; it becomes a column again, with the data type of its schema, as in `vec merge`. - Merging no longer overwrites the schemas of a collection that occurs in multiple - datasets, it unites them and rejects two versions of the same schema in one collection. + datasets, it unites them. - Merging warns about collection-only properties that differ between the datasets and have to be removed. - Merging fills in missing collection values of datasets that only list a single @@ -51,10 +53,13 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0. (invalid JSON) or failed for nullable integers. - `get_pyarrow_type` no longer modifies the schema of objects with `patternProperties`. - Validation checks that all features in GeoParquet files have a collection that is - listed in `schemas`, and the required properties per collection. + listed in `schemas`, and the required properties per collection instead of requiring + the properties of all collections for all features. - Validation checks the schemas of all collections, not only the first one. - Validation no longer fails with a `KeyError` for GeoParquet files with multiple collections, but without a collection column. +- Reading GeoJSON with a limit or a selection of properties no longer fails for features + without a geometry. ## [v0.3.1] - 2026-09-24 diff --git a/README.md b/README.md index 9688241..45537d1 100644 --- a/README.md +++ b/README.md @@ -152,9 +152,12 @@ the dataset is written anyway. Check it with `vec validate` afterwards. | An optional collection-only property differs between the datasets | Warning, the property is removed | Warning, the property is removed | | A feature has an empty or missing geometry | Error | Warning, the feature is kept | | A collection-level value doesn't fit the data type of its schema | Error | Warning, the value is left empty | -| A collection implements multiple versions of a schema, e.g. of an extension | Error | Error | +| The datasets use different versions of the Vecorel specification or of an extension | Error, before any data is read | Error, before any data is read | | The schemas of the datasets conflict | Error | Error | +Datasets with different versions of the Vecorel specification or of an extension can't be +merged, as the CLI can't upgrade them automatically yet. + The strict mode only checks what a merge can break or check with little effort. It doesn't validate the values against the schemas, e.g. patterns or value ranges, so a merged dataset is only valid if the values of the source datasets are valid. diff --git a/tests/test_merge.py b/tests/test_merge.py index 7f9b79d..0f7541e 100644 --- a/tests/test_merge.py +++ b/tests/test_merge.py @@ -666,3 +666,47 @@ def test_merge_writes_required_nested_properties_with_nulls_as_nullable(tmp_fold assert pq.read_schema(out).field("attrs").nullable rows, _ = _read(out) assert [r["attrs"] for r in rows] == [{"k": "v"}, None, {"k": "x"}, {"k": "y"}] + + +def test_merge_parquet_scans_all_parts_at_once_if_nothing_is_added(tmp_folder, monkeypatch): + from vecorel_cli.conversion.duckdb import DuckDBBaseConverter + + queries = [] + write_query = DuckDBBaseConverter.write_query + + def spy(self, con, source_query, *args, **kwargs): + queries.append(source_query) + return write_query(self, con, source_query, *args, **kwargs) + + monkeypatch.setattr(DuckDBBaseConverter, "write_query", spy) + + a = _part(tmp_folder, "a", "same", 2, collection={"region": "x"}) + b = _part(tmp_folder, "b", "same", 2, collection={"region": "x"}) + c = _part(tmp_folder, "c", "same", 2, collection={"region": "y"}) + + DuckDBBaseConverter().merge_parquet([a, b], tmp_folder / "shared.parquet") + assert queries[-1].count("read_parquet(") == 1 + assert len(pq.read_table(tmp_folder / "shared.parquet")) == 4 + + # the parts disagree on region, so it goes into the rows of each part + out = tmp_folder / "differing.parquet" + DuckDBBaseConverter().merge_parquet([a, c], out) + assert queries[-1].count("read_parquet(") == 2 + rows, _ = _read(out) + assert [r["region"] for r in rows] == ["x", "x", "y", "y"] + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_fails_for_different_versions_before_reading(tmp_folder, engine, monkeypatch): + admin2 = ADMIN.replace("/v0.1.0/", "/v0.2.0/") + columns = {"admin:country_code": ["DE", "DE"]} + a = _part(tmp_folder, "a", "a", 2, columns=columns, schemas=[CORE, ADMIN]) + b = _part(tmp_folder, "b", "b", 2, columns=columns, schemas=[CORE, ADMIN]) + b = _with_collection(b, {"schemas": {"b": [CORE, admin2]}, "collection": "b"}) + + def no_read(*args, **kwargs): + raise AssertionError("data was read") + + monkeypatch.setattr(GeoParquet, "read", no_read) + with pytest.raises(ValueError, match="different versions of a schema"): + _merge(tmp_folder, [a, b], engine, strict=False) diff --git a/tests/test_ops.py b/tests/test_ops.py index 580588c..a8c2335 100644 --- a/tests/test_ops.py +++ b/tests/test_ops.py @@ -1,7 +1,9 @@ +import re + import pytest from vecorel_cli.vecorel.collection import Collection -from vecorel_cli.vecorel.ops import merge_collections, warn_missing_required +from vecorel_cli.vecorel.ops import check_versions, merge_collections, warn_missing_required from vecorel_cli.vecorel.schemas import Schemas, VecorelSchema @@ -214,3 +216,30 @@ def warning(self, message): with pytest.raises(ValueError, match=message): merge_collections(collections, strict=True, schema_map=schema_map) + + +def test_check_versions(): + core1 = "https://vecorel.org/specification/v0.1.0/schema.yaml" + core2 = "https://vecorel.org/specification/v0.2.0/schema.yaml" + ext1 = "https://example.com/ext/v0.1.0/schema.yaml" + ext2 = "https://example.com/ext/v0.2.0/schema.yaml" + other = "https://example.com/other/v0.5.0/schema.yaml" + + # the same versions, and different extensions, are fine + check_versions( + [ + Collection({"schemas": {"a": [core1, ext1]}}), + Collection({"schemas": {"b": [core1, ext1, other]}}), + ] + ) + + # different versions fail, also across collections and datasets + with pytest.raises(ValueError, match=re.escape(f"{core1}; {core2}")): + check_versions([Collection({"schemas": {"a": [core1], "b": [core2]}})]) + with pytest.raises(ValueError, match=re.escape(f"{ext1}; {ext2}")): + check_versions( + [ + Collection({"schemas": {"a": [core1, ext1]}}), + Collection({"schemas": {"b": [core1, ext2]}}), + ] + ) diff --git a/tests/test_validate.py b/tests/test_validate.py index 6daaf19..2d3af22 100644 --- a/tests/test_validate.py +++ b/tests/test_validate.py @@ -173,3 +173,18 @@ def test_validate_checks_required_properties_per_collection(tmp_parquet_file): columns = {"admin:country_code": ["DE", "DE", None]} _write_collections(tmp_parquet_file, ["a", "a", "b"], columns, schemas) assert ValidateData().validate(tmp_parquet_file).errors == [] + + +def test_validate_scopes_missing_required_columns_to_the_collection(tmp_parquet_file): + core = "https://vecorel.org/specification/v0.1.0/schema.yaml" + admin = "https://vecorel.org/administrative-division-extension/v0.1.0/schema.yaml" + schemas = {"a": [core, admin], "b": [core]} + + # only collection a requires the property, and it has no rows + _write_collections(tmp_parquet_file, ["b", "b"], {}, schemas) + assert ValidateData().validate(tmp_parquet_file).errors == [] + + # reported once, for the collection that requires it + _write_collections(tmp_parquet_file, ["a", "b"], {}, schemas) + errors = [str(e) for e in ValidateData().validate(tmp_parquet_file).errors] + assert errors == ["admin:country_code: Required field is missing for collection 'a'"] diff --git a/vecorel_cli/conversion/duckdb.py b/vecorel_cli/conversion/duckdb.py index 9cca45b..139beee 100644 --- a/vecorel_cli/conversion/duckdb.py +++ b/vecorel_cli/conversion/duckdb.py @@ -14,7 +14,13 @@ from ..encoding.geoparquet import GeoParquet from ..parquet.types import constant_array from ..vecorel.hilbert import hilbert_keys_for_table, hilbert_reference_bounds -from ..vecorel.ops import get_collection_id, merge_collections, report, warn_missing_required +from ..vecorel.ops import ( + check_versions, + get_collection_id, + merge_collections, + report, + warn_missing_required, +) from ..vecorel.util import find_differing_crs from .base import BaseConverter @@ -667,13 +673,16 @@ def merge_parquet( if properties is not None: properties = set(properties) | {"geometry"} + # before any data is read + collections = [GeoParquet(path).get_collection() for path in paths] + check_versions(collections) + con = duckdb.connect() con.install_extension("spatial") con.load_extension("spatial") source_crs = self._common_crs(con, paths) - collections = [GeoParquet(path).get_collection() for path in paths] if collection is None: collection = merge_collections( collections, properties=properties, log=self, strict=strict @@ -686,6 +695,9 @@ def merge_parquet( # so two parts of one dataset can legitimately differ; the targets then come from # the union rather than from whichever part happens to be first selects = [] + # the selected columns of all parts, in order + union = [] + per_part = False for i, (path, part) in enumerate(zip(paths, collections)): names = pq.read_schema(path).names collection_only = set(part.get_collection_only_properties()) @@ -702,29 +714,33 @@ def merge_parquet( } keep_collection = properties is None or "collection" in properties fill = get_collection_id(part) if keep_collection else None + # the statistics tell whether the column has nulls without a scan + gaps = False + if keep_collection and "collection" in names: + with pq.ParquetFile(path) as pq_file: + gaps = bool(GeoParquet._columns_with_nulls(pq_file, {"collection"})) if fill is not None and "collection" not in names and "collection" not in collection: # Features of multiple collections must state their collection constants["collection"] = fill - elif fill is None and keep_collection: - # the statistics tell whether the column has nulls without a scan - if "collection" in names: - with pq.ParquetFile(path) as pq_file: - unknown = bool(GeoParquet._columns_with_nulls(pq_file, {"collection"})) - else: - unknown = "collection" not in collection - if unknown: - report( - f"Can't determine the collection of the features in {path}", self, strict - ) + elif ( + fill is None + and keep_collection + and (gaps if "collection" in names else "collection" not in collection) + ): + report(f"Can't determine the collection of the features in {path}", self, strict) + coalesce = fill is not None and gaps columns = [] for name in names: if name == "bbox" or (properties is not None and name not in properties): continue - if name == "collection" and fill is not None: + if name not in union: + union.append(name) + if name == "collection" and coalesce: columns.append(f'coalesce("collection", {_sql_literal(fill)}) AS "collection"') else: columns.append(_sql_name(name)) + per_part = per_part or coalesce or bool(constants) query = f"SELECT {', '.join(columns)}" source = f"read_parquet({_sql_path(path)})" if constants: @@ -738,7 +754,13 @@ def merge_parquet( query += f", {table}.*" source += f" CROSS JOIN {table}" selects.append(f"{query} FROM {source}") - source_query = " UNION ALL BY NAME ".join(selects) + if per_part: + source_query = " UNION ALL BY NAME ".join(selects) + else: + # nothing to add to any part, so a single scan over all of them + sources = "[" + ",".join(_sql_path(path) for path in paths) + "]" + columns = ", ".join(_sql_name(name) for name in union) + source_query = f"SELECT {columns} FROM read_parquet({sources}, union_by_name=true)" targets = [row[0] for row in con.execute(f"DESCRIBE {source_query}").fetchall()] return self.write_query( diff --git a/vecorel_cli/encoding/geojson.py b/vecorel_cli/encoding/geojson.py index 085e211..ea1a281 100644 --- a/vecorel_cli/encoding/geojson.py +++ b/vecorel_cli/encoding/geojson.py @@ -266,8 +266,8 @@ def _stream_json( if k == "id" and include: data["id"].append(json_stream.to_standard_types(v)) elif k == "geometry" and include: - geom = shape(json_stream.to_standard_types(v)) - data["geometry"].append(geom) + geom = json_stream.to_standard_types(v) + data["geometry"].append(shape(geom) if geom else None) elif k == "properties": for prop_key, prop_value in v.items(): # Property is not relevant diff --git a/vecorel_cli/validation/geoparquet.py b/vecorel_cli/validation/geoparquet.py index 99de5c7..6d46773 100644 --- a/vecorel_cli/validation/geoparquet.py +++ b/vecorel_cli/validation/geoparquet.py @@ -43,9 +43,12 @@ def _validate(self, num: Optional[int] = None, schema_map: SchemaMapping = {}): validator=self, ) - # Check that all required fields are present + # Check that all required fields are present; with multiple collections each + # collection only requires what its own schemas require, see validate_collection_column columns = data.columns required_props = schema.get("required", []) + if has_multiple_collections: + required_props = [key for key in required_props if key == "collection"] collection_props = schema.get("collection", {}) for key in required_props: if key not in columns and not collection_props.get(key, False): @@ -105,9 +108,12 @@ def _validate(self, num: Optional[int] = None, schema_map: SchemaMapping = {}): # Validate data of the column issues = [] - # Without a collection column the rows can't be grouped, which is reported above - if validate_data and (not has_multiple_collections or "collection" not in columns): + if validate_data and not has_multiple_collections: issues = validate_column(data[key], prop_schema) + elif validate_data and "collection" not in columns: + # the rows can't be grouped, which is reported above; the merged schema of + # all collections would report values that are valid for their own collection + pass elif validate_data: # Validate data for each collection separately for cid, cschema in schemas.items(): diff --git a/vecorel_cli/vecorel/ops.py b/vecorel_cli/vecorel/ops.py index 1ec4249..99991d4 100644 --- a/vecorel_cli/vecorel/ops.py +++ b/vecorel_cli/vecorel/ops.py @@ -1,3 +1,4 @@ +import re from typing import Optional import pandas as pd @@ -31,6 +32,8 @@ def merge( excludes: Optional[list[str]] = None, strict: bool = False, ) -> tuple[GeoDataFrame, Collection]: + # before any data is read + check_versions([item.get_collection() for item in encodings]) frames = [item.read(properties=properties, schema_map=schema_map) for item in encodings] collections = [item.get_collection() for item in encodings] if excludes: @@ -186,6 +189,26 @@ def get_collection_id(collection: Collection) -> Optional[str]: return None +def check_versions(collections: list[Collection]): + """ + Fails if the datasets use different versions of the Vecorel specification or of an + extension, in any of their collections, as they can't be upgraded automatically yet. + """ + versions = {} + for collection in collections: + for uris in collection.get_schemas().values(): + for uri in uris: + unversioned = re.sub(Schemas.version_pattern, "/", uri) + if unversioned != uri: + versions.setdefault(unversioned, set()).add(uri) + conflicts = ["; ".join(sorted(uris)) for uris in versions.values() if len(uris) > 1] + if conflicts: + raise ValueError( + "The datasets use different versions of a schema, which can't be merged: " + + " and ".join(sorted(conflicts)) + ) + + def warn_missing_required( collections: list[Collection], properties: list[str], From 2e636f7bfea30c53904604e795349d8ff0fa7037 Mon Sep 17 00:00:00 2001 From: Matthias Mohr Date: Mon, 28 Sep 2026 14:39:25 +0200 Subject: [PATCH 4/6] Code review --- tests/test_merge.py | 33 +++++++++++++++++++++++++++++ vecorel_cli/conversion/duckdb.py | 36 +++++++++++++++++++++++--------- vecorel_cli/vecorel/ops.py | 3 ++- 3 files changed, 61 insertions(+), 11 deletions(-) diff --git a/tests/test_merge.py b/tests/test_merge.py index 0f7541e..80dc849 100644 --- a/tests/test_merge.py +++ b/tests/test_merge.py @@ -710,3 +710,36 @@ def no_read(*args, **kwargs): monkeypatch.setattr(GeoParquet, "read", no_read) with pytest.raises(ValueError, match="different versions of a schema"): _merge(tmp_folder, [a, b], engine, strict=False) + + +@pytest.mark.parametrize("engine", ENGINES) +@pytest.mark.parametrize("cids", [("a", "a"), ("a", "b")]) +def test_merge_modes_required_property_that_no_part_has(tmp_folder, engine, log, cids): + # only the parts of collection a use the extension that requires the property + a = _part(tmp_folder, "a", cids[0], 2, schemas=[CORE, ADMIN]) + b = _part(tmp_folder, "b", cids[1], 2, schemas=[CORE, ADMIN] if cids[1] == "a" else None) + _check_modes( + tmp_folder, + [a, b], + engine, + log, + "Rows have no value for a required property: admin:country_code", + ) + + +def test_merge_modes_ids_that_only_differ_in_type(tmp_folder, log): + # the writer converts the ids to strings, so 1 and "1" are the same id + geojson = tmp_folder / "numeric.json" + feature = { + "type": "Feature", + "id": 1, + "geometry": {"type": "Point", "coordinates": [0, 0]}, + "properties": {}, + } + collection = {"schemas": {"same": [CORE]}, "collection": "same"} + geojson.write_text( + json.dumps({"type": "FeatureCollection", **collection, "features": [feature]}) + ) + b = _part(tmp_folder, "b", "same", 1, ids=["1"]) + + _check_modes(tmp_folder, [geojson, b], "geopandas", log, "repeat an id within their collection") diff --git a/vecorel_cli/conversion/duckdb.py b/vecorel_cli/conversion/duckdb.py index 139beee..29c9660 100644 --- a/vecorel_cli/conversion/duckdb.py +++ b/vecorel_cli/conversion/duckdb.py @@ -339,6 +339,8 @@ def write_query( suffix_duplicate_ids: bool = True, strict: bool = True, merging: bool = False, + # the selected properties of a merge, None for all + properties=None, ) -> str: """Write the rows a SELECT returns as a Vecorel GeoParquet file: drop rows that no required property or geometry survives, report an id that is not @@ -348,7 +350,8 @@ def write_query( run over. Ids only need to be unique per collection. Without `strict`, a missing required value is only reported and empty geometries are kept, as the in-memory merge does. - When `merging` strictly, empty geometries are an error instead of being dropped. + When `merging` strictly, empty geometries are an error instead of being dropped, + and so are required properties that no part has. """ compression = compression or "zstd" if compression == "zstd" and compression_level is None: @@ -377,15 +380,20 @@ def write_query( # The collections that require a property, None for all rows required_by = {} + # The same for required properties that no part has at all + absent_by = {} def null_condition(schema, skip=("geometry",), cid=None): required = [ - r - for r in schema.get("required", []) - if r not in skip and r not in collection_only and r in selected_targets + r for r in schema.get("required", []) if r not in skip and r not in collection_only ] for r in required: - required_by.setdefault(r, []).append(cid) + if r in selected_targets: + required_by.setdefault(r, []).append(cid) + elif merging and r not in collection and (properties is None or r in properties): + # not selected properties are reported by warn_missing_required + absent_by.setdefault(r, []).append(cid) + required = [r for r in required if r in selected_targets] return " OR ".join(f"{_sql_name(target)} IS NULL" for target in required) or None if per_collection: @@ -401,15 +409,22 @@ def null_condition(schema, skip=("geometry",), cid=None): null_cond = " OR ".join(conditions) else: null_cond = null_condition(collection.merge_schemas({})) + + def in_collections(cids): + return f'"collection" IN ({", ".join(map(_sql_literal, cids))})' + stats = ["count(*)"] if null_cond: - stats.append(f"count(*) FILTER (WHERE {null_cond})") # per property, so that the report can name what is missing for name, cids in required_by.items(): condition = f"{_sql_name(name)} IS NULL" if None not in cids: - condition += f' AND "collection" IN ({", ".join(map(_sql_literal, cids))})' + condition += f" AND {in_collections(cids)}" stats.append(f"count(*) FILTER (WHERE {condition})") + # a property that no part has is missing in all rows of the collections requiring it + for cids in absent_by.values(): + condition = "TRUE" if None in cids else in_collections(cids) + stats.append(f"count(*) FILTER (WHERE {condition})") # row numbers are unique by construction check_ids = "id" in selected_targets and not ids_are_generated if check_ids: @@ -435,8 +450,8 @@ def null_condition(schema, skip=("geometry",), cid=None): con.execute(f"SELECT {', '.join(stats)} FROM ({source_query})").fetchone() ) total = values.pop(0) - invalid = values.pop(0) if null_cond else 0 nulls = {name: values.pop(0) for name in required_by} if null_cond else {} + nulls.update({name: values.pop(0) for name in absent_by}) if check_ids: non_null = values.pop(0) distinct = values.pop(0) @@ -460,9 +475,9 @@ def null_condition(schema, skip=("geometry",), cid=None): for name, count in sorted(nulls.items()) if count ) - if invalid and (merging or not strict): + if missing and (merging or not strict): report(f"Rows have no value for a required property: {missing}", self, strict) - elif invalid: + elif missing: # A null in a required property is an error, whatever the count: # the writer rejects nulls in the non-nullable required fields # anyway, and silently dropping rows would make that data-quality @@ -772,6 +787,7 @@ def merge_parquet( source_crs=source_crs, strict=strict, merging=True, + properties=properties, **kwargs, ) diff --git a/vecorel_cli/vecorel/ops.py b/vecorel_cli/vecorel/ops.py index 99991d4..2dacce3 100644 --- a/vecorel_cli/vecorel/ops.py +++ b/vecorel_cli/vecorel/ops.py @@ -112,7 +112,8 @@ def merge( if "id" in merged.columns: with_id = merged[merged["id"].notna()] key = [c for c in ("collection", "id") if c in merged.columns] - duplicates = int(with_id.duplicated(subset=key).sum()) + # the writer converts the ids to strings, so 1 and "1" are the same id + duplicates = int(with_id[key].astype(str).duplicated().sum()) if duplicates: report(f"{duplicates} rows repeat an id within their collection", log, strict) From d021f93bb5500d29e0002a6d4a50c4a37cc07e96 Mon Sep 17 00:00:00 2001 From: Matthias Mohr Date: Mon, 28 Sep 2026 16:25:51 +0200 Subject: [PATCH 5/6] Further code review --- CHANGELOG.md | 13 ++++++++-- README.md | 8 ++++++ tests/test_encoding_geoparquet.py | 43 +++++++++++++++++++++++++++++++ tests/test_merge.py | 30 +++++++++++++++++++++ vecorel_cli/conversion/duckdb.py | 14 ++++++---- vecorel_cli/encoding/geojson.py | 5 +++- vecorel_cli/merge.py | 6 +++-- vecorel_cli/parquet/types.py | 9 ++++++- vecorel_cli/vecorel/ops.py | 19 ++++++++++---- 9 files changed, 131 insertions(+), 16 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 92551ed..2957aa3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -20,12 +20,19 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0. ### Changed - **Breaking:** `vec merge` keeps all properties by default, `--include` restricts them to - the core properties plus the given ones and `--exclude` removes any property. It reports - required properties that are not included, also collection-only ones. + 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 + ones. +- `vec merge` stores the collection in a column only if the datasets have multiple + collections or a dataset has a collection column, otherwise only in the collection + metadata. - **Breaking:** `vec merge` is strict by default: it fails if the merged dataset would be invalid, e.g. because an id repeats within a collection. - **Breaking:** Merging fails before reading any data if the datasets use different versions of the Vecorel specification or of an extension. +- **Breaking:** Merging doesn't drop features with an empty or missing geometry anymore + (`merge_parquet` did): they are an error in strict mode, otherwise they are kept with + a warning. - The message for missing required values names the missing properties with counts, instead of showing the SQL condition. @@ -60,6 +67,8 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0. collections, but without a collection column. - Reading GeoJSON with a limit or a selection of properties no longer fails for features without a geometry. +- Writing no longer fails for a date property with the same value for all features, the + value is stored as an ISO 8601 date in the collection. ## [v0.3.1] - 2026-09-24 diff --git a/README.md b/README.md index 45537d1..7092f01 100644 --- a/README.md +++ b/README.md @@ -133,7 +133,15 @@ All other datasets (e.g. GeoJSON or datasets that need to be reprojected) are me Use `--engine` to choose the engine explicitly. `-i` and `-e` apply to all properties, including collection-level metadata. +The geometry is required, so it can't be excluded. Constants that differ between the datasets are moved from the collection metadata to the features. +The collection of the features is stored in a column if the datasets have multiple collections +or if a dataset has a collection column already, otherwise only in the collection metadata. + +Geometries are merged as they are, except for the reprojection to the target CRS: +they are neither made valid nor converted to other geometry types (use `vec improve -g` for that), +and features with an empty or missing geometry are not dropped (see below). +The bounding boxes (`bbox`) are computed again for the merged GeoParquet file. #### Strict and non-strict mode diff --git a/tests/test_encoding_geoparquet.py b/tests/test_encoding_geoparquet.py index 0f3a1bb..567c8a6 100644 --- a/tests/test_encoding_geoparquet.py +++ b/tests/test_encoding_geoparquet.py @@ -187,6 +187,49 @@ def test_get_pyarrow_type_keeps_the_schema(): assert get_pyarrow_type(schema) == pa.map_(pa.string(), pa.string()) +def test_constant_array_parses_dates_completely(): + import datetime + + import pytest + + from vecorel_cli.parquet.types import constant_array + + schema = {"type": "date"} + for value in ["2024-01-01", "2024-01-01T00:00:00Z"]: + assert constant_array(value, schema).to_pylist() == [datetime.date(2024, 1, 1)] + for value in ["2024-01-01-invalid", "2024-01-01T12:00:00Z"]: + with pytest.raises(ValueError): + constant_array(value, schema) + + +def test_write_moves_a_constant_date_to_the_collection(tmp_parquet_file): + import datetime + + from geopandas import GeoDataFrame + from shapely.geometry import box + + gdf = GeoDataFrame( + { + "id": ["1", "2"], + "d": [datetime.date(2024, 1, 1)] * 2, + "geometry": [box(0, 0, 1, 1), box(1, 0, 2, 1)], + }, + crs="EPSG:4326", + ) + gp = GeoParquet(tmp_parquet_file) + gp.set_collection( + { + "schemas": {"c": ["https://vecorel.org/specification/v0.1.0/schema.yaml"]}, + "collection": "c", + "schemas:custom": {"properties": {"d": {"type": "date"}}}, + } + ) + gp.write(gdf) + + # dates are stored as ISO 8601 strings at the collection level + assert GeoParquet(tmp_parquet_file).get_collection()["d"] == "2024-01-01" + + def test_read_hydrates_array_and_object_constants(tmp_parquet_file): from geopandas import GeoDataFrame from shapely.geometry import box diff --git a/tests/test_merge.py b/tests/test_merge.py index 80dc849..d692833 100644 --- a/tests/test_merge.py +++ b/tests/test_merge.py @@ -743,3 +743,33 @@ def test_merge_modes_ids_that_only_differ_in_type(tmp_folder, log): b = _part(tmp_folder, "b", "same", 1, ids=["1"]) _check_modes(tmp_folder, [geojson, b], "geopandas", log, "repeat an id within their collection") + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_rejects_excluding_the_geometry(tmp_folder, engine): + a = _part(tmp_folder, "a", "a", 2) + b = _part(tmp_folder, "b", "b", 2) + with pytest.raises(ValueError, match="The geometry can't be excluded"): + _merge(tmp_folder, [a, b], engine, excludes=["geometry"], strict=False) + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_stores_a_shared_collection_in_the_metadata(tmp_folder, engine): + a = _part(tmp_folder, "a", "same", 2) + b = _part(tmp_folder, "b", "same", 2) + out = _merge(tmp_folder, [a, b], engine) + + rows, collection = _read(out) + assert collection["collection"] == "same" + assert all("collection" not in row for row in rows) + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_gives_all_features_a_collection_if_a_dataset_has_a_column(tmp_folder, engine): + a = _part(tmp_folder, "a", "same", 2, columns={"collection": ["same", "same"]}) + b = _part(tmp_folder, "b", "same", 2) + out = _merge(tmp_folder, [a, b], engine) + + rows, _ = _read(out) + assert [r["collection"] for r in rows] == ["same"] * 4 + assert _errors(out) == [] diff --git a/vecorel_cli/conversion/duckdb.py b/vecorel_cli/conversion/duckdb.py index 29c9660..be741b6 100644 --- a/vecorel_cli/conversion/duckdb.py +++ b/vecorel_cli/conversion/duckdb.py @@ -713,8 +713,13 @@ def merge_parquet( # the selected columns of all parts, in order union = [] per_part = False - for i, (path, part) in enumerate(zip(paths, collections)): - names = pq.read_schema(path).names + part_names = [pq.read_schema(path).names for path in paths] + # A column only if the merged collection can't carry the collection, or if a part + # has a column already, which all features then need, as in the in-memory merge + collection_column = "collection" not in collection or any( + "collection" in names for names in part_names + ) + for i, (path, part, names) in enumerate(zip(paths, collections, part_names)): collection_only = set(part.get_collection_only_properties()) # A part keeps its constants in its collection; one the merged collection does not # carry, because the parts disagree on it, goes back into the rows, as `vec merge` does @@ -734,13 +739,12 @@ def merge_parquet( if keep_collection and "collection" in names: with pq.ParquetFile(path) as pq_file: gaps = bool(GeoParquet._columns_with_nulls(pq_file, {"collection"})) - if fill is not None and "collection" not in names and "collection" not in collection: - # Features of multiple collections must state their collection + if fill is not None and "collection" not in names and collection_column: constants["collection"] = fill elif ( fill is None and keep_collection - and (gaps if "collection" in names else "collection" not in collection) + and (gaps if "collection" in names else collection_column) ): report(f"Can't determine the collection of the features in {path}", self, strict) coalesce = fill is not None and gaps diff --git a/vecorel_cli/encoding/geojson.py b/vecorel_cli/encoding/geojson.py index ea1a281..2efc186 100644 --- a/vecorel_cli/encoding/geojson.py +++ b/vecorel_cli/encoding/geojson.py @@ -1,3 +1,4 @@ +import datetime import json from pathlib import Path from typing import Optional, Union @@ -340,8 +341,10 @@ def _fix_omit_nulled_properties(obj): class VecorelJSONEncoder(json.JSONEncoder): def default(self, o): - if isinstance(o, pd.Timestamp): + if isinstance(o, (pd.Timestamp, datetime.datetime)): return to_iso8601(o) + elif isinstance(o, datetime.date): + return o.isoformat() elif isinstance(o, np.ndarray): return o.tolist() elif isinstance(o, set): diff --git a/vecorel_cli/merge.py b/vecorel_cli/merge.py index 09d61ff..dd916a8 100644 --- a/vecorel_cli/merge.py +++ b/vecorel_cli/merge.py @@ -26,7 +26,7 @@ class MergeDatasets(BaseCommand): This simply appends the datasets to each other. Ids that repeat within a collection are reported, but not changed. Each feature keeps its collection, which is stored in a column if the - datasets have multiple collections. + datasets have multiple collections or a dataset has a collection column. Local GeoParquet files that are all in the target CRS are merged with DuckDB, which doesn't need to fit the data into memory. All other datasets are merged in memory. @@ -65,7 +65,7 @@ def get_cli_args(): "excludes", type=click.STRING, multiple=True, - help="Properties to exclude.", + help="Properties to exclude, except for the geometry.", ), "engine": click.option( "--engine", @@ -99,6 +99,8 @@ def merge( raise ValueError("No source files provided") if engine not in self.engines: raise ValueError(f"Engine must be one of {', '.join(self.engines)}") + if excludes and "geometry" in excludes: + raise ValueError("The geometry can't be excluded") encodings = [create_encoding(s) for s in source] for encoding in encodings: if not encoding.exists(): diff --git a/vecorel_cli/parquet/types.py b/vecorel_cli/parquet/types.py index 99dbf43..c2dda47 100644 --- a/vecorel_cli/parquet/types.py +++ b/vecorel_cli/parquet/types.py @@ -325,7 +325,14 @@ def from_json_value(value, dtype=None): ts = ts.tz_localize("UTC") if ts.tzinfo is None else ts.tz_convert("UTC") return ts.to_pydatetime() if dtype == "date": - return datetime.date.fromisoformat(value[:10]) + try: + return datetime.date.fromisoformat(value) + except ValueError: + # a date-time at midnight, e.g. 2024-01-01T00:00:00Z + ts = pd.Timestamp(value) + if ts != ts.normalize(): + raise ValueError(f"'{value}' is not a date") + return ts.date() if dtype == "binary": return base64.b64decode(value, validate=True) return value diff --git a/vecorel_cli/vecorel/ops.py b/vecorel_cli/vecorel/ops.py index 2dacce3..469d53a 100644 --- a/vecorel_cli/vecorel/ops.py +++ b/vecorel_cli/vecorel/ops.py @@ -34,6 +34,8 @@ def merge( ) -> tuple[GeoDataFrame, Collection]: # before any data is read check_versions([item.get_collection() for item in encodings]) + # the geometry is required, as in merge_parquet + excludes = [key for key in excludes or [] if key != "geometry"] frames = [item.read(properties=properties, schema_map=schema_map) for item in encodings] collections = [item.get_collection() for item in encodings] if excludes: @@ -49,6 +51,11 @@ def merge( if properties is not None: warn_missing_required(collections, properties, schema_map, log, strict) props = merged_collection.merge_schemas(schema_map=schema_map).get("properties", {}) + # A column only if the merged collection can't carry the collection, or if a dataset + # has a column already, which all features then need, as in merge_parquet + collection_column = "collection" not in merged_collection or any( + "collection" in gdf.columns for gdf in frames + ) data = [] for item, gdf, collection in zip(encodings, frames, collections): @@ -83,12 +90,14 @@ def merge( cid = get_collection_id(collection) if keep_collection else None if cid is not None and "collection" in gdf.columns: gdf["collection"] = gdf["collection"].fillna(cid) - elif cid is not None: + elif cid is not None and collection_column: gdf["collection"] = cid - elif keep_collection and ( - gdf["collection"].isna().any() - if "collection" in gdf.columns - else "collection" not in merged_collection + elif ( + cid is None + and keep_collection + and ( + gdf["collection"].isna().any() if "collection" in gdf.columns else collection_column + ) ): report(f"Can't determine the collection of the features in {item.uri}", log, strict) From 8cefebd9e5c50ea5e60e299dc85ca608fcdda52f Mon Sep 17 00:00:00 2001 From: Matthias Mohr Date: Mon, 28 Sep 2026 17:11:47 +0200 Subject: [PATCH 6/6] Last code review? --- CHANGELOG.md | 2 ++ tests/test_merge.py | 24 ++++++++++++++++++++++++ vecorel_cli/vecorel/ops.py | 6 +----- 3 files changed, 27 insertions(+), 5 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 2957aa3..117796a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -26,6 +26,8 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0. - `vec merge` stores the collection in a column only if the datasets have multiple collections or a dataset has a collection column, otherwise only in the collection metadata. +- `vec merge` keeps properties without any value instead of removing them, so the + merged dataset has the same properties with both engines. - **Breaking:** `vec merge` is strict by default: it fails if the merged dataset would be invalid, e.g. because an id repeats within a collection. - **Breaking:** Merging fails before reading any data if the datasets use different diff --git a/tests/test_merge.py b/tests/test_merge.py index d692833..ec910fa 100644 --- a/tests/test_merge.py +++ b/tests/test_merge.py @@ -773,3 +773,27 @@ def test_merge_gives_all_features_a_collection_if_a_dataset_has_a_column(tmp_fol rows, _ = _read(out) assert [r["collection"] for r in rows] == ["same"] * 4 assert _errors(out) == [] + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_keeps_columns_without_values(tmp_folder, engine): + custom = {"properties": {"typed": {"type": "int32"}}} + parts = [] + for name in ("a", "b"): + part = _part( + tmp_folder, + name, + name, + 2, + columns={"typed": [1, 2], "untyped": ["x", "y"]}, + custom=custom, + ) + _with_nullable_column(part, "typed", [None, None], pa.int32()) + parts.append(_with_nullable_column(part, "untyped", [None, None])) + out = _merge(tmp_folder, parts, engine) + + schema = pq.read_schema(out) + assert str(schema.field("typed").type) == "int32" + assert str(schema.field("untyped").type) == "string" + rows, _ = _read(out) + assert all(r["typed"] is None and r["untyped"] is None for r in rows) diff --git a/vecorel_cli/vecorel/ops.py b/vecorel_cli/vecorel/ops.py index 469d53a..4764709 100644 --- a/vecorel_cli/vecorel/ops.py +++ b/vecorel_cli/vecorel/ops.py @@ -111,12 +111,8 @@ def merge( data.append(gdf) # Concatenate all GeoDataFrames to a single GeoDataFrame + # Columns without any value are kept, as in merge_parquet merged = GeoDataFrame(pd.concat(data, ignore_index=True)) - # Remove empty columns, except for the geometry, which is required - geometry = merged.geometry.name - merged = merged.drop( - columns=[c for c in merged.columns if c != geometry and merged[c].isna().all()] - ) if "id" in merged.columns: with_id = merged[merged["id"].notna()]