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