Skip to main content

Utilidades reutilizables para Spark y Delta Lake en arquitecturas Lakehouse.

Project description

gss-bi-udfs

Creo modulo para guardar UDFs comunes a todas las areas de BI.

contrato de carga Excel

gss_bi_udfs.io.load_latest_excel_file carga el ultimo .xls o .xlsx versionado desde /Volumes/bronze/files/{env}/{source_file}/. La lectura del archivo se realiza con pandas y la creacion del DataFrame final se realiza con Spark, por lo que son dos etapas distintas de error y validacion.

Firma publica:

load_latest_excel_file(
    spark=None,
    source_file=None,
    env=None,
    header=True,
    sheet_name=0,
    schema=None,
)

schema es opcional y debe ser pyspark.sql.types.StructType. Cuando se informa, se aplica en spark.createDataFrame(pdf, schema=schema) para fijar el contrato de salida Spark. Es util cuando el Excel tiene columnas con formatos o tipos mixtos y Spark/Arrow no puede inferir un tipo unico.

El mismo parametro esta disponible en load_latest_file(..., schema=None) para cargas excel, xls y xlsx, y en load_latest_excel(..., schema=None) por compatibilidad. load_latest_excel ya no oculta fallas de lectura o conversion retornando None: propaga el error original con contexto. None queda reservado para el caso en que no se encuentre ningun archivo Excel.

Ejemplo con schema Spark:

from pyspark.sql.types import StructType, StructField, StringType, DoubleType
from gss_bi_udfs import io

schema = StructType([
    StructField("cuenta", StringType(), True),
    StructField("descripcion", StringType(), True),
    StructField("importe", DoubleType(), True),
])

df_excel = io.load_latest_excel_file(
    spark=spark,
    source_file="ospysa-muspeg/finanzas/parametria_tesoreria_mutual",
    env="dev",
    header=True,
    sheet_name=0,
    schema=schema,
)

configuracion de metadata para gss_spark_flow

PipelineOrchestrator usa PostgreSQL como persistencia de metadata por defecto. La conexion puede recibirse como metadata_connection, como metadata_dsn o desde la variable GSS_SPARK_FLOW_POSTGRES_DSN:

export GSS_SPARK_FLOW_POSTGRES_DSN="postgresql://user:password@host:5432/db"
orchestrator = PipelineOrchestrator(schema="metadata")

PipelineOrchestrator ya no recibe spark ni cambia el namespace con USE CATALOG / USE. Si se necesita operar metadata legacy en Unity Catalog, se debe construir el store explicitamente y pasarlo al orquestador:

from gss_spark_flow.metadata_store import create_metadata_store

metadata_store = create_metadata_store(
    spark=spark,
    backend="unity_catalog",
    catalog="operaciones_dev",
    schema="metadata",
)
orchestrator = PipelineOrchestrator(metadata_store=metadata_store)

PipelineContext puede recibir spark explicitamente, pero tambien puede resolverlo bajo demanda desde gss_spark_flow.utils.get_spark() cuando una operacion lo necesita. Esto permite crear contextos para metadata PostgreSQL o stores inyectados sin forzar Spark durante la inicializacion:

from gss_spark_flow import PipelineContext

ctx = PipelineContext(
    pipeline_name="silver/dim_bup",
    task_name="dim_bup",
    dataset="dim_bup",
    auto_start=False,
)

Para ejecuciones Spark se puede seguir pasando la sesion de forma explicita. Esto no cambia el backend de metadata: si se necesita metadata Unity Catalog/Delta, debe declararse de forma explicita.

ctx = PipelineContext(
    spark,
    pipeline_name="silver/dim_bup",
    task_name="dim_bup",
    dataset="dim_bup",
    metadata_backend="unity_catalog",
    metadata_catalog="operaciones_dev",
    metadata_schema="metadata",
)

configuracion de ejecucion de pipelines

La metadata incluye metadata.pipeline_execution_config para alojar inicialmente pipelines que se ejecutan en Fabric o Databricks:

