Source code for oasislmf.preparation.summaries

__all__ = [
    'get_summary_mapping',
    'generate_summaryxref_files',
    'get_ri_inuring_priority_output_levels',
    'merge_oed_to_mapping',
    'write_exposure_summary',
    'get_exposure_summary',
    'get_useful_summary_cols',
    'get_xref_df',
    'write_summary_levels',
    'write_mapping_file',
    'calculated_summary_cols'
]

import logging
import pathlib
import shutil

import io
import json
import os

import numpy as np
import pandas as pd

from ..utils.coverages import SUPPORTED_COVERAGE_TYPES
from ..utils.data import (
    factorize_dataframe,
    get_dataframe,
    get_json,
    merge_dataframes,
    set_dataframe_column_dtypes,
    fill_na_with_categoricals
)
from ..utils.defaults import (
    SOURCE_IDX,
    SAR_ID,
    SUMMARY_MAPPING,
    SUMMARY_OUTPUT,
    SUMMARY_TOP_LEVEL_COLS,
    get_default_exposure_profile,
)
from ..utils.exceptions import OasisException
from ..utils.log import oasis_log
from ..utils.path import as_path
from ..utils.status import OASIS_KEYS_STATUS, OASIS_KEYS_STATUS_MODELLED
from oasislmf.pytools.common.data import gul_summary_xref_dtype, fm_summary_xref_dtype
from oasislmf.pytools.converters.csvtobin.utils.common import df_to_ndarray
from oasislmf.utils.fm import SUPPORTED_FM_LEVELS


MAP_SUMMARY_DTYPES = {
    'loc_id': 'int',
    SOURCE_IDX['loc']: 'int',
    SOURCE_IDX['acc']: 'int',
    'item_id': 'int',
    'layer_id': 'int',
    'coverage_id': 'int',
    'peril_id': 'category',
    'agg_id': 'int',
    'output_id': 'int',
    'coverage_type_id': 'int',
    'tiv': 'float',
    'building_id': 'int',
    'risk_id': 'int',
}

logger = logging.getLogger(__name__)


