diff --git a/notebooks/test_delete_items_from_collection.ipynb b/notebooks/test_delete_items_from_collection.ipynb index 13e9c9b..3cffbbe 100644 --- a/notebooks/test_delete_items_from_collection.ipynb +++ b/notebooks/test_delete_items_from_collection.ipynb @@ -9,8 +9,8 @@ { "metadata": { "ExecuteTime": { - "end_time": "2026-09-15T08:25:06.246567Z", - "start_time": "2026-09-15T08:24:59.582135700Z" + "end_time": "2026-09-22T07:07:48.821699400Z", + "start_time": "2026-09-22T07:07:45.460405400Z" } }, "cell_type": "code", @@ -29,8 +29,8 @@ "metadata": { "collapsed": true, "ExecuteTime": { - "end_time": "2026-09-15T08:25:18.518515200Z", - "start_time": "2026-09-15T08:25:18.506034800Z" + "end_time": "2026-09-22T07:07:48.847731200Z", + "start_time": "2026-09-22T07:07:48.821699400Z" } }, "source": [ @@ -45,8 +45,8 @@ { "metadata": { "ExecuteTime": { - "end_time": "2026-09-15T08:25:22.668653500Z", - "start_time": "2026-09-15T08:25:20.543609200Z" + "end_time": "2026-09-22T07:07:59.435229900Z", + "start_time": "2026-09-22T07:07:58.469442300Z" } }, "cell_type": "code", @@ -73,8 +73,8 @@ { "metadata": { "ExecuteTime": { - "end_time": "2026-09-15T08:25:25.181893800Z", - "start_time": "2026-09-15T08:25:25.131787900Z" + "end_time": "2026-09-22T07:08:01.640645Z", + "start_time": "2026-09-22T07:08:01.598193500Z" } }, "cell_type": "code", @@ -87,19 +87,19 @@ { "data": { "text/plain": [ - " item_id modelID typology \\\n", - "0 IUCNGET_global_v321_2024_IND IUCNGET_global_v321_2024_IND IUCNGET \n", - "1 IUCNGET_global_v315_2024_AFR IUCNGET_global_v315_2024_AFR IUCNGET \n", - "2 IUCNGET_global_v314_2024_NEO IUCNGET_global_v314_2024_NEO IUCNGET \n", - "3 IUCNGET_global_v314_2024_IND IUCNGET_global_v314_2024_IND IUCNGET \n", - "4 IUCNGET_global_v314_2024_AFR IUCNGET_global_v314_2024_AFR IUCNGET \n", + " item_id modelID \\\n", + "0 IUCNGET_global_v317_2024_NEO IUCNGET_global_v317_2024_NEO \n", + "1 IUCNGET_global_v317_2024_IND IUCNGET_global_v317_2024_IND \n", + "2 IUCNGET_global_v317_2024_HOL IUCNGET_global_v317_2024_HOL \n", + "3 IUCNGET_global_v317_2024_AFR IUCNGET_global_v317_2024_AFR \n", + "4 EUNIS2021plus_panEU_v311_2024_STE EUNIS2021plus_panEU_v311_2024_STE \n", "\n", - " training_year model_version spatial_region spatial_zone \n", - "0 2024 3.21 global IndoMalesian \n", - "1 2024 3.15 global African \n", - "2 2024 3.14 global Neotropical \n", - "3 2024 3.14 global IndoMalesian \n", - "4 2024 3.14 global African " + " typology training_year model_version spatial_region spatial_zone \n", + "0 IUCNGET 2024 3.17 global Neotropical \n", + "1 IUCNGET 2024 3.17 global IndoMalesian \n", + "2 IUCNGET 2024 3.17 global Holartic \n", + "3 IUCNGET 2024 3.17 global African \n", + "4 EUNIS2021plus 2024 3.11 panEU Steppic " ], "text/html": [ "
\n", @@ -132,53 +132,53 @@ " \n", " \n", " 0\n", - " IUCNGET_global_v321_2024_IND\n", - " IUCNGET_global_v321_2024_IND\n", + " IUCNGET_global_v317_2024_NEO\n", + " IUCNGET_global_v317_2024_NEO\n", " IUCNGET\n", " 2024\n", - " 3.21\n", + " 3.17\n", " global\n", - " IndoMalesian\n", + " Neotropical\n", " \n", " \n", " 1\n", - " IUCNGET_global_v315_2024_AFR\n", - " IUCNGET_global_v315_2024_AFR\n", + " IUCNGET_global_v317_2024_IND\n", + " IUCNGET_global_v317_2024_IND\n", " IUCNGET\n", " 2024\n", - " 3.15\n", + " 3.17\n", " global\n", - " African\n", + " IndoMalesian\n", " \n", " \n", " 2\n", - " IUCNGET_global_v314_2024_NEO\n", - " IUCNGET_global_v314_2024_NEO\n", + " IUCNGET_global_v317_2024_HOL\n", + " IUCNGET_global_v317_2024_HOL\n", " IUCNGET\n", " 2024\n", - " 3.14\n", + " 3.17\n", " global\n", - " Neotropical\n", + " Holartic\n", " \n", " \n", " 3\n", - " IUCNGET_global_v314_2024_IND\n", - " IUCNGET_global_v314_2024_IND\n", + " IUCNGET_global_v317_2024_AFR\n", + " IUCNGET_global_v317_2024_AFR\n", " IUCNGET\n", " 2024\n", - " 3.14\n", + " 3.17\n", " global\n", - " IndoMalesian\n", + " African\n", " \n", " \n", " 4\n", - " IUCNGET_global_v314_2024_AFR\n", - " IUCNGET_global_v314_2024_AFR\n", - " IUCNGET\n", + " EUNIS2021plus_panEU_v311_2024_STE\n", + " EUNIS2021plus_panEU_v311_2024_STE\n", + " EUNIS2021plus\n", " 2024\n", - " 3.14\n", - " global\n", - " African\n", + " 3.11\n", + " panEU\n", + " Steppic\n", " \n", " \n", "\n", @@ -244,8 +244,8 @@ { "metadata": { "ExecuteTime": { - "end_time": "2026-09-15T08:25:29.386166Z", - "start_time": "2026-09-15T08:25:29.308068100Z" + "end_time": "2026-09-22T07:10:52.829664300Z", + "start_time": "2026-09-22T07:10:52.499465400Z" } }, "cell_type": "code", @@ -255,7 +255,7 @@ ], "id": "4d5fe40f5ffc5142", "outputs": [], - "execution_count": 7 + "execution_count": 8 }, { "metadata": { @@ -298,13 +298,59 @@ ], "execution_count": 8 }, + { + "metadata": { + "ExecuteTime": { + "end_time": "2026-09-22T07:10:13.316345200Z", + "start_time": "2026-09-22T07:10:13.306016200Z" + } + }, + "cell_type": "code", + "source": [ + "# delete collections id\n", + "ldelcol = ['test-V311', 'test-V311-jobdb', 'IUCNGET-v3-MECE']" + ], + "id": "857e74f1803a5361", + "outputs": [], + "execution_count": 6 + }, + { + "metadata": { + "ExecuteTime": { + "end_time": "2026-09-22T07:10:59.695370600Z", + "start_time": "2026-09-22T07:10:56.402445700Z" + } + }, + "cell_type": "code", + "source": [ + "for item in ldelcol:\n", + " print(f'delete collection: {item}')\n", + " st.delete_collection(collection_name=item)" + ], + "id": "bc3b1c4b09ee83f6", + "outputs": [ + { + "name": "stdout", + "output_type": "stream", + "text": [ + "delete collection: test-V311\n", + "Collection test-V311 deleted successfully\n", + "delete collection: test-V311-jobdb\n", + "Collection test-V311-jobdb deleted successfully\n", + "delete collection: IUCNGET-v3-MECE\n", + "Collection IUCNGET-v3-MECE deleted successfully\n" + ] + } + ], + "execution_count": 9 + }, { "metadata": {}, "cell_type": "code", "outputs": [], "execution_count": null, "source": "", - "id": "857e74f1803a5361" + "id": "37795ddfde18e0f5" } ], "metadata": { diff --git a/src/eo_processing/config/settings.py b/src/eo_processing/config/settings.py index cc7de74..34d6612 100644 --- a/src/eo_processing/config/settings.py +++ b/src/eo_processing/config/settings.py @@ -93,11 +93,11 @@ } OPENEO_EXTRACT_CDSE_JOB_OPTIONS: Dict[str, Union[str, int, float, List]] = { - "driver-memory": "4G", - "driver-memoryOverhead": "2G", + "driver-memory": "3G", + "driver-memoryOverhead": "1G", "driver-cores": 1, - "executor-memory": "2G", - "executor-memoryOverhead": "1G", + "executor-memory": "4G", + "executor-memoryOverhead": "2G", "python-memory": "disable", "executor-cores": 1, "max-executors": 20, diff --git a/src/eo_processing/resources/udf_max_occurence_hierarchical_merger.py b/src/eo_processing/resources/udf_max_occurence_hierarchical_merger.py index 3348772..248f3df 100644 --- a/src/eo_processing/resources/udf_max_occurence_hierarchical_merger.py +++ b/src/eo_processing/resources/udf_max_occurence_hierarchical_merger.py @@ -19,17 +19,33 @@ def apply_metadata(metadata: CubeMetadata, context:Dict) -> CubeMetadata: def _select_highest_prob_class(cube: xr.DataArray, raster_codes) -> xr.DataArray: """ Select per model the highest probability of occurrence class - :param cube: data cube with probabilities for all classes of three levels + :param cube: data cube with probabilities for all classes per model (level) :param raster_codes: dataframe with raster code values - :return: data cube with highest probably of occurrence class per model (level) + :return: data cube with raster value of highest occurrence probabilities per pixel for each model """ # Create nodata mask before filling with 0 nodata_mask = cube.isnull().all(dim="bands") - # Identify the band with the highest probability for each pixel - cube= cube.fillna(0) #make sure argmax is not returning all slice N/A + # check + if cube.sizes["bands"] != len(raster_codes): + raise ValueError( + f"Expected one raster code per band, got " + f"{cube.sizes['bands']} bands and {len(raster_codes)} codes." + ) + + # we make sure we have no value 255 since that is officially the nodata value of PROBA results + # so set 255 to nan + cube = cube.where(cube != 255, np.nan) + + # make sure argmax is not returning all slice N/A + cube= cube.fillna(0) + + # All-zero pixels have no valid class; argmax would incorrectly select band 0. + nodata_mask = nodata_mask | (cube == 0).all(dim="bands") + try: - max_band = cube.dropna(dim="bands", how='all').argmax(axis=0) # Index of max value, OpenEO need bands ? + # Identify the band with the highest probability for each pixel + max_band = cube.argmax(axis=cube.get_axis_num("bands")) # Index of max value, OpenEO need bands ? except Exception as e: inspect(message=f"EXCEPTION {e} in argmax for {raster_codes}") @@ -64,14 +80,18 @@ def _merge_hierarchical(cube: xr.DataArray, df_high_prob) -> xr.DataArray: df_l2 = df_high_prob[(df_high_prob.level == '2')] if not df_l2.empty: - # now we run over each of this Level 2 raster to imprint into Level1 + # now we run over each Level 2 model cube (index in dataframe) to imprint into Level1 for row in df_l2.itertuples(): + # since the Index of the dataframe represents the band number in the cube aImprint = cube.isel(bands=[row.Index]) - nodata = [0 , -1] + nodata = [0 , -1, 255] # get the Level 1 habitat code from the level 2 data - lsub = np.unique(aImprint).tolist() - lsub = [x for x in lsub if x != nodata and not np.isnan(x)] + lsub = [x for x in np.unique(aImprint).tolist() if x not in nodata and not np.isnan(x)] + if len(lsub) == 0: + inspect(message=f"no data in level 2 for model {row.model}") + continue + # for valid members we can generate the level 1 class by flooring the level 2 class to the nearest 10000 lsub = [*set([int(np.floor(x / 10000) * 10000) for x in lsub])] if len(lsub) != 1: @@ -90,14 +110,16 @@ def _merge_hierarchical(cube: xr.DataArray, df_high_prob) -> xr.DataArray: df_l3 = df_high_prob[(df_high_prob.level == '3')] if not df_l3.empty: - # now we run over each of this Level 3 raster to imprint into Level2 + # now we run over each Level 3 model cube (index in dataframe) to imprint into Level2 for row in df_l3.itertuples(): aImprint = cube.isel(bands=[row.Index]) - nodata = [0 , -1] + nodata = [0 , -1, 255] # get the Level 2 habitat code from the level 3 data - lsub = np.unique(aImprint).tolist() - lsub = [x for x in lsub if x != nodata and not np.isnan(x)] + lsub = [x for x in np.unique(aImprint).tolist() if x not in nodata and not np.isnan(x)] + if len(lsub) == 0: + inspect(message=f"no data in level 3 for model {row.model}") + continue lsub = [*set([int(np.floor(x / 100) * 100) for x in lsub])] if len(lsub) != 1: @@ -138,7 +160,7 @@ def parse_prob_classes_fromStac(band_names: List[str]) -> pd.DataFrame: level, class_name, habitat, raster_code = match.groups() band_info.append((band_nr, level, class_name, habitat, int(raster_code))) else: - print('skipping {}'.format(band_name)) + inspect(message='skipping parsing of band names for band {}. Band name has wrong format.'.format(band_name)) # Create DataFrame df = pd.DataFrame(band_info, columns=["band_nr", "level", "model", "habitat", "raster_code"]) @@ -146,50 +168,76 @@ def parse_prob_classes_fromStac(band_names: List[str]) -> pd.DataFrame: def apply_datacube(cube: xr.DataArray, context:Dict) -> xr.DataArray: inspect(message=f"xarray dims {cube.dims}") - # important for filling by level - max_cube_initialized = False - ### get the list of classes as output from inference run # use returned metadata to build up the class dictionary input_band_names = cube.indexes["bands"].values inspect(message=f"input cube band names ({len(input_band_names)}): {input_band_names}") df = parse_prob_classes_fromStac(input_band_names) - inspect(message=f"## context parameters") + inspect(message=f"## parsed band names into dataframe") inspect(message=f"{df}") ### Determine first the highest probability per model (leveled) inspect(message=f"## determine highest probability per model/level") # read in the selected band names from the raster stack (per level and class) - for (level, class_name), group in df.groupby(["level", "model"]): + max_cube_initialized = False + high_prob_records = [] + for (level, class_name), group in df.groupby(["level", "model"]): + # get band_indices and raster_codes for this model band_indices = group["band_nr"].values - 1 # Convert to 0-based index raster_codes = group["raster_code"].values - # select the bands for this class + + # select the bands for this model (its classes) subset_cube = cube.isel(bands=list(band_indices)) # get a 2D array with the winning prob of the habitat classes max_probability = _select_highest_prob_class(subset_cube, raster_codes) + # Do not add groups that have no valid result anywhere. + if bool(max_probability.isnull().all().item()): + inspect( + message=( + f"Skipping level {level}, model {class_name}: " + "result contains only nodata." + ) + ) + continue + if not max_cube_initialized: - # Iniitialize the output cube only on the first iteration + # Iniitialize the output cube only on the first iteration which is mainly typology level 1 max_cube = max_probability max_cube_initialized = True else: # Append the result in the output cube max_cube = xr.concat([max_cube, max_probability], dim="bands") + # Append metadata only when the corresponding cube was appended. + high_prob_records.append( + {"level": level, "model": class_name, "count": len(group), } + ) + + if not max_cube_initialized: + raise RuntimeError( + "No valid model results found: all groups contain only nodata." + ) + # check if xaaray has band dimension - if only level1 is processed that can happen - if not "bands" in max_cube.dims: + if "bands" not in max_cube.dims: max_cube = max_cube.expand_dims(dim={"bands": [0]}) # create new dataframe with bands from highest_prob as LUT - df_high_prob = pd.DataFrame({'count':df.groupby(["level","model"]).size()}).reset_index(level=["model","level"]) + df_high_prob = pd.DataFrame(high_prob_records, columns=["level", "model", "count"]) + inspect(message=f"## dataframe with highest probabilities") + inspect(message=f"{df_high_prob}") + + if df_high_prob.iloc[0]["level"] != "1": + raise RuntimeError( + "Cannot merge hierarchy because no valid Level 1 result is available." + ) ### Merge highest probability classes in hierarchical way inspect(message=f"## merge highest probabilities") max_cube = _merge_hierarchical(max_cube, df_high_prob) - # make sure the output Xarray has set correct dtype (not float64) - max_cube = max_cube.astype("uint32") return max_cube \ No newline at end of file diff --git a/src/eo_processing/resources/udf_merge_proba_cubes.py b/src/eo_processing/resources/udf_merge_proba_cubes.py new file mode 100644 index 0000000..e4b1d7f --- /dev/null +++ b/src/eo_processing/resources/udf_merge_proba_cubes.py @@ -0,0 +1,73 @@ +import pandas as pd +import numpy as np +import xarray as xr +import re +from typing import Dict, List +from openeo.udf import inspect +from openeo.metadata import CubeMetadata + +def apply_metadata(metadata: CubeMetadata, context:Dict) -> CubeMetadata: + """ Rename the bands by using openeo apply_metadata function + :param metadata: Metadata of the input data cube + :param context: Context of the UDF + :return: renamed labels + """ + band_names = context.get('band_names') + return metadata.rename_labels(dimension="bands", target=band_names) + +def apply_datacube(cube: xr.DataArray, context:Dict) -> xr.DataArray: + """ merge proba cubes + + :param cube: data cube + :param context: dictionary to provide external data - key class_mapping is used to provide external re-mapping dict + :return: unique proba data cube + """ + sorted_band_names = context.get('band_names') + number_added_cubes = context.get('number_add_cubes') + + if not sorted_band_names: + raise ValueError("sorted_band_names must not be empty") + + inspect(message=f"merge proba cubes. add {number_added_cubes} cubes to base cube.") + inspect(message=f"output_band_names ({len(sorted_band_names)}): {sorted_band_names}") + + # first we create a new cube with the correct band names and all filled zeros + proba_cube = ( + xr.full_like( + cube.isel(bands=0, drop=True), + fill_value=np.nan, + dtype=np.float32, + ) + .expand_dims(bands=sorted_band_names) + .transpose(*cube.dims) + .copy() + ) + + # creating the list of band names of the cubes we need to split + input_band_names = cube.indexes["bands"].values + inspect(message=f"input cube band names ({len(input_band_names)}): {input_band_names}") + + # loop over the input bands.... inspect to which corresponding band it belongs and imprint in proba_cube + suffix_pattern = re.compile(r"_cube\d+$") + sorted_band_names_set = set(sorted_band_names) + + for cube_band_name in input_band_names: + # strip trailing "_cube{n}" suffix if present, then match exactly + match_band_name = suffix_pattern.sub("", cube_band_name) + + if match_band_name in sorted_band_names_set: + cube_values = cube.loc[dict(bands=cube_band_name)] + mask = np.logical_and(cube_values != 255, ~np.isnan(cube_values)) + proba_cube.loc[dict(bands=match_band_name)] = xr.where( + mask, cube_values, proba_cube.loc[dict(bands=match_band_name)] + ) + else: + inspect(message=f"no corresponding output band name found for input band name {cube_band_name}") + + # check + unfilled_mask = proba_cube.isnull().all(dim=[d for d in proba_cube.dims if d != "bands"]) + unfilled_bands = proba_cube["bands"].where(unfilled_mask, drop=True).values.tolist() + if unfilled_bands: + inspect(message=f"WARNING: {len(unfilled_bands)} output band(s) never filled: {unfilled_bands}") + + return proba_cube \ No newline at end of file diff --git a/src/eo_processing/utils/stac_helper.py b/src/eo_processing/utils/stac_helper.py index 3e43ba0..e8a032d 100644 --- a/src/eo_processing/utils/stac_helper.py +++ b/src/eo_processing/utils/stac_helper.py @@ -8,6 +8,47 @@ import pandas as pd import re + +TILE_PATTERNS = [ + r"E\d{3}N\d{3}", # EUNIS2021plus, e.g. E426N410 + r"\d{2}[^\W\d_][A-Z]{2}\d{2}", # IUCNGET, e.g. 48πXH44, 17λQC33 +] + +def extract_tile_id(basename: str) -> str: + for tile_pattern in TILE_PATTERNS: + match = re.search(rf"_(?P{tile_pattern})_", basename) + if match: + return match.group("tileID") + + raise ValueError(f"No supported tileID pattern found in basename: {basename}") + +def extract_real_tileid(tileid_variant): + """ + Extract the real tileID by removing trailing letter suffixes. + Ensures the tileID ends with a number. + Example: '48πXH34a' -> '48πXH34' + """ + # Remove any trailing letters after the last digit + return re.sub(r'[a-zA-Z]+$', '', tileid_variant) + +def combine_band_names(band_name_lists: List[List[str]]) -> List[str]: + """ + Combines multiple lists of band names into a single, sorted list. + + This function takes a list of band name collections, merges them into one + set to ensure uniqueness, and then sorts the combined names alphabetically. + + :param band_name_lists: list[list[str]] + A list containing multiple lists of band names. Each inner list represents + a collection of band names. + :return: list[str] + A sorted list of unique band names. + """ + combined = set() + for names in band_name_lists: + combined.update(names) + return sorted(combined) + def get_stac_collection_url(collection_id: str, catalog_url: str = "https://catalogue.weed.apex.esa.int/") -> str: """ Fetches the URL of a STAC (SpatioTemporal Asset Catalog) collection based on the provided @@ -102,29 +143,30 @@ def query_proba_results(df_AOI: gpd.GeoDataFrame, collection_id:str, processing_ stac_url:str = 'https://catalogue.weed.apex.esa.int', info_debug: bool = True, postprocess: bool = True) -> gpd.GeoDataFrame: """ - Queries and retrieves PROBA results intersecting with a given Area of Interest (AOI) from a STAC catalog. - - This function searches for PROBA results within the bounding box of the provided AOI and retrieves metadata - and asset information from the specified STAC catalog. The retrieved data is reformatted and returned - as a GeoDataFrame containing details about the intersecting PROBA tiles. - - Arguments: - :param df_AOI: A GeoDataFrame representing the Area of Interest (AOI). The GeoDataFrame should contain geometry - information and a coordinate reference system. If the coordinate reference system is not EPSG:4326, - the function will reproject it to EPSG:4326. - :param collection_id: A string specifying the collection identifier to search within the STAC catalog. - :param processing_year: An integer specifying the year for which to retrieve PROBA results. - :param stac_url: A string specifying the URL of the STAC catalog to query. Defaults to 'https://catalogue.weed.apex.esa.int'. - :param info_debug: A boolean flag to enable or disable debug logging information. Defaults to True. - :param postprocess: A boolean flag to determine whether to postprocess the retrieved data. Defaults to True. - - Returns: - A GeoDataFrame containing metadata and details about the PROBA results that intersect with the AOI. The GeoDataFrame - includes additional columns extracted from the metadata and asset information of the intersecting tiles, such as - datetime, bounding box, tile ID, and others. - - Raises: - ValueError: If no intersecting PROBA tiles are found in the STAC catalog for the specified collection. + Queries results from a STAC (SpatioTemporal Asset Catalog) service and + processes the output as per user requirements. + + This function retrieves geometries within a specified area of interest (AOI), + filters them based on the provided collection ID and processing year, and + optionally post-processes the results for further aggregation and refinement. + + :param df_AOI: A GeoDataFrame (gpd.GeoDataFrame) representing the area of interest (AOI). + Expected to have its coordinate system as EPSG:4326 or convertible to it. + :param collection_id: A string representing the STAC collection ID to query. + :param processing_year: An integer specifying the processing year for filtering results. + :param stac_url: A string representing the URL of the STAC service. + Default is 'https://catalogue.weed.apex.esa.int'. + :param info_debug: A boolean flag to enable or disable debug messages during processing. + Default is True. + :param postprocess: A boolean flag to activate post-processing of PROBA tiles. + Default is True. + + :return: A GeoDataFrame (gpd.GeoDataFrame) containing the results with the following attributes: + - 'tileID': The tile identifier. + - 'band_names': A unique, sorted list of band names for the tile. + - 'item_url': A list of item URLs related to the tile. + + :raises ValueError: If no intersecting PROBA tiles are found in the specified STAC collection. """ if info_debug: print(f"get_modelID_asset_geometry_from_STAC") # convert AOI into BBOX in 4326 @@ -144,83 +186,53 @@ def query_proba_results(df_AOI: gpd.GeoDataFrame, collection_id:str, processing_ search = client.search( collections=[collection_id], bbox=bbox_4326, - fields=["properties", "assets.openEO.href"], + fields=["properties", "assets", 'links'], ) results = [] for item in search.items_as_dicts(): - results.append( - [item['properties']['datetime'], item['properties']['proj:bbox'], item['properties']['proj:shape'], - item['properties']['proj:code'], item['assets']['openEO']['href']]) - - # build dataframe - df_result = pd.DataFrame(results, columns=['datetime', 'file_bbox', 'file_shape', 'file_epsg', 'file_url']) + item_result = [item['properties']['datetime'], item['assets']['openEO']['href'], + [band['name'] for band in item['assets']['openEO']['bands']], item['links'][0]['href'], + item['properties']['proj:shape']] + results.append(item_result) + # build Pandas dataframe + df_result = pd.DataFrame(results, columns=['datetime', 'file_url', 'band_names', 'item_url', 'file_shape']) if info_debug: print(f"- found {len(df_result)} intersecting PROBA tiles") # check if there are any results if df_result.empty: - ValueError(f"No intersecting PROBA tiles found in the STAC ({collection_id}).") + raise ValueError(f"No intersecting PROBA tiles found in the STAC ({collection_id}).") if postprocess: if info_debug: print(f"- postprocessing PROBA tiles") # split out from file_url important parts (file_name, tile_id, etc) + # split the tileID and processing year out of the file names df_result['basename'] = df_result['file_url'].apply(lambda x: os.path.basename(x)) - df_result[['project_typology', 'type', 'processing_year', 'tileID', 'model_short', 'inference_run_version', - 'procesisng_start']] = df_result['basename'].str.split('_', expand=True) - df_result['processing_year'] = df_result['processing_year'].str[-4:].astype(int) + df_result["tileID"] = df_result["basename"].apply(extract_tile_id) + + df_result["processing_year"] = ( + df_result["basename"] + .str.extract(r"year(\d{4})", expand=False) + .astype(int) + ) # first limit results to processing year df_result = df_result[df_result['processing_year'] == processing_year] # check if we have tiles smaller than our standard 20x20km grid - yes then make sure tile name is correct - def extract_real_tileid(tileid_variant): - """ - Extract the real tileID by removing trailing letter suffixes. - Ensures the tileID ends with a number. - Example: '48πXH34a' -> '48πXH34' - """ - # Remove any trailing letters after the last digit - return re.sub(r'[a-zA-Z]+$', '', tileid_variant) - for idx, row in df_result.iterrows(): if row.file_shape != [2000, 2000]: df_result.at[idx, 'tileID'] = extract_real_tileid(row.tileID) - # now we can filter out spatial duplicates for same used modelID_short name - # NOTE: that assumes that NEVER different inference runs of smae modelID were saved in same STAC catalog - df_result = df_result.drop_duplicates(subset=['tileID', 'model_short'], keep='first') - - # last step. we have to prepare the output file_name. - # Step 1: Check if we have duplicate tileIDs with different model_short values - duplicate_tiles = df_result.groupby('tileID')['model_short'].apply(lambda x: list(x.unique())).to_dict() - tiles_with_multiple_models = {k: v for k, v in duplicate_tiles.items() if len(v) > 1} - - if tiles_with_multiple_models: - if info_debug: print(f" -- Found {len(tiles_with_multiple_models)} tiles with multiple model_short values") - - # Step 2: For tiles with multiple models, condense model_short names - # Create a condensed model_short by combining unique values - for tile, models in tiles_with_multiple_models.items(): - # Sort models to ensure consistent naming - strata = [x.split('-')[0] for x in models] - condensed_name = '-'.join(sorted(strata)) + '-' + '-'.join(models[0].split('-')[1:]) - # Update all rows for this tileID with the condensed name - df_result.loc[df_result['tileID'] == tile, 'model_short'] = condensed_name - - # Step 3: Now remove duplicate tileIDs (keeping first occurrence) - df_result = df_result.drop_duplicates(subset=['tileID'], keep='first') - else: - if info_debug: print(" -- No duplicate tileIDs found with different model_short values") - # Still remove any exact duplicates - df_result = df_result.drop_duplicates(subset=['tileID'], keep='first') - - # Step 4: Create the file_prefix column properly - df_result['file_prefix'] = df_result.apply( - lambda - row: f"{row['project_typology']}_mece-cube_year{row['processing_year']}_{row['tileID']}_{row['model_short']}_{row['inference_run_version']}", - axis=1 - ) + ### for results with the same tileID (group them): combine band_names into one + ### unique, sorted list, and merge item_url into a list and for the rest take first entry + agg_dict = { + col: (list if col in ['band_names', 'item_url'] else 'first') + for col in df_result.columns if col != 'tileID' + } + df_result = df_result.groupby('tileID', as_index=False).agg(agg_dict) + # filter to final needed - df_result = df_result[['tileID', 'file_prefix']] + df_result = df_result[['tileID','band_names', 'item_url']] return df_result