diff --git a/CHANGELOG.md b/CHANGELOG.md index cd0e7db..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. @@ -46,15 +55,6 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0. - 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. -- The in-memory merge keeps constants that all datasets share in the collection, instead - of failing on or moving array and object constants into the rows. -- `vec merge` with DuckDB keeps rows without a required value (with a warning) and - rows with an empty geometry, like the in-memory merge. -- `merge_parquet` warns about a constant that doesn't fit the type of its schema. -- Merging rejects two versions of the same schema in one collection. -- `vec merge --exclude` no longer reads datasets other than GeoParquet twice. -- Merging warns when `--include` drops a required collection-only property. -- Merging accepts features again whose collection can't be determined. ## [v0.3.1] - 2026-09-24 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 934264d..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 @@ -730,10 +731,16 @@ 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) + table = _constants_table({"x": value}, {"x": {"type": dtype}}, log=log, source="part.parquet") - assert table.column("x").to_pylist() == [value] - assert f"'x' doesn't fit its type {dtype}" in capsys.readouterr().out + # 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): diff --git a/tests/test_merge.py b/tests/test_merge.py index 13a6603..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 @@ -129,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 @@ -243,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] == [ @@ -333,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() @@ -370,7 +371,7 @@ def test_merge_accepts_a_collection_it_cannot_determine(tmp_folder, engine, valu if values: a = _with_nullable_column(a, "collection", values) b = _part(tmp_folder, "b", "b", 2) - out = _merge(tmp_folder, [a, b], engine) + out = _merge(tmp_folder, [a, b], engine, strict=False) rows, _ = _read(out) assert [r["collection"] for r in rows][-2:] == ["b", "b"] @@ -381,7 +382,7 @@ def test_merge_accepts_a_collection_it_cannot_determine(tmp_folder, engine, valu 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() @@ -396,7 +397,7 @@ def test_merge_keeps_rows_without_a_required_value(tmp_folder, engine): b = _part( tmp_folder, "b", "b", 2, columns={"admin:country_code": ["FR", "FR"]}, schemas=[CORE, ADMIN] ) - out = _merge(tmp_folder, [a, b], engine) + 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"] @@ -410,7 +411,7 @@ def test_merge_keeps_empty_geometries(tmp_folder, engine): gdf = gp.read() gdf.loc[0, "geometry"] = shapely.Polygon() gp.write(gdf, dehydrate=False) - out = _merge(tmp_folder, [a, b], engine) + out = _merge(tmp_folder, [a, b], engine, strict=False) assert pq.read_metadata(out).num_rows == 4 @@ -418,7 +419,7 @@ def test_merge_keeps_empty_geometries(tmp_folder, engine): 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() @@ -474,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 4fa629a..580588c 100644 --- a/tests/test_ops.py +++ b/tests/test_ops.py @@ -175,7 +175,42 @@ def warning(self, message): assert "producer" not in merged log = Log() - warn_missing_required(merged, properties, schema_map, log) - assert log.warnings == [ - "Required properties are not included, the merged file will be invalid: producer" + 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 ec9f1d0..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,16 +6,15 @@ 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 @@ -57,43 +54,22 @@ 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, log: Optional[LoggerMixin] = None) -> 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 {} - dtype = schema.get("type") - try: - pa_type = get_pyarrow_type(schema) if dtype else None - except Exception: - # a schema that has no pyarrow type, e.g. an object with additionalProperties - pa_type = None - if pa_type is None: - arrays.append(pa.array([_to_arrow_value(value, dtype)])) - continue try: - arrays.append(pa.array([_to_arrow_value(value, dtype)], type=pa_type)) - except (ValueError, TypeError, OverflowError) as e: - if log: - log.warning(f"Constant '{key}' doesn't fit its type {dtype}, keeping it as is: {e}") - 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())) @@ -356,6 +332,7 @@ def write_query( 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 +342,7 @@ def write_query( run over. Ids only need to be unique per collection. Without `strict`, a missing required value is only reported and empty geometries are kept, as the in-memory merge does. + When `merging` strictly, empty geometries are an error instead of being dropped. """ compression = compression or "zstd" if compression == "zstd" and compression_level is None: @@ -391,21 +369,27 @@ def write_query( schema_groups = collection.get_schemas() per_collection = "collection" in selected_targets and len(schema_groups) > 1 - def null_condition(schema, skip=("geometry",)): + # The collections that require a property, None for all rows + required_by = {} + + def null_condition(schema, skip=("geometry",), cid=None): required = [ r for r in schema.get("required", []) if r not in skip and r not in collection_only and r in selected_targets ] + for r in required: + required_by.setdefault(r, []).append(cid) return " OR ".join(f"{_sql_name(target)} IS NULL" for target in required) or None if per_collection: # Each collection only requires what its own schemas require custom_schemas = collection.get_custom_schemas() conditions = ['"collection" IS NULL'] + required_by["collection"] = [None] for cid, group in schema_groups.items(): schema = group.merge_schemas(custom_schemas=custom_schemas) - cond = null_condition(schema, skip=("geometry", "collection")) + 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) @@ -414,6 +398,12 @@ def null_condition(schema, skip=("geometry",)): 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: @@ -440,6 +430,7 @@ def null_condition(schema, skip=("geometry",)): ) 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) @@ -451,28 +442,35 @@ def null_condition(schema, skip=("geometry",)): "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 and not strict: - self.warning( - f"{invalid} of {total} rows have no value for a required property, " - "the merged file will be invalid" - ) + missing = ", ".join( + f"{name} ({count:,} of {total:,} rows)" + for name, count in sorted(nulls.items()) + if count + ) + if invalid and (merging or not strict): + report(f"Rows have no value for a required property: {missing}", self, strict) elif invalid: # A null in a required property is an error, whatever the count: # the writer rejects nulls in the non-nullable required fields # anyway, and silently dropping rows would make that data-quality # decision for the user (vecorel/cli#33) raise ValueError( - f"{invalid} of {total} rows have no value for a required property " - f"({null_cond}). Handle them in the converter: fix the mapping, " - "fill the values in column_migrations, or exclude the rows with " - "a column_filters entry." + f"Rows have no value for a required property: {missing}. " + "Handle them in the converter: fix the mapping, fill the values in " + "column_migrations, or exclude the rows with a column_filters entry." ) - if blanks and strict: + 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: @@ -573,6 +571,7 @@ def null_condition(schema, skip=("geometry",)): compression_level=compression_level, geoparquet_version=geoparquet_version, crs=source_crs, + strict=strict, ) return output_file @@ -643,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. @@ -652,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") @@ -668,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, @@ -696,6 +705,17 @@ def merge_parquet( if fill is not None and "collection" not in names and "collection" not in collection: # Features of multiple collections must state their collection constants["collection"] = fill + elif fill is None and keep_collection: + # the statistics tell whether the column has nulls without a scan + if "collection" in names: + with pq.ParquetFile(path) as pq_file: + unknown = bool(GeoParquet._columns_with_nulls(pq_file, {"collection"})) + else: + unknown = "collection" not in collection + if unknown: + report( + f"Can't determine the collection of the features in {path}", self, strict + ) columns = [] for name in names: @@ -711,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, log=self)) + 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}") @@ -725,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/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 6456d00..09d61ff 100644 --- a/vecorel_cli/merge.py +++ b/vecorel_cli/merge.py @@ -30,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" @@ -71,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 @@ -82,6 +91,7 @@ def merge( includes=[], excludes=[], engine="auto", + strict=True, ): if not isinstance(source, list): raise ValueError("Source must be a list.") @@ -123,16 +133,21 @@ def merge( target.uri, properties=properties, suffix_duplicate_ids=False, - strict=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 + 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 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 8157df2..1ec4249 100644 --- a/vecorel_cli/vecorel/ops.py +++ b/vecorel_cli/vecorel/ops.py @@ -5,11 +5,23 @@ 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, @@ -17,6 +29,7 @@ def merge( schema_map: SchemaMapping = {}, log: Optional[LoggerMixin] = None, excludes: Optional[list[str]] = None, + strict: bool = False, ) -> tuple[GeoDataFrame, Collection]: frames = [item.read(properties=properties, schema_map=schema_map) for item in encodings] collections = [item.get_collection() for item in encodings] @@ -27,7 +40,12 @@ def merge( 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) + merged_collection = merge_collections( + collections, properties=properties, log=log, strict=strict, schema_map=schema_map + ) + if properties is not None: + warn_missing_required(collections, properties, schema_map, log, strict) + props = merged_collection.merge_schemas(schema_map=schema_map).get("properties", {}) data = [] for item, gdf, collection in zip(encodings, frames, collections): @@ -37,7 +55,26 @@ def merge( 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 keep_collection = properties is None or "collection" in properties cid = get_collection_id(collection) if keep_collection else None @@ -45,6 +82,12 @@ def merge( 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 @@ -57,22 +100,77 @@ def merge( # 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()] 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) - if log and properties is not None: - warn_missing_required(merged_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) + + 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) -> Optional[str]: """ The collection id of a dataset that doesn't state it for all features: @@ -89,20 +187,36 @@ def get_collection_id(collection: Collection) -> Optional[str]: 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)) - in_collection = collection_only & set(collection.keys()) - missing = set(schema.get("required", [])) - set(properties) - in_collection - {"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() @@ -130,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 @@ -142,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)