[docs] def get_useful_summary_cols(oed_hierarchy): return [ oed_hierarchy['accnum']['ProfileElementName'], oed_hierarchy['locnum']['ProfileElementName'], 'loc_id', oed_hierarchy['polnum']['ProfileElementName'], oed_hierarchy['portnum']['ProfileElementName'], SOURCE_IDX['loc'], SOURCE_IDX['acc'], 'item_id', 'layer_id', 'coverage_id', 'peril_id', 'agg_id', 'output_id', 'coverage_type_id', 'tiv', 'building_id', 'risk_id', 'intensity_adjustment', 'return_period' ]
def is_property_damage(summary_df): property_damage_cov_id = [SUPPORTED_COVERAGE_TYPES['buildings']['id'], SUPPORTED_COVERAGE_TYPES['contents']['id']] summary_df['is_property_damage'] = summary_df['coverage_type_id'].isin(property_damage_cov_id) return summary_df
[docs] calculated_summary_cols = {'is_property_damage': is_property_damage}
[docs] def get_xref_df(il_inputs_df): top_level_layers_df = il_inputs_df.loc[il_inputs_df['level_id'] == il_inputs_df['level_id'].max(), ['top_agg_id'] + SUMMARY_TOP_LEVEL_COLS].drop_duplicates() bottom_level_layers_df = il_inputs_df[il_inputs_df['level_id'] == 0].copy() bottom_level_layers_df.drop(columns=SUMMARY_TOP_LEVEL_COLS, inplace=True) return (merge_dataframes(bottom_level_layers_df, top_level_layers_df, join_on=['top_agg_id']) .drop_duplicates(subset=['gul_input_id', 'layer_id'], keep='first') .sort_values(['gul_input_id', 'layer_id'], kind='stable') )
@oasis_log
[docs] def get_summary_mapping(inputs_df, oed_hierarchy, is_fm_summary=False): """Create a DataFrame with linking information between Ktools `OasisFiles` And the Exposure data Args: inputs_df (pandas.DataFrame): datafame from gul_inputs.get_gul_input_items(..) / il_inputs.get_il_input_items(..) oed_hierarchy (dict): OED profile hierarchy, used to resolve the acc/loc/pol/port column names to keep is_fm_summary (bool): Indicates whether an FM summary mapping is required Returns: pandas.DataFrame: Subset of columns from gul_inputs_df / il_inputs_df """ # Case GUL+FM (based on il_inputs_df) if is_fm_summary: summary_mapping = inputs_df.rename(columns={'gul_input_id': 'agg_id'}) # GUL Only else: summary_mapping = inputs_df.copy(deep=True) summary_mapping['layer_id'] = 1 summary_mapping['agg_id'] = summary_mapping['item_id'] summary_mapping.drop( [c for c in summary_mapping.columns if c not in get_useful_summary_cols(oed_hierarchy)], axis=1, inplace=True ) acc_num = oed_hierarchy['accnum']['ProfileElementName'] loc_num = oed_hierarchy['locnum']['ProfileElementName'] policy_num = oed_hierarchy['polnum']['ProfileElementName'] portfolio_num = oed_hierarchy['portnum']['ProfileElementName'] dtypes = { **{t: 'str' for t in [portfolio_num, policy_num, acc_num, loc_num, 'peril_id']}, **{t: 'uint8' for t in ['coverage_type_id']}, **{t: 'uint32' for t in [SOURCE_IDX['loc'], SOURCE_IDX['acc'], 'loc_id', 'item_id', 'layer_id', 'coverage_id', 'agg_id', 'output_id', 'building_id', 'risk_id']}, **{t: 'float64' for t in ['tiv']} } summary_mapping = set_dataframe_column_dtypes(summary_mapping, dtypes) return summary_mapping
[docs] def merge_oed_to_mapping(summary_map_df, exposure_df, oed_column_join, oed_column_info): """Create a factorized col (summary ids) based on a list of oed column names {'Col_A': 0, 'Col_B': 1, 'Col_C': 2} Args: summary_map_df (pandas.DataFrame): dataframe return from get_summary_mapping exposure_df (pandas.DataFrame): Summary map file path oed_column_join (list): column to join on oed_column_info (dict): Dictionary of columns to pick from exposure_df and their default value Returns: pandas.DataFrame: New DataFrame of summary_map_df + exposure_df merged on exposure index """ column_set = set(oed_column_info) columns_found = [c for c in column_set if c in exposure_df.columns and c not in summary_map_df.columns] columns_missing = list(set(column_set) - set(columns_found)) new_summary_map_df = merge_dataframes(summary_map_df, exposure_df.loc[:, columns_found + oed_column_join], join_on=oed_column_join, how='inner') for col, default in oed_column_info.items(): if col in columns_missing: new_summary_map_df[col] = default fill_na_with_categoricals(new_summary_map_df, oed_column_info) return new_summary_map_df
def group_by_oed(oed_col_group, summary_map_df, exposure_df, sort_by, accounts_df=None): """Adds list of OED fields from `column_set` to summary map file Args: oed_col_group (list): OED column names the rows are grouped by to form the summaries summary_map_df (pandas.DataFrame): dataframe return from get_summary_mapping exposure_df (pandas.DataFrame): DataFrame loaded from location.csv sort_by (str): column the grouped rows are ordered by accounts_df (pandas.DataFrame): DataFrame loaded from accounts.csv Returns: tuple: a 3-tuple ``(summary_ids, summary_id_values, summary_tiv)``: - summary_ids (numpy.ndarray): integer summary id (1..n) for each row, e.g. ``array([1, 2, 1, 2, 1, 2, 1, 2, 1, 2, ...])`` - summary_id_values (numpy.ndarray): distinct values used to factorize the rows into summaries, e.g. ``array(['Layer1', 'Layer2'], dtype=object)`` - summary_tiv (pandas.DataFrame): total TIV aggregated per summary group """ oed_cols = oed_col_group # All required columns exposure_cols = [c for c in oed_cols if c not in summary_map_df.columns and c not in calculated_summary_cols] # columns which are in locations / Accounts file mapped_cols = [c for c in oed_cols + [SOURCE_IDX['loc'], SOURCE_IDX['acc'], sort_by] if c in summary_map_df.columns] # Columns already in summary_map_df to_calculate_column = [c for c in oed_cols if c in calculated_summary_cols] tiv_cols = ['tiv', 'loc_id', 'building_id', 'coverage_type_id'] # Extract mapped_cols from summary_map_df summary_group_df = summary_map_df.loc[:, list(set(tiv_cols).union(mapped_cols))] # Search Loc / Acc files and merge in remaining if exposure_cols: # Subject-at-risk (location / account) file columns exposure_cols_loc = [c for c in exposure_cols if c in exposure_df.columns] exposure_col_df = exposure_df.loc[:, exposure_cols_loc + [SAR_ID]].drop_duplicates(subset=SAR_ID, keep='first') summary_group_df = merge_dataframes(summary_group_df, exposure_col_df, join_on=SAR_ID, how='left') # Account file columns if isinstance(accounts_df, pd.DataFrame): accounts_cols = [c for c in exposure_cols if c in set(accounts_df.columns) - set(exposure_df.columns)] if accounts_cols: accounts_col_df = accounts_df.loc[:, accounts_cols + [SOURCE_IDX['acc']]] summary_group_df = merge_dataframes(summary_group_df, accounts_col_df, join_on=SOURCE_IDX['acc'], how='left') for col in to_calculate_column: summary_group_df = calculated_summary_cols[col](summary_group_df) fill_na_with_categoricals(summary_group_df, 0) summary_group_df.sort_values(by=[sort_by], inplace=True, kind='stable') summary_ids = factorize_dataframe(summary_group_df, by_col_labels=oed_cols) summary_tiv = summary_group_df.drop_duplicates(['loc_id', 'building_id', 'coverage_type_id'] + oed_col_group, keep='first').groupby(oed_col_group, observed=True).agg({'tiv': "sum"}) return summary_ids[0], summary_ids[1], summary_tiv @oasis_log
[docs] def write_summary_levels(exposure_df, accounts_df, exposure_data, target_dir): """Json file with list Available / Recommended columns for use in the summary reporting Available: Columns which exists in input files and has at least one non-zero / NaN value Recommended: Columns which are available + also in the list of `useful` groupings SUMMARY_LEVEL_LOC { 'GUL': { 'available': ['AccNumber', 'LocNumber', 'istenant', 'buildingid', 'countrycode', 'latitude', 'longitude', 'streetaddress', 'postalcode', 'occupancycode', 'constructioncode', 'locperilscovered', 'BuildingTIV', 'ContentsTIV', 'BITIV', 'PortNumber'], 'IL': { ... etc ... } } """ # Manage internal columns, (Non-OED exposure input) int_excluded_cols = ['loc_id', SOURCE_IDX['loc']] desc_non_oed = 'Not an OED field' int_oasis_cols = { 'coverage_type_id': 'Oasis coverage type', 'peril_id': 'OED peril code', 'coverage_id': 'Oasis coverage identifier', } # GUL perspective (loc columns only) l_col_list = (exposure_df.drop(columns=['peril_group_id'], errors='ignore') .replace(0, np.nan) .dropna(how='any', axis=1) .columns.to_list()) l_col_info = exposure_data.get_input_fields('Loc') gul_avail = {k: l_col_info[k.lower()]["Type & Description"] if k.lower() in l_col_info else desc_non_oed for k in set([c for c in l_col_list]).difference(int_excluded_cols)} # IL perspective (join of acc + loc col with no dups) il_avail = {} if accounts_df is not None: a_col_list = accounts_df.loc[:, ~accounts_df.isnull().all()].columns.to_list() a_col_info = exposure_data.get_input_fields('Acc') a_avail = set([c for c in a_col_list]) il_avail = {k: a_col_info[k.lower()]["Type & Description"] if k.lower() in a_col_info else desc_non_oed for k in a_avail.difference(gul_avail.keys())} # Write JSON gul_summary_lvl = {'GUL': {'available': {**gul_avail, **il_avail, **int_oasis_cols}}} il_summary_lvl = {'IL': {'available': {**gul_avail, **il_avail, **int_oasis_cols}}} if il_avail else {} with io.open(os.path.join(target_dir, 'exposure_summary_levels.json'), 'w', encoding='utf-8') as f: f.write(json.dumps({**gul_summary_lvl, **il_summary_lvl}, sort_keys=True, ensure_ascii=False, indent=4))
@oasis_log
[docs] def write_mapping_file(sum_inputs_df, target_dir, is_fm_summary=False): """Writes a summary map file, used to build summarycalc xref files. Args: sum_inputs_df (pandas.DataFrame): dataframe return from get_summary_mapping target_dir (str): directory the summary map file is written to is_fm_summary (bool): Indicates whether an FM summary mapping is required Returns: str: Summary xref file path """ target_dir = as_path( target_dir, 'Target IL input files directory', is_dir=True, preexists=False ) # Set chunk size for writing the CSV files - default is max 20K, min 1K chunksize = min(2 * 10 ** 5, max(len(sum_inputs_df), 1000)) if is_fm_summary: sum_mapping_fp = os.path.join(target_dir, SUMMARY_MAPPING['fm_map_fn']) else: sum_mapping_fp = os.path.join(target_dir, SUMMARY_MAPPING['gul_map_fn']) try: sum_inputs_df.to_csv( path_or_buf=sum_mapping_fp, encoding='utf-8', mode=('w' if os.path.exists(sum_mapping_fp) else 'a'), chunksize=chunksize, index=False ) except (IOError, OSError) as e: raise OasisException("Exception raised in 'write_mapping_file'", e) return sum_mapping_fp
def get_column_selection(summary_set): """Given a analysis_settings summary definition, return either 1. the set of OED columns requested to group by 2. If no information key 'oed_fields', then group all outputs into a single summary_set Args: summary_set (dict): summary group dictionary from the `analysis_settings.json` Returns: list: List of selected OED columns to create summary groups from """ if "oed_fields" not in summary_set: return [] if not summary_set["oed_fields"]: return [] # Use OED column list set in analysis_settings file elif isinstance(summary_set['oed_fields'], list) and len(summary_set['oed_fields']) > 0: return [c for c in summary_set['oed_fields']] elif isinstance(summary_set['oed_fields'], str) and len(summary_set['oed_fields']) > 0: return [summary_set['oed_fields']] else: raise OasisException( 'Error processing settings file: "oed_fields" ' 'is expected to be a list of strings, not {}'.format(type(summary_set['oed_fields'])) ) def get_ri_settings(run_dir): """Return the contents of ri_layers.json Example: { "1": { "inuring_priority": 1, "risk_level": "LOC", "directory": " ... /runs/ProgOasis-20190501145127/RI_1" } } Args: run_dir (str): The file path of the model run directory Returns: dict: metadata for the Reinsurance layers """ return get_json(src_fp=os.path.join(run_dir, 'ri_layers.json'))
[docs] def get_ri_inuring_priority_output_levels(run_dir): """Load the mapping from OED InuringPriority to RI output level from the input directory. The mapping is created during input generation and records, for each OED InuringPriority value, the index of the last RI layer (output level) that belongs to that priority. When an InuringPriority spans multiple risk levels (e.g. LOC and ACC), its output level is the highest RI layer index among those risk levels. Args: run_dir (str): Directory containing ``ri_inuring_priority_output_levels.json`` Returns: dict: mapping ``{inuring_priority: output_level}`` with integer keys/values """ raw = get_json(src_fp=os.path.join(run_dir, 'ri_inuring_priority_output_levels.json')) return {int(k): int(v) for k, v in raw.items()}
def write_df_to_csv_file(df, target_dir, filename): """Write a generated summary xref dataframe to disk in csv format. Args: df (pandas.DataFrame): The dataframe output of get_df( .. ) target_dir (str): Abs directory to write a summary_xref file filename (str): Name of output file """ target_dir = as_path(target_dir, 'Input files directory', is_dir=True, preexists=False) pathlib.Path(target_dir).mkdir(parents=True, exist_ok=True) chunksize = min(2 * 10 ** 5, max(len(df), 1000)) csv_fp = os.path.join(target_dir, filename) try: df.to_csv( path_or_buf=csv_fp, encoding='utf-8', mode=('w'), chunksize=chunksize, index=False ) except (IOError, OSError) as e: raise OasisException("Exception raised in 'write_df_to_csv_file'", e) return csv_fp def write_df_to_parquet_file(df, target_dir, filename): """Write a generated summary xref dataframe to disk in parquet format. Args: df (pandas.DataFrame): The dataframe output of get_df( .. ) target_dir (str): Abs directory to write a summary_xref file filename (str): Name of output file """ target_dir = as_path( target_dir, 'Output files directory', is_dir=True, preexists=False ) parquet_fp = os.path.join(target_dir, filename) try: df.to_parquet(path=parquet_fp, engine='pyarrow') except (IOError, OSError) as e: raise OasisException( "Exception raised in 'write_df_to_parquet_file'", e ) return parquet_fp def get_summary_xref_df( map_df, exposure_df, accounts_df, summaries_info_dict, summaries_type, id_set_index='output_id' ): """Create a Dataframe for either gul / il / ri based on a section from the analysis settings Args: map_df (pandas.DataFrame): Summary Map dataframe (GUL / IL) exposure_df (pandas.DataFrame): Location OED data accounts_df (pandas.DataFrame): Accounts OED data id_set_index (str): column of map_df the summary xref ids are taken from summaries_info_dict (list): list of dictionary definition for a summary group from the analysis_settings file, e.g.:: [{ "id": 1, "oed_fields": [], ... }, ... ] summaries_type (str): Text label to use as key in summary description either ['gul', 'il', 'ri'] Returns: summaryxref_df (pandas.DataFrame): Dataframe containing abstracted summary data for ktools summary_desc (dictionary): dictionary of dataFrames listing what summary_ids map to """ summaryxref_df = pd.DataFrame() summary_desc = {} all_cols = set(map_df.columns.to_list() + exposure_df.columns.to_list() + list(calculated_summary_cols.keys())) if isinstance(accounts_df, pd.DataFrame): all_cols.update(accounts_df.columns.to_list()) # Extract the summary id index column depending on id_set_index map_df.sort_values(id_set_index, inplace=True, kind='stable') ids_set_df = map_df.loc[:, [id_set_index]].rename(columns={'output_id': "output"}) # For each granularity build a set grouping for summary_set in summaries_info_dict: summary_set_df = ids_set_df cols_group_by = get_column_selection(summary_set) file_extension = 'csv' if summary_set.get('ord_output', {}).get('parquet_format'): file_extension = 'parquet' desc_key = '{}_S{}_summary-info.{}'.format( summaries_type, summary_set['id'], file_extension ) # an empty intersection means no selected columns from the input data if not set(cols_group_by).intersection(all_cols): # is the intersection empty because the columns don't exist? if set(cols_group_by).difference(all_cols): err_msg = 'Input error: Summary set columns missing from the input files: {}'.format( set(cols_group_by).difference(all_cols)) raise OasisException(err_msg) # Fall back to setting all in single group summary_set_df['summary_id'] = 1 summary_desc[desc_key] = pd.DataFrame(data=['All-Risks'], columns=['_not_set_']) summary_desc[desc_key].insert(loc=0, column='summary_id', value=1) summary_desc[desc_key].insert(loc=len(summary_desc[desc_key].columns), column='tiv', value=map_df.drop_duplicates(['building_id', 'loc_id', 'coverage_type_id'], keep='first').tiv.sum()) else: ( summary_set_df['summary_id'], set_values, tiv_values ) = group_by_oed(cols_group_by, map_df, exposure_df, id_set_index, accounts_df) # Build description file summary_desc_df = pd.DataFrame(data=list(set_values), columns=cols_group_by) summary_desc_df.insert(loc=0, column='summary_id', value=range(1, len(set_values) + 1)) summary_desc[desc_key] = pd.merge(summary_desc_df, tiv_values, left_on=cols_group_by, right_on=cols_group_by) # Appends summary set to '__summaryxref.csv' summary_set_df['summaryset_id'] = summary_set['id'] summaryxref_df = pd.concat([summaryxref_df, summary_set_df.drop_duplicates()], sort=True, ignore_index=True) dtypes = { t: 'uint32' for t in ['coverage_id', 'summary_id', 'summaryset_id'] } summaryxref_df = set_dataframe_column_dtypes(summaryxref_df, dtypes) return summaryxref_df, summary_desc @oasis_log
[docs] def generate_summaryxref_files( location_df, account_df, model_run_fp, analysis_settings, il=False, ri=False, rl=False, intermediary_csv=False ): """Top level function for creating the summaryxref files from the manager.py Args: location_df (pandas.DataFrame): Source locations, joined to the summary map to build the summary groupings and their description files account_df (pandas.DataFrame): Source accounts, joined the same way. Required for the il, ri and rl summary levels model_run_fp (str): Model run directory file path analysis_settings (dict): Model analysis settings file il (bool): Boolean to indicate the insured loss level mode - false if the source accounts file path not provided to Oasis files gen. ri (bool): Boolean to indicate the RI loss level mode - false if the source accounts file path not provided to Oasis files gen. rl (bool): Boolean to indicate the RL loss level mode - false if the source accounts file path not provided to Oasis files gen. intermediary_csv (bool): If True, also write a csv copy of each summaryxref file alongside the binary """ # Boolean checks for summary generation types (gul / il / ri) gul_summaries = all([ analysis_settings.get('gul_output'), analysis_settings.get('gul_summaries'), ]) il_summaries = all([ analysis_settings.get('il_output'), analysis_settings.get('il_summaries'), il, ]) ri_summaries = all([ analysis_settings.get('ri_output'), analysis_settings.get('ri_summaries'), ri, ]) rl_summaries = all([ analysis_settings.get('rl_output'), analysis_settings.get('rl_summaries'), rl, ]) # Verify account file + il_map file il_map_fp = os.path.join(model_run_fp, 'input', SUMMARY_MAPPING['fm_map_fn']) if il_summaries or ri_summaries or rl_summaries: if account_df is None: raise OasisException('No account file found.') if not os.path.exists(il_map_fp): raise OasisException('No summary map file found.') # Load il_map if present if os.path.exists(il_map_fp): il_map_df = get_dataframe( src_fp=il_map_fp, lowercase_cols=False, col_dtypes=MAP_SUMMARY_DTYPES, empty_data_error_msg='No summary map file found.', ) il_map_df = il_map_df[list(set(il_map_df).intersection(MAP_SUMMARY_DTYPES))] if gul_summaries: gul_map_df = il_map_df gul_map_df['item_id'] = gul_map_df['agg_id'] elif gul_summaries: gul_map_fp = os.path.join(model_run_fp, 'input', SUMMARY_MAPPING['gul_map_fn']) gul_map_df = get_dataframe( src_fp=gul_map_fp, lowercase_cols=False, col_dtypes=MAP_SUMMARY_DTYPES, empty_data_error_msg='No summary map file found.', ) gul_map_df = gul_map_df[list(set(gul_map_df).intersection(MAP_SUMMARY_DTYPES))] if gul_summaries: # Load GUL summary map id_set_index = 'item_id' gul_summaryxref_df, gul_summary_desc = get_summary_xref_df( gul_map_df, location_df, account_df, analysis_settings['gul_summaries'], 'gul', id_set_index ) # Write Xref file df_to_ndarray(gul_summaryxref_df, gul_summary_xref_dtype).tofile(os.path.join(model_run_fp, 'input', f"{SUMMARY_OUTPUT['gul']}.bin")) if intermediary_csv: write_df_to_csv_file(gul_summaryxref_df, os.path.join(model_run_fp, 'input'), f"{SUMMARY_OUTPUT['gul']}.csv") # Write summary_id description files for desc_key in gul_summary_desc: if desc_key.split('.')[-1] == 'parquet': write_df_to_parquet_file(gul_summary_desc[desc_key], os.path.join(model_run_fp, 'output'), desc_key) else: write_df_to_csv_file(gul_summary_desc[desc_key], os.path.join(model_run_fp, 'output'), desc_key) if il_summaries: # Load FM summary map il_map_fp = os.path.join(model_run_fp, 'input', SUMMARY_MAPPING['fm_map_fn']) il_map_df = get_dataframe( src_fp=il_map_fp, lowercase_cols=False, col_dtypes=MAP_SUMMARY_DTYPES, empty_data_error_msg='No summary map file found.', ) il_map_df = il_map_df[list(set(il_map_df).intersection(MAP_SUMMARY_DTYPES))] il_summaryxref_df, il_summary_desc = get_summary_xref_df( il_map_df, location_df, account_df, analysis_settings['il_summaries'], 'il' ) # Write Xref file df_to_ndarray(il_summaryxref_df, fm_summary_xref_dtype).tofile(os.path.join(model_run_fp, 'input', f"{SUMMARY_OUTPUT['il']}.bin")) if intermediary_csv: write_df_to_csv_file(il_summaryxref_df, os.path.join(model_run_fp, 'input'), f"{SUMMARY_OUTPUT['il']}.csv") # Write summary_id description files for desc_key in il_summary_desc: if desc_key.split('.')[-1] == 'parquet': write_df_to_parquet_file(il_summary_desc[desc_key], os.path.join(model_run_fp, 'output'), desc_key) else: write_df_to_csv_file(il_summary_desc[desc_key], os.path.join(model_run_fp, 'output'), desc_key) if ri_summaries or rl_summaries: if ('il_summaries' not in analysis_settings) or (not il_summaries): il_map_fp = os.path.join(model_run_fp, 'input', SUMMARY_MAPPING['fm_map_fn']) il_map_df = get_dataframe( src_fp=il_map_fp, lowercase_cols=False, col_dtypes=MAP_SUMMARY_DTYPES, empty_data_error_msg='No summary map file found.', ) il_map_df = il_map_df[list(set(il_map_df).intersection(MAP_SUMMARY_DTYPES))] ri_summaryxref_df, ri_summary_desc = get_summary_xref_df( il_map_df, location_df, account_df, analysis_settings['ri_summaries'], 'ri' ) # Write Xref file for each inuring priority where output has been requested. # analysis_settings['ri_inuring_priorities'] contains OED InuringPriority values; we use the # mapping file (written during input generation) to convert each to the RI output level (last # RI layer index for that OED priority). ri_settings = get_ri_settings(os.path.join(model_run_fp, 'input')) ri_layers = {int(x) for x in ri_settings} inuring_priority_to_output_level = get_ri_inuring_priority_output_levels( os.path.join(model_run_fp, 'input') ) valid_oed_priorities = set(inuring_priority_to_output_level.keys()) max_oed_priority = max(valid_oed_priorities) ri_inuring_priorities_oed = set(int(p) for p in analysis_settings.get('ri_inuring_priorities', [])) ri_inuring_priorities_oed.add(max_oed_priority) # final priority always gets output if not ri_inuring_priorities_oed.issubset(valid_oed_priorities): ri_missing = ri_inuring_priorities_oed.difference(valid_oed_priorities) ri_missing = [str(p) for p in sorted(ri_missing)] missing_str = ', '.join(ri_missing[:-1]) missing_str += ' and ' * (len(ri_missing) > 1) + ri_missing[-1] missing_str = ('priority ' if len(ri_missing) == 1 else 'priorities ') + missing_str raise OasisException(f'Requested outputs for inuring {missing_str} lie outside of scope.') # Convert OED priorities to the RI output levels that should receive summary xref files. # For gross RL output every RI layer is needed; for net RI output only the last layer of # each requested OED priority is needed. if rl_summaries: ri_output_levels = ri_layers else: ri_output_levels = {inuring_priority_to_output_level[p] for p in ri_inuring_priorities_oed} for output_level in ri_output_levels: summary_ri_fp = os.path.join( model_run_fp, 'input', os.path.basename(ri_settings[str(output_level)]['directory'])) df_to_ndarray(ri_summaryxref_df, fm_summary_xref_dtype).tofile(os.path.join(summary_ri_fp, f"{SUMMARY_OUTPUT['il']}.bin")) if intermediary_csv: write_df_to_csv_file(ri_summaryxref_df, summary_ri_fp, f"{SUMMARY_OUTPUT['il']}.csv") # Write summary_id description files output_dir = os.path.join(model_run_fp, 'output') for desc_key in ri_summary_desc: if desc_key.split('.')[-1] == 'parquet': write_df_to_parquet_file(ri_summary_desc[desc_key], output_dir, desc_key) else: write_df_to_csv_file(ri_summary_desc[desc_key], output_dir, desc_key) # Copy the inuring-priority-to-output-level mapping into the output directory so it is # available alongside the results for downstream consumers. pathlib.Path(output_dir).mkdir(parents=True, exist_ok=True) shutil.copyfile( os.path.join(model_run_fp, 'input', 'ri_inuring_priority_output_levels.json'), os.path.join(output_dir, 'ri_inuring_priority_output_levels.json'), )
def code_column(name): """Name of the coded column for the exposure summary field `name`. The codes live in their own namespace because a summary field is named after an OED column, and the summary is grouped by fields while summing columns of its own. `NumberOfBuildings` is an ordinary field to want a breakdown by, and its coded values would otherwise overwrite the `number_of_buildings` being summed. Args: name (str): name of the column being coded Returns: str: the name the coded column is held under """ return f'__code__{name}' def key_column(name): """Name of the coded column for `name` in its role as one of the summary's own group keys. The summary always groups by peril id and lookup status, and separately by each summary field. A field is named after an OED column, so it can perfectly well be named `Status` -- which `convert_col_name` turns into 'status'. The two roles are coded under names of their own because sharing one column would make the group key one element shorter for such a field than for every other, and the caller's lookups would then all miss and report the bucket as empty. Args: name (str): name of the column being coded Returns: str: the name the coded column is held under """ return f'__key__{name}' def encode_exposure_summary_keys(df, field_names): """Replace the columns the exposure summary groups by with integer codes. Every deduplication and grouping hashes its key columns, and for the object columns holding peril ids, lookup statuses and OED field values that hashing dominates the summary. Coding them up front means each pass hashes integers instead. Args: df (pandas.DataFrame): dataframe `df_summary_peril` from `get_exposure_summary` field_names (iterable): names of the exposure summary fields the rows are grouped by Returns: tuple: a 2-tuple ``(encoded, codes)``, where ``encoded`` holds the coded key columns, named by `key_column` for the summary's own group keys and by `code_column` for the summary fields, alongside the columns being summed, and ``codes`` maps each coded column name to its ``{value: code}`` lookup. A value missing from a lookup does not appear in `df`. """ # taken as arrays as the frame is a concatenation, so its index is not unique encoded = {col: df[col].to_numpy() for col in ['loc_id', 'coverage_type_id', 'tiv', 'number_of_buildings', 'number_of_risks']} # a field named 'peril_id' or 'status' is coded under both its names, so that the two roles stay # separate columns and the group key keeps the same shape for every field -- see `key_column` coded_as = {col: [key_column(col)] for col in ['peril_id', 'status']} for col in field_names: coded_as.setdefault(col, []).append(code_column(col)) codes = {} for col, column_names in coded_as.items(): column_codes, values = pd.factorize(df[col], sort=False) for column_name in column_names: encoded[column_name] = column_codes codes[col] = {value: code for code, value in enumerate(values.tolist())} return pd.DataFrame(encoded, copy=False), codes def get_exposure_summary_stats(encoded, field_name, group_by_status): """Aggregate the exposure summary statistics for every value of one OED field at once. The row-wise equivalent filters the frame down to one (field value, status, coverage type) combination at a time and sums it. As the deduplication it does only ever drops rows within such a combination, deduplicating on the combined key and grouping gives the same numbers from a single pass over the frame instead of one pass per combination. The field value is part of that key, so the deduplication is per field and cannot be hoisted out and shared. Two exposure rows sharing a `loc_id` -- which a repeated `LocNumber` produces, and which nothing upstream rejects -- may hold different field values, and sharing the key across fields would drop one of them and zero its bucket. Args: encoded (pandas.DataFrame): the coded frame from `encode_exposure_summary_keys` field_name (str): name of the exposure summary field the rows are grouped by group_by_status (bool): keep one row per lookup status, rather than treating the statuses as interchangeable the way the 'all' status does Returns: tuple: a 3-tuple ``(tiv, by_coverage, overall)`` of dicts keyed by the group key, which is ``(field_value[, status][, coverage_type_id])``, dropping the enclosing tuple for a single key. ``tiv`` holds the summed TIV, ``by_coverage`` and ``overall`` the location, building and risk counts per coverage type and over all coverage types respectively. """ counts = { 'number_of_locations': ('loc_id', 'size'), 'number_of_buildings': ('number_of_buildings', 'sum'), 'number_of_risks': ('number_of_risks', 'sum'), } group_keys = [code_column(field_name)] + ([key_column('status')] if group_by_status else []) coverage_keys = group_keys + ['coverage_type_id'] # each frame is deduplicated from the one above it, which is already the first row per group tiv_rows = encoded.drop_duplicates(group_keys + ['loc_id', key_column('peril_id'), 'coverage_type_id']) coverage_rows = tiv_rows.drop_duplicates(coverage_keys + ['loc_id']) location_rows = coverage_rows.drop_duplicates(group_keys + ['loc_id']) tiv = tiv_rows.groupby(coverage_keys, observed=True, sort=False)['tiv'].sum() by_coverage = coverage_rows.groupby(coverage_keys, observed=True, sort=False).agg(**counts) overall = location_rows.groupby(group_keys, observed=True, sort=False).agg(**counts) return tiv.to_dict(), by_coverage.to_dict('index'), overall.to_dict('index') @oasis_log def get_exposure_totals(df): """Return dictionary with total TIVs and number of locations Args: df (pandas.DataFrame): dataframe `df_summary_peril` from `get_exposure_summary` Returns: dict: totals section for exposure_summary dictionary """ dedupe_cols = ['loc_id', 'coverage_type_id'] within_scope = df['status'].isin(OASIS_KEYS_STATUS_MODELLED).to_numpy() # masks rather than frames, so only the scope being summed is held in memory scope_masks = { "modelled": within_scope, "not-modelled": ~within_scope, "portfolio": np.full(len(df), True), } totals = {} for scope, mask in scope_masks.items(): scope_df = df[mask] by_location = scope_df.drop_duplicates(subset='loc_id') totals[scope] = { "tiv": scope_df.drop_duplicates(subset=dedupe_cols)['tiv'].sum(), "number_of_locations": len(by_location), "number_of_buildings": int(by_location['number_of_buildings'].sum()), "number_of_risks": int(by_location['number_of_risks'].sum()), } return totals def convert_col_name(col_name): """Convert a column from OED format to exposure summary report format. For example `CountryCode` will be converted to `country_code`. Args: col_name (str): original OED column name Returns: str: exposure summary report field name """ col_list = [col_name[0].lower()] for i, c in enumerate(col_name[1:]): if c.isupper() and not col_name[i].isupper(): col_list += ['_'] col_list += c.lower() return ''.join(col_list) @oasis_log
[docs] def get_exposure_summary( exposure_df, keys_df, exposure_profile=get_default_exposure_profile(), additional_fields=[] ): """Create exposure summary as dictionary of TIVs and number of locations grouped by peril and validity respectively. returns a python dict(). Args: exposure_df (pandas.DataFrame): source exposure dataframe keys_df (pandas.DataFrame): dataFrame holding keys data (success and errors) exposure_profile (dict): profile defining exposure file additional_fields (list): extra exposure columns to group the summary by, on top of loc_id Returns: dict: Exposure summary dictionary """ # get location tivs by coveragetype df_summary = [] exposure_fields = ['loc_id'] exposure_col_names = ['loc_id'] for field_name in additional_fields: if field_name in exposure_df.columns: exposure_fields += [field_name] exposure_col_names += [convert_col_name(field_name)] else: logger.warn(f'exposure summary field not found: {field_name}') for field_name in exposure_profile: if 'FMTermType' in exposure_profile[field_name].keys(): if exposure_profile[field_name]['FMTermType'] == 'TIV' and exposure_profile[field_name]['ProfileElementName'] in exposure_df.columns: cov_name = exposure_profile[field_name]['ProfileElementName'] coverage_type_id = exposure_profile[field_name]['CoverageTypeID'] fields = exposure_fields + [cov_name] column_names = exposure_col_names + ['tiv'] tmp_df = exposure_df[fields].copy() tmp_df.columns = column_names tmp_df['coverage_type_id'] = coverage_type_id # Add number_of_buildings column if 'NumberOfBuildings' in exposure_df: tmp_df['number_of_buildings'] = exposure_df['NumberOfBuildings'] else: tmp_df['number_of_buildings'] = 1 # Add number_of_risks column if 'IsAggregate' in exposure_df: tmp_df['number_of_risks'] = tmp_df['number_of_buildings'] tmp_df.loc[exposure_df['IsAggregate'] == 0, 'number_of_risks'] = 1 else: tmp_df['number_of_risks'] = 1 df_summary.append(tmp_df) df_summary = pd.concat(df_summary) # fix 0 number_of_buildings and risks df_summary = df_summary.replace({'number_of_buildings': 0, 'number_of_risks': 0}, 1) # get all perils peril_list = keys_df['peril_id'].drop_duplicates().to_list() # Initialise sub categories oed_categories = {'peril_id': peril_list} for field_name in exposure_col_names[1:]: # ignore 'loc_id' field_list = df_summary[field_name].drop_duplicates().to_list() oed_categories[field_name] = field_list df_summary_peril = [] for peril_id in peril_list: tmp_df = df_summary.copy() tmp_df['peril_id'] = peril_id df_summary_peril.append(tmp_df) df_summary_peril = pd.concat(df_summary_peril) df_summary_peril = df_summary_peril.merge(keys_df, how='left', on=['loc_id', 'coverage_type_id', 'peril_id']) no_return = OASIS_KEYS_STATUS['noreturn']['id'] df_summary_peril['status'] = df_summary_peril['status'].fillna(no_return) # Compile summary of exposure data exposure_summary = {} # Create totals section exposure_summary['total'] = get_exposure_totals(df_summary_peril) exposure_summary.update(get_exposure_summary_fields(df_summary_peril, oed_categories)) return exposure_summary
def get_exposure_summary_fields(df, oed_categories): """Summarise the exposure by every value of every field the summary is grouped by. Args: df (pandas.DataFrame): dataframe `df_summary_peril` from `get_exposure_summary` oed_categories (dict): the values to summarise by, per exposure summary field name Returns: dict: the exposure summary sections for the fields in `oed_categories`, holding the TIV and the location, building and risk counts per field value and lookup status """ exposure_summary = {} for field_name, field_list in oed_categories.items(): exposure_summary[field_name] = {} for value in field_list: exposure_summary[field_name][value] = {} # Create dictionary structure for all and each validity status for status in ['all'] + list(OASIS_KEYS_STATUS.keys()): exposure_summary[field_name][value][status] = {} exposure_summary[field_name][value][status]['tiv'] = 0.0 exposure_summary[field_name][value][status]['tiv_by_coverage'] = {} exposure_summary[field_name][value][status]['number_of_locations'] = 0 exposure_summary[field_name][value][status]['number_of_locations_by_coverage'] = {} exposure_summary[field_name][value][status]['number_of_buildings'] = 0 exposure_summary[field_name][value][status]['number_of_buildings_by_coverage'] = {} exposure_summary[field_name][value][status]['number_of_risks'] = 0 exposure_summary[field_name][value][status]['number_of_risks_by_coverage'] = {} coverage_ids = [(coverage_type, info['id']) for coverage_type, info in SUPPORTED_COVERAGE_TYPES.items()] encoded, codes = encode_exposure_summary_keys(df, oed_categories) for field_name, field_list in oed_categories.items(): stats_by_status = get_exposure_summary_stats(encoded, field_name, group_by_status=True) stats_for_all = get_exposure_summary_stats(encoded, field_name, group_by_status=False) for status in ['all'] + list(OASIS_KEYS_STATUS.keys()): if status == 'all': (tiv, by_coverage, overall), group_key = stats_for_all, () else: # a status absent from the frame has no code, and so matches no group (tiv, by_coverage, overall), group_key = stats_by_status, (codes['status'].get(status),) for field_value in field_list: field_summary = exposure_summary[field_name][field_value][status] field_code = codes[field_name].get(field_value) for coverage_type, coverage_type_id in coverage_ids: coverage_key = (field_code,) + group_key + (coverage_type_id,) tiv_sum = float(tiv.get(coverage_key, 0.0)) field_summary['tiv_by_coverage'][coverage_type] = tiv_sum field_summary['tiv'] += tiv_sum counts = by_coverage.get(coverage_key) for count_name in ['number_of_locations', 'number_of_buildings', 'number_of_risks']: field_summary[f'{count_name}_by_coverage'][coverage_type] = int(counts[count_name]) if counts else 0 counts = overall.get((field_code,) + group_key if group_key else field_code) for count_name in ['number_of_locations', 'number_of_buildings', 'number_of_risks']: field_summary[count_name] = int(counts[count_name]) if counts else 0 return exposure_summary @oasis_log def write_gul_errors_map( target_dir, exposure_df, keys_errors_df, exposure_profile, ): """Create csv file to help map keys errors back to original exposures. Args: target_dir (str): directory on disk to write csv file exposure_df (pandas.DataFrame): source exposure dataframe keys_errors_df (pandas.DataFrame): keys errors dataframe exposure_profile (dict): profile defining exposure file """ cols = ['loc_id', 'PortNumber', 'AccNumber', 'LocNumber', 'peril_id', 'coverage_type_id', 'tiv', 'status', 'message'] gul_error_map_fp = os.path.join(target_dir, 'gul_errors_map.csv') exposure_id_cols = ['loc_id', 'PortNumber', 'AccNumber', 'LocNumber'] keys_error_cols = ['loc_id', 'peril_id', 'coverage_type_id', 'status', 'message'] cov_level_id = SUPPORTED_FM_LEVELS['site coverage']['id'] tiv_maps = {term_info['CoverageTypeID']: term_info['ProfileElementName'] for term_info in exposure_profile.values() if (term_info.get('FMTermType') == 'TIV' and term_info.get('FMLevel') == cov_level_id and term_info['ProfileElementName'] in exposure_df.columns)} exposure_cols = list(set(exposure_id_cols + list(tiv_maps.values())).intersection(exposure_df.columns)) keys_errors_df.columns = keys_error_cols gul_inputs_errors_df = exposure_df[exposure_cols].merge(keys_errors_df[keys_error_cols], on=['loc_id']) gul_inputs_errors_df['tiv'] = 0.0 for cov_type in tiv_maps: tiv_field = tiv_maps[cov_type] gul_inputs_errors_df['tiv'] = np.where( gul_inputs_errors_df['coverage_type_id'] == cov_type, gul_inputs_errors_df[tiv_field], gul_inputs_errors_df['tiv'] ) gul_inputs_errors_df['tiv'] = gul_inputs_errors_df['tiv'].fillna(0.0) gul_inputs_errors_df[list(set(cols).intersection(gul_inputs_errors_df.columns))].to_csv(gul_error_map_fp, index=False) @oasis_log
[docs] def write_exposure_summary( target_dir, exposure_df, keys_fp, keys_errors_fp, exposure_profile, additional_fields=[] ): """Create exposure summary as dictionary of TIVs and number of locations grouped by peril and validity respectively. Writes dictionary as json file to disk. Args: target_dir (str): directory on disk to write exposure summary file exposure_df (pandas.DataFrame): source exposure dataframe keys_fp (str): file path to keys file keys_errors_fp (str): file path to keys errors file exposure_profile (dict): profile defining exposure file additional_fields (list[str]): list of additional OED fields to add to exposure summary file Returns: str: Exposure summary file path """ keys_success_df = keys_errors_df = None # get keys success if keys_fp: try: keys_success_df = get_dataframe(src_fp=keys_fp, lowercase_cols=False)[['LocID', 'PerilID', 'CoverageTypeID']] except OasisException: # Assume empty file on read error. keys_success_df = pd.DataFrame(columns=['locid']) else: keys_success_df['status'] = OASIS_KEYS_STATUS['success']['id'] keys_success_df.columns = ['loc_id', 'peril_id', 'coverage_type_id', 'status'] # get keys errors if keys_errors_fp: try: keys_errors_df = get_dataframe(src_fp=keys_errors_fp, lowercase_cols=False)[['LocID', 'PerilID', 'CoverageTypeID', 'Status', 'Message']] except OasisException: # Assume empty file on read error. keys_errors_df = pd.DataFrame(columns=['locid']) else: keys_errors_df.columns = ['loc_id', 'peril_id', 'coverage_type_id', 'status', 'message'] if not keys_errors_df.empty: write_gul_errors_map(target_dir, exposure_df, keys_errors_df, exposure_profile) # concatinate keys responses & run df_keys = pd.concat([keys_success_df, keys_errors_df]) exposure_summary = get_exposure_summary(exposure_df, df_keys, exposure_profile, additional_fields) # write exposure summary as json fileV fp = os.path.join(target_dir, 'exposure_summary_report.json') with io.open(fp, 'w', encoding='utf-8') as f: f.write(json.dumps(exposure_summary, ensure_ascii=False, indent=4)) return fp