|
11 | 11 | import dve.core_engine.backends.implementations.duckdb # pylint: disable=unused-import |
12 | 12 | import dve.core_engine.backends.implementations.spark # pylint: disable=unused-import |
13 | 13 | import dve.parser.file_handling as fh |
| 14 | +from dve.core_engine.backends.exceptions import MessageBearingError |
14 | 15 | from dve.core_engine.backends.readers import _READER_REGISTRY |
15 | 16 | from dve.core_engine.configuration.v1 import SchemaName, V1EngineConfig, _ModelConfig |
16 | 17 | from dve.core_engine.loggers import get_logger |
| 18 | +from dve.core_engine.message import FeedbackMessage |
17 | 19 | from dve.core_engine.type_hints import URI, SubmissionResult |
18 | 20 | from dve.metadata_parser.model_generator import JSONtoPyd |
19 | 21 |
|
@@ -52,11 +54,34 @@ def load_reader( |
52 | 54 | backend_reader_kwargs: Optional[dict[str, Any]] = None, |
53 | 55 | ): |
54 | 56 | """Loads the readers for the diven feed, model name and file extension""" |
55 | | - reader_config = dataset[model_name].reader_config[f".{file_extension.lower()}"] |
56 | | - reader = _READER_REGISTRY[reader_config.reader]( |
57 | | - **reader_config.kwargs_, **backend_reader_kwargs if backend_reader_kwargs else {} |
58 | | - ) |
59 | | - return reader |
| 57 | + try: |
| 58 | + 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 |
| 63 | + except KeyError as exc: |
| 64 | + if file_extension: |
| 65 | + err_msg = ( |
| 66 | + f"The supplied file extension `{file_extension if file_extension else None}`" |
| 67 | + +f" is not a supported file format for {model_name}." |
| 68 | + ) |
| 69 | + else: |
| 70 | + err_msg = "No supplied file extension. Unable to parse file without a file extension." |
| 71 | + |
| 72 | + raise MessageBearingError( |
| 73 | + "The file extension provided is not supported.", |
| 74 | + messages=[ |
| 75 | + FeedbackMessage( |
| 76 | + entity=model_name, |
| 77 | + record=None, |
| 78 | + failure_type="submission", |
| 79 | + error_location="Whole File", |
| 80 | + error_code="InvalidFileExtension", |
| 81 | + error_message=err_msg, |
| 82 | + ) |
| 83 | + ], |
| 84 | + ) from exc |
60 | 85 |
|
61 | 86 |
|
62 | 87 | def unpersist_all_rdds(spark: SparkSession): |
|
0 commit comments