Metadata-Version: 2.4
Name: gss-bi-udfs
Version: 0.2.2
Summary: Utilidades reutilizables para Spark y Delta Lake en arquitecturas Lakehouse.
Author: Geronimo Forconi
Requires-Python: >=3.8
Description-Content-Type: text/markdown

# 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:

```bash
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:

```python
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.

```python
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:

```bash
export GSS_SPARK_FLOW_FABRIC_TOKEN="..."
export GSS_SPARK_FLOW_DATABRICKS_TOKEN="..."
```

Tambien podes pasar parametros runtime por corrida:

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

Para validar el payload sin ejecutar:

```python
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:

```python
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:

```json
{
  "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:

```python
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`.
   ```bash
   git checkout main
   git pull
   git status
   ```
2. Defini la nueva version semantica (`X.Y.Z`) y crea el tag `vX.Y.Z`.
   ```bash
   git tag -a v0.1.5 -m "Release v0.1.5"
   ```
3. Publica rama y tag en GitHub.
   ```bash
   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`).
