Documentación de `gss_bi_udfs`

Referencia funcional de los módulos utils, io, transforms y merges. Cada función incluye definición, parámetros, retorno y ejemplo de uso.

Módulo `utils`

set_spark

set_spark(spark)

Configura una sesión Spark global para uso interno del módulo.

ParámetroTipoDescripción
sparkSparkSessionSesión Spark a registrar globalmente.
Retorna: None.
from pyspark.sql import SparkSession
from gss_bi_udfs import utils

spark = SparkSession.builder.getOrCreate()
utils.set_spark(spark)

get_spark

get_spark(spark=None)

Obtiene Spark en este orden: parámetro explícito, sesión configurada por set_spark, builtins.spark (Databricks), o builder.getOrCreate().

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional explícita.
Retorna: SparkSession activa.

get_env

get_env(default="dev")

Lee el entorno de ejecución desde la variable ENV.

ParámetroTipoDescripción
defaultstrValor por defecto si ENV no existe.
Retorna: str con el entorno (ej. dev, qa, pro).

get_env_catalog

get_env_catalog(catalog)

Ajusta el catálogo según ambiente. Si ENV=pro devuelve el catálogo base; en otros casos agrega sufijo.

ParámetroTipoDescripción
catalogstrNombre base del catálogo.
Retorna: str con catálogo ajustado por ambiente.

get_env_table_path

get_env_table_path(catalog, table_path)

Compone el identificador final de tabla: catálogo ajustado + ruta de tabla.

ParámetroTipoDescripción
catalogstrCatálogo base.
table_pathstrRuta esquema.tabla.
Retorna: str con ruta completa (catalogo_env.schema.tabla).

get_schema_root_location

get_schema_root_location(spark=None, catalog=None, schema=None)

Consulta DESCRIBE SCHEMA EXTENDED y extrae RootLocation.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional.
catalogstrCatálogo base.
schemastrEsquema dentro del catálogo.
Retorna: str con path físico del esquema.

get_table_info

get_table_info(spark=None, *, full_table_name=None, catalog=None, schema=None, table=None)

Resuelve información integral de tabla (nombre, path, existencia y metadatos de provider/type).

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional.
full_table_namestr | NoneFormato catalog.schema.table.
catalog/schema/tablestr | NoneAlternativa al parámetro anterior.
Retorna: dict con catalog, schema, table, full_table_name, path, exists, provider, table_type.
info = utils.get_table_info(full_table_name="fi_comunes.silver.dim_cliente")
print(info["full_table_name"], info["path"], info["exists"])

get_default_value_by_type

get_default_value_by_type(dtype)

Devuelve un valor default de Spark (Column) según tipo de dato.

ParámetroTipoDescripción
dtypepyspark.sql.types.DataTypeTipo de dato del campo.
Retorna: Column con default (-999, "N/A", False, 1900-01-01 o null).

Módulo `io`

load_latest_parquet

load_latest_parquet(spark=None, data_base=None, schema=None, table=None, env=None)

Carga el parquet más reciente para la tabla solicitada desde /Volumes/bronze/{db}_{schema}/{env}/{table}/.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional.
data_basestrNombre de base lógica.
schemastrNombre de esquema.
tablestrNombre de tabla.
envstr | NoneAmbiente; si no se pasa, usa ENV o dev.
Retorna: DataFrame de Spark o None si no encuentra archivos.
df_clientes = io.load_latest_parquet(
    spark=spark,
    data_base="timepro",
    schema="insudb",
    table="clientes",
    env="dev"
)

return_parquets_and_register_temp_views

return_parquets_and_register_temp_views(spark=None, tables_load=None, verbose=False, env=None)

Lee parquets según configuración y crea vistas temporales. Devuelve los DataFrames en un diccionario.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional.
tables_loaddictMapa de carga por base/esquema/lista de tablas y vistas.
verboseboolImprime mensajes de estado.
envstr | NoneAmbiente de lectura.
Retorna: dict con claves db.schema.table y valores DataFrame.

parquets_register_temp_views

parquets_register_temp_views(spark=None, tables_load=None, verbose=False, env=None)

Misma carga de parquets que la función anterior, pero sin devolver diccionario.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional.
tables_loaddictConfiguración de tablas/vistas.
verboseboolMensajes de materialización.
envstr | NoneAmbiente de lectura.
Retorna: None.

load_latest_excel

load_latest_excel(spark=None, source_file=None, env=None)

