from pyspark.sql import functions as F
# ===========================
# Catalog & Schema Names
# ===========================
CATALOG = "ecommerce"
BRONZE_SCHEMA = "bronze"
SILVER_SCHEMA = "silver"
# ===========================
# Helper Functions
# ===========================
def bronze_table(table_name):
return f"{CATALOG}.{BRONZE_SCHEMA}.{table_name}"
def silver_table(table_name):
return f"{CATALOG}.{SILVER_SCHEMA}.{table_name}"
def show_distinct_values(df, column_name): """ Displays all distinct values from the specified column. """ df.select(column_name).distinct().show(truncate=False)
def replace_column_values(df, column_name=None, replacement_dict=None): """ Replaces values in the specified column using a dictionary. Parameters: df (DataFrame): Input DataFrame. column_name (str): Column in which values need to be replaced. replacement_dict (dict): Dictionary containing old and new values. Returns: DataFrame: Updated DataFrame. """ return df.replace( to_replace=replacement_dict, subset=[column_name] )
def write_delta_table(df, catalog_name, schema_name, table_name,
overwriteSchema='overwriteSchema', mode="overwrite", overwrite_schema=True):
"""
Writes a DataFrame to a Delta table.
Parameters:
df (DataFrame): Input DataFrame.
catalog_name (str): Unity Catalog name.
schema_name (str): Schema name (bronze/silver/gold).
table_name (str): Target table name.
mode (str): Write mode (overwrite, append, etc.).
overwrite_schema (bool): Whether to overwrite the schema.
Returns:
None
"""
(
df.write
.format("delta")
.mode(mode)
.option(overwriteSchema, str(overwrite_schema).lower())
.saveAsTable(f"{catalog_name}.{schema_name}.{table_name}")
)
def read_delta_table(spark, catalog_name, schema_name, table_name):
"""
Reads a Delta table and returns it as a DataFrame.
"""
return spark.table(f"{catalog_name}.{schema_name}.{table_name}")
def find_duplicates(df, columns):
"""
Finds duplicate records based on one or more columns.
Parameters:
df (DataFrame): Input DataFrame.
columns (list): List of column names to check for duplicates.
Returns:
DataFrame: Duplicate values with their occurrence count.
"""
return (
df.groupBy(columns)
.count()
.filter(F.col("count") > 1)
)