diff --git a/CHANGELOG.md b/CHANGELOG.md index ed5131f0..7b012d15 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,6 +8,7 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0. ## [Unreleased] ### Added +- Added `PerFileConverterMixin`, which converts each source file on its own and merges the parts; ES uses it, as its 50 provinces do not fit in memory together. - Added `pixi run check-hcat`, which compares the HCAT mapping tables the converters read online with the taxonomy. - Added `FiboaDuckDBBaseConverter` for SQL-based conversion of large Parquet sources. - Added `PerFileBaseConverter` to process multi-file sources incrementally. diff --git a/fiboa_cli/conversion/per_file.py b/fiboa_cli/conversion/per_file.py new file mode 100644 index 00000000..d0161f19 --- /dev/null +++ b/fiboa_cli/conversion/per_file.py @@ -0,0 +1,34 @@ +import os +from tempfile import TemporaryDirectory + +from vecorel_cli.conversion.duckdb import DuckDBBaseConverter + + +class PerFileConverterMixin: + """Converts each source file on its own and merges the parts, for sources too large to read at once.""" + + def convert(self, output_file, cache=None, input_files=None, variant=None, **kwargs) -> str: + self.select_variant(variant) + urls = input_files or self.get_urls() + if not urls or len(urls) == 1: + return super().convert( + output_file, cache=cache, input_files=input_files, variant=variant, **kwargs + ) + + directory = os.path.dirname(output_file) or "." + os.makedirs(directory, exist_ok=True) + with TemporaryDirectory(dir=directory) as parts_folder: + parts = [] + for index, (uri, target) in enumerate(urls.items()): + self.info(f"Converting source {index + 1}/{len(urls)}: {uri}") + part = os.path.join(parts_folder, f"part_{index}.parquet") + super().convert( + part, cache=cache, input_files={uri: target}, variant=variant, **kwargs + ) + parts.append(part) + + self.info(f"Merging {len(parts)} parts into {output_file}") + merge_options = ("compression", "compression_level", "geoparquet_version") + return DuckDBBaseConverter().merge_parquet( + parts, output_file, **{k: kwargs[k] for k in merge_options if k in kwargs} + ) diff --git a/fiboa_cli/datasets/es.py b/fiboa_cli/datasets/es.py index 742c803e..d88da6fe 100644 --- a/fiboa_cli/datasets/es.py +++ b/fiboa_cli/datasets/es.py @@ -4,10 +4,12 @@ from vecorel_cli.vecorel.extensions import ADMIN_DIVISION from ..conversion.fiboa_converter import FiboaBaseConverter +from ..conversion.per_file import PerFileConverterMixin from .commons.hcat import AddHCATMixin -class Converter(AddHCATMixin, FiboaBaseConverter): +# 50 provincial GeoPackages that do not fit in memory together +class Converter(PerFileConverterMixin, AddHCATMixin, FiboaBaseConverter): id = "es" short_name = "Spain" title = "Spain Declared Crops (Cultivos Declarados SIGPAC)" diff --git a/tests/test_per_file.py b/tests/test_per_file.py new file mode 100644 index 00000000..e30f91d5 --- /dev/null +++ b/tests/test_per_file.py @@ -0,0 +1,24 @@ +import shutil + +import geopandas as gpd + +from fiboa_cli.converters import Converters + +FIXTURE = "tests/data-files/convert/es/1501_ALAVA_cd_2025_20250105.gpkg.zip" + + +def test_each_source_is_converted_on_its_own_and_merged(tmp_path): + from fiboa_cli import Registry # noqa + + copy = tmp_path / "copy.gpkg.zip" + shutil.copy(FIXTURE, copy) + output = tmp_path / "es.parquet" + + Converters().load("es").convert( + str(output), input_files={FIXTURE: ["*.gpkg"], str(copy): ["*.gpkg"]} + ) + + merged = gpd.read_parquet(output) + assert len(merged) == 20 # 10 rows per source + assert merged["id"].is_unique + assert set(tmp_path.iterdir()) == {copy, output} # the parts are gone