campo uso
pipeline_name nombre logico del pipeline dentro del orquestador
medallion_layer capa del medallero a la que pertenece el pipeline: bronze, silver o gold
execution_engine fabric o databricks
fabric_group_id workspace/group id de Fabric
fabric_pipeline_id id del pipeline de Fabric
databricks_workspace workspace de Databricks
databricks_path path del job/notebook/pipeline en Databricks
parameters_json parametros serializados como JSON
is_active permite activar/desactivar sin borrar

Ejemplo:

orchestrator = PipelineOrchestrator(schema="metadata")

orchestrator.add_fabric_pipeline_config(
    pipeline_name="finanzas_fabric",
    group_id="fabric-group-id",
    pipeline_id="fabric-pipeline-id",
    medallion_layer="bronze",
    parameters={"periodo": "2026-04"},
)

orchestrator.add_databricks_pipeline_config(
    pipeline_name="finanzas_databricks",
    workspace="https://adb-xxx.azuredatabricks.net",
    path="/Shared/pipelines/finanzas",
    medallion_layer="silver",
    parameters={"periodo": "2026-04"},
)

configs = orchestrator.list_pipeline_execution_configs()

Con esa configuracion, cada step ya trae el pipeline_name. El orquestador usa ese nombre para buscar en metadata la configuracion de ejecucion del pipeline: motor, ids, workspace/path y parametros base.

medallion_layer se define a nivel de pipeline_name, no de dependencia. La regla esperada es que un mismo pipeline pertenezca siempre a una unica capa del medallero. Por ejemplo, 0.Master_Sancor_Load_Table puede quedar marcado como bronze, silver/dim_bup como silver y gold/dim_bup como gold.

plan = orchestrator.get_group_execution_plan(group_name)

for step in plan:
    decision = orchestrator.get_dataset_execution_decision(
        pipeline_name=step["pipeline_name"],
        dataset=step["dataset"],
    )

    if decision["should_process"]:
        result = orchestrator.execute_dataset(step)
        print(result)
    else:
        ctx = PipelineContext(
            spark,
            pipeline_name=step["pipeline_name"],
            task_name=step["dataset"],
            dataset=step["dataset"],
            metadata_backend="postgres",
        )
        ctx.skip(
            reason=decision["reason"],
            decision=decision["decision"],
            reused_run_id=decision.get("reused_run_id"),
            reused_task_run_id=decision.get("reused_task_run_id"),
        )

execute_dataset dispara:

  • Fabric: POST https://api.fabric.microsoft.com/v1/workspaces/{group_id}/items/{pipeline_id}/jobs/instances?jobType=Pipeline
  • Databricks con job existente: POST {workspace}/api/2.1/jobs/run-now
  • Databricks con notebook one-shot: POST {workspace}/api/2.1/jobs/runs/submit

Para Databricks, parameters_json debe incluir una de estas opciones de infraestructura:

  • databricks_job_id o job_id
  • existing_cluster_id
  • new_cluster

Esos campos se usan para disparar la ejecucion y no se pasan como parametros al notebook. El resto de los parametros si se envian al notebook.

Tokens soportados por variable de entorno:

export GSS_SPARK_FLOW_FABRIC_TOKEN="..."
export GSS_SPARK_FLOW_DATABRICKS_TOKEN="..."

Tambien podes pasar parametros runtime por corrida:

orchestrator.execute_dataset(
    step,
    runtime_parameters={"periodo": "2026-04"},
)

Para validar el payload sin ejecutar:

orchestrator.execute_dataset(step, dry_run=True)

configuracion de ingesta bronze con Copy Data

La metadata separa la configuracion funcional del procesamiento (metadata.dataset_config) de la configuracion fisica de ingesta hacia bronze. Para cargas desde onpremise u origenes equivalentes, la configuracion del elemento Copy Data de Fabric vive en metadata.bronze_ingestion_config.

