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.

configuracion de catalog/schema para gss_spark_flow

PipelineOrchestrator permite definir defaults por entorno:

export GSS_SPARK_FLOW_CATALOG=operaciones_dev
export GSS_SPARK_FLOW_SCHEMA=metadata

Con eso, PipelineOrchestrator(spark) aplica USE CATALOG y USE <schema> automaticamente. Si se pasan catalog/schema en el constructor, esos valores tienen prioridad sobre las variables de entorno.

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
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(spark)

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

orchestrator.add_databricks_pipeline_config(
    pipeline_name="finanzas_databricks",
    workspace="https://adb-xxx.azuredatabricks.net",
    path="/Shared/pipelines/finanzas",
    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.

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"],
        )
        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)

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
    }
  ]
}

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

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.

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(spark)

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"],
        )
        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-0.2.3.tar.gz (110.8 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-0.2.3-py3-none-any.whl (41.9 kB view details)

Uploaded Python 3

File details

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

File metadata

  • Download URL: gss_bi_udfs-0.2.3.tar.gz
  • Upload date:
  • Size: 110.8 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-0.2.3.tar.gz
Algorithm Hash digest
SHA256 7b3db277858e0d3798155ee155564a8cc81b98f31aa6973eea9661dab51bcc7d
MD5 0311e87f2ee73756da54d57209a1801d
BLAKE2b-256 7812a61774f0520bb5d248e8d157d9dd323c2e97e410faa177f0db1f7df1105d

See more details on using hashes here.

File details

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

File metadata

  • Download URL: gss_bi_udfs-0.2.3-py3-none-any.whl
  • Upload date:
  • Size: 41.9 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-0.2.3-py3-none-any.whl
Algorithm Hash digest
SHA256 58e8aeeeaf7ac48e5f6355ad1cda303361605eecfe0ee85c361527863a8331a7
MD5 a24fde8dd77cf3513fe0027fe97839df
BLAKE2b-256 3e062b15282c3d921e1062601d71ab47fc218cca986a801fd6e7140c1f45cf95

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