Skip to content
Closed
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
42 changes: 30 additions & 12 deletions src/dve/core_engine/backends/base/reader.py
Original file line number Diff line number Diff line change
@@ -1,23 +1,23 @@
"""Abstract implementation of the file parser."""

from abc import ABC, abstractmethod
from collections.abc import Iterator
from inspect import ismethod
from typing import Any, ClassVar, Optional, TypeVar
from typing import Any, ClassVar, Iterator, Optional, TypeVar

from pydantic import BaseModel
from typing_extensions import Protocol

from dve.core_engine.backends.exceptions import (
CriticalMessageBearingError,
MessageBearingError,
ReaderLacksEntityTypeSupport
ReaderLacksEntityTypeSupport,
)
from dve.core_engine.backends.types import EntityName, EntityType
from dve.core_engine.configuration.v1 import (
AllowedAdditionalReaderChecks,
_ReaderAdditionalChecksConfig,
)
from dve.core_engine.constants import RECORD_INDEX_COLUMN_NAME
from dve.core_engine.message import FeedbackMessage
from dve.core_engine.type_hints import URI, ArbitraryFunction, WrapDecorator
from dve.parser.file_handling.service import open_stream
Expand Down Expand Up @@ -62,6 +62,7 @@ class BaseFileReader(ABC):
"""An abstract representation of a reader for some file type."""

__read_methods__: ClassVar[_ReadFunctions] = {}

"""
A dictionary mapping implemented entity types to their read functions.

Expand Down Expand Up @@ -93,7 +94,6 @@ class variable for the subclass.
entity_type: Optional[type] = getattr(method, _ENTITY_TYPE_ATTR_NAME, None)
if entity_type is None:
continue

cls.__read_methods__[entity_type] = method # type: ignore

@abstractmethod
Expand Down Expand Up @@ -137,8 +137,10 @@ def read_to_entity_type(
self.raise_if_not_sensible_file(resource, entity_name)

if entity_type == Iterator[dict[str, Any]]:
entity = self.read_to_py_iterator(
resource, entity_name, schema, all_model_fields # type: ignore
entity = self.filter_null_records_py_iterator(
self.read_to_py_iterator(
resource, entity_name, schema, all_model_fields # type: ignore
)
)

else:
Expand All @@ -148,12 +150,14 @@ def read_to_entity_type(
except KeyError as err:
raise ReaderLacksEntityTypeSupport(entity_type=entity_type) from err

entity = reader_func(
self,
resource,
entity_name,
schema,
all_model_fields=all_model_fields, # type: ignore
entity = self.filter_null_records(
reader_func(
self,
resource,
entity_name,
schema,
all_model_fields=all_model_fields, # type: ignore
)
)

if config := additional_checks.get("check_empty"):
Expand Down Expand Up @@ -235,3 +239,17 @@ def raise_if_not_sensible_file(
error_message=self.ft_error_message,
),
)

def filter_null_records_py_iterator(
self, records: Iterator[dict[str, Any]]
) -> Iterator[dict[str, Any]]:
"""Strip null records from py iterator"""

def _is_non_null_record(record: dict[str, Any]) -> bool:
return any(v is not None for k, v in record.items() if k != RECORD_INDEX_COLUMN_NAME)

yield from filter(_is_non_null_record, records)

def filter_null_records(self, entity: EntityType) -> EntityType:
"""Strip null records from entity"""
raise NotImplementedError()
5 changes: 4 additions & 1 deletion src/dve/core_engine/backends/exceptions.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ def __init__(self, *args: object, messages: Messages) -> None:
self.messages = messages
"""The messages to be returned as part of the error."""


class CriticalMessageBearingError(BackendError):
"""
A backend error that comes with a pre-created message.
Expand All @@ -44,6 +45,7 @@ def __init__(self, *args: object, message: FeedbackMessage) -> None:
self.message = message
"""The message to be returned as part of the error."""


class UnableToParseCSVError(CriticalMessageBearingError):
"""An error raised when unable to parse a CSV file"""

Expand All @@ -60,7 +62,8 @@ def __init__(
failure_type="submission",
is_informational=False,
error_type="csv read",
error_message=error_message or "Unable to parse the CSV file. Please check the structure of your CSV.", # pylint: disable=C0301
error_message=error_message
or "Unable to parse the CSV file. Please check the structure of your CSV.", # pylint: disable=C0301
error_code=error_code or "MalformedCSV",
)
)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,11 @@