bronze_ingestion_config se define por dataset, no por pipeline_name + dataset: un dataset bronze debe tener una unica definicion canonica de ingesta. Si varios pipelines consumen ese dataset, deben depender del mismo dataset ya ingestado, no redefinir como se carga.

Para evitar repetir ids de Fabric en cada dataset, las conexiones se centralizan en metadata.connection_config:

campo uso
connection_key alias logico usado por otras tablas
fabric_connection_id id de la conexion configurada en Fabric
connection_type tipo de conexion, por ejemplo SQL Server o Azure Data Lake Storage Gen2
description descripcion funcional
is_active permite activar/desactivar sin borrar

Campos principales de metadata.bronze_ingestion_config:

campo uso
dataset dataset bronze ingestable
source_connection_key conexion origen definida en connection_config
source_connection_type tipo de origen esperado por Fabric; para onpremise, SQL Server
source_database base de datos origen
source_schema schema origen
source_table tabla origen
source_query consulta usada por Copy Data; puede ser select * o incluir filtro CDC
sink_connection_key conexion destino definida en connection_config
sink_connection_type tipo de destino esperado por Fabric; para bronze, Azure Data Lake Storage Gen2
sink_path_template ruta parametrizada de destino
sink_file_format formato de archivo destino; para bronze, Parquet
is_cdc indica si la consulta aplica criterio CDC
is_active permite activar/desactivar sin borrar

Para resolver el payload de procesamiento desde el orquestador:

bronze_config = orchestrator.get_bronze_ingestion_config(dataset)
source_connection = orchestrator.get_connection_config(
    bronze_config["source_connection_key"]
)
sink_connection = orchestrator.get_connection_config(
    bronze_config["sink_connection_key"]
)

if bronze_config["is_cdc"]:
    cdc_config = orchestrator.get_cdc_ingestion_config(
        pipeline_name=pipeline_name,
        dataset=dataset,
    )

get_cdc_ingestion_config(...) agrega la configuracion funcional de metadata.dataset_config y el estado de metadata.watermark_state. Devuelve:

campo uso
watermark_column columna que ordena los cambios
watermark_type tipo logico usado para serializar el watermark
lower_watermark ultimo valor confirmado; es exclusivo en la siguiente lectura
has_previous_watermark indica si existe un limite inferior confirmado
upper_bound_query consulta que fija el limite superior del lote
initial_query_template lectura inicial hasta el limite superior inclusivo
incremental_query_template lectura (lower_watermark, upper_watermark]

La consulta de origen se selecciona con esta regla:

condicion consulta
has_previous_watermark = false initial_query_template
has_previous_watermark = true incremental_query_template

El valor de source_schema no interviene en esta decision. El pipeline debe:

  1. Ejecutar upper_bound_query y conservar el upper_watermark obtenido.
  2. Seleccionar la plantilla segun has_previous_watermark.
  3. Reemplazar {lower_watermark} y {upper_watermark} por literales SQL compatibles con watermark_type.
  4. Ejecutar Copy Data con la consulta resultante.
  5. Confirmar el límite superior solamente después de una copia exitosa.

No debe recalcularse el maximo durante la copia, porque eso puede dejar huecos entre lotes.

Las plantillas usan los placeholders {lower_watermark} y {upper_watermark}. Despues de una copia exitosa, confirmar el lote y cerrar la tarea con:

orchestrator.complete_cdc_ingestion(
    pipeline_name=pipeline_name,
    dataset=dataset,
    run_id=run_id,
    task_run_id=task_run_id,
    upper_watermark=upper_watermark,
)

Esta operacion actualiza watermark_state y registra watermark_in / watermark_out en task_run. Ante una copia fallida no debe invocarse, por lo que el watermark confirmado permanece sin cambios.

Es obligatorio ejecutar complete_cdc_ingestion(...) despues de cada Copy Data CDC exitoso, pasando exactamente el upper_watermark obtenido antes de iniciar la copia. Si no se ejecuta:

  • metadata.watermark_state no avanza.
  • La siguiente llamada devuelve el mismo lower_watermark.
  • Si nunca hubo una confirmacion previa, lower_watermark sigue siendo null y has_previous_watermark sigue siendo false, por lo que se vuelve a seleccionar initial_query_template.

