Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
34 changes: 34 additions & 0 deletions fiboa_cli/conversion/per_file.py
Original file line number Diff line number Diff line change
@@ -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}
)
4 changes: 3 additions & 1 deletion fiboa_cli/datasets/es.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)"
Expand Down
24 changes: 24 additions & 0 deletions tests/test_per_file.py
Original file line number Diff line number Diff line change
@@ -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
Loading