Skip to main content

Librería de utilidades comunes para pipelines ELT orientados a Databricks (DLT).

Project description

DBX ELT Utilities v2.0

dbx-elt-utils es la libreria estandar para crear pipelines Lakeflow Spark Declarative Pipelines (SDP) en Databricks. Proporciona metodos para ingesta hibrida, setup automatico de entorno (Local vs Nube), limpieza de datos, y validacion de calidad.

v2.0 — Soporte completo para la nueva API pyspark.pipelines (dp), Liquid Clustering, ExpectationsFactory, 5 funciones nuevas de limpieza, y checkpoint en UC Volumes.

Instalacion

# Produccion (Databricks Runtime)
pip install dbx-elt-utils

# Desarrollo local (VSCode + Databricks Connect)
pip install dbx-elt-utils[local]

La libreria no fuerza pyspark como dependencia, evitando conflictos con el Databricks Runtime.


Quick Start

from dbx_elt_utils.notebook_utils import init_notebook

notebook = init_notebook()

env   = notebook.env     # '_dev' o '_prod'
spark = notebook.spark   # Sesion PySpark
dp    = notebook.dp      # API SDP (pyspark.pipelines o mock local)

Modulos

1. notebook_utils — Setup y Mock Local

La funcion init_notebook() detecta automaticamente si estas en Databricks o en VSCode y configura todo:

Atributo Descripcion
notebook.env Sufijo de entorno: _dev o _prod
notebook.spark Sesion PySpark (nativa o Connect)
notebook.dp Modulo SDP (real o _SdpMock)
notebook.dlt Alias legacy (= notebook.dp)
notebook.is_local True si corre en VSCode

Mock SDP completo — en local, notebook.dp es un _SdpMock que soporta todos los decoradores sin error:

@dp.materialized_view(name="mi_tabla", comment="...")
@dp.expect("id_not_null", "id IS NOT NULL")
def mi_tabla():
    return spark.read.table("source")

Funciones de testing:

Funcion Descripcion
get_test_spark() Sesion Spark (Workspace o Databricks Connect)
get_local_source_table(spark, fqn) Resuelve tabla oficial vs temporal de test
clean_local_test_table(spark, fqn) Elimina tabla _tmp_sql despues del test
display_test_results(spark, data) Muestra DataFrame adaptado al entorno
stop_local_spark() Libera streams activos sin matar la sesion
get_checkpoint_base(env) Ruta checkpoints en UC Volumes
get_schema_location_base(env) Ruta schema location (hash-based)

2. ingest_utils — Ingesta Hibrida (Bronze)

from dbx_elt_utils.ingest_utils import ingesta_hibrida

df = ingesta_hibrida(spark, SOURCE_LANDING, tipo="auto_detect")

Detecta automaticamente el tipo de origen y retorna el DataFrame adecuado:

Tipo detectado Origen Metodo
external_table Tabla Unity Catalog spark.readStream.table()
delta_path Ruta Delta (dbfs/abfss) spark.readStream.format("delta")
auto_loader Archivos (CSV/JSON/Parquet) cloudFiles con schema evolution
batch_files Ruta sin streaming spark.read.format()

Deteccion inteligente:

  • Rutas con / -> archivos (Auto Loader o Delta path segun contenido _delta_log)
  • Formato catalog.schema.table -> tabla Unity Catalog
  • schemaEvolutionMode=addNewColumns activado automaticamente en Auto Loader

3. clean_utils — Limpieza y Transformacion (Silver)

from dbx_elt_utils.clean_utils import (
    extraer_valor_array_string, parse_json_array_column,
    normalize_string, safe_cast, dedup_by_key, add_surrogate_key
)
Funcion Que hace Ejemplo
extraer_valor_array_string(col) ["12345"] a 12345, [] a NULL Limpia arrays JSON de APIs
parse_json_array_column(df, col, alias) ["a","b"] a array real de Spark Explode-ready
normalize_string(col) Trim + lowercase + sin acentos Estandarizacion de texto
safe_cast(col, tipo) Cast seguro sin error (NULL si falla) safe_cast(col("edad"), "int")
dedup_by_key(df, keys, order_col) Dedup por clave con row_number CDC / merge patterns
add_surrogate_key(df, cols, alias) SHA-256 surrogate key add_surrogate_key(df, ["id","src"])

4. expectations_utils — Validacion de Calidad (Nuevo v2.0)

from dbx_elt_utils.expectations_utils import ExpectationsFactory as EF

rules = EF.combine(
    EF.not_null("id", "nombre"),
    EF.in_range("edad", 0, 150),
    EF.freshness("updated_at", max_hours=48),
    EF.in_set("estado", ["ACTIVO", "INACTIVO"]),
)