from dve.common.error_utils import get_feedback_errors_uri
from dve.core_engine.backends.base.utilities import _get_non_heterogenous_type
from dve.core_engine.backends.utilities import DEFAULT_ISO_FORMATS, datetime_format_to_regex
from dve.core_engine.backends.utilities import (
DEFAULT_ISO_FORMATS,
datetime_format_to_regex,
polars_filter_null_records,
)
from dve.core_engine.constants import RECORD_INDEX_COLUMN_NAME
from dve.core_engine.type_hints import URI, EntityName
from dve.metadata_parser.utilities import resilient_get
Expand Down Expand Up @@ -500,3 +504,14 @@
stmt = f"TRIM({quoted_name})"
return _cast_as_ddb_type(stmt, type_) if parent_element else stmt
raise ValueError(f"No equivalent DuckDB type for {type_annotation!r}")


def _ddb_filter_null_records(self, entity: DuckDBPyRelation):
df = polars_filter_null_records(entity.pl()) # pylint: disable=W0612

Check warning on line 510 in src/dve/core_engine/backends/implementations/duckdb/duckdb_helpers.py

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Remove the unused local variable "df".

See more on https://sonarcloud.io/project/issues?id=NHSDigital_data-validation-engine&issues=AaEhJb5WRP6Qm1DADDxu&open=AaEhJb5WRP6Qm1DADDxu&pullRequest=180
return self._connection.sql("SELECT * from df")