El notebook workspace/notebooks/complete_copy_data.py permite finalizar ambos tipos de carga con el mismo flujo:

  • Siempre registra las metricas devueltas por Copy Data.
  • Ante errores marca el task_run como FAILED y no modifica el watermark.
  • Para una carga no CDC ejecuta PipelineContext.success_task().
  • Para una carga CDC extrae upper_watermark de upper_bound_output y ejecuta complete_cdc_ingestion(...).

Los widgets esperados son:

widget contenido
validation_output respuesta generada por prepare_copy_data.py
copydata_output salida JSON de la actividad Copy Data
upper_bound_output salida de Lookup con upper_watermark; obligatorio solamente para CDC

En Fabric, upper_bound_output puede tener la forma {"firstRow": {"upper_watermark": "..."}} o {"value": [{"upper_watermark": "..."}]}.

Ejemplo de ruta destino:

@concat(pipeline().parameters.DataBase, '/', pipeline().parameters.Schema, '/', pipeline().parameters.Table)

plan de ejecucion para Fabric como orquestador unico

Cuando Fabric coordina ejecuciones Fabric + Databricks, conviene generar un plan enriquecido una sola vez y pasarlo como JSON al Data Pipeline:

plan_json = orchestrator.get_group_execution_plan_json(group_name)

Ese JSON incluye, por cada step:

  • pipeline_name y dataset
  • execution_engine
  • execution_config
  • execution_parameters
  • dependency_decision
  • dataset_execution_decision

Ejemplo conceptual:

{
  "group_name": "modelo_estrella_finanzas_daily",
  "steps": [
    {
      "step": 1,
      "pipeline_name": "bronze/clientes",
      "dataset": "timepro.insudb.client",
      "execution_engine": "fabric",
      "can_execute": true,
      "should_process": true,
      "has_dataset_execution_policy": true
    },
    {
      "step": 2,
      "pipeline_name": "silver/clientes",
      "dataset": "fi_comunes.silver.clientes",
      "execution_engine": "databricks",
      "can_execute": true,
      "should_process": true,
      "has_dataset_execution_policy": false
    }
  ]
}

Para orquestar por medallero, el metodo recomendado es:

plan = orchestrator.get_group_medallion_execution_plan(group_name)

Ese plan devuelve, para cada target del grupo, tres etapas:

  • bronze: lista de pipeline_name + dataset que puede dispararse en paralelo.
  • silver: lista de pipeline_name + dataset en orden topologico.
  • gold: lista de pipeline_name + dataset en orden topologico.

Ejemplo conceptual:

[
  {
    "step": 1,
    "pipeline_name": "gold/dim_bup",
    "dataset": "dim_bup",
    "steps": [
      {
        "step": 1,
        "layer": "bronze",
        "parallel": true,
        "items": [
          {"pipeline_name": "0.Master_Sancor_Load_Table", "dataset": "bup.bup.persons"},
          {"pipeline_name": "0.Master_Sancor_Load_Table", "dataset": "bup.bup.addresses"}
        ]
      },
      {
        "step": 2,
        "layer": "silver",
        "parallel": false,
        "items": [
          {"pipeline_name": "silver/dim_bup", "dataset": "dim_bup"}
        ]
      },
      {
        "step": 3,
        "layer": "gold",
        "parallel": false,
        "items": [
          {"pipeline_name": "gold/dim_bup", "dataset": "dim_bup"}
        ]
      }
    ]
  }
]

El plan trae solo los identificadores de ejecucion (pipeline_name y dataset) por item. Al momento de disparar cada item, se debe validar si puede ejecutarse con assert_can_execute(...) o inspeccionar get_dependency_diagnostics(...).

que representa cada tabla de metadata

