Skip to main content

easyfabric.loaders.table_utils

spark_operation_with_retries​

def spark_operation_with_retries(operation_func,
config_manager: ConfigManager = None,
operation_name: str = "Spark Operation")

Executes a Spark/Delta operation with retry logic for concurrency conflicts.

Arguments:

  • operation_func - A callable that performs the Spark/Delta operation.
  • config_manager ConfigManager - Configuration manager for retry settings.
  • operation_name str - Name of the operation for logging.

refresh_table​

def refresh_table(table_config: TableConfig,
config_manager: ConfigManager = None,
layer: str = "bronze",
return_table: bool = False,
history: bool = False) -> DataFrame | None

Refresh the metadata of a Spark table and optionally return it as a DataFrame, with retry logic on failure.

Arguments:

  • table_config TableConfig - Table configuration object.
  • config_manager ConfigManager - Configuration manager for environment/layer.
  • layer str - Data layer name.
  • return_table bool, optional - If True, return the refreshed Spark DataFrame. Defaults to False.
  • history bool, optional - If True, get historical version of table. Defaults to False.

Returns:

DataFrame | None: Spark DataFrame if return_table is True, else None.

truncate_bronze_table​

def truncate_bronze_table(table_config: TableConfig,
config_manager: ConfigManager = None)

Deletes all data from a specified bronze lakehouseA place where you store both "raw" data (like files) and "organized" data (like tables). It combines the best of a File Cabinet and a Database. table by using Truncate and refreshes the table cache.

Arguments:

  • table_config TableConfig - Configuration object containing the table details.
  • config_manager ConfigManager - Configuration manager to retrieve lakehouse and schema details.

Raises:

  • Exception - If the TRUNCATE TABLE operation fails.

delete_silver​

def delete_silver(table_config: TableConfig,
config_manager: ConfigManager = None)

Deletes all data from the specified silver table and refreshes the cache to ensure the latest state of the table is visible for subsequent operations. This function retrieves the silver lakehouse configuration, constructs the necessary table name and performs deletion using Spark.

Arguments:

  • table_config TableConfig - Configuration object that provides table-related information such as table names. It is used to determine the silver table name linked to the provided configuration.
  • config_manager ConfigManager - Configuration manager instance that provides access to various lakehouse configurations, including the silver lakehouse details required for this operation.

Raises:

  • Exception - If the TRUNCATE TABLE operation fails, wrapping the original error.

align_filter_columns​

def align_filter_columns(df_filter, target_columns)

Rename the filter query's columns to the target's spelling and return the join keys.

Column names are matched case-insensitively, so a filter query written as SELECT 'x' AS CustomerName joins on a bronze column stored as customername.

Raises:

  • ValueError - When the filter query shares no column with the target.

apply_filter_query​

def apply_filter_query(df_target, df_filter)

Keep only the target rows that match the filter query on their shared columns.

The result keeps the target's own columns in their original order, so positional writes downstream are unaffected by the join.

DateConversion Objects​

@dataclass(frozen=True)
class DateConversion()

A generated date parse, whichever form it takes, split into its parts.

measure_expression​

@property
def measure_expression() -> str

The same parse, in the one form that cannot abort the job.

try_to_timestamp yields NULL for a value the pattern cannot read, also with spark.sql.ansi.enabled. It does not escape one case: under spark.sql.legacy.timeParserPolicy EXCEPTION, a value the pre-3.0 parser would still have read -- a date with a time behind it, a day without its leading zero -- raises SparkUpgradeException here as it does in the conversion. report_date_conversion_losses turns that into a stop that names column, pattern and value.

Measuring a to_date with try_to_timestamp is the same question: both run the value through one formatter, and a pattern that cannot read a value fails for both. Only the parse matters here, not the result type.

DateConversionLoss Objects​

@dataclass(frozen=True)
class DateConversionLoss()

A typed column whose every filled source value converted to NULL.

parse_date_conversion​

def parse_date_conversion(expression: str) -> Optional[DateConversion]

Split a generated date expression into its source column and pattern.

None voor al het andere: zonder patroon valt er niets te melden.

Een date-doel komt binnen als cast(try_to_timestamp(...) as date). De parse zit in de kern, dus de cast eromheen gaat er eerst af.

date_conversion_aggregations​

def date_conversion_aggregations(
conversions: list[DateConversion]) -> list[tuple[str, str]]

Build the (alias, SQL) pairs that measure every conversion in one pass.

date_conversion_losses​

def date_conversion_losses(
measurements: dict,
conversions: list[DateConversion]) -> list[DateConversionLoss]

Keep the conversions that emptied an entirely filled column.

measure_date_conversions​

def measure_date_conversions(
df: DataFrame, aggregations: list[tuple[str, str]]) -> Optional[dict]

Run the one aggregation, or report that it could not run and give up.

A measurement is diagnostics: it may cost a scan, it may find nothing, but it may never be the reason a load fails. Whatever Spark raises here -- a runtime without try_to_timestamp, a column type that will not cast, a failing executor -- the load continues with exactly the behaviour it had before the measurement existed. None says so; an empty dict would read as measured.

One error is passed on: a value Spark refuses to parse under timeParserPolicy EXCEPTION. The conversion raises it as well, so the load fails on it regardless; passing it on lets the caller name the column.

unparseable_value​

def unparseable_value(error: Exception) -> Optional[str]

The value Spark refused under timeParserPolicy EXCEPTION, if that is the error.

Spark raises this for a value the pre-3.0 parser would have read and the current one cannot -- a date with a time behind it, a day without its leading zero -- because the two would disagree on what it means. The conversion raises it too, so the load cannot get past it either way.

conversion_rejecting​

def conversion_rejecting(
df: DataFrame,
conversions: list[DateConversion]) -> Optional[DateConversion]

The conversion whose pattern raised, found by parsing one column at a time.

stop_on_unparseable_value​

def stop_on_unparseable_value(df: DataFrame, conversions: list[DateConversion],
error: Exception) -> None

Turn Spark's parser complaint into one that names column, pattern and value.

report_date_conversion_losses​

def report_date_conversion_losses(
df: DataFrame,
expressions: list[tuple[str, str]],
violation_policy: Optional[ViolationPolicy] = None) -> None

Report every date conversion that turns a filled column entirely to NULL.

Meet op de bron-dataframe: een expressie overschrijft de kolom die ze leest.

date_conversion_policy​

def date_conversion_policy(table_config: TableConfig) -> ViolationPolicy

The policy that decides what a fully emptied date column costs.

apply_silver_expressions​

def apply_silver_expressions(
df: DataFrame,
expressions: list[tuple[str, str]],
violation_policy: Optional[ViolationPolicy] = None) -> DataFrame

Apply every SilverExpression, after measuring what the date ones cost.