def duckdb_filter_null_recs(cls):
"""Add method to class to filter null records for duckdb relations"""
setattr(cls, "filter_null_records", _ddb_filter_null_records)
return cls
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
)
from dve.core_engine.backends.implementations.duckdb.duckdb_helpers import (
duckdb_check_entity_empty,
duckdb_filter_null_recs,
duckdb_record_index,
duckdb_write_parquet,
get_duckdb_type_from_annotation,
Expand All @@ -37,6 +38,7 @@
from dve.parser.file_handling import get_content_length


@duckdb_filter_null_recs
@duckdb_check_entity_empty
@duckdb_record_index
@duckdb_write_parquet
Expand Down Expand Up @@ -123,7 +125,8 @@ def read_to_relation( # pylint: disable=unused-argument
raise UnableToParseCSVError(
entity_name="csv_structure",
error_code=self.ft_error_code,
error_message=self.ft_error_message or "Unable to parse CSV file. Structure is likely malformed.", # pylint: disable=C0301
error_message=self.ft_error_message
or "Unable to parse CSV file. Structure is likely malformed.", # pylint: disable=C0301
) from exc

if self.null_empty_strings:
Expand All @@ -135,6 +138,7 @@ def read_to_relation( # pylint: disable=unused-argument
return rel


@duckdb_check_entity_empty
@polars_record_index
class PolarsToDuckDBCSVReader(DuckDBCSVReader):
"""
Expand Down Expand Up @@ -183,7 +187,8 @@ def read_to_relation( # pylint: disable=unused-argument
raise UnableToParseCSVError(
entity_name="csv_structure",
error_code=self.ft_error_code,
error_message=self.ft_error_message or "Unable to parse CSV file. Structure is likely malformed.", # pylint: disable=C0301
error_message=self.ft_error_message
or "Unable to parse CSV file. Structure is likely malformed.", # pylint: disable=C0301
) from exc

if self.null_empty_strings:
Expand All @@ -200,7 +205,8 @@ def read_to_relation( # pylint: disable=unused-argument
raise UnableToParseCSVError(
entity_name="csv_structure",
error_code=self.ft_error_code,
error_message=self.ft_error_message or "Found zero records after loading CSV. File is likely malformed.", # pylint: disable=C0301
error_message=self.ft_error_message
or "Found zero records after loading CSV. File is likely malformed.", # pylint: disable=C0301
)

return entity
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
from dve.core_engine.backends.base.reader import BaseFileReader, read_function
from dve.core_engine.backends.implementations.duckdb.duckdb_helpers import (
duckdb_check_entity_empty,
duckdb_filter_null_recs,
duckdb_record_index,
duckdb_write_parquet,
get_duckdb_type_from_annotation,
Expand All @@ -19,6 +20,7 @@
from dve.core_engine.type_hints import URI, EntityName


@duckdb_filter_null_recs
@duckdb_check_entity_empty
@duckdb_record_index
@duckdb_write_parquet
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@
from dve.core_engine.backends.exceptions import CriticalMessageBearingError
from dve.core_engine.backends.implementations.duckdb.duckdb_helpers import (
duckdb_check_entity_empty,
duckdb_filter_null_recs,
duckdb_write_parquet,
)
from dve.core_engine.backends.readers.xml import XMLStreamReader
Expand All @@ -23,6 +24,7 @@
from dve.core_engine.type_hints import URI


@duckdb_filter_null_recs
@duckdb_check_entity_empty
@polars_record_index
@duckdb_write_parquet
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
from dve.core_engine.backends.implementations.spark.spark_helpers import (
get_type_from_annotation,
spark_check_entity_empty,
spark_filter_null_recs,
spark_record_index,
spark_write_parquet,
)
Expand All @@ -21,6 +22,7 @@
from dve.parser.file_handling import get_content_length


@spark_filter_null_recs
@spark_check_entity_empty
@spark_record_index
@spark_write_parquet
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -12,13 +12,15 @@
from dve.core_engine.backends.implementations.spark.spark_helpers import (
get_type_from_annotation,
spark_check_entity_empty,
spark_filter_null_recs,
spark_record_index,
spark_write_parquet,
)
from dve.core_engine.type_hints import URI, EntityName
from dve.parser.file_handling import get_content_length


@spark_filter_null_recs
@spark_check_entity_empty
@spark_record_index
@spark_write_parquet
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
df_is_empty,
get_type_from_annotation,
spark_check_entity_empty,
spark_filter_null_recs,
spark_record_index,
spark_write_parquet,
)
Expand All @@ -30,6 +31,7 @@
"""The mode to use when parsing XML files with Spark."""


@spark_filter_null_recs
@spark_check_entity_empty
@spark_record_index
@spark_write_parquet
Expand Down Expand Up @@ -57,6 +59,7 @@ def read_to_dataframe(
)


@spark_filter_null_recs
@spark_check_entity_empty
@spark_record_index
@spark_write_parquet
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -638,3 +638,17 @@ def get_spark_cast_statement_from_annotation(
stmt = f"TRIM({quoted_name})"
return _cast_as_spark_type(stmt, type_) if parent_element else stmt
raise ValueError(f"No equivalent Spark type for {type_annotation!r}")


def _spark_filter_null_records(self, entity: DataFrame) -> DataFrame: # pylint: disable=W0613
return entity.dropna(
how="all", subset=[cl for cl in entity.columns if cl != RECORD_INDEX_COLUMN_NAME]
)


def spark_filter_null_recs(cls):
"""Class decorator to add spark method for filtering records where values
are all null (aside from record_index).
"""
setattr(cls, "filter_null_records", _spark_filter_null_records)
return cls
2 changes: 1 addition & 1 deletion src/dve/core_engine/backends/readers/utilities.py
Original file line number Diff line number Diff line change
Expand Up @@ -66,7 +66,7 @@ def raise_message_bearing_error_on_header_differences(
reporting_field="csv_header",
error_code=field_check_error_code,
error_message=field_check_error_message,
)
),
)


Expand Down
9 changes: 9 additions & 0 deletions src/dve/core_engine/backends/utilities.py
Original file line number Diff line number Diff line change
Expand Up @@ -272,3 +272,12 @@ def polars_record_index(cls):
setattr(cls, "add_record_index", _add_polars_record_index)
setattr(cls, "drop_record_index", _drop_polars_record_index)
return cls


def polars_filter_null_records(entity: pl.DataFrame) -> pl.DataFrame:
"""Strip records where all values (aside from record index) are null"""
return entity.filter(
~pl.all_horizontal(
pl.col([cl for cl in entity.columns if cl != RECORD_INDEX_COLUMN_NAME]).is_null()
)
)
24 changes: 24 additions & 0 deletions tests/test_core_engine/test_backends/fixtures.py
Original file line number Diff line number Diff line change
Expand Up @@ -589,4 +589,28 @@ def nested_parquet_custom_dc_err_details(temp_dir):

yield file_path

@pytest.fixture
def temp_json_file_w_null_recs(temp_dir: Path):

class SimpleModel(BaseModel):
varchar_field: str
bigint_field: int
date_field: date
timestamp_field: datetime

field_names: list[str] = ["varchar_field","bigint_field","date_field","timestamp_field"]
typed_data = [
["hi", 1, date(2023, 1, 3), datetime(2023, 1, 3, 12, 0, 3)],
[None, None, None, None],
["bye", 3, date(2023, 3, 7), datetime(2023, 5, 9, 15, 21, 53)],
[None, None, None, None]
]

test_data = [dict(zip(field_names, rw)) for rw in typed_data]

with open(temp_dir.joinpath("test.json"), mode="w") as json_file:
json.dump(test_data, json_file, default=str)

yield temp_dir.joinpath("test.json"), test_data, SimpleModel


Loading
Loading