tabla representa
run Ejecucion logica de un pipeline. Guarda inicio, fin, estado, disparador y datos operativos generales.
task_run Detalle por tarea/dataset dentro de una ejecucion. Guarda estado, conteos, watermarks y decisiones de skip/reuso.
error_log Errores asociados a una corrida o tarea.
watermark_state Ultimo watermark confirmado por pipeline_name + dataset.
dataset_config Configuracion funcional del procesamiento de un dataset: origen, destino, watermark y claves de merge.
connection_config Catalogo centralizado de conexiones logicas usadas por configs tecnicas.
bronze_ingestion_config Configuracion fisica de ingesta bronze por dataset fuente.
pipeline_dependency DAG entre nodos ejecutables parent pipeline/dataset -> child pipeline/dataset.
pipeline_dependency_group Targets que un grupo quiere materializar, normalmente nodos gold, con orden de prioridad.
pipeline_dependency_group_schedule Calendarios asociados a grupos de ejecucion.
dataset_execution_policy Politica de frescura propia de datasets fuente, principalmente bronze no-CDC.
pipeline_execution_config Configuracion tecnica de ejecucion por pipeline_name y capa del medallero (medallion_layer).
lineage Relacion input/output registrada por ejecucion.
data_quality Resultados de controles de calidad.
schema_version Versiones/migraciones aplicadas sobre la metadata.

Hay dos decisiones distintas:

decision metadata pregunta que responde
dependency_decision metadata.pipeline_dependency si el pipeline_name + dataset puede correr segun sus padres del DAG
dataset_execution_decision metadata.dataset_execution_policy si un dataset con politica configurada debe recargarse o puede reutilizar una carga fresca

Decision de modelado: cuando una dependencia declara parent_dataset, la validacion se resuelve contra el ultimo metadata.task_run de ese dataset. Esto permite que una capa dependiente avance cuando el dataset requerido termino en SUCCESS, aunque el metadata.run que agrupa varios datasets siga abierto o termine con otro estado por errores en datasets no relacionados. Cuando la dependencia no declara dataset, la validacion conserva semantica de pipeline y se resuelve contra metadata.run.

dataset_execution_policy debe configurarse solamente para datasets cuya frescura se decide por politica propia, por ejemplo datasets onpremise cargados en Fabric. Para datasets silver/gold sin fila en esa tabla, el plan devuelve has_dataset_execution_policy = false y should_process = true; su ejecucion queda gobernada por pipeline_dependency.

politicas de recarga por dataset

La metadata incluye metadata.dataset_execution_policy para definir cada cuanto conviene volver a cargar un dataset fuente, independientemente de los grupos o pipelines que lo consuman. Esta tabla no define dependencias entre pipelines: esas dependencias viven en metadata.pipeline_dependency.

En la practica, dataset_execution_policy debe usarse para datasets onpremise o fuentes equivalentes cuya recarga se decide por frescura propia. Para datasets derivados como silver/gold, si no existe una fila en esta tabla, el orquestador no aplica politica de recarga y deja que pipeline_dependency gobierne si el step puede ejecutarse.

Antes de evaluar la frescura por ultima carga exitosa, el orquestador revisa si ya existe un metadata.task_run en RUNNING para el mismo pipeline_name + dataset. Si existe, get_dataset_execution_decision devuelve should_process = false, decision = SKIP_ALREADY_RUNNING y reason = dataset_has_running_task. Esto permite usar Databricks como punto consistente de decision/reserva antes de delegar el procesamiento real a Fabric.

Para datasets CDC, no debe cargarse una fila en dataset_execution_policy cuando su actualizacion se gobierna por metadata.watermark_state y por dataset_config.watermark_column / dataset_config.watermark_type. En ese caso, la frecuencia efectiva de recarga del dataset queda determinada por el watermark y no por una politica de frescura independiente.

campo uso
dataset nombre del dataset al que aplica la politica
min_interval_minutes intervalo minimo entre cargas exitosas
schedule_expr expresion opcional de calendario del dataset
timezone zona horaria asociada al calendario
skip_if_fresh permite registrar SKIPPED si el dataset sigue fresco
is_active permite activar/desactivar sin borrar