@dp.table(name="mi_tabla")
@dp.expect_all_or_drop(rules)
def mi_tabla():
    ...
Metodo Genera
EF.not_null("col") col IS NOT NULL
EF.not_empty("col") col IS NOT NULL AND trim(col) != ''
EF.positive("col") col > 0
EF.in_range("col", 0, 100) col >= 0 AND col <= 100
EF.in_set("col", [...]) col IN ('a', 'b', 'c')
EF.freshness("col", 24) col >= current_timestamp() - INTERVAL 24 HOURS
EF.regex_match("col", pattern) col RLIKE 'pattern'
EF.combine(...) Fusiona multiples dicts en uno

Ejemplo Completo: Bronze + Silver

# -- Bronze --
from dbx_elt_utils.notebook_utils import init_notebook
from dbx_elt_utils.ingest_utils import ingesta_hibrida

notebook = init_notebook()
env, spark, dp = notebook.env, notebook.spark, notebook.dp

SOURCE = f"landing{env}.schema.mi_tabla"

@dp.materialized_view(
    name="mi_tabla",
    comment=f"Bronze desde {SOURCE}",
    cluster_by=["id"]
)
def bronze_mi_tabla():
    df = ingesta_hibrida(spark, SOURCE)
    df.createOrReplaceTempView("v_raw")
    return spark.sql("SELECT *, current_timestamp() as _ingested_at FROM v_raw")
# -- Silver (CDC SCD1) --
from dbx_elt_utils.expectations_utils import ExpectationsFactory as EF
from pyspark.sql.functions import col

BRONZE_SOURCE = f"bronze{env}.schema.mi_tabla"

dp.create_streaming_table(
    name="mi_tabla",
    cluster_by=["id"],
    expect_all_or_drop=EF.combine(
        EF.not_null("id"),
        EF.not_empty("nombre")
    )
)

dp.apply_changes(
    target="mi_tabla",
    source=BRONZE_SOURCE,
    keys=["id"],
    sequence_by=col("_ingested_at"),
    stored_as_scd_type=1
)

Changelog

v2.0.0

  • notebook.dp — Alias para pyspark.pipelines (nueva API SDP)
  • _SdpMock completo: materialized_view, table, expect*, apply_changes, create_streaming_table, append_flow
  • Checkpoint path en UC Volumes (/Volumes/bronze{env}/temporary/checkpoints)
  • Schema location hash-based para Auto Loader
  • ExpectationsFactory — Modulo nuevo de validacion de calidad
  • clean_utils: 5 funciones nuevas (normalize_string, safe_cast, dedup_by_key, add_surrogate_key, parse_json_array_column)
  • ingest_utils: Deteccion delta_path y tipo batch_files
  • Dual import: pyspark.pipelines con fallback dlt

v1.3.1

  • Version estable anterior (API @dlt.table)

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

dbx_elt_utils-2.2.8.tar.gz (33.8 kB view details)

Uploaded Source

Built Distribution

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

dbx_elt_utils-2.2.8-py3-none-any.whl (23.3 kB view details)

Uploaded Python 3

File details

Details for the file dbx_elt_utils-2.2.8.tar.gz.

File metadata

  • Download URL: dbx_elt_utils-2.2.8.tar.gz
  • Upload date:
  • Size: 33.8 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.2.0 CPython/3.13.7

File hashes

Hashes for dbx_elt_utils-2.2.8.tar.gz
Algorithm Hash digest
SHA256 b0722d282e04badc5ce7eb1756aff8950e8b2a2cf5fdf1b39f6e931c86e12544
MD5 22fda0b8736c8f8c71b1314b6a96b0e0
BLAKE2b-256 0a0fb8272d59e6ca5ada0984becd9e5b358a610007efc71691523aeef0b85e0b

See more details on using hashes here.

File details

Details for the file dbx_elt_utils-2.2.8-py3-none-any.whl.

File metadata

  • Download URL: dbx_elt_utils-2.2.8-py3-none-any.whl
  • Upload date:
  • Size: 23.3 kB
  • Tags: Python 3
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.2.0 CPython/3.13.7

File hashes

Hashes for dbx_elt_utils-2.2.8-py3-none-any.whl
Algorithm Hash digest
SHA256 7ddd84b282ce5ff1943e111acdb61bc8adeb7342c31028c7c9ed43ceadf8ae63
MD5 3b94af7cda9bff8238e9bc56c1c0d146
BLAKE2b-256 8120aa6a4ccec5cee4fc57ba67895de4653a58c5ceb47d098c1a0d5ac72d78fc

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