Carga el último archivo Excel desde /Volumes/bronze/excel/{env}/{source_file}/ y lo convierte a DataFrame Spark.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional.
source_filestrRuta relativa del archivo dentro de bronze/excel.
envstr | NoneAmbiente de lectura.
Retorna: DataFrame de Spark o None si falla/no hay archivos.

return_excels_and_register_temp_views

return_excels_and_register_temp_views(spark=None, files_load=None, verbose=False, env=None)

Carga Excels configurados, crea vistas temporales y devuelve diccionario de DataFrames.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional.
files_loaddictConfiguración dominio/subdominio/lista de archivos-vistas.
verboseboolMensajes de materialización.
envstr | NoneAmbiente de lectura.
Retorna: dict con claves dominio.subdominio.archivo y valores DataFrame.

excels_register_temp_views

excels_register_temp_views(spark=None, files_load=None, verbose=False, env=None)

Carga Excels y materializa vistas temporales sin retornar estructura adicional.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional.
files_loaddictConfiguración de archivos y vistas.
verboseboolMensajes de materialización.
envstr | NoneAmbiente de lectura.
Retorna: None.

load_and_materialize_views

load_and_materialize_views(action, **kwargs)

Dispatcher para acciones de carga: parquets/excels con o sin retorno.

ParámetroTipoDescripción
actionstrNombre de acción soportada.
kwargsdictParámetros de la acción seleccionada.
Retorna: dict resultante de la acción; {} si no existe.
tables_load = {
    "timepro": {"insudb": [{"table": "clientes", "view": "vw_clientes"}]}
}

out = io.load_and_materialize_views(
    action="return_parquets_and_register_temp_views",
    spark=spark,
    tables_load=tables_load,
    env="dev"
)

save_table_to_delta

save_table_to_delta(spark=None, df=None, catalog=None, schema=None, table_name=None)

Guarda un DataFrame en formato Delta con overwrite y registración en metastore.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional.
dfDataFrameDataset de entrada a persistir.
catalogstrCatálogo destino base.
schemastrEsquema destino.
table_namestrNombre de tabla destino.
Retorna: None.
io.save_table_to_delta(
    spark=spark,
    df=df_clientes,
    catalog="fi_comunes",
    schema="silver",
    table_name="dim_cliente"
)

Módulo `transforms`

add_hashid

add_hashid(df, columns, new_col_name="hashid")

Concatena columnas, calcula hash con xxhash64 y ubica la nueva PK técnica al inicio del DataFrame.

ParámetroTipoDescripción
dfDataFrameDataFrame de entrada.
columnslist[str]Columnas para construir el hash.
new_col_namestrNombre de la nueva columna hash.
Retorna: DataFrame con columna hash agregada y reordenada.
df_hash = transforms.add_hashid(
    df=df_clientes,
    columns=["id_cliente", "tipo_doc", "nro_doc"],
    new_col_name="sk_cliente"
)

get_default_record

get_default_record(spark, df)

Genera un único registro default usando el esquema de referencia.

ParámetroTipoDescripción
sparkSparkSessionSesión Spark activa.
dfDataFrameDataFrame cuyo esquema se usa para defaults.
Retorna: DataFrame de una fila con valores default por tipo.
df_default = transforms.get_default_record(spark, df_hash)
df_dim = df_hash.unionByName(df_default)

Módulo `merges`

merge_scd2

merge_scd2(spark=None, df_dim_src=None, table_name=None, business_keys=None, surrogate_key=None, eow_date="9999-12-31")

Implementa Slowly Changing Dimension Tipo 2 sobre una tabla Delta, cerrando vigencias antiguas e insertando nuevas versiones.

ParámetroTipoDescripción
sparkSparkSession | NoneSesión opcional.
df_dim_srcDataFrameFuente dimensional con columnas de negocio (sin surrogate key).
table_namestrTabla destino en formato catalog.schema.table.
business_keysstr | list[str]Clave de negocio simple o compuesta.
surrogate_keystrNombre de la clave técnica a generar.
eow_datestrFecha fin de vigencia para registros activos.
Retorna: None. Ejecuta el merge directamente sobre Delta.
merges.merge_scd2(
    spark=spark,
    df_dim_src=df_clientes.select("codigo_cliente", "segmento", "estado"),
    table_name="fi_comunes.silver.dim_cliente",
    business_keys="codigo_cliente",
    surrogate_key="sk_dim_cliente",
    eow_date="9999-12-31"
)

Incluye comportamiento validado contra el código y tests actuales del repositorio.