From ad41c13127aa2cd3f9ec92ee22de8faf220adcab Mon Sep 17 00:00:00 2001 From: Ashwani Srivastav Date: Mon, 7 Sep 2026 09:44:42 +0000 Subject: [PATCH 1/5] adding deletion threshold --- .../oecd/regional_education/manifest.json | 4 +++- .../oecd/regional_education/validation_config.json | 13 +++++++++++++ 2 files changed, 16 insertions(+), 1 deletion(-) create mode 100644 statvar_imports/oecd/regional_education/validation_config.json diff --git a/statvar_imports/oecd/regional_education/manifest.json b/statvar_imports/oecd/regional_education/manifest.json index 520c86b868..db31e8de6c 100644 --- a/statvar_imports/oecd/regional_education/manifest.json +++ b/statvar_imports/oecd/regional_education/manifest.json @@ -15,10 +15,12 @@ "import_inputs": [ { "template_mcf": "output/oecd_regional_education.tmcf", - "cleaned_csv": "output/oecd_regional_education.csv" + "cleaned_csv": "output/oecd_regional_education.csv", + "node_mcf": "output/*.mcf" } ], "cron_schedule": "0 10 1,15 * *", + "validation_config_file": "validation_config.json", "source_files": [ "gcs_output/source_files/*.csv" ] diff --git a/statvar_imports/oecd/regional_education/validation_config.json b/statvar_imports/oecd/regional_education/validation_config.json new file mode 100644 index 0000000000..f03e4c8ebf --- /dev/null +++ b/statvar_imports/oecd/regional_education/validation_config.json @@ -0,0 +1,13 @@ +{ + "schema_version": "1.0", + "rules": [ + { + "rule_id": "check_deleted_records_percent", + "description": "Allow up to 5% deleted records due to OECD dataflow 2.5 NUTS 2024 regional restructuring and historical series revisions.", + "validator": "DELETED_RECORDS_PERCENT", + "params": { + "threshold": 4 + } + } + ] +} From 277154dcc6b00571ced8c8dd64c4fbc8cfbb99e8 Mon Sep 17 00:00:00 2001 From: Ashwani Srivastav Date: Mon, 7 Sep 2026 09:56:46 +0000 Subject: [PATCH 2/5] adding deletion threshold --- statvar_imports/oecd/regional_education/validation_config.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/statvar_imports/oecd/regional_education/validation_config.json b/statvar_imports/oecd/regional_education/validation_config.json index f03e4c8ebf..15b4ef71dc 100644 --- a/statvar_imports/oecd/regional_education/validation_config.json +++ b/statvar_imports/oecd/regional_education/validation_config.json @@ -3,7 +3,7 @@ "rules": [ { "rule_id": "check_deleted_records_percent", - "description": "Allow up to 5% deleted records due to OECD dataflow 2.5 NUTS 2024 regional restructuring and historical series revisions.", + "description": "Allow up to 4% deleted records due to OECD dataflow 2.5 NUTS 2024 regional restructuring and historical series revisions.", "validator": "DELETED_RECORDS_PERCENT", "params": { "threshold": 4 From e6dcb36e2113bdf58881def0e7d97ea30d2a7228 Mon Sep 17 00:00:00 2001 From: Ashwani Srivastav Date: Mon, 7 Sep 2026 10:29:30 +0000 Subject: [PATCH 3/5] adding checks --- .../regional_education/validation_config.json | 18 ++++++++++++++++-- 1 file changed, 16 insertions(+), 2 deletions(-) diff --git a/statvar_imports/oecd/regional_education/validation_config.json b/statvar_imports/oecd/regional_education/validation_config.json index 15b4ef71dc..3dcdd01ba5 100644 --- a/statvar_imports/oecd/regional_education/validation_config.json +++ b/statvar_imports/oecd/regional_education/validation_config.json @@ -3,10 +3,24 @@ "rules": [ { "rule_id": "check_deleted_records_percent", - "description": "Allow up to 4% deleted records due to OECD dataflow 2.5 NUTS 2024 regional restructuring and historical series revisions.", + "description": "Allow up to 5% deleted records due to OECD dataflow 2.5 NUTS 2024 regional restructuring and historical series revisions.", "validator": "DELETED_RECORDS_PERCENT", "params": { - "threshold": 4 + "threshold": 5 + } + }, + { + "rule_id": "check_max_date_consistent", + "description": "Ensure MaxDate is uniform across all StatVars", + "validator": "MAX_DATE_CONSISTENT" + }, + { + "rule_id": "check_max_date_freshness", + "description": "Ensure MaxDate is fresh within allowable 2-year reporting lag for OECD annual data", + "validator": "SQL_VALIDATOR", + "params": { + "query": "SELECT StatVar, MaxDate FROM stats", + "condition": "CAST(MaxDate AS INTEGER) >= date_part('year', current_date) - 2" } } ] From 244270cba117940b9bfa18f20a01c9c3c33ed29e Mon Sep 17 00:00:00 2001 From: Ashwani Srivastav Date: Mon, 7 Sep 2026 11:00:12 +0000 Subject: [PATCH 4/5] update check condition --- statvar_imports/oecd/regional_education/validation_config.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/statvar_imports/oecd/regional_education/validation_config.json b/statvar_imports/oecd/regional_education/validation_config.json index 3dcdd01ba5..b54bfda467 100644 --- a/statvar_imports/oecd/regional_education/validation_config.json +++ b/statvar_imports/oecd/regional_education/validation_config.json @@ -20,7 +20,7 @@ "validator": "SQL_VALIDATOR", "params": { "query": "SELECT StatVar, MaxDate FROM stats", - "condition": "CAST(MaxDate AS INTEGER) >= date_part('year', current_date) - 2" + "condition": "CAST(SUBSTRING(CAST(MaxDate AS VARCHAR), 1, 4) AS INTEGER) >= date_part('year', current_date) - 2" } } ] From 18db5cd12909b5a5bce7a90523cc8b2b613b76e0 Mon Sep 17 00:00:00 2001 From: Ashwani Srivastav Date: Wed, 9 Sep 2026 13:01:20 +0000 Subject: [PATCH 5/5] adding mcf file --- .../oecd/regional_education/manifest.json | 5 +- .../oecd_regional_education_custom_schema.mcf | 4 + .../oecd/regional_education/preprocess.py | 173 +++++++++++++----- 3 files changed, 132 insertions(+), 50 deletions(-) create mode 100644 statvar_imports/oecd/regional_education/oecd_regional_education_custom_schema.mcf diff --git a/statvar_imports/oecd/regional_education/manifest.json b/statvar_imports/oecd/regional_education/manifest.json index db31e8de6c..6b95ca6db4 100644 --- a/statvar_imports/oecd/regional_education/manifest.json +++ b/statvar_imports/oecd/regional_education/manifest.json @@ -10,7 +10,7 @@ "scripts": [ "../../../util/download_util_script.py --download_url='https://sdmx.oecd.org/public/rest/data/OECD.CFE.EDS,DSD_REG_EDU@DF_ATTAIN,/A.........?dimensionAtObservation=AllDimensions&format=csvfilewithlabels' --output_folder=gcs_output/source_files", "preprocess.py", - "../../../tools/statvar_importer/stat_var_processor.py --input_data=gcs_output/source_files/oecd_regional_education_data.csv --pv_map=oecd_regional_education_pvmap.csv --config_file=oecd_regional_education_metadata.csv --places_resolved_csv=oecd_regional_education_places_resolved.csv --existing_statvar_mcf=gs://unresolved_mcf/scripts/statvar/stat_vars.mcf --output_path=output/oecd_regional_education" + "../../../tools/statvar_importer/stat_var_processor.py --input_data=gcs_output/source_files/oecd_regional_education_data.csv --pv_map=oecd_regional_education_pvmap.csv --config_file=oecd_regional_education_metadata.csv --places_resolved_csv=oecd_regional_education_places_resolved.csv --existing_statvar_mcf=gs://unresolved_mcf/scripts/statvar/stat_vars.mcf --output_path=output/oecd_regional_education --output_counters=counters/oecd_regional_education_counters.csv" ], "import_inputs": [ { @@ -22,7 +22,8 @@ "cron_schedule": "0 10 1,15 * *", "validation_config_file": "validation_config.json", "source_files": [ - "gcs_output/source_files/*.csv" + "gcs_output/source_files/*.csv", + "counters/*.csv" ] } ] diff --git a/statvar_imports/oecd/regional_education/oecd_regional_education_custom_schema.mcf b/statvar_imports/oecd/regional_education/oecd_regional_education_custom_schema.mcf new file mode 100644 index 0000000000..8d9bb338e6 --- /dev/null +++ b/statvar_imports/oecd/regional_education/oecd_regional_education_custom_schema.mcf @@ -0,0 +1,4 @@ +Node: dcid:PostSecondaryNonTertiaryEducation__UpperSecondaryEducation +typeOf: dcs:SchoolGradeLevelEnum +name: "PostSecondaryNonTertiaryEducation__UpperSecondaryEducation" +description: "Upper secondary and post-secondary non-tertiary education (ISCED 2011 levels 3 and 4)." diff --git a/statvar_imports/oecd/regional_education/preprocess.py b/statvar_imports/oecd/regional_education/preprocess.py index b8e4405230..77b43f5b1d 100644 --- a/statvar_imports/oecd/regional_education/preprocess.py +++ b/statvar_imports/oecd/regional_education/preprocess.py @@ -1,55 +1,132 @@ +import csv import os import re -from absl import logging +import shutil -# --- Add this line to set verbosity --- -logging.set_verbosity(logging.INFO) -# For even more detail if you have debug messages: -# logging.set_verbosity(logging.DEBUG) -# -------------------------------------- +try: + from absl import logging + logging.set_verbosity(logging.INFO) +except ImportError: + import logging as std_logging -def rename_target_file(base_path='.'): + class _CompatLogger: + def __init__(self): + self._logger = std_logging.getLogger(__name__) + self._logger.setLevel(std_logging.INFO) + if not self._logger.handlers: + handler = std_logging.StreamHandler() + handler.setFormatter( + std_logging.Formatter('%(levelname)s:%(message)s')) + self._logger.addHandler(handler) + + def set_verbosity(self, level): + self._logger.setLevel(level) + + def info(self, msg, *args, **kwargs): + self._logger.info(msg, *args, **kwargs) + + def warning(self, msg, *args, **kwargs): + self._logger.warning(msg, *args, **kwargs) + + def error(self, msg, *args, **kwargs): + self._logger.error(msg, *args, **kwargs) + + logging = _CompatLogger() + logging.set_verbosity(std_logging.INFO) + + +def preprocess(base_path='.'): folder_name = 'gcs_output/source_files' target_folder = os.path.join(base_path, folder_name) + counters_folder = os.path.join(base_path, 'counters') + os.makedirs(counters_folder, exist_ok=True) + output_folder = os.path.join(base_path, 'output') + os.makedirs(output_folder, exist_ok=True) + + custom_schema_file = os.path.join( + base_path, 'oecd_regional_education_custom_schema.mcf') + if os.path.isfile(custom_schema_file): + shutil.copyfile( + custom_schema_file, + os.path.join(output_folder, 'oecd_regional_education_custom_schema.mcf')) + logging.info(f"Copied custom schema to {output_folder}") + + places_resolved_file = os.path.join( + base_path, 'oecd_regional_education_places_resolved.csv') + valid_places = set() + if os.path.isfile(places_resolved_file): + with open(places_resolved_file, 'r', encoding='utf-8') as f: + reader = csv.DictReader(f) + for row in reader: + if row.get('dcid', '').strip(): + valid_places.add(row['place_name'].strip()) + logging.info(f"Loaded {len(valid_places)} valid places from {places_resolved_file}") + else: + logging.warning(f"Places resolved file not found: {places_resolved_file}") + + if not os.path.isdir(target_folder): + logging.error(f"Folder '{folder_name}' not found in '{base_path}'") + return + + pattern = re.compile(r'^A.*$', re.IGNORECASE) + raw_file = None + for filename in os.listdir(target_folder): + if pattern.match(filename): + raw_file = filename + break + + target_csv = os.path.join(target_folder, 'oecd_regional_education_data.csv') + + if raw_file: + src_path = os.path.join(target_folder, raw_file) + tmp_path = os.path.join(target_folder, 'filtered_tmp.csv') + logging.info(f"Filtering '{raw_file}' into 'oecd_regional_education_data.csv'...") + _filter_csv(src_path, tmp_path, valid_places) + if os.path.exists(target_csv): + os.remove(target_csv) + os.rename(tmp_path, target_csv) + if src_path != target_csv and os.path.exists(src_path): + os.remove(src_path) + logging.info("Preprocessing and filtering completed successfully.") + elif os.path.isfile(target_csv) and valid_places: + tmp_path = os.path.join(target_folder, 'filtered_tmp.csv') + logging.info(f"Checking and filtering existing '{target_csv}'...") + _filter_csv(target_csv, tmp_path, valid_places) + os.replace(tmp_path, target_csv) + logging.info("Filtering completed successfully.") + else: + logging.info("No matching source data file found to process.") + + +def _filter_csv(src_path: str, dst_path: str, valid_places: set): + with open(src_path, 'r', encoding='utf-8', errors='replace') as fin, \ + open(dst_path, 'w', encoding='utf-8', newline='') as fout: + reader = csv.reader(fin) + writer = csv.writer(fout) + + header = next(reader, None) + if not header: + return + writer.writerow(header) + + ref_area_idx = header.index('REF_AREA') if 'REF_AREA' in header else None + if ref_area_idx is None: + logging.warning("REF_AREA column not found in header, copying all rows.") + for row in reader: + writer.writerow(row) + return + + kept = 0 + dropped = 0 + for row in reader: + if len(row) > ref_area_idx and row[ref_area_idx].strip() in valid_places: + writer.writerow(row) + kept += 1 + else: + dropped += 1 + + logging.info(f"Filtered source data: {kept} rows kept, {dropped} rows with unresolved places dropped.") + - try: - # Check if the folder exists - if not os.path.isdir(target_folder): - logging.fatal(f"Folder '{folder_name}' not found in '{base_path}'") # Changed to error for non-fatal issues - return # Exit function if folder not found - - # Pattern to match any file starting with 'A' - pattern = re.compile(r'^A.*$', re.IGNORECASE) - renamed = False - - # Search through files in the folder - for filename in os.listdir(target_folder): - logging.info(f"Checking file: {filename}") - if pattern.match(filename): - old_path = os.path.join(target_folder, filename) - new_path = os.path.join(target_folder, 'oecd_regional_education_data.csv') - - try: - os.rename(old_path, new_path) - logging.info(f"Renamed '{filename}' to 'oecd_regional_education_data.csv'") - renamed = True - except PermissionError: - logging.warning(f"Permission denied while renaming '{filename}'.") # Changed to warning - except OSError as e: - logging.fatal(f"OS error while renaming '{filename}': {e}") # Changed to error - break # Rename only the first match - - if not renamed: - logging.info("No matching file starting with 'A' found to rename.") - - except FileNotFoundError as e: - # This block might not be hit if caught earlier, but good for other FileNotFoundError - logging.fatal(f"File system error: {e}") - except PermissionError as e: - # This block might not be hit if caught earlier, but good for other PermissionError - logging.fatal(f"Global permission error: {e}") - except Exception as e: - logging.fatal(f"An unexpected critical error occurred: {e}") # Changed to critical for unexpected errors - -# Run it -rename_target_file() \ No newline at end of file +if __name__ == '__main__': + preprocess() \ No newline at end of file