From 9c5de944729dc12def6029e8080630d0bbe6faf7 Mon Sep 17 00:00:00 2001 From: Alex Galarce Date: Tue, 29 Jul 2025 22:05:51 +0000 Subject: [PATCH] Merged in feature/improved-error-out-logging (pull request #637) Feature/improved error out logging * Improve error logging in file processing and JSON parsing * Enhance JSON parsing error handling for CSV compatibility * black formatting * isort * Merge remote-tracking branch 'origin/main' into feature/improved-error-out-logging Approved-by: Katon Minhas --- fieldExtraction/src/investment/main.py | 79 ++++++++++++++++------- fieldExtraction/src/utils/string_utils.py | 38 +++++++---- 2 files changed, 83 insertions(+), 34 deletions(-) diff --git a/fieldExtraction/src/investment/main.py b/fieldExtraction/src/investment/main.py index 6d97612..b68b220 100644 --- a/fieldExtraction/src/investment/main.py +++ b/fieldExtraction/src/investment/main.py @@ -1,20 +1,23 @@ - import concurrent.futures +import logging import os import random import traceback -import pandas as pd from datetime import datetime +import pandas as pd + random.seed(42) +# from src.tracking.batch_tracking import BatchTracker +import traceback + import src.investment.file_processing as file_processing import src.tracking.tracking_utils as tracking_utils import src.utils.io_utils as io_utils import src.utils.string_utils as string_utils from src import config -# from src.tracking.batch_tracking import BatchTracker -import traceback + def safe_process_file(item, all_dataset, run_timestamp): """call process_file and catch any exceptions that occur @@ -27,41 +30,71 @@ def safe_process_file(item, all_dataset, run_timestamp): Returns: pd.DataFrame: DataFrame containing the results or an error message """ + file_id = ( + item[0] if isinstance(item, (list, tuple)) and len(item) > 0 else item + ) # unpack file_id from item (passed in as a tuple below) try: return file_processing.process_file(item, all_dataset, run_timestamp) - except Exception as e: # When there's an issue with the processing inside the future - print(f"Error processing file {item[0]}: {str(e)}") - print(traceback.format_exc()) - return pd.DataFrame([{"error": str(e), "file_id": item[0]}]) # Return a single-row dataframe so we can still concat it later + except ( + Exception + ) as e: # When there's an issue with the processing inside the future + error_type = type(e).__name__ + error_message = str(e) + full_traceback = traceback.format_exc() + + logging.error(f"Error processing file {file_id}: {error_message}") + logging.error(f"Error type: {error_type}") + logging.error(f"Error Message: {error_message}") + + logging.error(f"Full traceback:\n{full_traceback}") + return pd.DataFrame( + [ + { + "file_id": file_id, + "error": error_message.replace(",", ";").replace( + "\n", " " + ), # Replace problematic characters for CSV compatibility + "error_type": error_type, + "traceback": full_traceback.replace(",", ";").replace( + "\n", " | " + ), # keep line breaks as separators + } + ] + ) # Return a single-row dataframe so we can still concat it later -def main(testing=False, test_params = {}): +def main(testing=False, test_params={}): run_timestamp = datetime.now().strftime(f"run_%Y%m%d_%H:%M_{config.BATCH_ID}") - + # Read and process input if not testing: input_dict = io_utils.read_input() max_workers = config.MAX_WORKERS - else: - input_dict = io_utils.read_input(test_params['local_input_dir']) - max_workers = test_params['max_workers'] - if 'input_files' in test_params: - input_dict = {k: v for k, v in input_dict.items() if k in test_params['input_files']} - + else: + input_dict = io_utils.read_input(test_params["local_input_dir"]) + max_workers = test_params["max_workers"] + if "input_files" in test_params: + input_dict = { + k: v for k, v in input_dict.items() if k in test_params["input_files"] + } + print(f"Total Input Files : {len(input_dict)}") # load all crosswalk files, mappings and embeddings all_dataset = io_utils.load_all_dataset() io_utils.load_embeddings() - + # Process files with run_timestamp successful_results = [] error_results = [] with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as executor: # use executor to process files concurrently; use submit() instead of map() # so that we can catch exceptions and continue processing other files - futures = [executor.submit(safe_process_file, item, all_dataset, run_timestamp) for item in input_dict.items()] - + futures = [ + executor.submit(safe_process_file, item, all_dataset, run_timestamp) + for item in input_dict.items() + ] + # collect results as they complete for future in concurrent.futures.as_completed(futures): try: @@ -70,11 +103,13 @@ def main(testing=False, test_params = {}): error_results.append(results) else: successful_results.append(results) - except Exception as e: # When there's an issue with the future itself (not the processing inside the future) + except ( + Exception + ) as e: # When there's an issue with the future itself (not the processing inside the future) print(f"Error processing future: {str(e)}") # Coerce this to the format expected by pd.concat (so a single-row dataframe) error_results.append(pd.DataFrame([{"error": str(e)}])) - + # only concat if we have results if len(successful_results) > 0: FINAL_RESULT_DF = pd.concat(successful_results) @@ -85,7 +120,7 @@ def main(testing=False, test_params = {}): ERROR_RESULT_DF = pd.concat(error_results) else: ERROR_RESULT_DF = pd.DataFrame() - + if not testing: if config.WRITE_TO_S3: io_utils.write_s3(FINAL_RESULT_DF, "", run_timestamp, "final") diff --git a/fieldExtraction/src/utils/string_utils.py b/fieldExtraction/src/utils/string_utils.py index 272cae0..33239eb 100644 --- a/fieldExtraction/src/utils/string_utils.py +++ b/fieldExtraction/src/utils/string_utils.py @@ -6,15 +6,13 @@ from datetime import datetime import numpy as np import pandas as pd + import src.utils.llm_utils as llm_utils from src import config from src.enums.delimiters import Delimiter from src.prompts import preprocessing_prompts -from src.regex.regex_patterns import ( - BACKTICK_PATTERN, - PIPE_PATTERN, - TRIPLE_BACKTICK_PATTERN, -) +from src.regex.regex_patterns import (BACKTICK_PATTERN, PIPE_PATTERN, + TRIPLE_BACKTICK_PATTERN) def extract_text_from_delimiters( @@ -274,19 +272,26 @@ def universal_json_load(string_dict: str): # Try to match dictionary or list of dictionaries first dict_matches = re.findall(r"\{.*\}|\[\s*?\{.*\}\s*?\]", string_dict, re.DOTALL) if dict_matches: - matched_str = dict_matches[-1] # Take the last match in the string + matched_str = dict_matches[-1] # Take the last match in the string try: return json.loads(matched_str) except json.JSONDecodeError as e: print(f"Failed to parse JSON: {e}") print(f"Original string: {string_dict}") print(f"Traceback: {traceback.format_exc()}") - raise ValueError(f"JSON parsing failed: {e.msg} at position {e.pos}") from e - + # Use placeholders to preserve JSON structure characters for CSV compatibility + csv_safe_string = ( + string_dict.replace(",", "[COMMA]") + .replace("\n", "[NEWLINE]") + .replace('"', "[QUOTE]") + ) + raise ValueError( + f"JSON parsing failed: {e.msg} at position {e.pos}; Original string:\n{csv_safe_string}" + ) from e # If no dict pattern found, try to match list of strings list_matches = re.findall(r"\[(.*?)\]", string_dict, re.DOTALL) if list_matches: - matched_str = f"[{list_matches[-1]}]" # Take the last match in the string + matched_str = f"[{list_matches[-1]}]" # Take the last match in the string try: # Clean up the string list before parsing cleaned_str = re.sub(r'(? str: # If no match found, return original value return state_value + def flatten_to_strings(items) -> list[str]: """Helper function to flatten and stringify any nested structure @@ -608,5 +622,5 @@ def flatten_to_strings(items) -> list[str]: result.append(str(item).strip()) else: result.append(str(items).strip()) - - return result \ No newline at end of file + + return result