diff --git a/CHANGELOG.md b/CHANGELOG.md index 4583d47..61a8766 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -13,30 +13,39 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0. without loading them into memory; `--engine` chooses the engine explicitly. - `vec merge --crs first` keeps the CRS of the first dataset, the default is still EPSG:4326. - `merge_parquet` accepts `properties` to restrict the merged properties. +- `vec merge --no-strict` warns instead of failing if the merged dataset would be invalid + and writes it anyway; `merge_parquet` has a `strict` parameter (default: `True`). + See the README for the differences between the modes. ### Changed - **Breaking:** `vec merge` keeps all properties by default, `--include` restricts them to - the core properties plus the given ones and `--exclude` removes any property. It warns if - a required property is not included. + the core properties plus the given ones and `--exclude` removes any property. It reports + required properties that are not included, also collection-only ones. +- **Breaking:** `vec merge` is strict by default: it fails if the merged dataset would be + invalid, e.g. because an id repeats within a collection. +- The message for missing required values names the missing properties with counts, + instead of showing the SQL condition. ### Fixed - `merge_parquet` no longer drops a property that each part kept as a constant in - its collection but on which the parts disagree; it becomes a column again, as in - `vec merge`. + 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. + datasets, it unites them and rejects two versions of the same schema in one collection. - Merging warns about collection-only properties that differ between the datasets and have to be removed. - Merging fills in missing collection values of datasets that only list a single collection in `schemas`. -- Merging hydrates array and object constants correctly, instead of spreading them - over the rows, dropping them (`vec merge`) or failing (`merge_parquet`). +- `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` hydrates date-time, date and binary constants with their types and - NaN as null. - `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. diff --git a/README.md b/README.md index 6a13780..9688241 100644 --- a/README.md +++ b/README.md @@ -132,6 +132,33 @@ Local GeoParquet files that are all in the target CRS are merged with DuckDB, so All other datasets (e.g. GeoJSON or datasets that need to be reprojected) are merged in memory. Use `--engine` to choose the engine explicitly. +`-i` and `-e` apply to all properties, including collection-level metadata. +Constants that differ between the datasets are moved from the collection metadata to the features. + +#### Strict and non-strict mode + +By default, `vec merge` is strict: problems that make the merged dataset invalid are errors +and no dataset is written. +With `--no-strict`, `vec merge` is fail-safe: these problems are reported as warnings and +the dataset is written anyway. Check it with `vec validate` afterwards. + +| Problem | Strict (default) | `--no-strict` | +| ------- | ---------------- | ------------- | +| A required property has no value for some features | Error | Warning, the property is written as nullable | +| An id repeats within a collection | Error | Warning | +| The collection of a feature can't be determined | Error | Warning, the collection is left empty | +| `-i` or `-e` removes a required property | Error | Warning | +| A required collection-only property differs between the datasets | Error | Warning, the property is removed | +| An optional collection-only property differs between the datasets | Warning, the property is removed | Warning, the property is removed | +| A feature has an empty or missing geometry | Error | Warning, the feature is kept | +| A collection-level value doesn't fit the data type of its schema | Error | Warning, the value is left empty | +| A collection implements multiple versions of a schema, e.g. of an extension | Error | Error | +| The schemas of the datasets conflict | Error | Error | + +The strict mode only checks what a merge can break or check with little effort. +It doesn't validate the values against the schemas, e.g. patterns or value ranges, +so a merged dataset is only valid if the values of the source datasets are valid. + Check `vec merge --help` for more details. ### Create JSON Schema from Vecorel Schema diff --git a/tests/test_convert_duckdb.py b/tests/test_convert_duckdb.py index 373f984..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 @@ -723,6 +724,25 @@ def test_merge_parquet_hydrates_nan_as_null(tmp_folder): } +@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_merge.py b/tests/test_merge.py index f38db58..7f9b79d 100644 --- a/tests/test_merge.py +++ b/tests/test_merge.py @@ -1,4 +1,5 @@ import json +import re import sys from pathlib import Path @@ -11,6 +12,7 @@ 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 @@ -62,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() @@ -106,14 +130,14 @@ def _part( return path -def _with_nullable_column(path, name, values): +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.string()), pa.array(values, pa.string())) + 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 @@ -182,6 +206,20 @@ def test_merge_hydrates_array_and_object_constants(tmp_folder, engine): 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"}}} @@ -206,7 +244,7 @@ def test_merge_keeps_ids_that_repeat_in_other_collections(tmp_folder, engine, lo 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) + out = _merge(tmp_folder, [a, b, c], engine, strict=False) rows, _ = _read(out) assert [(r["collection"], r["id"]) for r in rows] == [ @@ -296,7 +334,7 @@ def test_merge_warns_when_includes_drop_required_properties(tmp_folder, engine, b = _part( tmp_folder, "b", "b", 2, columns={"admin:country_code": ["FR", "FR"]}, schemas=[CORE, ADMIN] ) - _merge(tmp_folder, [a, b], engine, includes=["foo"]) + _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() @@ -325,20 +363,63 @@ def test_merge_fills_missing_collection_values(tmp_folder, engine): 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"]) + 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") + _merge(tmp_folder, [a, b], "geopandas", strict=False) assert "repeat an id" not in log() @@ -394,3 +475,194 @@ def crs_of(path): assert "Merging in memory, as the datasets must be reprojected" in log() assert crs_of(_merge(tmp_folder, [c, d], crs="EPSG:3857")) == 3857 assert "Merging with DuckDB" in log() + + +def _with_collection(path, collection): + """Replaces the collection metadata of a file.""" + table = pq.read_table(path) + metadata = dict(table.schema.metadata) + metadata[b"collection"] = json.dumps(collection).encode("utf-8") + pq.write_table(table.replace_schema_metadata(metadata), path) + return path + + +def _check_modes(folder, parts, engine, log, message, **kwargs): + """With --no-strict the merge warns and writes the file, by default it fails, both with the message.""" + out = _merge(folder, parts, engine, strict=False, **kwargs) + assert message in log() + assert out.exists() + out.unlink() + with pytest.raises(ValueError, match=re.escape(message)): + _merge(folder, parts, engine, **kwargs) + return out + + +@pytest.mark.parametrize("engine", ENGINES) +@pytest.mark.parametrize("cids", [("a", "a"), ("a", "b")]) +def test_merge_modes_null_in_required_property(tmp_folder, engine, log, cids): + a = _part( + tmp_folder, + "a", + cids[0], + 2, + columns={"admin:country_code": ["DE", "DE"]}, + schemas=[CORE, ADMIN], + ) + a = _with_nullable_column(a, "admin:country_code", ["DE", None]) + b = _part( + tmp_folder, + "b", + cids[1], + 2, + columns={"admin:country_code": ["FR", "FR"]}, + schemas=[CORE, ADMIN], + ) + _check_modes( + tmp_folder, + [a, b], + engine, + log, + "Rows have no value for a required property: admin:country_code", + ) + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_modes_repeating_ids(tmp_folder, engine, log): + a = _part(tmp_folder, "a", "a", 2, ids=["1", "2"]) + b = _part(tmp_folder, "b", "a", 2, ids=["1", "3"]) + out = _check_modes(tmp_folder, [a, b], engine, log, "repeat an id within their collection") + assert out.name.startswith("merged") + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_modes_unknown_collection(tmp_folder, engine, log): + a = _with_nullable_column(_part(tmp_folder, "a", "a", 2), "collection", [None, None]) + a = _with_collection(a, {"schemas": {"a": [CORE], "x": [CORE]}}) + b = _part(tmp_folder, "b", "b", 2) + _check_modes(tmp_folder, [a, b], engine, log, "Can't determine the collection of the features") + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_modes_includes_remove_required_properties(tmp_folder, engine, log): + # -i also applies to collection-only properties + custom = { + "properties": {"producer": {"type": "string"}}, + "collection": {"producer": True}, + "required": ["producer"], + } + a = _part(tmp_folder, "a", "a", 2, collection={"producer": "X"}, custom=custom) + b = _part(tmp_folder, "b", "b", 2, collection={"producer": "X"}, custom=custom) + out = _check_modes( + tmp_folder, + [a, b], + engine, + log, + "Required properties are not included, the merged file will be invalid: producer", + includes=["foo"], + ) + assert out.name.startswith("merged") + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_modes_required_collection_only_property_differs(tmp_folder, engine, log): + custom = { + "properties": {"producer": {"type": "string"}}, + "collection": {"producer": True}, + "required": ["producer"], + } + a = _part(tmp_folder, "a", "a", 2, collection={"producer": "X"}, custom=custom) + b = _part(tmp_folder, "b", "b", 2, collection={"producer": "Y"}, custom=custom) + _check_modes( + tmp_folder, + [a, b], + engine, + log, + "Collection-only properties differ between the datasets and are removed: producer", + ) + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_modes_empty_geometries(tmp_folder, engine, log): + gdf = gpd.GeoDataFrame( + {"id": ["e0", "e1"], "geometry": [shapely.box(0, 0, 1, 1), shapely.Polygon()]}, + crs="EPSG:4326", + ) + a = tmp_folder / "empty.parquet" + gp = GeoParquet(a) + gp.set_collection({"schemas": {"a": [CORE]}, "collection": "a"}) + gp.write(gdf, dehydrate=False) + b = _part(tmp_folder, "b", "b", 2) + + _check_modes(tmp_folder, [a, b], engine, log, "1 of 4 rows have an empty or missing geometry") + # --no-strict keeps the rows + rows, _ = _read(_merge(tmp_folder, [a, b], engine, strict=False)) + assert len(rows) == 4 + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_modes_constant_that_does_not_fit_its_schema(tmp_folder, engine, log): + custom = {"properties": {"n": {"type": "uint8"}}} + a = _part(tmp_folder, "a", "a", 2, collection={"n": 300}, custom=custom) + b = _part(tmp_folder, "b", "b", 2, collection={"n": 1}, custom=custom) + out = _check_modes(tmp_folder, [a, b], engine, log, "n: Value 300 doesn't fit data type uint8") + + # --no-strict writes a null instead + rows, _ = _read(_merge(tmp_folder, [a, b], engine, strict=False)) + assert [(r["collection"], r["n"]) for r in rows] == [ + ("a", None), + ("a", None), + ("b", 1), + ("b", 1), + ] + assert out.name.startswith("merged") + + +@pytest.mark.parametrize("engine", ENGINES) +def test_merge_hydrates_binary_constants(tmp_folder, engine): + custom = {"properties": {"blob": {"type": "binary"}}} + a = _part(tmp_folder, "a", "a", 2, collection={"blob": "aGVsbG8="}, custom=custom) + b = _part(tmp_folder, "b", "b", 2, collection={"blob": "d29ybGQ="}, custom=custom) + + rows, _ = _read(_merge(tmp_folder, [a, b], engine)) + assert [(r["collection"], r["blob"]) for r in rows] == [ + ("a", b"hello"), + ("a", b"hello"), + ("b", b"world"), + ("b", b"world"), + ] + + +def test_merge_modes_all_geometries_missing(tmp_folder, log): + paths = [] + for cid in ("a", "b"): + path = tmp_folder / f"{cid}.json" + features = [ + {"type": "Feature", "id": f"{cid}{i}", "geometry": None, "properties": {}} + for i in range(2) + ] + collection = {"schemas": {cid: [CORE]}, "collection": cid} + path.write_text( + json.dumps({"type": "FeatureCollection", **collection, "features": features}) + ) + paths.append(path) + + _check_modes( + tmp_folder, paths, "geopandas", log, "4 of 4 rows have an empty or missing geometry" + ) + + +def test_merge_writes_required_nested_properties_with_nulls_as_nullable(tmp_folder, log): + struct = pa.struct([("k", pa.string())]) + custom = { + "required": ["attrs"], + "properties": {"attrs": {"type": "object", "properties": {"k": {"type": "string"}}}}, + } + a = _part(tmp_folder, "a", "a", 2, columns={"attrs": [{"k": "v"}, {"k": "w"}]}, custom=custom) + a = _with_nullable_column(a, "attrs", [{"k": "v"}, None], struct) + b = _part(tmp_folder, "b", "a", 2, columns={"attrs": [{"k": "x"}, {"k": "y"}]}, custom=custom) + + out = _merge(tmp_folder, [a, b], "duckdb", strict=False) + assert "Rows have no value for a required property: attrs" in log() + assert pq.read_schema(out).field("attrs").nullable + rows, _ = _read(out) + assert [r["attrs"] for r in rows] == [{"k": "v"}, None, {"k": "x"}, {"k": "y"}] diff --git a/tests/test_ops.py b/tests/test_ops.py index 2e6912a..580588c 100644 --- a/tests/test_ops.py +++ b/tests/test_ops.py @@ -1,7 +1,7 @@ import pytest from vecorel_cli.vecorel.collection import Collection -from vecorel_cli.vecorel.ops import merge_collections +from vecorel_cli.vecorel.ops import merge_collections, warn_missing_required from vecorel_cli.vecorel.schemas import Schemas, VecorelSchema @@ -110,6 +110,22 @@ def test_merge_collections_unites_schemas_of_a_collection(): [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: @@ -131,3 +147,70 @@ def warning(self, message): assert log.warnings == [ "Collection-only properties differ between the datasets and are removed: producer" ] + + +def test_warn_missing_required_collection_only_properties(tmp_path): + class Log: + warnings = [] + + def warning(self, message): + self.warnings.append(message) + + core = "https://vecorel.org/specification/v0.1.0/schema.yaml" + ext = "https://example.com/producer/v0.1.0/schema.yaml" + path = tmp_path / "schema.yaml" + path.write_text( + "$schema: https://vecorel.org/sdl/v0.2.0/schema.json\n" + "required: [producer]\n" + "collection: {producer: true}\n" + "properties: {producer: {type: string}}\n" + ) + schema_map = {ext: path} + properties = ["id", "geometry", "collection"] + collections = [ + Collection({"schemas": {"c1": [core, ext]}, "producer": "A"}), + Collection({"schemas": {"c2": [core, ext]}, "producer": "A"}), + ] + merged = merge_collections(collections, properties=properties) + assert "producer" not in merged + + log = Log() + warn_missing_required(collections, properties, schema_map, log) + message = "Required properties are not included, the merged file will be invalid: producer" + assert log.warnings == [message] + + with pytest.raises(ValueError, match=message): + warn_missing_required(collections, properties, schema_map, strict=True) + + +def test_merge_collections_uses_the_schema_map(tmp_path): + class Log: + warnings = [] + + def warning(self, message): + self.warnings.append(message) + + # the extension is only available through the schema map + core = "https://vecorel.org/specification/v0.1.0/schema.yaml" + ext = "https://example.com/producer/v0.1.0/schema.yaml" + path = tmp_path / "schema.yaml" + path.write_text( + "$schema: https://vecorel.org/sdl/v0.2.0/schema.json\n" + "required: [producer]\n" + "collection: {producer: true}\n" + "properties: {producer: {type: string}}\n" + ) + schema_map = {ext: path} + collections = [ + Collection({"schemas": {"c1": [core, ext]}, "producer": "A"}), + Collection({"schemas": {"c2": [core, ext]}, "producer": "B"}), + ] + message = "Collection-only properties differ between the datasets and are removed: producer" + + log = Log() + merged = merge_collections(collections, log=log, schema_map=schema_map) + assert "producer" not in merged + assert log.warnings == [message] + + with pytest.raises(ValueError, match=message): + merge_collections(collections, strict=True, schema_map=schema_map) diff --git a/vecorel_cli/conversion/duckdb.py b/vecorel_cli/conversion/duckdb.py index 10f8ca4..9cca45b 100644 --- a/vecorel_cli/conversion/duckdb.py +++ b/vecorel_cli/conversion/duckdb.py @@ -1,5 +1,3 @@ -import base64 -import datetime import json import os from pathlib import Path @@ -8,15 +6,16 @@ import duckdb import numpy as np -import pandas as pd 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 get_pyarrow_type +from ..parquet.types import constant_array from ..vecorel.hilbert import hilbert_keys_for_table, hilbert_reference_bounds -from ..vecorel.ops import get_collection_id, merge_collections, warn_missing_required +from ..vecorel.ops import get_collection_id, merge_collections, report, warn_missing_required +from ..vecorel.util import find_differing_crs from .base import BaseConverter @@ -55,49 +54,25 @@ def _sql_literal(value) -> str: return f"'{escaped}'" -# Collection metadata is JSON, so temporal and binary values are encoded as strings -def _to_arrow_value(value, dtype): - 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) - return ( - ts.tz_localize("UTC") if ts.tzinfo is None else ts.tz_convert("UTC") - ).to_pydatetime() - if dtype == "date": - return datetime.date.fromisoformat(value[:10]) - if dtype == "binary": - return base64.b64decode(value) - return value - - -def _constants_table(constants: dict, props: dict) -> pa.Table: +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(): - schema = props.get(key) or {} try: - pa_type = get_pyarrow_type(schema) if schema.get("type") else None - arrays.append(pa.array([_to_arrow_value(value, schema.get("type"))], type=pa_type)) - except Exception: - arrays.append(pa.array([value])) + 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())) -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) - - # This converter is experimental, use with caution. # Results may not be fully compliant yet. # Use this primarily for datasets that are too large to be processed by the default converter. @@ -321,9 +296,8 @@ def _common_crs(self, con, sources: list): They are combined without reprojection, and DuckDB before 1.5 drops the CRS from the metadata, so it has to be read from the files themselves. """ - source_crs = None - reference = None - for i, source in enumerate(sources): + crs_values = [] + for source in sources: crs = None row = con.execute( "SELECT value FROM parquet_kv_metadata(?) WHERE key = 'geo'", [source] @@ -332,15 +306,14 @@ def _common_crs(self, con, sources: list): geo = json.loads(bytes(row[0])) primary = geo.get("primary_column", "") crs = geo.get("columns", {}).get(primary, {}).get("crs") - if i == 0: - source_crs = crs - reference = _normalize_crs(crs) - elif not _equal_crs(_normalize_crs(crs), reference): - raise ValueError( - f"The sources use different coordinate reference systems: {source} " - "differs from the first source. Reproject the sources to a common CRS." - ) - return source_crs + crs_values.append(crs) + differing = find_differing_crs(crs_values) + if differing is not None: + raise ValueError( + f"The sources use different coordinate reference systems: {sources[differing]} " + "differs from the first source. Reproject the sources to a common CRS." + ) + return crs_values[0] if crs_values else None def write_query( self, @@ -358,6 +331,8 @@ def write_query( geoparquet_version: Optional[str] = None, original_geometries: bool = False, suffix_duplicate_ids: bool = True, + strict: bool = True, + merging: bool = False, ) -> str: """Write the rows a SELECT returns as a Vecorel GeoParquet file: drop rows that no required property or geometry survives, report an id that is not @@ -365,6 +340,9 @@ def write_query( `targets` names the properties the query returns, which is what the checks run over. Ids only need to be unique per collection. + Without `strict`, a missing required value is only reported and empty + geometries are kept, as the in-memory merge does. + When `merging` strictly, empty geometries are an error instead of being dropped. """ compression = compression or "zstd" if compression == "zstd" and compression_level is None: @@ -388,23 +366,30 @@ def write_query( # Same required-value check, empty-geometry drop and id uniqueness # check as in the GeoDataFrame-based codepath, in one scan collection_only = set(collection.get_collection_only_properties()) - has_collection_column = "collection" in selected_targets + schema_groups = collection.get_schemas() + per_collection = "collection" in selected_targets and len(schema_groups) > 1 + + # The collections that require a property, None for all rows + required_by = {} - def null_condition(schema): + def null_condition(schema, skip=("geometry",), cid=None): required = [ r for r in schema.get("required", []) - if r != "geometry" and r not in collection_only and r in selected_targets + if r not in skip and r not in collection_only and r in selected_targets ] - return " OR ".join(f'"{target}" IS NULL' for target in required) or None + for r in required: + required_by.setdefault(r, []).append(cid) + return " OR ".join(f"{_sql_name(target)} IS NULL" for target in required) or None - schema_groups = collection.get_schemas() - if len(schema_groups) > 1 and has_collection_column: + 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(): - cond = null_condition(group.merge_schemas(custom_schemas=custom_schemas)) + 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) @@ -413,12 +398,16 @@ def null_condition(schema): stats = ["count(*)"] if null_cond: stats.append(f"count(*) FILTER (WHERE {null_cond})") + # per property, so that the report can name what is missing + for name, cids in required_by.items(): + condition = f"{_sql_name(name)} IS NULL" + if None not in cids: + condition += f' AND "collection" IN ({", ".join(map(_sql_literal, cids))})' + 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 has_collection_column else '"id"' - ) + id_key = 'struct_pack(c := "collection", i := "id")' if per_collection else '"id"' stats.append('count("id")') stats.append(f'count(DISTINCT {id_key}) FILTER (WHERE "id" IS NOT NULL)') blank_cond = None @@ -441,6 +430,7 @@ def null_condition(schema): ) total = values.pop(0) invalid = values.pop(0) if null_cond else 0 + nulls = {name: values.pop(0) for name in required_by} if null_cond else {} if check_ids: non_null = values.pop(0) distinct = values.pop(0) @@ -452,23 +442,35 @@ def null_condition(schema): "The repeating ids get a ~ suffix in the output." ) elif distinct < non_null: - self.warning( - f"{non_null - distinct:,} of {non_null:,} rows repeat an id within their collection" + report( + f"{non_null - distinct:,} of {non_null:,} rows repeat an id within their collection", + self, + strict, ) blanks = values.pop(0) if blank_cond else 0 repairs = values.pop(0) if repair_cond else 0 - if invalid: + missing = ", ".join( + f"{name} ({count:,} of {total:,} rows)" + for name, count in sorted(nulls.items()) + if count + ) + if invalid and (merging or not strict): + report(f"Rows have no value for a required property: {missing}", self, strict) + elif invalid: # A null in a required property is an error, whatever the count: # the writer rejects nulls in the non-nullable required fields # anyway, and silently dropping rows would make that data-quality # decision for the user (vecorel/cli#33) raise ValueError( - f"{invalid} of {total} rows have no value for a required property " - f"({null_cond}). Handle them in the converter: fix the mapping, " - "fill the values in column_migrations, or exclude the rows with " - "a column_filters entry." + f"Rows have no value for a required property: {missing}. " + "Handle them in the converter: fix the mapping, fill the values in " + "column_migrations, or exclude the rows with a column_filters entry." ) - if blanks: + 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: @@ -531,7 +533,7 @@ def null_condition(schema): 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 @@ -569,12 +571,13 @@ def null_condition(schema): 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 of a collection (id~1, id~2, ...), like the GeoDataFrame-based codepath, when the source does not provide unique ids @@ -583,8 +586,7 @@ def _suffix_duplicate_ids_in_file( 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.""" - names = pq.read_schema(output_file).names - key = '"collection", "id"' if "collection" in names else '"id"' + key = '"collection", "id"' if per_collection else '"id"' reported = False while True: total, duplicated = con.execute( @@ -640,7 +642,13 @@ def _suffix_duplicate_ids_in_file( raise def merge_parquet( - self, paths: list, output_file, collection=None, properties=None, **kwargs + 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. @@ -649,6 +657,8 @@ def merge_parquet( their geometries are already valid polygons; pass `original_geometries=False` 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") @@ -665,9 +675,11 @@ def merge_parquet( collections = [GeoParquet(path).get_collection() for path in paths] if collection is None: - collection = merge_collections(collections, properties=properties, log=self) + collection = merge_collections( + collections, properties=properties, log=self, strict=strict + ) if properties is not None: - warn_missing_required(collection, properties, {}, self) + 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, @@ -689,16 +701,21 @@ def merge_parquet( and (properties is None or key in properties) } keep_collection = properties is None or "collection" in properties - fill = None - if keep_collection and "collection" in names: - # Fill gaps like the in-memory merge does, if the part's collection is known - try: - fill = get_collection_id(part, path) - except ValueError: - pass - elif keep_collection and "collection" not in collection: + fill = get_collection_id(part) if keep_collection else None + if fill is not None and "collection" not in names and "collection" not in collection: # Features of multiple collections must state their collection - constants["collection"] = get_collection_id(part, path) + constants["collection"] = fill + elif fill is None and keep_collection: + # the statistics tell whether the column has nulls without a scan + if "collection" in names: + with pq.ParquetFile(path) as pq_file: + unknown = bool(GeoParquet._columns_with_nulls(pq_file, {"collection"})) + else: + unknown = "collection" not in collection + if unknown: + report( + f"Can't determine the collection of the features in {path}", self, strict + ) columns = [] for name in names: @@ -714,7 +731,10 @@ def merge_parquet( # 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)) + 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}") @@ -728,6 +748,8 @@ def merge_parquet( collection, targets=targets, source_crs=source_crs, + strict=strict, + merging=True, **kwargs, ) diff --git a/vecorel_cli/create_geoparquet.py b/vecorel_cli/create_geoparquet.py index 7aefd76..fedf455 100644 --- a/vecorel_cli/create_geoparquet.py +++ b/vecorel_cli/create_geoparquet.py @@ -59,7 +59,9 @@ def create( # Read source data encodings = [create_encoding(s) for s in source] # Merge encodings into a single GeoDataFrame - geodata, collection = merge(encodings, properties=properties, schema_map=schema_map) + geodata, collection = merge( + encodings, properties=properties, schema_map=schema_map, log=self + ) # Write to target target_encoding = GeoParquet(target) diff --git a/vecorel_cli/encoding/base.py b/vecorel_cli/encoding/base.py index fd46909..84ee33a 100644 --- a/vecorel_cli/encoding/base.py +++ b/vecorel_cli/encoding/base.py @@ -112,15 +112,23 @@ 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 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 c022d08..09d61ff 100644 --- a/vecorel_cli/merge.py +++ b/vecorel_cli/merge.py @@ -14,6 +14,7 @@ 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): @@ -29,6 +30,9 @@ class MergeDatasets(BaseCommand): 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" @@ -70,6 +74,12 @@ def get_cli_args(): 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.", + ), } @runnable @@ -81,6 +91,7 @@ def merge( includes=[], excludes=[], engine="auto", + strict=True, ): if not isinstance(source, list): raise ValueError("Source must be a list.") @@ -103,10 +114,6 @@ def merge( properties = None if includes: properties = list(set(Registry.core_properties) | set(includes)) - if excludes: - if properties is None: - properties = self.get_available_properties(encodings) - properties = list(set(properties) - set(excludes)) blocker = self.get_duckdb_blocker(encodings, target, crs) if engine == "duckdb" and blocker: @@ -115,19 +122,32 @@ def merge( 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) + gdf, collection = merge_( + encodings, + crs=crs, + properties=properties, + log=self, + excludes=excludes, + strict=strict, + ) target.set_collection(collection) - target.write(gdf, properties=properties) + target.write(gdf, properties=properties, strict=strict) return target @@ -150,20 +170,16 @@ def get_duckdb_blocker( 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. """ - from .conversion.duckdb import _equal_crs, _normalize_crs - for encoding in [*encodings, target]: if not isinstance(encoding, GeoParquet) or not isinstance(encoding.uri, Path): return "DuckDB only merges local GeoParquet files" - reference = _normalize_crs(crs) if crs else None + crs_values = [] for encoding in encodings: geo = encoding.get_geoparquet_metadata() or {} column = geo.get("columns", {}).get(geo.get("primary_column"), {}) - source_crs = _normalize_crs(column.get("crs")) - if reference is None: - reference = source_crs - elif not _equal_crs(source_crs, reference): - return "the datasets must be reprojected to a common CRS" + 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 e43e9e8..99dbf43 100644 --- a/vecorel_cli/parquet/types.py +++ b/vecorel_cli/parquet/types.py @@ -1,4 +1,7 @@ +import base64 +import binascii import datetime +from typing import Optional import numpy as np import pandas as pd @@ -307,3 +310,40 @@ def get_pyarrow_type_for_geopandas(dtype): } SUPPORTED_PROTOCOLS = ["http", "https", "s3", "gs"] + + +def from_json_value(value, dtype=None): + """ + A collection-level value as a Python value for the data type. + Collection metadata is JSON, so temporal and binary values are encoded as strings. + """ + if value is None or (isinstance(value, float) and np.isnan(value)): + return None + if isinstance(value, str): + if dtype == "date-time": + ts = pd.Timestamp(value) + ts = ts.tz_localize("UTC") if ts.tzinfo is None else ts.tz_convert("UTC") + return ts.to_pydatetime() + if dtype == "date": + return datetime.date.fromisoformat(value[:10]) + if dtype == "binary": + return base64.b64decode(value, validate=True) + return value + + +def constant_array(value, schema: Optional[dict] = None, length: int = 1) -> pa.Array: + """ + An array that repeats a collection-level value, typed by its schema if given. + Raises a ValueError if the value doesn't fit the schema. + """ + dtype = (schema or {}).get("type") + try: + pa_type = get_pyarrow_type(schema) if dtype else None + except Exception: + pa_type = None + if pa_type is None: + return pa.array([from_json_value(value)] * length) + try: + return pa.array([from_json_value(value, dtype)] * length, type=pa_type) + except (pa.ArrowException, ValueError, TypeError, OverflowError, binascii.Error) as e: + raise ValueError(f"Value {value!r} doesn't fit data type {dtype}: {e}") from e diff --git a/vecorel_cli/vecorel/ops.py b/vecorel_cli/vecorel/ops.py index 8e10f38..1ec4249 100644 --- a/vecorel_cli/vecorel/ops.py +++ b/vecorel_cli/vecorel/ops.py @@ -5,32 +5,89 @@ 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]: - data = [] - collections = [] + 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", {}) - for item in encodings: - # Load the dataset - gdf = item.read(hydrate=True, properties=properties, schema_map=schema_map) - collection = item.get_collection() + data = [] + 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 - if "collection" not in gdf.columns or gdf["collection"].isna().any(): - cid = get_collection_id(collection, item.uri) - if "collection" in gdf.columns: - gdf["collection"] = gdf["collection"].fillna(cid) - else: - gdf["collection"] = cid + keep_collection = properties is None or "collection" in properties + cid = get_collection_id(collection) if keep_collection else None + if cid is not None and "collection" in gdf.columns: + gdf["collection"] = gdf["collection"].fillna(cid) + elif cid is not None: + gdf["collection"] = cid + elif keep_collection and ( + gdf["collection"].isna().any() + if "collection" in gdf.columns + else "collection" not in merged_collection + ): + report(f"Can't determine the collection of the features in {item.uri}", log, strict) if not crs: # If no CRS is given, use the first CRS that is available as the base CRS @@ -39,33 +96,86 @@ def merge( # Change the CRS if necessary gdf.to_crs(crs=crs, inplace=True) - # Add data to lists data.append(gdf) - collections.append(collection) # Concatenate all GeoDataFrames to a single GeoDataFrame merged = GeoDataFrame(pd.concat(data, ignore_index=True)) - # Remove empty columns - merged.dropna(axis=1, how="all", inplace=True) + # Remove empty columns, except for the geometry, which is required + geometry = merged.geometry.name + merged = merged.drop( + columns=[c for c in merged.columns if c != geometry and merged[c].isna().all()] + ) - if log and "id" in merged.columns: + if "id" in merged.columns: with_id = merged[merged["id"].notna()] - duplicates = int(with_id.duplicated(subset=["collection", "id"]).sum()) + key = [c for c in ("collection", "id") if c in merged.columns] + duplicates = int(with_id.duplicated(subset=key).sum()) if duplicates: - log.warning(f"{duplicates} rows repeat an id within their collection") + report(f"{duplicates} rows repeat an id within their collection", log, strict) - # Merge all collections - collection = merge_collections(collections, properties=properties, log=log) - if log and properties is not None: - warn_missing_required(collection, properties, schema_map, log) + check_merged_data(merged, merged_collection, schema_map, log, strict, properties) + + return merged, merged_collection + + +def check_merged_data( + merged: GeoDataFrame, + collection: Collection, + schema_map: SchemaMapping = {}, + log: Optional[LoggerMixin] = None, + strict: bool = False, + properties=None, +): + """ + Reports empty geometries and missing required values, per collection. + Properties that are not selected are reported by warn_missing_required. + """ + geometries = merged.geometry + empty = int((geometries.isna() | geometries.is_empty).sum()) + if empty: + report(f"{empty} of {len(merged)} rows have an empty or missing geometry", log, strict) - return merged, collection + 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 get_collection_id(collection: Collection, source=None) -> str: +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: @@ -73,23 +183,40 @@ def get_collection_id(collection: Collection, source=None) -> str: schemas = collection.get_schemas() if len(schemas) == 1: return next(iter(schemas.keys())) - raise ValueError(f"Can't determine the collection of the features in {source}") + return None def warn_missing_required( - collection: Collection, properties: list[str], schema_map: SchemaMapping, log: LoggerMixin + collections: list[Collection], + properties: list[str], + schema_map: SchemaMapping, + log: Optional[LoggerMixin] = None, + strict: bool = False, ): - schema = collection.merge_schemas(schema_map=schema_map) - collection_only = set(collection.get_collection_only_properties(schema_map=schema_map)) - missing = set(schema.get("required", [])) - set(properties) - collection_only - {"geometry"} + """ + 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: - log.warning( - f"Required properties are not included, the merged file will be invalid: {', '.join(sorted(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 + collections: list[Collection], + properties=None, + log: Optional[LoggerMixin] = None, + strict: bool = False, + schema_map: SchemaMapping = {}, ) -> Collection: schemas = Schemas() custom_schemas = VecorelSchema() @@ -117,9 +244,10 @@ def merge_collections( other_props = {k: v for k, v in other_props.items() if k in properties} custom_schemas = custom_schemas.pick(properties) - if log: + 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 @@ -129,12 +257,13 @@ def merge_collections( and (properties is None or k in properties) } if keys: - dropped |= keys & set(c.get_collection_only_properties()) - if dropped: - log.warning( - "Collection-only properties differ between the datasets and are removed: " - + ", ".join(sorted(dropped)) - ) + 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 00f5862..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: @@ -341,6 +342,14 @@ def add_all(self, schemas: Union[RawSchemas, "Schemas"]): 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: 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