From d388e98a56312cc6478503c4ea4b4e8d125e5874 Mon Sep 17 00:00:00 2001 From: Ivor Bosloper Date: Sat, 26 Sep 2026 23:22:14 +0200 Subject: [PATCH 01/12] Merge: hydrate only what the parts disagree on (review #1) The in-memory merge hydrated every constant and relied on the dehydration on write to move them back, which fails for arrays and objects. It now merges the collections first and hydrates only the keys the merged collection drops, as merge_parquet does. Co-Authored-By: Claude Opus 5.5 --- tests/test_merge.py | 14 ++++++++++++++ vecorel_cli/encoding/base.py | 12 ++++++++++-- vecorel_cli/vecorel/ops.py | 26 ++++++++++++++------------ 3 files changed, 38 insertions(+), 14 deletions(-) diff --git a/tests/test_merge.py b/tests/test_merge.py index f38db58..8ea307b 100644 --- a/tests/test_merge.py +++ b/tests/test_merge.py @@ -182,6 +182,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"}}} 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/vecorel/ops.py b/vecorel_cli/vecorel/ops.py index 8e10f38..cb436dc 100644 --- a/vecorel_cli/vecorel/ops.py +++ b/vecorel_cli/vecorel/ops.py @@ -17,13 +17,19 @@ def merge( schema_map: SchemaMapping = {}, log: Optional[LoggerMixin] = None, ) -> 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] + merged_collection = merge_collections(collections, properties=properties, log=log) - 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) + ] + gdf = item.hydrate_from_collection(gdf, schema_map=schema_map, keys=keys) if "collection" not in gdf.columns or gdf["collection"].isna().any(): cid = get_collection_id(collection, item.uri) @@ -39,9 +45,7 @@ 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)) @@ -54,12 +58,10 @@ def merge( if duplicates: log.warning(f"{duplicates} rows repeat an id within their collection") - # 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) + warn_missing_required(merged_collection, properties, schema_map, log) - return merged, collection + return merged, merged_collection def get_collection_id(collection: Collection, source=None) -> str: From 325f87475ba88d80d45daa3f34eac5ff371e0b32 Mon Sep 17 00:00:00 2001 From: Ivor Bosloper Date: Sat, 26 Sep 2026 23:24:51 +0200 Subject: [PATCH 02/12] Merge: keep rows as they are with DuckDB, like in memory (review #2) vec merge only appends datasets, so the DuckDB engine no longer applies the converter checks there: a missing required value is reported with a merge-appropriate warning instead of an error with converter advice, and empty geometries are kept, as the in-memory merge does. This is a new strict flag on write_query that vec merge turns off; converters and merge_parquet itself keep the checks. The per-collection null condition no longer repeats the collection check. Co-Authored-By: Claude Opus 5.5 --- tests/test_merge.py | 28 ++++++++++++++++++++++++++++ vecorel_cli/conversion/duckdb.py | 19 ++++++++++++++----- vecorel_cli/merge.py | 1 + 3 files changed, 43 insertions(+), 5 deletions(-) diff --git a/tests/test_merge.py b/tests/test_merge.py index 8ea307b..e0a2714 100644 --- a/tests/test_merge.py +++ b/tests/test_merge.py @@ -349,6 +349,34 @@ def test_merge_excludes_the_collection(tmp_folder, engine, log): 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) + + 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) + + 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) diff --git a/vecorel_cli/conversion/duckdb.py b/vecorel_cli/conversion/duckdb.py index 10f8ca4..373a935 100644 --- a/vecorel_cli/conversion/duckdb.py +++ b/vecorel_cli/conversion/duckdb.py @@ -358,6 +358,7 @@ def write_query( geoparquet_version: Optional[str] = None, original_geometries: bool = False, suffix_duplicate_ids: bool = True, + strict: bool = True, ) -> 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 +366,8 @@ 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. """ compression = compression or "zstd" if compression == "zstd" and compression_level is None: @@ -390,11 +393,11 @@ def write_query( collection_only = set(collection.get_collection_only_properties()) has_collection_column = "collection" in selected_targets - def null_condition(schema): + def null_condition(schema, skip=("geometry",)): 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 @@ -404,7 +407,8 @@ def null_condition(schema): custom_schemas = collection.get_custom_schemas() conditions = ['"collection" IS NULL'] 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")) if cond: conditions.append(f'("collection" = {_sql_literal(cid)} AND ({cond}))') null_cond = " OR ".join(conditions) @@ -457,7 +461,12 @@ def null_condition(schema): ) blanks = values.pop(0) if blank_cond else 0 repairs = values.pop(0) if repair_cond else 0 - if invalid: + 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" + ) + 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 @@ -468,7 +477,7 @@ def null_condition(schema): "fill the values in column_migrations, or exclude the rows with " "a column_filters entry." ) - if blanks: + if 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: diff --git a/vecorel_cli/merge.py b/vecorel_cli/merge.py index c022d08..3fa6f2b 100644 --- a/vecorel_cli/merge.py +++ b/vecorel_cli/merge.py @@ -121,6 +121,7 @@ def merge( target.uri, properties=properties, suffix_duplicate_ids=False, + strict=False, ) else: if engine == "auto": From 6856d80431575833242c2b4e8a7816585d7960c0 Mon Sep 17 00:00:00 2001 From: Ivor Bosloper Date: Sat, 26 Sep 2026 23:26:50 +0200 Subject: [PATCH 03/12] Report constants that don't fit their schema type (review #3) _constants_table caught every exception and fell back to an inferred type silently. It now only catches the conversion errors and warns with the property and its type; a schema without a pyarrow type still falls back without a warning. Co-Authored-By: Claude Opus 5.5 --- tests/test_convert_duckdb.py | 15 ++++++++++++++- vecorel_cli/conversion/duckdb.py | 19 +++++++++++++++---- 2 files changed, 29 insertions(+), 5 deletions(-) diff --git a/tests/test_convert_duckdb.py b/tests/test_convert_duckdb.py index 373f984..934264d 100644 --- a/tests/test_convert_duckdb.py +++ b/tests/test_convert_duckdb.py @@ -9,7 +9,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 +723,19 @@ 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) + + assert table.column("x").to_pylist() == [value] + assert f"'x' doesn't fit its type {dtype}" in capsys.readouterr().out + + def test_merge_parquet_checks_ids_that_convert_generated(tmp_folder, capsys): """Each part numbered its own rows from zero, so the merge has to check what convert() is allowed to take for granted.""" diff --git a/vecorel_cli/conversion/duckdb.py b/vecorel_cli/conversion/duckdb.py index 373a935..1d7f1c1 100644 --- a/vecorel_cli/conversion/duckdb.py +++ b/vecorel_cli/conversion/duckdb.py @@ -12,6 +12,7 @@ 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 @@ -72,15 +73,25 @@ def _to_arrow_value(value, dtype): return value -def _constants_table(constants: dict, props: dict) -> pa.Table: +def _constants_table(constants: dict, props: dict, log: Optional[LoggerMixin] = None) -> 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 schema.get("type") else None - arrays.append(pa.array([_to_arrow_value(value, schema.get("type"))], type=pa_type)) + 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])) return pa.table(arrays, names=list(constants.keys())) @@ -723,7 +734,7 @@ 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)) query += f", {table}.*" source += f" CROSS JOIN {table}" selects.append(f"{query} FROM {source}") From 62d12dfee2b6080269e22314f727649635cca8ca Mon Sep 17 00:00:00 2001 From: Ivor Bosloper Date: Sat, 26 Sep 2026 23:27:49 +0200 Subject: [PATCH 04/12] Merge: reject two versions of one schema in a collection (review #4) add_all only rejected conflicting core versions; two versions of the same extension (e.g. fiboa v0.2.0 and v0.3.0) in one collection were kept side by side. Schemas that only differ in their version segment are now rejected the same way. Co-Authored-By: Claude Opus 5.5 --- tests/test_ops.py | 16 ++++++++++++++++ vecorel_cli/vecorel/schemas.py | 9 +++++++++ 2 files changed, 25 insertions(+) diff --git a/tests/test_ops.py b/tests/test_ops.py index 2e6912a..b31a7fd 100644 --- a/tests/test_ops.py +++ b/tests/test_ops.py @@ -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: 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: From 9f5109883c2b3e028f14548f9e890a4868564c8e Mon Sep 17 00:00:00 2001 From: Ivor Bosloper Date: Sat, 26 Sep 2026 23:28:36 +0200 Subject: [PATCH 05/12] DuckDB: key ids by collection only with several collections (review #5) The id uniqueness check and the suffixing used a (collection, id) key whenever there is a collection column, which slows down the common single-collection conversion for nothing. Co-Authored-By: Claude Opus 5.5 --- vecorel_cli/conversion/duckdb.py | 17 +++++++---------- 1 file changed, 7 insertions(+), 10 deletions(-) diff --git a/vecorel_cli/conversion/duckdb.py b/vecorel_cli/conversion/duckdb.py index 1d7f1c1..05e43fb 100644 --- a/vecorel_cli/conversion/duckdb.py +++ b/vecorel_cli/conversion/duckdb.py @@ -402,7 +402,8 @@ 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 def null_condition(schema, skip=("geometry",)): required = [ @@ -412,8 +413,7 @@ def null_condition(schema, skip=("geometry",)): ] return " OR ".join(f'"{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'] @@ -431,9 +431,7 @@ def null_condition(schema, skip=("geometry",)): # 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 @@ -551,7 +549,7 @@ def null_condition(schema, skip=("geometry",)): 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 @@ -594,7 +592,7 @@ def null_condition(schema, skip=("geometry",)): 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 @@ -603,8 +601,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( From 5427a332417665b722a23bc4a536ccbc92ef6d66 Mon Sep 17 00:00:00 2001 From: Ivor Bosloper Date: Sat, 26 Sep 2026 23:30:54 +0200 Subject: [PATCH 06/12] Merge: apply excludes to the loaded data in memory (review #6) With --exclude, every dataset without a cheap property listing (e.g. GeoJSON) was read in full to get its columns and then read again by the merge. The in-memory merge now takes the excludes and resolves them on the data it has loaded; the DuckDB path still lists the GeoParquet columns, which only reads the schema. A merge without the collection no longer adds it back or needs it for the duplicate check. Co-Authored-By: Claude Opus 5.5 --- tests/test_merge.py | 23 +++++++++++++++++++++++ vecorel_cli/merge.py | 13 ++++++++----- vecorel_cli/vecorel/ops.py | 14 ++++++++++++-- 3 files changed, 43 insertions(+), 7 deletions(-) diff --git a/tests/test_merge.py b/tests/test_merge.py index e0a2714..9779915 100644 --- a/tests/test_merge.py +++ b/tests/test_merge.py @@ -11,6 +11,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 +63,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() diff --git a/vecorel_cli/merge.py b/vecorel_cli/merge.py index 3fa6f2b..ef82978 100644 --- a/vecorel_cli/merge.py +++ b/vecorel_cli/merge.py @@ -103,10 +103,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,6 +111,11 @@ 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], @@ -126,7 +127,9 @@ def merge( 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 + ) target.set_collection(collection) target.write(gdf, properties=properties) diff --git a/vecorel_cli/vecorel/ops.py b/vecorel_cli/vecorel/ops.py index cb436dc..6736cdc 100644 --- a/vecorel_cli/vecorel/ops.py +++ b/vecorel_cli/vecorel/ops.py @@ -16,9 +16,17 @@ def merge( properties=None, schema_map: SchemaMapping = {}, log: Optional[LoggerMixin] = None, + excludes: Optional[list[str]] = None, ) -> tuple[GeoDataFrame, Collection]: frames = [item.read(properties=properties, schema_map=schema_map) for item in encodings] collections = [item.get_collection() for item in encodings] + if excludes: + if properties is None: + properties = set() + for gdf, collection in zip(frames, collections): + properties |= set(gdf.columns) | set(collection.keys()) + properties = list(set(properties) - set(excludes)) + frames = [gdf.drop(columns=[c for c in excludes if c in gdf.columns]) for gdf in frames] merged_collection = merge_collections(collections, properties=properties, log=log) data = [] @@ -31,7 +39,8 @@ def merge( ] gdf = item.hydrate_from_collection(gdf, schema_map=schema_map, keys=keys) - if "collection" not in gdf.columns or gdf["collection"].isna().any(): + keep_collection = properties is None or "collection" in properties + if keep_collection and ("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) @@ -54,7 +63,8 @@ def merge( if log and "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") From be2fed7824d15dd7cfcb71c092c9f31b7f606215 Mon Sep 17 00:00:00 2001 From: Ivor Bosloper Date: Sat, 26 Sep 2026 23:31:15 +0200 Subject: [PATCH 07/12] DuckDB: quote required property names with _sql_name (review #7) Co-Authored-By: Claude Opus 5.5 --- vecorel_cli/conversion/duckdb.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/vecorel_cli/conversion/duckdb.py b/vecorel_cli/conversion/duckdb.py index 05e43fb..9d514b2 100644 --- a/vecorel_cli/conversion/duckdb.py +++ b/vecorel_cli/conversion/duckdb.py @@ -411,7 +411,7 @@ def null_condition(schema, skip=("geometry",)): for r in schema.get("required", []) 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 + 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 From 7c2dc1a766e69ecaa1de08882f662527cbfce3ec Mon Sep 17 00:00:00 2001 From: Ivor Bosloper Date: Sat, 26 Sep 2026 23:32:26 +0200 Subject: [PATCH 08/12] Share the CRS comparison of merge and DuckDB (review #8) get_duckdb_blocker re-implemented the loop of DuckDBBaseConverter._common_crs with its private helpers. Both now use find_differing_crs from vecorel.util. Co-Authored-By: Claude Opus 5.5 --- vecorel_cli/conversion/duckdb.py | 36 ++++++++++---------------------- vecorel_cli/merge.py | 13 +++++------- vecorel_cli/vecorel/util.py | 19 +++++++++++++++++ 3 files changed, 35 insertions(+), 33 deletions(-) diff --git a/vecorel_cli/conversion/duckdb.py b/vecorel_cli/conversion/duckdb.py index 9d514b2..bc4344d 100644 --- a/vecorel_cli/conversion/duckdb.py +++ b/vecorel_cli/conversion/duckdb.py @@ -18,6 +18,7 @@ from ..parquet.types import get_pyarrow_type 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.util import find_differing_crs from .base import BaseConverter @@ -96,19 +97,6 @@ def _constants_table(constants: dict, props: dict, log: Optional[LoggerMixin] = 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. @@ -332,9 +320,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] @@ -343,15 +330,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, diff --git a/vecorel_cli/merge.py b/vecorel_cli/merge.py index ef82978..6456d00 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): @@ -154,20 +155,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/vecorel/util.py b/vecorel_cli/vecorel/util.py index 1a9d30d..4e0812d 100644 --- a/vecorel_cli/vecorel/util.py +++ b/vecorel_cli/vecorel/util.py @@ -16,6 +16,25 @@ file_cache = {} +def find_differing_crs(crs_values: list, reference=None) -> Optional[int]: + """The index of the first GeoParquet crs value that differs from the reference + (by default the first value), None if they all agree. A missing crs is OGC:CRS84.""" + from pyproj import CRS + + def normalize(crs): + return CRS.from_user_input(crs if crs is not None else "OGC:CRS84") + + reference = normalize(reference) if reference is not None else None + for i, crs in enumerate(crs_values): + crs = normalize(crs) + if reference is None: + reference = crs + # GeoParquet coordinates are always x, y regardless of the CRS axis order + elif not crs.equals(reference, ignore_axis_order=True): + return i + return None + + def suffix_duplicate_ids(gdf): """Number the ids that appear on several rows (id~1, id~2, ...), which turns the id column into strings. Null ids are left alone. Returns the frame and From 2fb1d4c8f83a0830aed31d0c3ad40f90f82be6b2 Mon Sep 17 00:00:00 2001 From: Ivor Bosloper Date: Sat, 26 Sep 2026 23:33:55 +0200 Subject: [PATCH 09/12] Merge: warn when includes drop a required collection-only property (review #9) warn_missing_required skipped every collection-only property, but merge_collections drops those that are not in the included properties, so a required one disappeared without a warning. Only the ones the merged collection still carries are skipped now. Co-Authored-By: Claude Opus 5.5 --- tests/test_ops.py | 34 +++++++++++++++++++++++++++++++++- vecorel_cli/vecorel/ops.py | 3 ++- 2 files changed, 35 insertions(+), 2 deletions(-) diff --git a/tests/test_ops.py b/tests/test_ops.py index b31a7fd..4fa629a 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 @@ -147,3 +147,35 @@ 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(merged, properties, schema_map, log) + assert log.warnings == [ + "Required properties are not included, the merged file will be invalid: producer" + ] diff --git a/vecorel_cli/vecorel/ops.py b/vecorel_cli/vecorel/ops.py index 6736cdc..b0595f3 100644 --- a/vecorel_cli/vecorel/ops.py +++ b/vecorel_cli/vecorel/ops.py @@ -93,7 +93,8 @@ def warn_missing_required( ): 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"} + in_collection = collection_only & set(collection.keys()) + missing = set(schema.get("required", [])) - set(properties) - in_collection - {"geometry"} if missing: log.warning( f"Required properties are not included, the merged file will be invalid: {', '.join(sorted(missing))}" From c74f79e3033e237d7a9f653dcaf7bff0364e88a1 Mon Sep 17 00:00:00 2001 From: Ivor Bosloper Date: Sat, 26 Sep 2026 23:35:26 +0200 Subject: [PATCH 10/12] Merge: accept a collection that can't be determined (review #10) get_collection_id raised for rows without a collection when the metadata doesn't name one, so create-geoparquet and the in-memory merge failed on input they used to accept, while merge_parquet caught it in one place only. It now returns None and both merges leave such rows as they are. Co-Authored-By: Claude Opus 5.5 --- tests/test_merge.py | 15 +++++++++++++++ vecorel_cli/conversion/duckdb.py | 12 +++--------- vecorel_cli/vecorel/ops.py | 16 ++++++++-------- 3 files changed, 26 insertions(+), 17 deletions(-) diff --git a/tests/test_merge.py b/tests/test_merge.py index 9779915..13a6603 100644 --- a/tests/test_merge.py +++ b/tests/test_merge.py @@ -362,6 +362,21 @@ 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) + + 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) diff --git a/vecorel_cli/conversion/duckdb.py b/vecorel_cli/conversion/duckdb.py index bc4344d..ec9f1d0 100644 --- a/vecorel_cli/conversion/duckdb.py +++ b/vecorel_cli/conversion/duckdb.py @@ -692,16 +692,10 @@ 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 columns = [] for name in names: diff --git a/vecorel_cli/vecorel/ops.py b/vecorel_cli/vecorel/ops.py index b0595f3..8157df2 100644 --- a/vecorel_cli/vecorel/ops.py +++ b/vecorel_cli/vecorel/ops.py @@ -40,12 +40,11 @@ def merge( gdf = item.hydrate_from_collection(gdf, schema_map=schema_map, keys=keys) keep_collection = properties is None or "collection" in properties - if keep_collection and ("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 + 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 if not crs: # If no CRS is given, use the first CRS that is available as the base CRS @@ -74,10 +73,11 @@ def merge( return merged, merged_collection -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: @@ -85,7 +85,7 @@ 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( From 9b05e198c692366b8d95159dde85ed824ff16da9 Mon Sep 17 00:00:00 2001 From: Ivor Bosloper Date: Sat, 26 Sep 2026 23:35:50 +0200 Subject: [PATCH 11/12] Changelog for the review fixes Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 4583d47..cd0e7db 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -46,6 +46,15 @@ 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 From daaa30fd4e0da98fca4f58f9e5588072f549f892 Mon Sep 17 00:00:00 2001 From: Matthias Mohr Date: Mon, 28 Sep 2026 13:39:42 +0200 Subject: [PATCH 12/12] Add strict mode to vec merge, minor improvements and bug fixes (#60) * Add strict mode to vec merge, minor improvements and bug fixes * Fix typed constants, missing geometries, nested nulls and schema map in merge --- CHANGELOG.md | 36 ++--- README.md | 27 ++++ tests/test_convert_duckdb.py | 13 +- tests/test_merge.py | 210 +++++++++++++++++++++++++++-- tests/test_ops.py | 41 +++++- vecorel_cli/conversion/duckdb.py | 137 +++++++++++-------- vecorel_cli/create_geoparquet.py | 4 +- vecorel_cli/encoding/geoparquet.py | 37 +++++ vecorel_cli/merge.py | 21 ++- vecorel_cli/parquet/types.py | 40 ++++++ vecorel_cli/vecorel/ops.py | 160 +++++++++++++++++++--- 11 files changed, 611 insertions(+), 115 deletions(-) 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)