Ejemplo:

orchestrator = PipelineOrchestrator(schema="metadata")

orchestrator.add_dataset_execution_policy(
    dataset="timepro.insudb.intermedia",
    min_interval_minutes=1440,
    schedule_expr="0 0 6 * * ?",
    timezone="America/Argentina/Cordoba",
)

plan = orchestrator.get_group_execution_plan("modelo_estrella_finanzas_daily")

for step in plan:
    if not step["should_process"]:
        ctx = PipelineContext(
            spark,
            pipeline_name=step["pipeline_name"],
            task_name=step["dataset"],
            dataset=step["dataset"],
            metadata_backend="postgres",
        )
        ctx.skip(
            reason=step["decision_reason"],
            decision=step["decision"],
            reused_run_id=step["reused_run_id"],
            reused_task_run_id=step["reused_task_run_id"],
        )

para compilar local

python3 -m build

para publicar local (manual)

python3 -m twine upload dist/*

publicar nueva version en pypi con github actions

Prerequisitos:

  • Secret del repo configurado: PYPI_API_TOKEN
  • Workflow: .github/workflows/python-publish.yml

Pasos:

  1. Asegurate de tener los cambios listos en main.
    git checkout main
    git pull
    git status
    
  2. Defini la nueva version semantica (X.Y.Z) y crea el tag vX.Y.Z.
    git tag -a v0.1.5 -m "Release v0.1.5"
    
  3. Publica rama y tag en GitHub.
    git push origin main
    git push origin v0.1.5
    
  4. Crea el release asociado al tag.
    • GitHub -> Releases -> Draft a new release
    • Seleccionar tag v0.1.5
    • Publicar (Publish release)
  5. Verifica la corrida del workflow.
    • GitHub -> Actions -> Publish Python Package to PyPI
    • Debe finalizar en verde.
  6. Verifica la version publicada en PyPI.
    • https://pypi.org/project/gss-bi-udfs/

Notas:

  • El workflow valida que el tag tenga formato vX.Y.Z.
  • La version del paquete se toma del tag (sin la v).
  • Si falla el release, podes reintentar desde Actions con Run workflow y el input tag (ejemplo: v0.1.5).

Project details


Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

gss_bi_udfs-2.1.0.tar.gz (216.2 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

gss_bi_udfs-2.1.0-py3-none-any.whl (78.2 kB view details)

Uploaded Python 3

File details

Details for the file gss_bi_udfs-2.1.0.tar.gz.

File metadata

  • Download URL: gss_bi_udfs-2.1.0.tar.gz
  • Upload date:
  • Size: 216.2 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.1.0 CPython/3.13.12

File hashes

Hashes for gss_bi_udfs-2.1.0.tar.gz
Algorithm Hash digest
SHA256 876267f61679c9784dae5028dc20adc4a37909539fe191edeee89ddb83845577
MD5 905b722dfb6dbceff6057889ffb28cf6
BLAKE2b-256 ec194ea4d945db17e0280695366f4415c27dfe5f3d13b51a49ca09fe067291c4

See more details on using hashes here.

File details

Details for the file gss_bi_udfs-2.1.0-py3-none-any.whl.

File metadata

  • Download URL: gss_bi_udfs-2.1.0-py3-none-any.whl
  • Upload date:
  • Size: 78.2 kB
  • Tags: Python 3
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.1.0 CPython/3.13.12

File hashes

Hashes for gss_bi_udfs-2.1.0-py3-none-any.whl
Algorithm Hash digest
SHA256 73cf772c7fa82df91dbe9a2ffeeb5e85c5c8685cafc1a222a5d5321ef7106085
MD5 24cc4c7cba1f0f63d618c2c84fd59ebd
BLAKE2b-256 ede6826bfd46ec8c9e80f0ea34936d01dbaf7dde0c9a7131b4717a3960ce0fcf

See more details on using hashes here.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page