Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion src/dve/core_engine/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -82,7 +82,9 @@ def _ensure_just_file_stem(
@property
def file_name_with_ext(self):
"""Return file name with extension."""
return f"{self.file_name}.{self.file_extension}"
if self.file_extension:
return f"{self.file_name}.{self.file_extension}"
return self.file_name

@classmethod
def from_metadata_file(cls, submission_id: str, metadata_uri: Location):
Expand Down
6 changes: 3 additions & 3 deletions src/dve/pipeline/pipeline.py
Original file line number Diff line number Diff line change
Expand Up @@ -215,10 +215,10 @@ def write_file_to_parquet(

for model_name, model in models.items():
self._logger.info(f"Transforming {model_name} to stringified parquet")
reader: BaseFileReader = load_reader(
dataset, model_name, ext, self.backend_reader_kwargs
)
try:
reader: BaseFileReader = load_reader(
dataset, model_name, ext, self.backend_reader_kwargs
)
if not entity_type:
reader.write_parquet(
reader.read_to_py_iterator(
Expand Down
35 changes: 30 additions & 5 deletions src/dve/pipeline/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,9 +11,11 @@
import dve.core_engine.backends.implementations.duckdb # pylint: disable=unused-import
import dve.core_engine.backends.implementations.spark # pylint: disable=unused-import
import dve.parser.file_handling as fh
from dve.core_engine.backends.exceptions import MessageBearingError
from dve.core_engine.backends.readers import _READER_REGISTRY
from dve.core_engine.configuration.v1 import SchemaName, V1EngineConfig, _ModelConfig
from dve.core_engine.loggers import get_logger
from dve.core_engine.message import FeedbackMessage
from dve.core_engine.type_hints import URI, SubmissionResult
from dve.metadata_parser.model_generator import JSONtoPyd

Expand Down Expand Up @@ -52,11 +54,34 @@ def load_reader(
backend_reader_kwargs: Optional[dict[str, Any]] = None,
):
"""Loads the readers for the diven feed, model name and file extension"""
reader_config = dataset[model_name].reader_config[f".{file_extension.lower()}"]
reader = _READER_REGISTRY[reader_config.reader](
**reader_config.kwargs_, **backend_reader_kwargs if backend_reader_kwargs else {}
)
return reader
try:
reader_config = dataset[model_name].reader_config[f".{file_extension.lower()}"]
reader = _READER_REGISTRY[reader_config.reader](
**reader_config.kwargs_, **backend_reader_kwargs if backend_reader_kwargs else {}
)
return reader
except KeyError as exc:
if file_extension:
err_msg = (
f"The supplied file extension `{file_extension if file_extension else None}`"
+f" is not a supported file format for {model_name}."
)
else:
err_msg = "No supplied file extension. Unable to parse file without a file extension."

raise MessageBearingError(
"The file extension provided is not supported.",
messages=[
FeedbackMessage(
entity=model_name,
record=None,
failure_type="submission",
error_location="Whole File",
error_code="InvalidFileExtension",
error_message=err_msg,
)
],
) from exc


def unpersist_all_rdds(spark: SparkSession):
Expand Down
4 changes: 3 additions & 1 deletion tests/features/planets.feature
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,9 @@ Feature: Pipeline tests using the planets dataset
And I add initial audit entries for the submission
Then the latest audit record for the submission is marked with processing status file_transformation
When I run the file transformation phase
Then the latest audit record for the submission is marked with processing status failed
Then the latest audit record for the submission is marked with processing status error_report
When I run the error report phase
Then An error report is produced

Scenario: Handle a file with duplicated extension provided (spark)
Given I submit the planets file planets.csv.csv for processing
Expand Down
36 changes: 36 additions & 0 deletions tests/test_pipeline/test_pipeline_utils.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
from dve.core_engine.backends.exceptions import MessageBearingError
from dve.core_engine.configuration.v1 import _ModelConfig, _ReaderConfig
from dve.pipeline.utils import load_reader

import pytest


class TestLoadReader:
test_model_config = _ModelConfig(
fields={"test": "str"},
reporting_fields=["test"],
key_field="test",
reader_config={
".csv": _ReaderConfig(reader="TestCsvReader"),
}
)

def test_invalid_load_reader_with_file_ext(self):
with pytest.raises(MessageBearingError) as exc_info:
load_reader(
{"test": self.test_model_config},
"test_model",
"jpeg"
)

assert exc_info.value.messages[0].error_message == "The supplied file extension `jpeg` is not a supported file format for test_model."

def test_invalid_load_reader_missing_file_ext(self):
with pytest.raises(MessageBearingError) as exc_info:
load_reader(
{"test": self.test_model_config},
"test_model",
""
)

assert exc_info.value.messages[0].error_message == "No supplied file extension. Unable to parse file without a file extension."
Loading