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_managerConfigManager - Configuration manager for retry settings.operation_namestr - 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_configTableConfig - Table configuration object.config_managerConfigManager - Configuration manager for environment/layer.layerstr - Data layer name.return_tablebool, optional - If True, return the refreshed Spark DataFrame. Defaults to False.historybool, 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_configTableConfig - Configuration object containing the table details.config_managerConfigManager - 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_configTableConfig - 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_managerConfigManager - 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.