Skip to content

Commit cb36f39

Browse files
refactor: add additional logging
also moved READER_CONFIG out of the try except block as it's setup error of the DVE rather than user/submitter error.
1 parent 93cb81a commit cb36f39

2 files changed

Lines changed: 8 additions & 6 deletions

File tree

src/dve/pipeline/pipeline.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -241,6 +241,7 @@ def write_file_to_parquet(
241241
f"{out}{model_name}",
242242
)
243243
except MessageBearingError as exc:
244+
self._logger.error(f"Unable to process {model_name}", exc_info=exc)
244245
errors.extend(exc.messages)
245246

246247
return list(dict.fromkeys(errors)) # remove any duplicate errors

src/dve/pipeline/utils.py

Lines changed: 7 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -56,21 +56,17 @@ def load_reader(
5656
"""Loads the readers for the diven feed, model name and file extension"""
5757
try:
5858
reader_config = dataset[model_name].reader_config[f".{file_extension.lower()}"]
59-
reader = _READER_REGISTRY[reader_config.reader](
60-
**reader_config.kwargs_, **backend_reader_kwargs if backend_reader_kwargs else {}
61-
)
62-
return reader
6359
except KeyError as exc:
6460
if file_extension:
6561
err_msg = (
66-
f"The supplied file extension `{file_extension if file_extension else None}`"
62+
f"The supplied file extension `{file_extension}`"
6763
+f" is not a supported file format for {model_name}."
6864
)
6965
else:
7066
err_msg = "No supplied file extension. Unable to parse file without a file extension."
7167

7268
raise MessageBearingError(
73-
"The file extension provided is not supported.",
69+
f"The file extension provided ({file_extension}) is not supported for this collection.",
7470
messages=[
7571
FeedbackMessage(
7672
entity=model_name,
@@ -83,6 +79,11 @@ def load_reader(
8379
],
8480
) from exc
8581

82+
reader = _READER_REGISTRY[reader_config.reader](
83+
**reader_config.kwargs_, **backend_reader_kwargs if backend_reader_kwargs else {}
84+
)
85+
return reader
86+
8687

8788
def unpersist_all_rdds(spark: SparkSession):
8889
"""Unpersist any checkpointed or cached rdds to avoid memory leaks"""

0 commit comments

Comments
 (0)