diff --git a/geoprob_pipe/app_object.py b/geoprob_pipe/app_object.py index 549c6692..755987db 100644 --- a/geoprob_pipe/app_object.py +++ b/geoprob_pipe/app_object.py @@ -5,8 +5,7 @@ from datetime import datetime from typing import TYPE_CHECKING, Optional, List import pandas as pd - -from geoprob_pipe.cmd_app.parameter_input.expand_input_tables import run_expand_input_tables +from geoprob_pipe.cmd_app.parameter_input.expand.expand_input_tables import run_expand_input_tables from geoprob_pipe.cmd_app.parameter_input.input_parameter_figures import InputParameterFigures try: diff --git a/geoprob_pipe/calculations/systems/base_objects/base_system_build.py b/geoprob_pipe/calculations/systems/base_objects/base_system_build.py index 0a8c1f8e..64365cad 100644 --- a/geoprob_pipe/calculations/systems/base_objects/base_system_build.py +++ b/geoprob_pipe/calculations/systems/base_objects/base_system_build.py @@ -3,7 +3,7 @@ from typing import List, Tuple from pandas import DataFrame, Series import sqlite3 -from geoprob_pipe.cmd_app.parameter_input.expand_input_tables import run_expand_input_tables +from geoprob_pipe.cmd_app.parameter_input.expand.expand_input_tables import run_expand_input_tables def _gather_variable_correlations(geopackage_filepath: str diff --git a/geoprob_pipe/cmd_app/parameter_input/added_input_parameters.py b/geoprob_pipe/cmd_app/parameter_input/added_input_parameters.py index f5b556a9..0ec86458 100644 --- a/geoprob_pipe/cmd_app/parameter_input/added_input_parameters.py +++ b/geoprob_pipe/cmd_app/parameter_input/added_input_parameters.py @@ -2,7 +2,7 @@ from InquirerPy import inquirer import sqlite3 from datetime import datetime -from geoprob_pipe.cmd_app.parameter_input.expand_input_tables import run_expand_input_tables +from geoprob_pipe.cmd_app.parameter_input.expand.expand_input_tables import run_expand_input_tables from geoprob_pipe.cmd_app.parameter_input.initiate_input_excel_tables import initiate_input_excel_tables from geoprob_pipe.cmd_app.parameter_input.input_parameter_figures import InputParameterFigures from geoprob_pipe.cmd_app.parameter_input.export_input_parameter_excel import export_input_parameter_tables diff --git a/geoprob_pipe/cmd_app/parameter_input/expand/__init__.py b/geoprob_pipe/cmd_app/parameter_input/expand/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/geoprob_pipe/cmd_app/parameter_input/expand/expand_dataframes.py b/geoprob_pipe/cmd_app/parameter_input/expand/expand_dataframes.py new file mode 100644 index 00000000..5b3e4b47 --- /dev/null +++ b/geoprob_pipe/cmd_app/parameter_input/expand/expand_dataframes.py @@ -0,0 +1,256 @@ +from typing import Dict, List +from pandas import DataFrame, Series +import copy +import sqlite3 +from geoprob_pipe.calculations.systems.mappers.initial_input import INITIAL_INPUT_MAPPER + + +def _gather_required_input_parameters(geopackage_filepath: str) -> List[str]: + + conn = sqlite3.connect(geopackage_filepath) + cursor = conn.cursor() + cursor.execute(""" + SELECT geoprob_pipe_metadata."values" + FROM geoprob_pipe_metadata + WHERE metadata_type='geohydrologisch_model'; + """) + result = cursor.fetchone() + if not result: + raise ValueError + model_string = result[0] + conn.close() + df_dummy_data = DataFrame(INITIAL_INPUT_MAPPER[model_string]["input"]) + + _ = df_dummy_data.sort_values(by=["name"]) + df_dummy_data = df_dummy_data.sort_values(by=["name"]) + return df_dummy_data["name"].unique().tolist() + + +def merge_into_df(df: DataFrame, df_gather: DataFrame, on_cols: list[str]) -> DataFrame: + """If row (unique combination/calculation) does not have a parameter_input yet, + then combine first will gather the new value, if provided. + """ + + # Combine two dataframe to ensure index is equal + df = df.merge( + df_gather, + on=on_cols, + how="left", + suffixes=("", "_new"), + ) + + # Get value out of _new-col, if no value exists yet + df["parameter_input"] = df["parameter_input"].combine_first( + df["parameter_input_new"] + ) + + # Remove _new-col (not needed anymore) + df.drop(columns=["parameter_input_new"], inplace=True) + + return df + + +def _1a_excel_uittredepunt_en_scenario( + df: DataFrame, df_parameter_invoer_combined: DataFrame, parameter_name: str +) -> DataFrame: + + mask = ( + (df_parameter_invoer_combined["parameter"] == parameter_name) + & (df_parameter_invoer_combined["scope"] == "uittredepunt") + & (df_parameter_invoer_combined["ondergrondscenario_naam"].notna()) + ) + + df_gather = ( + df_parameter_invoer_combined.loc[ + mask, + ["scope_referentie", "ondergrondscenario_naam", "parameter_input"], + ] + .rename( + columns={ + "scope_referentie": "uittredepunt_id", + "ondergrondscenario_naam": "naam", + } + ) + .drop_duplicates(["uittredepunt_id", "naam"]) + ) + + return merge_into_df(df, df_gather, ["uittredepunt_id", "naam"]) + + +def _1b_gis_uittredepunt_en_scenario( + df: DataFrame, df_parameter_invoer_combined: DataFrame, parameter_name: str +) -> DataFrame: + + mask = ( + (df_parameter_invoer_combined["parameter"] == parameter_name) + & (df_parameter_invoer_combined["scope"] == "gis_uittredepunt") + & (df_parameter_invoer_combined["ondergrondscenario_naam"].notna()) + ) + + df_gather = ( + df_parameter_invoer_combined.loc[ + mask, + ["scope_referentie", "ondergrondscenario_naam", "parameter_input"], + ] + .rename( + columns={ + "scope_referentie": "uittredepunt_id", + "ondergrondscenario_naam": "naam", + } + ) + .drop_duplicates(["uittredepunt_id", "naam"]) + ) + + return merge_into_df(df, df_gather, ["uittredepunt_id", "naam"]) + + +def _2a_excel_uittredepunt( + df: DataFrame, df_parameter_invoer_combined: DataFrame, parameter_name: str +) -> DataFrame: + + mask = ( + (df_parameter_invoer_combined["parameter"] == parameter_name) + & (df_parameter_invoer_combined["scope"] == "uittredepunt") + & (df_parameter_invoer_combined["ondergrondscenario_naam"].isna()) + & (df_parameter_invoer_combined["parameter_input"] != {}) + ) + + df_gather = ( + df_parameter_invoer_combined.loc[ + mask, + ["scope_referentie", "parameter_input"], + ] + .rename(columns={"scope_referentie": "uittredepunt_id"}) + .drop_duplicates(["uittredepunt_id"]) + ) + + return merge_into_df(df, df_gather, ["uittredepunt_id"]) + + +def _2b_gis_uittredepunt( + df: DataFrame, df_parameter_invoer_combined: DataFrame, parameter_name: str +) -> DataFrame: + + mask = ( + (df_parameter_invoer_combined["parameter"] == parameter_name) + & (df_parameter_invoer_combined["scope"] == "gis_uittredepunt") + & (df_parameter_invoer_combined["parameter_input"] != {}) + ) + + df_gather = ( + df_parameter_invoer_combined.loc[ + mask, + ["scope_referentie", "parameter_input"], + ] + .rename(columns={"scope_referentie": "uittredepunt_id"}) + .drop_duplicates(["uittredepunt_id"]) + ) + + return merge_into_df(df, df_gather, ["uittredepunt_id"]) + + +def _3_excel_vak_en_scenario( + df: DataFrame, df_parameter_invoer_combined: DataFrame, parameter_name: str +) -> DataFrame: + + mask = ( + (df_parameter_invoer_combined["parameter"] == parameter_name) + & (df_parameter_invoer_combined["scope"] == "vak") + & (df_parameter_invoer_combined["ondergrondscenario_naam"].notna()) + ) + + df_gather = ( + df_parameter_invoer_combined.loc[ + mask, + ["scope_referentie", "ondergrondscenario_naam", "parameter_input"], + ] + .rename( + columns={ + "scope_referentie": "vak_id", + "ondergrondscenario_naam": "naam", + } + ) + .drop_duplicates(["vak_id", "naam"]) + ) + + return merge_into_df(df, df_gather, ["vak_id", "naam"]) + + +def _4_excel_vak( + df: DataFrame, df_parameter_invoer_combined: DataFrame, parameter_name: str +) -> DataFrame: + + mask = ( + (df_parameter_invoer_combined["parameter"] == parameter_name) + & (df_parameter_invoer_combined["scope"] == "vak") + & (df_parameter_invoer_combined["ondergrondscenario_naam"].isna()) + ) + + df_gather = ( + df_parameter_invoer_combined.loc[ + mask, + ["scope_referentie", "parameter_input"], + ] + .rename(columns={"scope_referentie": "vak_id"}) + .drop_duplicates(["vak_id"]) + ) + + return merge_into_df(df, df_gather, ["vak_id"]) + + +def _5_excel_traject( + df: DataFrame, df_parameter_invoer_combined: DataFrame, parameter_name: str +) -> DataFrame: + + df_gather = df_parameter_invoer_combined.loc[ + (df_parameter_invoer_combined["parameter"] == parameter_name) + & (df_parameter_invoer_combined["scope"] == "traject") + ] + + if df_gather.empty: + return df + + value = df_gather.iloc[0]["parameter_input"] + mask_na = df["parameter_input"].isna() + df.loc[mask_na, "parameter_input"] = [ + copy.deepcopy(value) for _ in range(mask_na.sum()) + ] + + return df + + +def expand_input_dataframes( + df_parameter_invoer_combined: DataFrame, + df_identifiers: DataFrame, + geopackage_filepath: str, +) -> Dict[str, DataFrame]: + + required_input_parameters = _gather_required_input_parameters( + geopackage_filepath=geopackage_filepath + ) + + collection_of_dfs: Dict[str, DataFrame] = {} + + for parameter_name in required_input_parameters: + df = df_identifiers.copy() + df["parameter_input"] = Series([None] * len(df), dtype="object") + + df = _1a_excel_uittredepunt_en_scenario( + df, df_parameter_invoer_combined, parameter_name + ) + df = _1b_gis_uittredepunt_en_scenario( + df, df_parameter_invoer_combined, parameter_name + ) + + df = _2a_excel_uittredepunt(df, df_parameter_invoer_combined, parameter_name) + df = _2b_gis_uittredepunt(df, df_parameter_invoer_combined, parameter_name) + + df = _3_excel_vak_en_scenario(df, df_parameter_invoer_combined, parameter_name) + + df = _4_excel_vak(df, df_parameter_invoer_combined, parameter_name) + + df = _5_excel_traject(df, df_parameter_invoer_combined, parameter_name) + + collection_of_dfs[parameter_name] = df + + return collection_of_dfs diff --git a/geoprob_pipe/cmd_app/parameter_input/expand/expand_input_tables.py b/geoprob_pipe/cmd_app/parameter_input/expand/expand_input_tables.py new file mode 100644 index 00000000..734e16e4 --- /dev/null +++ b/geoprob_pipe/cmd_app/parameter_input/expand/expand_input_tables.py @@ -0,0 +1,358 @@ +import os +import sqlite3 +from typing import Dict, List, Optional + +import numpy as np +import scipy.stats as sct +from geopandas import GeoDataFrame, read_file +from pandas import DataFrame, concat, notna, read_csv, read_sql, set_option +from probabilistic_library import FragilityValue + +from geoprob_pipe.cmd_app.parameter_input.input_parameter_tables import ( + InputParameterTables, +) +from geoprob_pipe.cmd_app.validation.validate_expanded_table import ( + validate_expand_tables, +) + +from .expand_dataframes import expand_input_dataframes + + +def _combine_parameter_invoer_sources(tables: InputParameterTables) -> DataFrame: + """Combineert de geo-gerefereerde parameter invoer met de handmatige invoer die oorspronkelijk uit de Excel kwam. + Zodoende kan vanuit één dataframe de invoer geëxplodeerd worden naar invoer per uittredepunt.""" + + # Gather raw data + df_gis_join_parameter_invoer = tables.df_gis_join_parameter_invoer + df_gis_join_parameter_invoer["scope"] = "gis_uittredepunt" + df_parameter_invoer = tables.df_parameter_invoer + + # Concatenate + import warnings + + with warnings.catch_warnings(): + warnings.simplefilter(action="ignore", category=FutureWarning) + df_parameter_invoer_combined = concat( + [df_gis_join_parameter_invoer, df_parameter_invoer], ignore_index=True + ) + + return df_parameter_invoer_combined + + +def _gather_hrd_frag_line_from_geopackage(ref: str, geopackage_filepath: str): + # Read database + conn = sqlite3.connect(geopackage_filepath) + df_frag_line: DataFrame = read_sql( + sql=f"SELECT * FROM fragility_values_invoer_hrd WHERE fragility_values_ref = '{ref}' AND kans < 1.0;", + con=conn, + ) + conn.close() + + # Filter beta > 8 (probabilistic library cannot work with that in 26.1.1) + df_frag_line["beta"] = -sct.norm.ppf(df_frag_line["kans"]) + df_frag_line = df_frag_line[df_frag_line["beta"] < 8.0].copy(deep=True) + df_frag_line = df_frag_line.drop(columns=["beta"]) + + # Validate + assert df_frag_line.__len__() > 2 + # It should be validated beforehand that all added fragility lines are at least 3 points. So if this assert triggers + # something should be improved earlier in validation. + + # Construct Fragility Values + df_frag_line = df_frag_line.sort_values(by=["waarde"]) + frag_points = [] + for index, row in df_frag_line.iterrows(): + fc = FragilityValue() + fc.x = row["waarde"] + fc.probability_of_failure = row["kans"] + frag_points.append(fc) + + return frag_points + + +def _gather_other_freq_line(ref: str, tables: InputParameterTables): + """Gather a freq line that is not from the HRD-database, and not from a .csv-file. It originates from the Excel, + either directly from the Excel, or imported into the 'fragility_values_invoer_hrd'-database table. + + :param ref: + :param tables: + :return: + """ + + df_data = tables.df_fragility_values_invoer + df_frag_line = df_data[df_data["fragility_values_ref"] == ref].copy(deep=True) + + # Filter beta > 8 (probabilistic library cannot work with that in 26.1.1) + df_frag_line["beta"] = -sct.norm.ppf(df_frag_line["kans"]) + df_frag_line = df_frag_line[df_frag_line["beta"] < 8.0].copy(deep=True) + df_frag_line = df_frag_line.drop(columns=["beta"]) + + # Validate + assert df_frag_line.__len__() > 2 + # It should be validated beforehand that all added fragility lines are at least 3 points. So if this assert triggers + # something should be improved earlier in validation. + + # Construct Fragility Values + df_frag_line = df_frag_line.sort_values(by=["waarde"]) + frag_points = [] + for index, row in df_frag_line.iterrows(): + fc = FragilityValue() + fc.x = row["waarde"] + fc.probability_of_failure = row["kans"] + frag_points.append(fc) + + return frag_points + + +def _gather_frag_line_from_csv(csv_file_name: str, geopackage_filepath: str): + + # Read csv-file + csv_dir = os.path.join(os.path.dirname(geopackage_filepath), "frag_csv_files") + os.makedirs(csv_dir, exist_ok=True) + path_to_csv = os.path.join(csv_dir, csv_file_name) + if not os.path.exists(path_to_csv): + raise FileNotFoundError( + f"CSV-file with fragility curve not found for reference '{csv_file_name}'. Please make sure to place your " + f"csv-files at the following location: \n{csv_dir}" + ) + df_frag_line = read_csv(path_to_csv, sep=",") + + # Validate + assert df_frag_line.__len__() > 2 + # It should be validated beforehand that all added fragility lines are at least 3 points. So if this assert triggers + # something should be improved earlier in validation. + + # Construct Fragility Values + df_frag_line = df_frag_line.sort_values(by=["waarde"]) + frag_points = [] + for index, row in df_frag_line.iterrows(): + fc = FragilityValue() + fc.x = row["waarde"] + fc.probability_of_failure = row["kans"] + frag_points.append(fc) + + return frag_points + + +def _collect_fragility_values( + tables: InputParameterTables, + fragility_refs: List[str], + geopackage_filepath: str, +) -> DataFrame: + """Collects the fragility values from the different sources. The sources are (a) the HRD-file, (b) the Excel input + file, and (c) the csv folder. + + :param tables: + :param fragility_refs: + :param geopackage_filepath: + :return: Returns a dataframe with columns fragility_values_ref and fragility_values. + """ + + # Non-HRD and non-csv available freqs + df_frag_invoer = tables.df_fragility_values_invoer + available_frag_invoer_refs = df_frag_invoer["fragility_values_ref"].unique() + + # Collect freqs + return_array = [] + for fragility_ref in fragility_refs: + # From .csv-file (not stored in GeoPackage) + if fragility_ref.endswith(".csv"): + return_array.append( + { + "fragility_values_ref": fragility_ref, + "fragility_values": _gather_frag_line_from_csv( + csv_file_name=fragility_ref, + geopackage_filepath=geopackage_filepath, + ), + } + ) + + # From HRD-database (previously stored in GeoPackage) + elif fragility_ref not in available_frag_invoer_refs: + return_array.append( + { + "fragility_values_ref": fragility_ref, + "fragility_values": _gather_hrd_frag_line_from_geopackage( + ref=fragility_ref, geopackage_filepath=geopackage_filepath + ), + } + ) + + # Otherwise, retrieve custom curve from Excel (previously stored in GeoPackage) + else: + return_array.append( + { + "fragility_values_ref": fragility_ref, + "fragility_values": _gather_other_freq_line( + ref=fragility_ref, tables=tables + ), + } + ) + + # Build dataframe + if return_array.__len__() == 0: + df = DataFrame(data=[], columns=["fragility_values_ref", "fragility_values"]) + return df + return DataFrame(return_array) + + +def _add_fragility_values_to_combined_parameter_invoer( + df_parameter_invoer_combined: DataFrame, + tables: InputParameterTables, + geopackage_filepath: str, + drop_ref: bool = True, +) -> DataFrame: + """Haalt uit de fragility values Excel de arrays op en vervang in de df_parameter_invoer_combined de referentie + met de daadwerkelijke fragility values.""" + + df = df_parameter_invoer_combined.copy(deep=True) + + # Replace empty values with NaN + set_option("future.no_silent_downcasting", True) + new_series = ( + df["fragility_values_ref"].replace("", np.nan).infer_objects(copy=False) + ) + df["fragility_values_ref"] = new_series + + # Gather referenced fragility value refs + fragility_refs = df["fragility_values_ref"].dropna().unique() + + # Collect fragility lines + df_frag_lines = _collect_fragility_values( + tables=tables, + fragility_refs=fragility_refs, + geopackage_filepath=geopackage_filepath, + ) + + # Attach to parameter invoer df + df = df.merge( + df_frag_lines, + left_on="fragility_values_ref", + right_on="fragility_values_ref", + how="left", + ) + + if drop_ref: + df = df.drop(columns=["fragility_values_ref"]) + + return df + + +def _collect_right_columns_combined_parameter_invoer( + df_parameter_invoer_combined: DataFrame, +) -> DataFrame: + """Parameter tabel omzetten naar juiste kolommen. Enkel per uittredepunt, scenario en parameter de + parameterinvoer.""" + + # Add parameter invoer: op uittredepunten niveau + df_parameter_invoer_combined = df_parameter_invoer_combined.drop( + columns=["bronnen", "opmerking"] + ) + parameter_description_columns = [ + "distribution_type", + "mean", + "variation", + "deviation", + "minimum", + "maximum", + "fragility_values", + ] + if "fragility_values_ref" in df_parameter_invoer_combined.columns: + parameter_description_columns.append("fragility_values_ref") + df_parameter_invoer_combined["parameter_input"] = ( + df_parameter_invoer_combined.apply( + lambda row: { + k: row[k] + for k in parameter_description_columns + if isinstance(row[k], list) or notna(row[k]) + }, + axis=1, + ) + ) + df_parameter_invoer_combined = df_parameter_invoer_combined.drop( + columns=parameter_description_columns + ) + return df_parameter_invoer_combined + + +def _construct_df_identifiers(geopackage_filepath: str, tables: InputParameterTables): + """Create identifiers table. + + The identifiers table is a table with unique rows (unique uittredepunt, vak and scenario combo). It is used as a + base to explode to input per scenario and uittredepunt.""" + + gdf_uittredepunten: GeoDataFrame = read_file( + geopackage_filepath, layer="uittredepunten" + ) + df_identifiers: DataFrame = gdf_uittredepunten[["uittredepunt_id", "vak_id"]] + df_scenario_invoer: DataFrame = tables.df_scenario_invoer + df_identifiers = df_identifiers.merge( + df_scenario_invoer, left_on="vak_id", right_on="vak_id" + ) + df_identifiers = df_identifiers.drop(columns=["kans"]) + return df_identifiers + + +def _concat_collection(collection: Dict[str, DataFrame]): + for parameter_name, df in collection.items(): + collection[parameter_name]["parameter_name"] = parameter_name + return_df = concat([df for _, df in collection.items()], ignore_index=True) + return_df = return_df.rename(columns={"naam": "ondergrondscenario_naam"}) + return return_df[ + [ + "parameter_name", + "vak_id", + "uittredepunt_id", + "ondergrondscenario_naam", + "parameter_input", + ] + ] + + +def run_expand_input_tables( + geopackage_filepath: str, + add_frag_ref: bool = False, + tables: Optional[InputParameterTables] = None, +) -> DataFrame: + """Performs logic to expand all parameter input from different levels to input on uittredepunt-level. + + :param geopackage_filepath: + :param add_frag_ref: Indien True, dan wordt voor parameter input met een distribution_type 'cdf_curve' de referentie + naar de invoer behouden. Deze is niet nodig voor de berekeningen, maar wordt wel gebruikt voor de visualisaties. + :param tables: + :return: DataFrame with columns parameter_name, vak_id, uittredepunt_id, ondergrondscenario_naam and + parameter_input. + """ + + if tables is None: + tables = InputParameterTables(geopackage_filepath=geopackage_filepath) + + # Construct df_parameter_invoer_combined + df_parameter_invoer_combined1 = _combine_parameter_invoer_sources(tables=tables) + df_parameter_invoer_combined2 = _add_fragility_values_to_combined_parameter_invoer( + df_parameter_invoer_combined=df_parameter_invoer_combined1, + tables=tables, + geopackage_filepath=geopackage_filepath, + drop_ref=add_frag_ref == False, + ) + df_parameter_invoer_combined3 = _collect_right_columns_combined_parameter_invoer( + df_parameter_invoer_combined=df_parameter_invoer_combined2 + ) + + # Expand + df_identifiers = _construct_df_identifiers( + geopackage_filepath=geopackage_filepath, tables=tables + ) + collection: Dict[str, DataFrame] = expand_input_dataframes( + df_parameter_invoer_combined=df_parameter_invoer_combined3, + df_identifiers=df_identifiers, + geopackage_filepath=geopackage_filepath, + ) + collected_input = _concat_collection(collection=collection) + # Temporary location of validation check. + validate_expand_tables( + tables=tables, + geopackage_filepath=geopackage_filepath, + df_expanded=collected_input, + ) + return collected_input diff --git a/geoprob_pipe/cmd_app/parameter_input/expand_input_tables.py b/geoprob_pipe/cmd_app/parameter_input/expand_input_tables.py deleted file mode 100644 index a16d32c5..00000000 --- a/geoprob_pipe/cmd_app/parameter_input/expand_input_tables.py +++ /dev/null @@ -1,390 +0,0 @@ -from typing import Dict, List, Optional -from pandas import DataFrame, isna, notna, concat, read_sql, read_csv, set_option -import sqlite3 -import numpy as np -import os -import scipy.stats as sct -from geopandas import GeoDataFrame, read_file -from geoprob_pipe.calculations.systems.mappers.initial_input import INITIAL_INPUT_MAPPER -from geoprob_pipe.cmd_app.parameter_input.input_parameter_tables import InputParameterTables -from probabilistic_library import FragilityValue - - -def _combine_parameter_invoer_sources(tables: InputParameterTables) -> DataFrame: - """ Combineert de geo-gerefereerde parameter invoer met de handmatige invoer die oorspronkelijk uit de Excel kwam. - Zodoende kan vanuit één dataframe de invoer geëxplodeerd worden naar invoer per uittredepunt. """ - - # Gather raw data - df_gis_join_parameter_invoer = tables.df_gis_join_parameter_invoer - df_gis_join_parameter_invoer['scope'] = 'gis_uittredepunt' - df_parameter_invoer = tables.df_parameter_invoer - - # Concatenate - import warnings - with warnings.catch_warnings(): - warnings.simplefilter(action='ignore', category=FutureWarning) - df_parameter_invoer_combined = concat( - [df_gis_join_parameter_invoer, df_parameter_invoer], ignore_index=True) - - return df_parameter_invoer_combined - - -def _gather_hrd_frag_line_from_geopackage(ref: str, geopackage_filepath: str): - # Read database - conn = sqlite3.connect(geopackage_filepath) - df_frag_line: DataFrame = read_sql( - sql=f"SELECT * FROM fragility_values_invoer_hrd WHERE fragility_values_ref = '{ref}' AND kans < 1.0;", - con=conn) - conn.close() - - # Filter beta > 8 (probabilistic library cannot work with that in 26.1.1) - df_frag_line['beta'] = -sct.norm.ppf(df_frag_line['kans']) - df_frag_line = df_frag_line[df_frag_line['beta'] < 8.0].copy(deep=True) - df_frag_line = df_frag_line.drop(columns=["beta"]) - - # Validate - assert df_frag_line.__len__() > 2 - # It should be validated beforehand that all added fragility lines are at least 3 points. So if this assert triggers - # something should be improved earlier in validation. - - # Construct Fragility Values - df_frag_line = df_frag_line.sort_values(by=["waarde"]) - frag_points = [] - for index, row in df_frag_line.iterrows(): - fc = FragilityValue() - fc.x = row["waarde"] - fc.probability_of_failure = row["kans"] - frag_points.append(fc) - - return frag_points - - -def _gather_other_freq_line(ref: str, tables: InputParameterTables): - """ Gather a freq line that is not from the HRD-database, and not from a .csv-file. It originates from the Excel, - either directly from the Excel, or imported into the 'fragility_values_invoer_hrd'-database table. - - :param ref: - :param tables: - :return: - """ - - df_data = tables.df_fragility_values_invoer - df_frag_line = df_data[df_data["fragility_values_ref"] == ref].copy(deep=True) - - # Filter beta > 8 (probabilistic library cannot work with that in 26.1.1) - df_frag_line['beta'] = -sct.norm.ppf(df_frag_line['kans']) - df_frag_line = df_frag_line[df_frag_line['beta'] < 8.0].copy(deep=True) - df_frag_line = df_frag_line.drop(columns=["beta"]) - - # Validate - assert df_frag_line.__len__() > 2 - # It should be validated beforehand that all added fragility lines are at least 3 points. So if this assert triggers - # something should be improved earlier in validation. - - # Construct Fragility Values - df_frag_line = df_frag_line.sort_values(by=["waarde"]) - frag_points = [] - for index, row in df_frag_line.iterrows(): - fc = FragilityValue() - fc.x = row["waarde"] - fc.probability_of_failure = row["kans"] - frag_points.append(fc) - - return frag_points - - -def _gather_frag_line_from_csv(csv_file_name: str, geopackage_filepath: str): - - # Read csv-file - csv_dir = os.path.join(os.path.dirname(geopackage_filepath), "frag_csv_files") - os.makedirs(csv_dir, exist_ok=True) - path_to_csv = os.path.join(csv_dir, csv_file_name) - if not os.path.exists(path_to_csv): - raise FileNotFoundError( - f"CSV-file with fragility curve not found for reference '{csv_file_name}'. Please make sure to place your " - f"csv-files at the following location: \n{csv_dir}") - df_frag_line = read_csv(path_to_csv, sep=",") - - # Validate - assert df_frag_line.__len__() > 2 - # It should be validated beforehand that all added fragility lines are at least 3 points. So if this assert triggers - # something should be improved earlier in validation. - - # Construct Fragility Values - df_frag_line = df_frag_line.sort_values(by=["waarde"]) - frag_points = [] - for index, row in df_frag_line.iterrows(): - fc = FragilityValue() - fc.x = row["waarde"] - fc.probability_of_failure = row["kans"] - frag_points.append(fc) - - return frag_points - - -def _collect_fragility_values( - tables: InputParameterTables, fragility_refs: List[str], geopackage_filepath: str, -) -> DataFrame: - """ Collects the fragility values from the different sources. The sources are (a) the HRD-file, (b) the Excel input - file, and (c) the csv folder. - - :param tables: - :param fragility_refs: - :param geopackage_filepath: - :return: Returns a dataframe with columns fragility_values_ref and fragility_values. - """ - - # Non-HRD and non-csv available freqs - df_frag_invoer = tables.df_fragility_values_invoer - available_frag_invoer_refs = df_frag_invoer['fragility_values_ref'].unique() - - # Collect freqs - return_array = [] - for fragility_ref in fragility_refs: - - # From .csv-file (not stored in GeoPackage) - if fragility_ref.endswith(".csv"): - return_array.append({ - "fragility_values_ref": fragility_ref, - "fragility_values": _gather_frag_line_from_csv( - csv_file_name=fragility_ref, geopackage_filepath=geopackage_filepath)}) - - # From HRD-database (previously stored in GeoPackage) - elif fragility_ref not in available_frag_invoer_refs: - return_array.append({ - "fragility_values_ref": fragility_ref, - "fragility_values": _gather_hrd_frag_line_from_geopackage( - ref=fragility_ref, geopackage_filepath=geopackage_filepath)}) - - # Otherwise, retrieve custom curve from Excel (previously stored in GeoPackage) - else: - return_array.append({ - "fragility_values_ref": fragility_ref, - "fragility_values": _gather_other_freq_line(ref=fragility_ref, tables=tables)}) - - # Build dataframe - if return_array.__len__() == 0: - df = DataFrame(data=[], columns=["fragility_values_ref", "fragility_values"]) - return df - return DataFrame(return_array) - - -def _add_fragility_values_to_combined_parameter_invoer( - df_parameter_invoer_combined: DataFrame, tables: InputParameterTables, - geopackage_filepath: str, drop_ref: bool = True) -> DataFrame: - """ Haalt uit de fragility values Excel de arrays op en vervang in de df_parameter_invoer_combined de referentie - met de daadwerkelijke fragility values. """ - - df = df_parameter_invoer_combined.copy(deep=True) - - # Replace empty values with NaN - set_option('future.no_silent_downcasting', True) - new_series = df['fragility_values_ref'].replace('', np.nan).infer_objects(copy=False) - df['fragility_values_ref'] = new_series - - # Gather referenced fragility value refs - fragility_refs = df['fragility_values_ref'].dropna().unique() - - # Collect fragility lines - df_frag_lines = _collect_fragility_values( - tables=tables, fragility_refs=fragility_refs, geopackage_filepath=geopackage_filepath) - - # Attach to parameter invoer df - df = df.merge( - df_frag_lines, left_on="fragility_values_ref", right_on="fragility_values_ref", how="left") - - if drop_ref: - df = df.drop(columns=["fragility_values_ref"]) - - return df - - -def _collect_right_columns_combined_parameter_invoer(df_parameter_invoer_combined: DataFrame) -> DataFrame: - """ Parameter tabel omzetten naar juiste kolommen. Enkel per uittredepunt, scenario en parameter de - parameterinvoer. """ - - # Add parameter invoer: op uittredepunten niveau - df_parameter_invoer_combined = df_parameter_invoer_combined.drop(columns=["bronnen", "opmerking"]) - parameter_description_columns = [ - "distribution_type", "mean", "variation", "deviation", "minimum", "maximum", "fragility_values"] - if "fragility_values_ref" in df_parameter_invoer_combined.columns: - parameter_description_columns.append("fragility_values_ref") - df_parameter_invoer_combined['parameter_input'] = df_parameter_invoer_combined.apply(lambda row: { - k: row[k] for k in parameter_description_columns if isinstance(row[k], list) or notna(row[k]) - }, axis=1) - df_parameter_invoer_combined = df_parameter_invoer_combined.drop(columns=parameter_description_columns) - return df_parameter_invoer_combined - - -def _construct_df_identifiers(geopackage_filepath: str, tables: InputParameterTables): - """ Create identifiers table. - - The identifiers table is a table with unique rows (unique uittredepunt, vak and scenario combo). It is used as a - base to explode to input per scenario and uittredepunt. """ - - gdf_uittredepunten: GeoDataFrame = read_file(geopackage_filepath, layer="uittredepunten") - df_identifiers: DataFrame = gdf_uittredepunten[["uittredepunt_id", "vak_id"]] - df_scenario_invoer: DataFrame = tables.df_scenario_invoer - df_identifiers = df_identifiers.merge(df_scenario_invoer, left_on="vak_id", right_on="vak_id") - df_identifiers = df_identifiers.drop(columns=["kans"]) - return df_identifiers - - -def _gather_required_input_parameters(geopackage_filepath: str) -> List[str]: - - conn = sqlite3.connect(geopackage_filepath) - cursor = conn.cursor() - cursor.execute(""" - SELECT geoprob_pipe_metadata."values" - FROM geoprob_pipe_metadata - WHERE metadata_type='geohydrologisch_model'; - """) - result = cursor.fetchone() - if not result: - raise ValueError - model_string = result[0] - conn.close() - df_dummy_data = DataFrame(INITIAL_INPUT_MAPPER[model_string]['input']) - - _ = df_dummy_data.sort_values(by=["name"]) - df_dummy_data = df_dummy_data.sort_values(by=["name"]) - return df_dummy_data['name'].unique().tolist() - # TODO: Return this to use in iteration to retrieve data - - # Method to use the system_function. Not necessary because of the dummy input. But for now kept. - # from geoprob_pipe.calculations.system_calculations.piping_moria.reliability_calculation import \ - # PipingMORIASystemReliabilityCalculation - # from geoprob_pipe.calculations.system_calculations.piping_moria.dummy_input import PIPING_DUMMY_INPUT - # import inspect - # - # reliability_calculation = PipingMORIASystemReliabilityCalculation(PIPING_DUMMY_INPUT) - # sig = inspect.signature(reliability_calculation.given_system_variables_setup_function) - # _ = [ - # name for name, param in sig.parameters.items() - # if param.default == inspect.Parameter.empty and param.kind in (param.POSITIONAL_OR_KEYWORD, param.KEYWORD_ONLY) - # ] - - # return ["mv_exit", "gamma_sat_deklaag"] # TODO - # return ["gamma_sat_deklaag"] # TODO - - -def _expand( - df_parameter_invoer_combined: DataFrame, df_identifiers: DataFrame, geopackage_filepath: str -) -> Dict[str, DataFrame]: - # Add parameter invoer: op uittredepunten niveau - required_input_parameters = _gather_required_input_parameters(geopackage_filepath=geopackage_filepath) - collection_of_dfs: Dict[str, DataFrame] = {} - - for parameter_name in required_input_parameters: - - # Gather and merge input on uittredepunten / scenario-niveau - df_gather = df_parameter_invoer_combined[ - (df_parameter_invoer_combined['parameter'] == parameter_name) & - (df_parameter_invoer_combined['scope'] == 'uittredepunt') & - (df_parameter_invoer_combined['ondergrondscenario_naam'].notna())] - df_gather = df_gather[["scope_referentie", "ondergrondscenario_naam", "parameter_input"]] - df_gather = df_gather.rename(columns={ - "scope_referentie": "uittredepunt_id", - "ondergrondscenario_naam": "naam"}) - df = df_identifiers.copy(deep=True) - df["parameter_input"] = np.nan - df['parameter_input'] = df['parameter_input'].combine_first( - df_identifiers.copy(deep=True).merge(df_gather, on=["uittredepunt_id", "naam"], how="left")[ - 'parameter_input']) - - # Uittredepunt niveau - df_gather = df_parameter_invoer_combined[ - (df_parameter_invoer_combined['parameter'] == parameter_name) & - (df_parameter_invoer_combined['scope'] == 'uittredepunt') & - (df_parameter_invoer_combined['ondergrondscenario_naam']).isna()] - df_gather = df_gather[["scope_referentie", "parameter_input"]] - df_gather = df_gather.rename(columns={"scope_referentie": "uittredepunt_id"}) - - df['parameter_input'] = df['parameter_input'].combine_first( - df_identifiers.copy(deep=True).merge(df_gather, on=["uittredepunt_id"], how="left")['parameter_input']) - - # Vak / scenario niveau - df_gather = df_parameter_invoer_combined[ - (df_parameter_invoer_combined['parameter'] == parameter_name) & - (df_parameter_invoer_combined['scope'] == 'vak') & - (df_parameter_invoer_combined['ondergrondscenario_naam'].notna())] - df_gather = df_gather[["scope_referentie", "ondergrondscenario_naam", "parameter_input"]] - df_gather = df_gather.rename(columns={ - "scope_referentie": "vak_id", - "ondergrondscenario_naam": "naam"}) - df['parameter_input'] = df['parameter_input'].combine_first( - df_identifiers.copy(deep=True).merge(df_gather, on=["vak_id", "naam"], how="left")['parameter_input']) - - # Vak niveau - df_gather = df_parameter_invoer_combined[ - (df_parameter_invoer_combined['parameter'] == parameter_name) & - (df_parameter_invoer_combined['scope'] == 'vak') & - (df_parameter_invoer_combined['ondergrondscenario_naam'].isna())] - df_gather = df_gather[["scope_referentie", "parameter_input"]] - df_gather = df_gather.rename(columns={"scope_referentie": "vak_id"}) - df['parameter_input'] = df['parameter_input'].combine_first( - df_identifiers.copy(deep=True).merge(df_gather, on=["vak_id"], how="left")['parameter_input']) - - # Traject niveau - df_gather = df_parameter_invoer_combined[ - (df_parameter_invoer_combined['parameter'] == parameter_name) & - (df_parameter_invoer_combined['scope'] == 'traject')] - if df_gather.__len__() >= 1: - traject_value = df_gather['parameter_input'].values[0] - df['parameter_input'] = df['parameter_input'].apply(lambda x: traject_value if isna(x) else x) - - # GIS spatial joins - df_gather = df_parameter_invoer_combined[ - (df_parameter_invoer_combined['parameter'] == parameter_name) & - (df_parameter_invoer_combined['scope'] == 'gis_uittredepunt')].copy(deep=True) - df_gather = df_gather[["scope_referentie", "parameter_input"]] - df_gather = df_gather.rename(columns={"scope_referentie": "uittredepunt_id"}) - df['parameter_input'] = df['parameter_input'].combine_first( - df_identifiers.copy(deep=True).merge(df_gather, on=["uittredepunt_id"], how="left")['parameter_input']) - - # Add to collection - collection_of_dfs[parameter_name] = df - - return collection_of_dfs - - -def _concat_collection(collection: Dict[str, DataFrame]): - for parameter_name, df in collection.items(): - collection[parameter_name]['parameter_name'] = parameter_name - return_df = concat([df for _, df in collection.items()], ignore_index=True) - return_df = return_df.rename(columns={"naam": "ondergrondscenario_naam"}) - return return_df[["parameter_name", "vak_id", "uittredepunt_id", "ondergrondscenario_naam", "parameter_input"]] - - -def run_expand_input_tables( - geopackage_filepath: str, add_frag_ref: bool = False, tables: Optional[InputParameterTables] = None, -) -> DataFrame: - """ Performs logic to expand all parameter input from different levels to input on uittredepunt-level. - - :param geopackage_filepath: - :param add_frag_ref: Indien True, dan wordt voor parameter input met een distribution_type 'cdf_curve' de referentie - naar de invoer behouden. Deze is niet nodig voor de berekeningen, maar wordt wel gebruikt voor de visualisaties. - :param tables: - :return: DataFrame with columns parameter_name, vak_id, uittredepunt_id, ondergrondscenario_naam and - parameter_input. - """ - - if tables is None: - tables = InputParameterTables(geopackage_filepath=geopackage_filepath) - - # Construct df_parameter_invoer_combined - df_parameter_invoer_combined1 = _combine_parameter_invoer_sources(tables=tables) - df_parameter_invoer_combined2 = _add_fragility_values_to_combined_parameter_invoer( - df_parameter_invoer_combined=df_parameter_invoer_combined1, tables=tables, - geopackage_filepath=geopackage_filepath, drop_ref=add_frag_ref==False) - df_parameter_invoer_combined3 = _collect_right_columns_combined_parameter_invoer( - df_parameter_invoer_combined=df_parameter_invoer_combined2) - - # Expand - df_identifiers = _construct_df_identifiers(geopackage_filepath=geopackage_filepath, tables=tables) - collection: Dict[str, DataFrame] = _expand( - df_parameter_invoer_combined=df_parameter_invoer_combined3, - df_identifiers=df_identifiers, - geopackage_filepath=geopackage_filepath) - - return _concat_collection(collection=collection) diff --git a/geoprob_pipe/cmd_app/validation/__init__.py b/geoprob_pipe/cmd_app/validation/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/geoprob_pipe/cmd_app/validation/validate_expanded_table.py b/geoprob_pipe/cmd_app/validation/validate_expanded_table.py new file mode 100644 index 00000000..b74ef127 --- /dev/null +++ b/geoprob_pipe/cmd_app/validation/validate_expanded_table.py @@ -0,0 +1,457 @@ +from __future__ import annotations + +import ast +import os +from datetime import datetime +from typing import TYPE_CHECKING + +import numpy as np +import pandas as pd + +from geoprob_pipe.utils.validation_messages import BColors + +if TYPE_CHECKING: + from geoprob_pipe.cmd_app.parameter_input.input_parameter_tables import ( + InputParameterTables, + ) + + +def get_mean(value) -> float | None: + """Find the mean value used from the excel or gis input. + + :param value: Data from the source tables + :return: Mean value as float or None if not found. + """ + if isinstance(value, dict): + return value.get("mean") + elif isinstance(value, str): + d = ast.literal_eval(value) # string → dict + return d.get("mean", np.nan) + return None + + +def _check_match( + df_check: pd.DataFrame, df_compare: pd.DataFrame, key: str +) -> pd.DataFrame: + """Check of the mean in the expanded table is close to the value in the source table. + + :param df_check: DataFrame with bools of checks. + :param df_compare: DataFrame with values of alle sources + :param key: Column name for the check. + :return: DataFrame with completed checks for a column. + """ + rel_tol = 1e-4 + col = key.replace("mean", "match") + mask_nan = df_compare[key].isna() + + df_check[col] = np.isclose(df_compare["mean"], df_compare[key], rtol=rel_tol) + + df_check[col] = df_check[col].mask(mask_nan, None) + + return df_check + + +def _merge_traject(df_input: pd.DataFrame, df_compare: pd.DataFrame) -> pd.DataFrame: + """Merge the input at the traject level into df_compare. + + :param df_input: Dataframe with inputs for expansion + :param df_compare: Dataframe with comparison results + :return: Expanded df_compare + """ + df_traject: pd.DataFrame = df_input.loc[df_input["scope"] == "traject"] + df_compare = df_compare.merge( + df_traject[["parameter", "mean"]], + on=["parameter"], + how="left", + suffixes=["", "_traject"], + ) + return df_compare + + +def _merge_vak_wild(df_input: pd.DataFrame, df_compare: pd.DataFrame) -> pd.DataFrame: + """Merge the input at the section wildcard level into df_compare. + + :param df_input: Dataframe with inputs for expansion + :param df_compare: Dataframe with comparison results + :return: Expanded df_compare + """ + df_vak_wild = df_input.loc[ + (df_input["scope"] == "vak") & (df_input["ondergrondscenario_naam"].isna()) + ] + df_vak_wild: pd.DataFrame = df_vak_wild.rename( + columns={"scope_referentie": "vak_id"} + ) + df_compare = df_compare.merge( + df_vak_wild[["parameter", "vak_id", "mean"]], + on=["parameter", "vak_id"], + how="left", + suffixes=["", "_vak_wild"], + ) + return df_compare + + +def _merge_vak_scen(df_input: pd.DataFrame, df_compare: pd.DataFrame) -> pd.DataFrame: + """Merge the input at the section scenario level into df_compare. + + :param df_input: Dataframe with inputs for expansion + :param df_compare: Dataframe with comparison results + :return: Expanded df_compare + """ + df_vak_scen: pd.DataFrame = df_input.loc[ + (df_input["scope"] == "vak") & (df_input["ondergrondscenario_naam"].notna()) + ] + df_vak_scen = df_vak_scen.rename(columns={"scope_referentie": "vak_id"}) + df_compare = df_compare.merge( + df_vak_scen[["parameter", "vak_id", "ondergrondscenario_naam", "mean"]], + on=["parameter", "vak_id", "ondergrondscenario_naam"], + how="left", + suffixes=["", "_vak_scen"], + ) + return df_compare + + +def _merge_utp_gis_wild( + df_input: pd.DataFrame, df_compare: pd.DataFrame +) -> pd.DataFrame: + """Merge the input at the exit-point gis wildcard level into df_compare. + + :param df_input: Dataframe with inputs for expansion + :param df_compare: Dataframe with comparison results + :return: Expanded df_compare + """ + df_gis_utp_wild: pd.DataFrame = df_input.loc[ + (df_input["scope"] == "gis_uittredepunt") + & (df_input["ondergrondscenario_naam"].isna()) + ] + df_gis_utp_wild = df_gis_utp_wild.rename( + columns={"scope_referentie": "uittredepunt_id"} + ) + df_compare = df_compare.merge( + df_gis_utp_wild[["parameter", "uittredepunt_id", "mean"]], + on=["parameter", "uittredepunt_id"], + how="left", + suffixes=["", "_gis_utp_wild"], + ) + return df_compare + + +def _merge_utp_excel_wild( + df_input: pd.DataFrame, df_compare: pd.DataFrame +) -> pd.DataFrame: + """Merge the input at the exit-point excel wildcard level into df_compare. + + :param df_input: Dataframe with inputs for expansion + :param df_compare: Dataframe with comparison results + :return: Expanded df_compare + """ + df_excel_utp_wild: pd.DataFrame = df_input.loc[ + (df_input["scope"] == "uittredepunt") + & (df_input["ondergrondscenario_naam"].isna()) + ] + df_excel_utp_wild = df_excel_utp_wild.rename( + columns={"scope_referentie": "uittredepunt_id"} + ) + df_compare = df_compare.merge( + df_excel_utp_wild[["parameter", "uittredepunt_id", "mean"]], + on=["parameter", "uittredepunt_id"], + how="left", + suffixes=["", "_excel_utp_wild"], + ) + return df_compare + + +def _merge_utp_gis_scen( + df_input: pd.DataFrame, df_compare: pd.DataFrame +) -> pd.DataFrame: + """Merge the input at the exit-point gis scenario level into df_compare. + + :param df_input: Dataframe with inputs for expansion + :param df_compare: Dataframe with comparison results + :return: Expanded df_compare + """ + df_gis_utp_scen: pd.DataFrame = df_input.loc[ + (df_input["scope"] == "gis_uittredepunt") + & (df_input["ondergrondscenario_naam"].notna()) + ] + df_gis_utp_scen = df_gis_utp_scen.rename( + columns={"scope_referentie": "uittredepunt_id"} + ) + df_compare = df_compare.merge( + df_gis_utp_scen[ + ["parameter", "uittredepunt_id", "ondergrondscenario_naam", "mean"] + ], + on=["parameter", "uittredepunt_id", "ondergrondscenario_naam"], + how="left", + suffixes=["", "_gis_utp_scen"], + ) + return df_compare + + +def _merge_utp_excel_scen( + df_input: pd.DataFrame, df_compare: pd.DataFrame +) -> pd.DataFrame: + """Merge the input at the exit-point excel scenario level into df_compare. + + :param df_input: Dataframe with inputs for expansion + :param df_compare: Dataframe with comparison results + :return: Expanded df_compare + """ + df_excel_utp_scen: pd.DataFrame = df_input.loc[ + (df_input["scope"] == "uittredepunt") + & (df_input["ondergrondscenario_naam"].notna()) + ] + df_excel_utp_scen = df_excel_utp_scen.rename( + columns={"scope_referentie": "uittredepunt_id"} + ) + df_compare = df_compare.merge( + df_excel_utp_scen[["parameter", "uittredepunt_id", "mean"]], + on=["parameter", "uittredepunt_id"], + how="left", + suffixes=["", "_excel_utp_scen"], + ) + return df_compare + + +def has_priority_error(row) -> bool: + """This function will check for all rows where the first True-value + in the row is, starting from the right to the left. + Values can be True, False or None. + + In the case of None; no value was provided at this level. + True means the value provided matches the value in the expanded table + And False means the value does not match the value in the expanded table. + If true is found it will check if any of the values to the right are False. + If one of these is found, there was a value found higher in the hierarchy that + was unused. + + :param row: Row from the DataFrame + :return: True if error found, else false + """ + # Find de column index van True met de hoogste plaats in de hierarchy. + last_true_idx: int = max( + (i for i, value in enumerate(row) if value is True), + default=-1, + ) + + values_after_last_true = row[last_true_idx + 1 :] + + return any(value is False for value in values_after_last_true) + + +def _clean_and_filter_dataframes( + df_input: pd.DataFrame, df_expanded: pd.DataFrame +) -> tuple[pd.DataFrame, pd.DataFrame]: + """Cleaup and filter the dataframes for futher use. + + :param df_input: DataFrame with input for expansion. + :param df_expanded: Exopanded DataFrame + :return: Prepared DataFrames + """ + df_expanded = df_expanded.rename(columns={"parameter_name": "parameter"}) + df_input["ondergrondscenario_naam"] = df_input["ondergrondscenario_naam"].replace( + "", pd.NA + ) + df_expanded["ondergrondscenario_naam"] = df_expanded[ + "ondergrondscenario_naam" + ].replace("", pd.NA) + df_expanded["uittredepunt_id"] = df_expanded["uittredepunt_id"].astype("Int64") + + # Filter Dataframes + df_expanded["mean"] = df_expanded["parameter_input"].apply(get_mean) + df_expanded = df_expanded[ + [ + "parameter", + "vak_id", + "uittredepunt_id", + "ondergrondscenario_naam", + "mean", + ] + ] + df_input = df_input[ + [ + "parameter", + "scope", + "scope_referentie", + "ondergrondscenario_naam", + "mean", + ] + ] + return df_input, df_expanded + + +def _setup_df_compare( + df_input: pd.DataFrame, df_expanded: pd.DataFrame +) -> pd.DataFrame: + """Setup a DataFrame with all possible value found in the input tables. + + :param df_input: DataFrame with input + :param df_expanded: Expanded DataFrame + :return: Expanded DataFrame with all possible values added in separate columns + per step. + """ + df_compare: pd.DataFrame = df_expanded.copy() + + df_compare = _merge_traject(df_input=df_input, df_compare=df_compare) + df_compare = _merge_vak_wild(df_input=df_input, df_compare=df_compare) + df_compare = _merge_vak_scen(df_input=df_input, df_compare=df_compare) + df_compare = _merge_utp_gis_wild(df_input=df_input, df_compare=df_compare) + df_compare = _merge_utp_excel_wild(df_input=df_input, df_compare=df_compare) + df_compare = _merge_utp_excel_scen(df_input=df_input, df_compare=df_compare) + df_compare = _merge_utp_gis_scen(df_input=df_input, df_compare=df_compare) + + return df_compare + + +def _setup_df_match( + df_compare: pd.DataFrame, df_expanded: pd.DataFrame +) -> pd.DataFrame: + """Setup a DataFrame with checks of the values are equal to the chosen value + in the expanded DataFrame. True means the same, False different and None means + no value found. + + :param df_compare: Expanded DataFrame with all possible values added in separate columns + :param df_expanded: Expanded DataFrame + :return: Expanded DataFrame with the check results in separate columns per step. + """ + df_check: pd.DataFrame = df_expanded.copy() + + df_check = _check_match(df_check, df_compare, "mean_traject") + df_check = _check_match(df_check, df_compare, "mean_vak_wild") + df_check = _check_match(df_check, df_compare, "mean_vak_scen") + df_check = _check_match(df_check, df_compare, "mean_gis_utp_wild") + df_check = _check_match(df_check, df_compare, "mean_excel_utp_wild") + df_check = _check_match(df_check, df_compare, "mean_gis_utp_scen") + df_check = _check_match(df_check, df_compare, "mean_excel_utp_scen") + + return df_check + + +def _setup_df_errors( + df_compare: pd.DataFrame, df_check: pd.DataFrame, df_expanded: pd.DataFrame +) -> pd.DataFrame: + """Collect the result of both checks in a single dataframe. + + :param df_check: Expanded DataFrame with the check results in separate columns + per step. + :param df_expanded: Expanded DataFrame + :return: DataFrame with the result of the error check for all rows. + """ + df_errors: pd.DataFrame = df_expanded.copy() + df_errors = df_errors.rename(columns={"mean": "used_value"}) + + cols = [ + "match_traject", + "match_vak_wild", + "match_vak_scen", + "match_gis_utp_wild", + "match_excel_utp_wild", + "match_gis_utp_scen", + "match_excel_utp_scen", + ] + df_errors["value_correct"] = df_check[cols].any(axis=1) + + # --- Check highest hierarchy --- + df_errors["uses"] = ( + df_check[cols].idxmax(axis=1).where(cond=df_errors["value_correct"], other=None) + ) + df_errors["should_use"] = ( + df_check[cols] + .idxmax(axis=1) + .where(cond=~df_errors["value_correct"], other=df_errors["uses"]) + ) + clean_names = { + "match_traject": "traject", + "match_vak_wild": "vak wildcard", + "match_vak_scen": "vak scenario", + "match_gis_utp_wild": "gis uittredepunt wildcard", + "match_excel_utp_wild": "excel uittredepunt wildcard", + "match_gis_utp_scen": "gis uittredepunt scenario", + "match_excel_utp_scen": "excel uittredepunt scenario", + } + df_errors["uses"] = df_errors["uses"].astype(str).map(clean_names) + df_errors["should_use"] = df_errors["should_use"].astype(str).map(clean_names) + arr = df_check[cols].to_numpy(dtype=object) + + df_errors["wrong_step"] = [has_priority_error(row) for row in arr] + + # FIXME Temporary fix for the 'buitenwaterstand' parameter + df_errors = df_errors.loc[df_expanded["parameter"] != "buitenwaterstand"] + # kolomnamem: used_value, is_correct, should_use, uses, + # zonder match als nette string + return df_errors + + +def validate_expand_tables( + tables: InputParameterTables, geopackage_filepath: str, df_expanded: pd.DataFrame +) -> None: + """Validation function for the expanded parameters table. Checks of the values in + the table correspond the values in de source tables. And checks of a value higher + in the hierarchy exist than was used. + + :param tables: Source tables. + :param df_expanded: expanded parameter table. + """ + df_excel: pd.DataFrame = tables.df_parameter_invoer.copy() + df_gis: pd.DataFrame = tables.df_gis_join_parameter_invoer.copy() + + # Remove empty columns for pandas.concat. + df_excel = df_excel.dropna(axis=1, how="all") + df_gis = df_gis.dropna(axis=1, how="all") + + # Prepare dataframes + df_excel["scope_referentie"] = df_excel["scope_referentie"].astype("Int64") + df_gis["scope_referentie"] = df_gis["scope_referentie"].astype("Int64") + df_gis["mean"] = df_gis["mean"].astype(float) + df_gis["scope"] = "gis_uittredepunt" + df_input: pd.DataFrame = pd.concat([df_excel, df_gis], ignore_index=True) + + df_input, df_expanded = _clean_and_filter_dataframes( + df_input=df_input, df_expanded=df_expanded + ) + + # Merge per step in hierarchy + df_compare: pd.DataFrame = _setup_df_compare( + df_input=df_input, df_expanded=df_expanded + ) + # Check of values are equal + df_match: pd.DataFrame = _setup_df_match( + df_compare=df_compare, df_expanded=df_expanded + ) + # Collect the errors + df_errors: pd.DataFrame = _setup_df_errors( + df_compare=df_compare, df_check=df_match, df_expanded=df_expanded + ) + + # Result: collect all rows without a match or wrong hierarchy + df_validation_output: pd.DataFrame = df_errors.loc[ + (df_errors["value_correct"]) | (df_errors["wrong_step"]) + ] + + if len(df_validation_output) > 0: + export_dir = os.path.join( + os.path.dirname(geopackage_filepath), + "exports", + str(datetime.now().strftime("%Y-%m-%d_%H%M%S")), + "parameter_input_process", + ) + os.makedirs(export_dir, exist_ok=True) + export_path = os.path.join(export_dir, "validation_expanded_table_input.xlsx") + if os.path.exists(export_path): + os.remove(export_path) + df_validation_output.to_excel(export_path, index=False) + print( + f"""{BColors.WARNING} + Er is een mismatch tussen de invoer parameters en de expanded parameter + tabel. De gedetailleerde lijst is geëxporteerd naar: \n + {export_path} + {BColors.ENDC}""" + ) + if os.environ.get("GEOPROB_DEBUG", False): + export_path = os.path.join(export_dir, "validation_df_compare.xlsx") + if os.path.exists(export_path): + os.remove(export_path) + df_compare.to_excel(export_path, index=False) + export_path = os.path.join(export_dir, "validation_df_match.xlsx") + if os.path.exists(export_path): + os.remove(export_path) + df_match.to_excel(export_path, index=False) diff --git a/geoprob_pipe/results/assemblage/objects.py b/geoprob_pipe/results/assemblage/objects.py index b26b5c12..e314e3da 100644 --- a/geoprob_pipe/results/assemblage/objects.py +++ b/geoprob_pipe/results/assemblage/objects.py @@ -1,6 +1,5 @@ from __future__ import annotations from dataclasses import dataclass -import math from typing import Optional, List, Tuple, cast from geoprob_pipe.results.assemblage.functions import bepaal_N_vak import scipy.stats as stats # importeer de scipy.stats module diff --git a/geoprob_pipe/visualizations/graphs/hfreq.py b/geoprob_pipe/visualizations/graphs/hfreq.py index b5462860..794b8937 100644 --- a/geoprob_pipe/visualizations/graphs/hfreq.py +++ b/geoprob_pipe/visualizations/graphs/hfreq.py @@ -1,16 +1,12 @@ from __future__ import annotations import os -# import matplotlib.pyplot as plt from typing import TYPE_CHECKING, List import plotly.graph_objects as go from datetime import datetime -# from geopandas import GeoDataFrame, read_file -# from probabilistic_library import FragilityValue import pydra_core as pydra from pandas import Series, concat, DataFrame, read_sql -from geoprob_pipe.cmd_app.parameter_input.expand_input_tables import run_expand_input_tables +from geoprob_pipe.cmd_app.parameter_input.expand.expand_input_tables import run_expand_input_tables import sqlite3 - if TYPE_CHECKING: from geoprob_pipe import GeoProbPipe @@ -237,6 +233,3 @@ def _optionally_export(self, export: bool = False, add_timestamp: bool = False): timestamp = datetime.now().strftime("%Y-%m-%d_%H%M") timestamp_str = f"{timestamp}_" self.fig.write_html(os.path.join(export_dir, f"{timestamp_str}hfreq.html"), include_plotlyjs='cdn') - # if self.geoprob_pipe.software_requirements.chrome_is_installed: - # self.fig.write_image(os.path.join(export_dir, f"{timestamp_str}hfreq.png"), format="png") - # Note CP: Export van PNG uitgezet. Meerwaarde beperkt denk ik aangezien je de HTML al hebt. diff --git a/geoprob_pipe/visualizations/graphs/river_waterlevel.py b/geoprob_pipe/visualizations/graphs/river_waterlevel.py index 51586f3c..c6ca923d 100644 --- a/geoprob_pipe/visualizations/graphs/river_waterlevel.py +++ b/geoprob_pipe/visualizations/graphs/river_waterlevel.py @@ -6,7 +6,7 @@ import plotly.colors as pc from plotly.graph_objects import Figure, Scatter from typing import TYPE_CHECKING, Dict, List, Tuple -from geoprob_pipe.cmd_app.parameter_input.expand_input_tables import run_expand_input_tables +from geoprob_pipe.cmd_app.parameter_input.expand.expand_input_tables import run_expand_input_tables if TYPE_CHECKING: from geoprob_pipe import GeoProbPipe diff --git a/tests/calculations/systems/test_build_and_run.py b/tests/calculations/systems/test_build_and_run.py index bf851b6c..9df05351 100644 --- a/tests/calculations/systems/test_build_and_run.py +++ b/tests/calculations/systems/test_build_and_run.py @@ -17,7 +17,7 @@ def test_worker(): geohydrologisch_model=model, geopackage_filepath=filepath, to_run_vakken_ids=None) - result = _worker(row_unique={'uittredepunt_id': 1, 'ondergrondscenario_naam': 'scenario1', 'vak_id': 4}) + _ = _worker(row_unique={'uittredepunt_id': 1, 'ondergrondscenario_naam': 'scenario1', 'vak_id': 4}) ##