From 131dff77dec4e72382f2eb3429ad9ca32a61eec6 Mon Sep 17 00:00:00 2001 From: Alex Galarce Date: Mon, 24 Mar 2025 18:56:50 +0000 Subject: [PATCH] Merged in feature/split-out-error-log (pull request #450) DRAFT: Feature/split out error log * Refactor result handling to separate successful and error results in processing * Add forced exception for testing in safe_process_file function * Merged main into feature/split-out-error-log * main.py edited online with Bitbucket * Add support for error output type in write_s3 function Approved-by: Katon Minhas --- fieldExtraction/src/investment/main.py | 25 +++++++++++++++++++------ fieldExtraction/src/utils/io_utils.py | 5 +++++ 2 files changed, 24 insertions(+), 6 deletions(-) diff --git a/fieldExtraction/src/investment/main.py b/fieldExtraction/src/investment/main.py index 20d2341..9ed3873 100644 --- a/fieldExtraction/src/investment/main.py +++ b/fieldExtraction/src/investment/main.py @@ -53,7 +53,8 @@ def main(testing=False, test_params = {}): io_utils.load_embeddings() # Process files with run_timestamp - all_results = [] + 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 @@ -63,25 +64,37 @@ def main(testing=False, test_params = {}): for future in concurrent.futures.as_completed(futures): try: results = future.result() - all_results.append(results) + if "error" in results.columns: + 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) print(f"Error processing future: {str(e)}") # Coerce this to the format expected by pd.concat (so a single-row dataframe) - all_results.append(pd.DataFrame([{"error": str(e)}])) + error_results.append(pd.DataFrame([{"error": str(e)}])) # only concat if we have results - if len(all_results) > 0: - FINAL_RESULT_DF = pd.concat(all_results) + if len(successful_results) > 0: + FINAL_RESULT_DF = pd.concat(successful_results) else: FINAL_RESULT_DF = pd.DataFrame() + + if len(error_results) > 0: + 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") + if not ERROR_RESULT_DF.empty: + io_utils.write_s3(ERROR_RESULT_DF, "", run_timestamp, "error") else: io_utils.write_local(FINAL_RESULT_DF, "", run_timestamp, "final") + if not ERROR_RESULT_DF.empty: + io_utils.write_local(ERROR_RESULT_DF, "", run_timestamp, "error") else: - return FINAL_RESULT_DF + return FINAL_RESULT_DF, ERROR_RESULT_DF if __name__ == "__main__": diff --git a/fieldExtraction/src/utils/io_utils.py b/fieldExtraction/src/utils/io_utils.py index 04bc8d7..9361b70 100644 --- a/fieldExtraction/src/utils/io_utils.py +++ b/fieldExtraction/src/utils/io_utils.py @@ -247,6 +247,11 @@ def write_s3(df, filename, run_timestamp, output_type): output_path = f"{config.BATCH_ID}/{run_timestamp}/{config.BATCH_ID}-RESULTS.csv" elif output_type == "individual": output_path = f"{config.BATCH_ID}/{run_timestamp}/individual/{filename}-RESULTS.csv" + elif output_type == "error": + output_path = f"{config.BATCH_ID}/{run_timestamp}/{config.BATCH_ID}-ERRORS.csv" + else: + # Fallback path for unknown output types + output_path = f"{config.BATCH_ID}/{run_timestamp}/{config.BATCH_ID}-unknown-{output_type}.csv" csv_buffer = df.to_csv(index=False).encode() try: