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, desde
un secret Databricks configurado por variables de entorno o desde
GSS_SPARK_FLOW_POSTGRES_DSN como fallback:
export GSS_SPARK_FLOW_POSTGRES_DSN_SECRET_SCOPE="kvgsssbx001"
export GSS_SPARK_FLOW_POSTGRES_DSN_SECRET_KEY="connection-string-metadata"
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_idojob_idexisting_cluster_idnew_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:
- Ejecutar
upper_bound_queryy conservar elupper_watermarkobtenido. - Seleccionar la plantilla segun
has_previous_watermark. - Reemplazar
{lower_watermark}y{upper_watermark}por literales SQL compatibles conwatermark_type. - Ejecutar Copy Data con la consulta resultante.
- 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_stateno avanza.- La siguiente llamada devuelve el mismo
lower_watermark. - Si nunca hubo una confirmacion previa,
lower_watermarksigue siendonullyhas_previous_watermarksigue siendofalse, por lo que se vuelve a seleccionarinitial_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_runcomoFAILEDy no modifica el watermark. - Para una carga no CDC ejecuta
PipelineContext.success_task(). - Para una carga CDC extrae
upper_watermarkdeupper_bound_outputy ejecutacomplete_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_nameydatasetexecution_engineexecution_configexecution_parametersdependency_decisiondataset_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 depipeline_name + datasetque puede dispararse en paralelo.silver: lista depipeline_name + dataseten orden topologico.gold: lista depipeline_name + dataseten 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:
- Asegurate de tener los cambios listos en
main.git checkout main git pull git status
- Defini la nueva version semantica (
X.Y.Z) y crea el tagvX.Y.Z.git tag -a v0.1.5 -m "Release v0.1.5"
- Publica rama y tag en GitHub.
git push origin main git push origin v0.1.5
- Crea el release asociado al tag.
- GitHub ->
Releases->Draft a new release - Seleccionar tag
v0.1.5 - Publicar (
Publish release)
- GitHub ->
- Verifica la corrida del workflow.
- GitHub ->
Actions->Publish Python Package to PyPI - Debe finalizar en verde.
- GitHub ->
- 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
ActionsconRun workflowy el inputtag(ejemplo:v0.1.5).
Project details
Release history Release notifications | RSS feed
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file gss_bi_udfs-4.2.1.tar.gz.
File metadata
- Download URL: gss_bi_udfs-4.2.1.tar.gz
- Upload date:
- Size: 314.7 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/7.0.0 CPython/3.13.14
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
8b416826d14cc22ba812265f4cae99f4917c8c42e379b16b8f6292d1ebcfd557
|
|
| MD5 |
e701c9582159f4523a535548a76295b9
|
|
| BLAKE2b-256 |
827fe9c1cfb5dbe135f62fa2c629b41905aabb749a41f6ec241c30128df2c8f6
|
File details
Details for the file gss_bi_udfs-4.2.1-py3-none-any.whl.
File metadata
- Download URL: gss_bi_udfs-4.2.1-py3-none-any.whl
- Upload date:
- Size: 95.8 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/7.0.0 CPython/3.13.14
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
8b4538548e2d8c684e935ac371057da2be2f5a2fdfafb6337f557063fca02a39
|
|
| MD5 |
cc7faa3efa76ae1839db456ff5475d85
|
|
| BLAKE2b-256 |
56748235b06a7736e5b43fb2a6346495bd2c0443461451246c5115e06d0d6882
|