diff --git a/CHANGELOG.md b/CHANGELOG.md index 776a999..117796a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,25 +7,97 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0. ## [Unreleased] +### 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 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. +- `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 + 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. + +### 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. +- 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 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. +- 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 -- 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 @@ -35,17 +107,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 @@ -53,30 +134,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. @@ -84,7 +179,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 @@ -93,6 +187,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, @@ -102,7 +198,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`. @@ -111,100 +212,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..7092f01 100644 --- a/README.md +++ b/README.md @@ -121,7 +121,54 @@ 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. +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 + +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 | +| 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. Check `vec merge --help` for more details. diff --git a/tests/test_convert_duckdb.py b/tests/test_convert_duckdb.py index ba182cf..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 @@ -621,6 +622,127 @@ 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 _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..567c8a6 100644 --- a/tests/test_encoding_geoparquet.py +++ b/tests/test_encoding_geoparquet.py @@ -148,3 +148,107 @@ 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_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 + + 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..ec910fa 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,707 @@ 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"}] + + +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) + + +@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") + + +@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) == [] + + +@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/tests/test_ops.py b/tests/test_ops.py index 44aaaed..a8c2335 100644 --- a/tests/test_ops.py +++ b/tests/test_ops.py @@ -1,5 +1,9 @@ +import re + +import pytest + from vecorel_cli.vecorel.collection import Collection -from vecorel_cli.vecorel.ops import merge_collections +from vecorel_cli.vecorel.ops import check_versions, merge_collections, warn_missing_required from vecorel_cli.vecorel.schemas import Schemas, VecorelSchema @@ -92,3 +96,150 @@ 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) + + +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 e0c65d1..2d3af22 100644 --- a/tests/test_validate.py +++ b/tests/test_validate.py @@ -108,3 +108,83 @@ 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 == [] + + +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 80d63ee..be741b6 100644 --- a/vecorel_cli/conversion/duckdb.py +++ b/vecorel_cli/conversion/duckdb.py @@ -9,9 +9,19 @@ 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 ( + check_versions, + get_collection_id, + merge_collections, + report, + warn_missing_required, +) +from ..vecorel.util import find_differing_crs from .base import BaseConverter @@ -22,6 +32,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 +60,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 +302,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 +312,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 +336,22 @@ 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, + # 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 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, + and so are required properties that no part has. """ compression = compression or "zstd" if compression == "zstd" and compression_level is None: @@ -346,23 +374,63 @@ 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 = {} + # 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 + ] + for r in required: + 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: + # 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({})) + + def in_collections(cids): + return f'"collection" IN ({", ".join(map(_sql_literal, cids))})' + stats = ["count(*)"] - null_cond = None - if required: - null_cond = " OR ".join(f'"{target}" IS NULL' for target in required) - stats.append(f"count(*) FILTER (WHERE {null_cond})") + if 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 {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: + 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: @@ -382,31 +450,48 @@ def write_query( 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) - 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 missing and (merging or not strict): + report(f"Rows have no value for a required property: {missing}", self, strict) + 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 # 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: + if blanks and merging: + report( + f"{blanks:,} of {total:,} rows have an empty or missing geometry", self, strict + ) + 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 +552,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 +592,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 +615,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 +638,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 +662,35 @@ 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"} + + # before any data is read + collections = [GeoParquet(path).get_collection() for path in paths] + check_versions(collections) con = duckdb.connect() con.install_extension("spatial") @@ -594,18 +698,90 @@ 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, + if collection is None: + 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 - sources = "[" + ",".join(_sql_path(path) for path in paths) + "]" - source_query = f"SELECT * FROM read_parquet({sources}, union_by_name=true)" + selects = [] + # the selected columns of all parts, in order + union = [] + per_part = False + 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 + 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 + # 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_column: + constants["collection"] = fill + elif ( + fill is None + and keep_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 + + columns = [] + for name in names: + if name == "bbox" or (properties is not None and name not in properties): + continue + 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: + # 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), + ) + query += f", {table}.*" + source += f" CROSS JOIN {table}" + selects.append(f"{query} FROM {source}") + 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()] - if collection is None: - from ..vecorel.ops import merge_collections - - collection = merge_collections([GeoParquet(path).get_collection() for path in paths]) - return self.write_query( con, source_query, @@ -613,6 +789,9 @@ def merge_parquet(self, paths: list, output_file, collection=None, **kwargs) -> collection, targets=targets, source_crs=source_crs, + strict=strict, + merging=True, + properties=properties, **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/geojson.py b/vecorel_cli/encoding/geojson.py index 085e211..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 @@ -266,8 +267,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 @@ -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/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..dd916a8 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 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. + + 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, except for the geometry.", + ), + "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,98 @@ 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)}") + 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(): + 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..c2dda47 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,47 @@ 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": + 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 + + +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..6d46773 100644 --- a/vecorel_cli/validation/geoparquet.py +++ b/vecorel_cli/validation/geoparquet.py @@ -43,14 +43,20 @@ 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): 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", {}) @@ -104,7 +110,11 @@ def _validate(self, num: Optional[int] = None, schema_map: SchemaMapping = {}): issues = [] if validate_data and not has_multiple_collections: issues = validate_column(data[key], prop_schema) - elif validate_data and has_multiple_collections: + 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(): vecorel_schema = cschema.merge_schemas( @@ -132,6 +142,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..4764709 100644 --- a/vecorel_cli/vecorel/ops.py +++ b/vecorel_cli/vecorel/ops.py @@ -1,24 +1,105 @@ +import re +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]: + # 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: + 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", {}) + # 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 = [] - 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 and collection_column: + gdf["collection"] = cid + 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) if not crs: # If no CRS is given, use the first CRS that is available as the base CRS @@ -27,22 +108,145 @@ 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 + # Columns without any value are kept, as in merge_parquet merged = GeoDataFrame(pd.concat(data, ignore_index=True)) - # Remove empty columns - merged.dropna(axis=1, how="all", inplace=True) - # Merge all collections - collection = merge_collections(collections, properties=properties) + if "id" in merged.columns: + with_id = merged[merged["id"].notna()] + key = [c for c in ("collection", "id") if c in merged.columns] + # 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) + + check_merged_data(merged, merged_collection, schema_map, log, strict, properties) - return merged, collection + 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, + ) -def merge_collections(collections: list[Collection], properties=None) -> Collection: +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 + + +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], + 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, + log: Optional[LoggerMixin] = None, + strict: bool = False, + schema_map: SchemaMapping = {}, +) -> Collection: schemas = Schemas() custom_schemas = VecorelSchema() handled_separately = ("schemas", "schemas:custom") @@ -69,6 +273,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