diff --git a/scripts/world_bank/wdi/README.md b/scripts/world_bank/wdi/README.md index ef1f0f5dc6..8bab931dd6 100644 --- a/scripts/world_bank/wdi/README.md +++ b/scripts/world_bank/wdi/README.md @@ -146,5 +146,18 @@ If you want to perform "only download", run the below command: python3 worldbank.py --mode=download ``` +### Historical Data Merge and Import Validation + +During processing (`--mode=process`), `worldbank.py` merges deleted historical +observations from GCS (`--historical_gcs_path`, defaulting to +`gs://unresolved_mcf/world_bank/wdi/deleted_rows_07_2026.csv`) with fresh World +Bank API data. Records are deduplicated across composite keys +(`StatisticalVariable`, `ISO3166Alpha3`, `Year`, `observationPeriod`, `unit`, +`measurementMethod`, `scalingFactor`), prioritizing fresh observations +(`keep='first'`). + +Automated import validation rules (deleted records threshold and `MaxDate` +freshness check) are configured in `validation_config.json`. + We highly recommend the use of the import validation tool for this import which you can find in https://github.com/datacommonsorg/tools/tree/master/import-validation-helper. diff --git a/scripts/world_bank/wdi/manifest.json b/scripts/world_bank/wdi/manifest.json index bc3927141e..35edea215f 100644 --- a/scripts/world_bank/wdi/manifest.json +++ b/scripts/world_bank/wdi/manifest.json @@ -20,7 +20,8 @@ "WorldBankCountries.csv", "schema_csvs/WorldBankIndicators_prod.csv" ], - "cron_schedule": "0 11 * * 2" + "cron_schedule": "0 11 * * 2", + "validation_config_file": "validation_config.json" } ] -} \ No newline at end of file +} diff --git a/scripts/world_bank/wdi/validation_config.json b/scripts/world_bank/wdi/validation_config.json new file mode 100644 index 0000000000..607e87a958 --- /dev/null +++ b/scripts/world_bank/wdi/validation_config.json @@ -0,0 +1,22 @@ +{ + "schema_version": "1.0", + "rules": [ + { + "rule_id": "check_deleted_records_percent", + "description": "Checks that deleted observation points do not exceed 0.1%.", + "validator": "DELETED_RECORDS_PERCENT", + "params": { + "threshold": 0.1 + } + }, + { + "rule_id": "check_max_date_freshness", + "description": "Verifies that latest MaxDate is within 3 years of current year for active StatVars.", + "validator": "SQL_VALIDATOR", + "params": { + "query": "SELECT StatVar, CAST(MaxDate AS INT) AS latest_year FROM stats WHERE StatVar NOT IN ('Count_Person_ResidingLessThan5MetersAboveSeaLevel_AsFractionOf_Count_Person', 'Count_Person_7To14Years_Employed_AsFractionOf_Count_Person_7To14Years', 'Count_Person_7To14Years_Female_Employed_AsFractionOf_Count_Person_7To14Years_Female', 'Count_Person_7To14Years_Male_Employed_AsFractionOf_Count_Person_7To14Years_Male', 'Amount_EconomicActivity_ExpenditureActivity_TertiaryEducationExpenditure_Government_AsFractionOf_Amount_EconomicActivity_ExpenditureActivity_EducationExpenditure_Government', 'Amount_Consumption_Alcohol_15OrMoreYears_AsFractionOf_Count_Person_15OrMoreYears', 'Count_Death_IntentionalSelfHarm_AsFractionOf_Count_Person', 'Count_Death_IntentionalSelfHarm_Female_AsFractionOf_Count_Person_Female', 'Count_Death_IntentionalSelfHarm_Male_AsFractionOf_Count_Person_Male', 'Amount_Consumption_RenewableEnergy_AsFractionOf_Amount_Consumption_Energy')", + "condition": "latest_year >= (EXTRACT(YEAR FROM CURRENT_DATE) - 3)" + } + } + ] +} diff --git a/scripts/world_bank/wdi/worldbank.py b/scripts/world_bank/wdi/worldbank.py index d864bfeaa2..5f274116c7 100644 --- a/scripts/world_bank/wdi/worldbank.py +++ b/scripts/world_bank/wdi/worldbank.py @@ -39,6 +39,10 @@ "indicatorSchemaFile", os.path.join(_MODULE_DIR, "schema_csvs/WorldBankIndicators_prod.csv"), "") flags.DEFINE_string('mode', '', 'Options: download or process') +flags.DEFINE_string( + 'historical_gcs_path', + 'gs://unresolved_mcf/world_bank/wdi/deleted_rows_07_2026.csv', + 'GCS path to the deleted historical data CSV file') # Remaps the columns provided by World Bank API. WORLDBANK_COL_REMAP = { @@ -472,7 +476,9 @@ def output_csv_and_tmcf_by_grouping(worldbank_dataframe, if saveOutput: TMCF_PATH = 'output/WorldBank.tmcf' else: - TMCF_PATH = 'test_data/output/output_generated.tmcf' + TMCF_PATH = os.path.join(_MODULE_DIR, + 'test_data/output/output_generated.tmcf') + os.makedirs(os.path.dirname(TMCF_PATH), exist_ok=True) with open(TMCF_PATH, 'w', newline='') as f_out: for index, enum in enumerate(tmcfs_for_stat_vars): tmcf, stat_var_obs_cols, stat_vars_in_group = enum @@ -504,15 +510,42 @@ def output_csv_and_tmcf_by_grouping(worldbank_dataframe, df = df.replace({'StatisticalVariable': RESOLUTION_TO_EXISTING_DCID}) if saveOutput: logging.info("Writing output csv") - df.drop('IndicatorCode', axis=1).to_csv('output/WorldBank.csv', - float_format='%.10f', - index=False) + output_file_path = 'output/WorldBank.csv' + final_df = merge_historical_data(df.drop('IndicatorCode', axis=1), + _FLAGS.historical_gcs_path) + final_df.to_csv(output_file_path, float_format='%.10f', index=False) else: return df except Exception as e: logging.fatal(f"Error generating output {e}") +def merge_historical_data(df, historical_gcs_path): + """Merges and deduplicates historical deleted data from GCS into df.""" + if not historical_gcs_path: + return df + try: + composite_keys = [ + 'StatisticalVariable', 'ISO3166Alpha3', 'Year', 'observationPeriod', + 'unit', 'measurementMethod', 'scalingFactor' + ] + deleted_df = retry_call( + pd.read_csv, + fargs=[historical_gcs_path], + fkwargs={'dtype': { + k: str for k in composite_keys + }}, + tries=3, + delay=5, + backoff=2) + df = pd.concat([df, deleted_df], ignore_index=True) + for col in composite_keys: + df[col] = df[col].fillna('').astype(str).str.removesuffix('.0') + return df.drop_duplicates(subset=composite_keys, keep='first') + except Exception as e: + logging.fatal(f"Could not read historical deleted data from GCS: {e}") + + def source_scaling_remap(row, scaling_factor_lookup, existing_stat_var_lookup): """ Scales values by sourceScalingFactor and inputs exisiting stat vars. @@ -543,6 +576,8 @@ def source_scaling_remap(row, scaling_factor_lookup, existing_stat_var_lookup): def process(indicator_codes, worldbank_dataframe, saveOutput=True): logging.info("Processing the input files") try: + os.makedirs('output', exist_ok=True) + # Add source description to note. def add_source_to_description(row): if not pd.isna(row['Source']): diff --git a/scripts/world_bank/wdi/worldbank_test.py b/scripts/world_bank/wdi/worldbank_test.py index b76d1a3ea4..2de1edbfd3 100644 --- a/scripts/world_bank/wdi/worldbank_test.py +++ b/scripts/world_bank/wdi/worldbank_test.py @@ -20,7 +20,7 @@ OUTPUT_PATH = "test_data/output" if not os.path.exists( os.path.join(_MODULE_DIR, OUTPUT_PATH, "output_generated.csv")): - os.mkdir(os.path.join(_MODULE_DIR, OUTPUT_PATH)) + os.makedirs(os.path.join(_MODULE_DIR, OUTPUT_PATH), exist_ok=True) GENERATED_CSV_PATH = os.path.join(_MODULE_DIR, OUTPUT_PATH, "output_generated.csv") GENERATED_TMCF_PATH = os.path.join(_MODULE_DIR, OUTPUT_PATH, @@ -77,6 +77,33 @@ def test_WDI(self): self.assertEqual(expected_tmcf_data.strip(), generated_tmcf_data.strip()) + def test_merge_historical_data(self): + fresh_df = pd.DataFrame({ + 'StatisticalVariable': ['dcid:SV1', 'dcid:SV2'], + 'ISO3166Alpha3': ['dcid:country/USA', 'dcid:country/IRQ'], + 'Year': ['2020', '1991'], + 'observationPeriod': ['P1Y', 'P1Y'], + 'Value0': [10.5, 200.0], + 'unit': ['USD', 'USD'], + 'measurementMethod': ['', ''], + 'scalingFactor': ['100', ''] + }) + historical_csv = io.StringIO( + "StatisticalVariable,ISO3166Alpha3,Year,observationPeriod,Value0," + "unit,measurementMethod,scalingFactor\n" + "dcid:SV2,dcid:country/IRQ,1991,P1Y,150.0,USD,,\n" + "dcid:SV2,dcid:country/IRQ,1992,P1Y,175.0,USD,,100\n") + merged = merge_historical_data(fresh_df, historical_csv) + self.assertEqual(len(merged), 3) + # Verify fresh observation (200.0) takes precedence over historical (150.0) + irq_1991 = merged[(merged['ISO3166Alpha3'] == 'dcid:country/IRQ') & + (merged['Year'] == '1991')] + self.assertEqual(float(irq_1991['Value0'].iloc[0]), 200.0) + # Verify scalingFactor remains string '100' without float decimal formatting + csv_out = merged.to_csv(float_format='%.10f', index=False) + self.assertIn(',100\n', csv_out) + self.assertNotIn('100.0000000000', csv_out) + if __name__ == '__main__': unittest.main()