Saltar al contenido
Lakehouse LabLakehouse LabPreparación Databricks Data Engineer
Módulo 18 · Professional

Declarativo

Contenido abierto

Puedes leer todo sin registrarte. Solo crearemos un perfil anónimo cuando decidas guardar tu progreso.

Professional

Spark Declarative Pipelines en Lakeflow

Declara datasets y dependencias para que Lakeflow gestione el grafo y la ejecución incremental sobre el framework actual de Spark.

Lectura pública
Al terminar podrás
  • Crear tablas streaming y materialized views
  • Comparar declarativo con Structured Streaming
  • Distinguir Spark Declarative Pipelines de la oferta gestionada Lakeflow
Ver fuentes y revisión

Metadatos editoriales

Última revisión
21 jul 2026
Nivel
Professional
Ruta relacionada
pipelines
Dominios blueprint
Spark Declarative Pipelines · Lakeflow
Estado
Revisión editorial interna
Fuentes principales
Spark Declarative Pipelines flows · Databricks · Materialized views · Databricks
Reportar un error
01
Modelo mental

Framework Spark Declarative Pipelines y Lakeflow

Spark Declarative Pipelines define datasets y dependencias; Lakeflow extiende el framework y gestiona el grafo, las actualizaciones, el linaje y los eventos.

Objetivo
Spark Declarative Pipelines define datasets y dependencias; Lakeflow extiende el framework y gestiona el grafo, las actualizaciones, el linaje y los eventos.
Duración estimada
17 min aprox.
Dificultad
Professional
Prerrequisitos
m12
Reportar un error en esta lección

En lugar de iniciar manualmente varios `writeStream`, cada función devuelve un DataFrame que define un dataset. Las lecturas entre datasets establecen dependencias y el pipeline determina el orden. Esto reduce código operativo, pero no elimina decisiones sobre contrato, incrementabilidad, calidad o coste.

El grafo debe expresar transformaciones de datos, no pasos imperativos con efectos externos. Crear archivos, llamar APIs o mutar tablas arbitrariamente dentro de una función declarativa rompe reevaluación y dificulta optimización. Esos efectos pertenecen a tareas de Jobs alrededor del pipeline.

Modelo mental

Spark Declarative Pipelines cambia la unidad de razonamiento desde una secuencia de comandos hacia un grafo de datasets y flows. El autor declara qué representa cada streaming table, materialized view o sink y sus dependencias se deducen de las lecturas; el motor construye el DAG, elige el orden válido y administra actualizaciones incrementales. En Databricks, Lakeflow pipelines es la oferta gestionada que extiende e interopera con el framework Apache Spark Declarative Pipelines sobre un runtime optimizado, añadiendo operación, event log, gobernanza y capacidades específicas. No debe confundirse el framework con el antiguo nombre comercial Delta Live Tables: código o exámenes previos pueden usar DLT, pero el modelo vigente se expresa como pipelines, flows y datasets. Declarativo no significa automático sin contrato: claves, semántica temporal, calidad, costes y compatibilidad siguen perteneciendo al diseño humano.

Pipeline

Unidad gestionada de desarrollo y ejecución que contiene datasets, flows, sinks, configuración y el grafo de dependencias que los relaciona.

Define la frontera de actualización, observabilidad y despliegue que el equipo opera como un producto coherente.
Flow

Relación declarativa que procesa una fuente mediante una consulta y escribe sus resultados en un destino administrado por el pipeline.

Separa la lógica de movimiento de datos del objeto persistente y permite varias entradas controladas hacia un target.
Evaluación declarativa

Fase en la que el runtime interpreta definiciones para descubrir objetos y dependencias antes de ejecutar el procesamiento efectivo de datos.

Explica por qué las funciones deben ser deterministas y no contener acciones, llamadas externas ni efectos dependientes del orden del archivo.
PySparkDos datasets declarativos conectados
from pyspark import pipelines as dp
from pyspark.sql import functions as F

@dp.table(name="orders_bronze")
def orders_bronze():
    return spark.readStream.table("main.raw.orders")

@dp.materialized_view(name="daily_order_totals")
def daily_order_totals():
    return (
        spark.read.table("orders_bronze")
          .groupBy(F.to_date("event_ts").alias("order_date"))
          .agg(F.sum("amount").alias("revenue"))
    )

El nombre lógico `orders_bronze` crea la dependencia; el pipeline administra actualización y metadatos.

Puntos clave

  • Las funciones declarativas devuelven DataFrames y no deben ejecutar acciones como `collect()` o escrituras manuales.
  • Las dependencias proceden de lecturas, no del orden físico de funciones en el archivo.
  • El event log ofrece progreso, calidad, linaje y errores del grafo.

Evita

  • Llamar `display`, `count` o `saveAsTable` dentro de una función decorada y mezclar declaración con ejecución.
  • Depender del orden del archivo en vez de leer explícitamente el dataset upstream.

Recuerdo activo

¿Qué crea una arista entre dos datasets del pipeline?

Borrador privado · solo en este navegador
02
Implementación

Modelo declarativo y grafo

Una streaming table procesa nuevas filas de una fuente streaming y conserva semántica incremental apropiada para bronze y silver append/CDC.

Objetivo
Una streaming table procesa nuevas filas de una fuente streaming y conserva semántica incremental apropiada para bronze y silver append/CDC.
Duración estimada
17 min aprox.
Dificultad
Professional
Prerrequisitos
m12
Reportar un error en esta lección

Una definición con `spark.readStream` produce un flujo continuo o por triggers dentro del pipeline. La tabla persiste resultados y el servicio gestiona checkpoints internos. Es adecuada cuando cada registro nuevo puede procesarse incrementalmente sin recalcular todo el resultado.

No debe leerse una fuente cambiante con `spark.read` y esperar semántica streaming. Tampoco se usa una streaming table para una consulta que necesita revisar cambios arbitrarios en ambos lados de un join batch; una materialized view puede permitir actualización incremental gestionada.

Modelo mental

Una streaming table representa un dataset cuyo estado crece o cambia mediante uno o más flows incrementales alimentados por fuentes streaming. No es simplemente una tabla Delta a la que alguien ejecuta append: el pipeline administra la consulta, el checkpoint y la relación entre definición y actualización. Es apropiada para ingestión bronze, transformaciones silver continuas y CDC cuando la semántica se puede expresar incrementalmente. Leerla como stream transmite cambios nuevos a consumidores; leerla como relación batch observa su estado materializado. La entrada debe ser realmente streaming (`readStream`, `STREAM(...)` o una fuente equivalente); declarar una tabla streaming no convierte por sí sola una consulta batch completa en incremental. También importa distinguir append puro de actualizaciones: AUTO CDC usa un flow gestionado para aplicar cambios a una streaming table destino, mientras un simple append conservaría múltiples versiones como filas independientes.

Streaming table

Dataset persistente administrado por un pipeline cuyos flows procesan incrementalmente una o varias fuentes streaming y conservan progreso recuperable.

Es el objeto principal para ingestión y transformaciones continuas sin gestionar manualmente cada `writeStream` y checkpoint.
Lectura streaming

Lectura que observa únicamente nuevos cambios disponibles desde una fuente y mantiene una posición incremental entre actualizaciones sucesivas.

Evita escanear el estado completo, pero exige que la semántica upstream pueda propagarse correctamente como cambios.
Full refresh

Reconstrucción del contenido de un dataset desde sus fuentes en lugar de continuar únicamente con el progreso incremental existente.

Puede ser necesaria tras cambios incompatibles y requiere considerar coste, retención y efectos sobre consumidores.
PySparkStreaming table desde archivos con Auto Loader
from pyspark import pipelines as dp

@dp.table(
    name="orders_bronze",
    comment="Pedidos crudos con metadatos de origen",
    table_properties={"quality": "bronze"},
)
def orders_bronze():
    return (
        spark.readStream.format("cloudFiles")
          .option("cloudFiles.format", "json")
          .option("cloudFiles.schemaLocation", "/Volumes/main/ops/schemas/orders")
          .load("/Volumes/main/landing/orders")
    )

La ubicación de esquema de Auto Loader sigue siendo necesaria aunque el pipeline gestione su propio estado.

Puntos clave

  • `spark.readStream` señala una entrada incremental.
  • El pipeline administra estado operativo; no se define `checkpointLocation` dentro del dataset.
  • La compatibilidad de transformaciones streaming sigue aplicando, incluidos watermarks para estado acotado.

Evita

  • Añadir manualmente `writeStream` o `checkpointLocation` dentro de una definición de pipeline.
  • Aplicar una agregación stateful sin watermark y trasladar al servicio un estado ilimitado.

Recuerdo activo

¿Quién gestiona el checkpoint de una streaming table en un pipeline?

Borrador privado · solo en este navegador
03
Operación

Streaming tables y materialized views

Una materialized view almacena el resultado de una consulta batch declarativa y el servicio intenta actualizarla incrementalmente cuando cambian sus dependencias.

Objetivo
Una materialized view almacena el resultado de una consulta batch declarativa y el servicio intenta actualizarla incrementalmente cuando cambian sus dependencias.
Duración estimada
17 min aprox.
Dificultad
Professional
Prerrequisitos
m12
Reportar un error en esta lección

A diferencia de una vista lógica, la salida se materializa para servir consultas rápidas. La definición suele usar `spark.read.table` porque describe el resultado completo correcto. El motor decide si puede aplicar cambios incrementales o si necesita recomputar según consulta y origen.

Es apropiada para joins, agregaciones y modelos gold donde el resultado puede cambiar por actualizaciones upstream. No promete que toda consulta sea siempre incremental; diseño de claves, filtros y operaciones influye en el plan de refresh y debe observarse en el event log.

Modelo mental

Una materialized view almacena el resultado de una consulta declarativa y lo actualiza cuando cambian sus dependencias. A diferencia de una vista lógica, no recalcula para cada lector; a diferencia de una streaming table, su consulta se formula sobre relaciones batch y describe el estado completo deseado. El motor intenta mantenerla incrementalmente cuando el plan y las fuentes lo permiten, pero el contrato no promete que todas las transformaciones eviten recomputación. Esto la hace adecuada para agregados, joins y productos gold cuya semántica es una instantánea consistente. El autor debe razonar sobre frescura del refresh, coste de actualización y capacidad de incrementalización. Una materialized view independiente creada desde SQL sigue usando un pipeline administrado por detrás, mientras un proyecto Lakeflow agrupa muchos objetos bajo una misma frontera operativa. Cambiar la definición puede alterar el plan y desencadenar refresh más amplio.

Materialized view

Resultado persistido de una consulta declarativa batch que el pipeline refresca para mantenerlo sincronizado con sus dependencias de datos.

Ofrece lecturas rápidas y consistentes para productos complejos sin recalcular toda la consulta por consumidor.
Incrementalización

Capacidad del motor para transformar cambios upstream en cambios equivalentes del resultado sin recomputar completamente la consulta declarada.

Determina coste y duración de refresh, pero depende del plan y no debe asumirse como garantía universal.
Query fingerprint

Identidad derivada de la definición y plan que ayuda a detectar cuándo cambió la lógica mantenida por una vista materializada.

Permite explicar refresh completos y relacionar variaciones de coste con despliegues concretos del código.
PySparkVista materializada gold
from pyspark import pipelines as dp
from pyspark.sql import functions as F

@dp.materialized_view(name="customer_order_metrics")
def customer_order_metrics():
    orders = spark.read.table("orders_silver")
    return (
        orders.groupBy("customer_id")
          .agg(
              F.countDistinct("order_id").alias("orders"),
              F.sum("amount").alias("lifetime_value"),
              F.max("event_ts").alias("last_order_at"),
          )
    )

Consulta el event log para confirmar si las actualizaciones concretas usan refresh incremental.

Puntos clave

  • La definición expresa el resultado completo, aunque el refresh pueda ser incremental.
  • Una materialized view almacena datos; una vista estándar recalcula al consultar.
  • El event log permite verificar modo y coste del refresh en vez de asumirlo.

Evita

  • Usar `readStream` por reflejo en una materialized view que describe un resultado completo cambiante.
  • Prometer refresh incremental para cualquier UDF o consulta sin observar el plan real.

Recuerdo activo

¿Por qué la definición de una materialized view puede usar una lectura batch y seguir actualizándose incrementalmente?

Borrador privado · solo en este navegador
04
Diagnóstico

Python y SQL en pipelines

Los flows permiten varias entradas hacia un target y distinguen append, AUTO CDC y cargas `ONCE` dentro del mismo modelo declarativo.

Objetivo
Los flows permiten varias entradas hacia un target y distinguen append, AUTO CDC y cargas `ONCE` dentro del mismo modelo declarativo.
Duración estimada
17 min aprox.
Dificultad
Professional
Prerrequisitos
m12
Reportar un error en esta lección

Un default flow acompaña la definición normal de un dataset. Flows adicionales pueden unir feeds regionales en una misma streaming table sin construir un `union` monolítico. Un append flow añade filas; AUTO CDC aplica cambios ordenados; `ONCE` ejecuta una carga batch una sola vez salvo full refresh.

Cada flow debe tener identidad y semántica compatibles con el target. Un target de AUTO CDC solo recibe flows AUTO CDC. Para backfill histórico, un append flow `ONCE` aislado puede coexistir con la ingesta continua, siempre que no duplique rangos ya procesados.

Modelo mental

Un flow es la unidad que describe cómo datos de una fuente llegan a un target. Separar flow y tabla permite que varios orígenes alimenten el mismo dataset bajo reglas claras: un append flow para regiones, un AUTO CDC flow para cambios de clientes o un flow `ONCE` para una carga histórica. La multiplicidad no equivale a permitir escrituras arbitrarias concurrentes; todos los flows forman parte del grafo y del protocolo gestionado del pipeline. `ONCE` ejecuta una carga una sola vez dentro del ciclo de vida registrado y puede volver a ejecutarse en un refresh completo, por lo que su lógica debe ser determinista. AUTO CDC es el nombre vigente recomendado; `APPLY CHANGES` conserva la misma sintaxis y sigue disponible, además de aparecer en material de certificación Professional. El estudiante debe reconocer ambos términos sin presentar el nombre anterior como la opción nueva.

Append flow

Flow declarativo que incorpora registros de una fuente al target sin interpretar cada nueva fila como una actualización de clave existente.

Permite unir fuentes append-only manteniendo progreso y observabilidad separados dentro del mismo pipeline.
Flow ONCE

Flow destinado a una carga finita que se ejecuta una vez en actualizaciones normales y puede repetirse durante full refresh.

Sirve para bootstrap o backfill, pero obliga a escribir lógica determinista y compatible con reconstrucción.
APPLY CHANGES

Nombre anterior todavía disponible para la API cuya opción recomendada actual se denomina AUTO CDC y conserva la misma sintaxis.

Puede aparecer en código legado y en el blueprint Professional, por lo que hay que reconocerlo sin confundir la recomendación vigente.
SQLFlows regionales hacia una tabla
CREATE OR REFRESH STREAMING TABLE main.bronze.orders_all;

CREATE FLOW orders_eu AS INSERT INTO main.bronze.orders_all
BY NAME SELECT *, 'eu' AS region
FROM STREAM(main.raw.orders_eu);

CREATE FLOW orders_us AS INSERT INTO main.bronze.orders_all
BY NAME SELECT *, 'us' AS region
FROM STREAM(main.raw.orders_us);

Valida claves y esquema comunes; `BY NAME` evita depender del orden físico de columnas.

Puntos clave

  • Varios append flows pueden escribir en una misma streaming table.
  • AUTO CDC targets solo aceptan flows AUTO CDC.
  • `ONCE` sirve para una carga finita y vuelve a ejecutarse en un full refresh.

Evita

  • Mezclar append y AUTO CDC sobre el mismo target sin respetar las restricciones del flow.
  • Usar `ONCE` para datos que seguirán llegando y dejar de ingerir silenciosamente después de la primera actualización.

Recuerdo activo

¿Cuándo volvería a ejecutarse un flow marcado `ONCE`?

Borrador privado · solo en este navegador
05
Decisión de diseño

Serverless y modos de ejecución

Un pipeline mantenible separa datasets por dominio y capa, parametriza catálogos/rutas y mantiene efectos operativos fuera de las definiciones.

Objetivo
Un pipeline mantenible separa datasets por dominio y capa, parametriza catálogos/rutas y mantiene efectos operativos fuera de las definiciones.
Duración estimada
17 min aprox.
Dificultad
Professional
Prerrequisitos
m12
Reportar un error en esta lección

Los archivos Python pueden agruparse por bronze, silver y gold o por dominio, siempre que los nombres de dataset sean únicos. Configuración como catálogo fuente, paths y umbrales entra mediante parámetros del pipeline, no constantes replicadas. Las funciones de transformación puras se prueban fuera de los decoradores.

Dev, test y prod ejecutan el mismo código con targets y permisos distintos. Una actualización se valida en un catálogo aislado y con datos representativos antes de promover. El owner revisa event log, lineage y cambios de esquema tras desplegar.

Modelo mental

Un proyecto declarativo mantenible organiza contratos, dominios y capas antes que archivos gigantes. Cada dataset tiene nombre estable, owner, comentario, claves, expectativas y dependencia clara; las funciones de definición son pequeñas y puras. La configuración de entorno —catálogo, schema, ubicaciones, tamaño o modo de ejecución— se inyecta desde el pipeline o bundle y no se codifica en cada notebook. El código compartido puede normalizar columnas y reglas, pero no debe generar dinámicamente un grafo imposible de revisar. Bronze conserva fidelidad de origen, silver aplica contratos y gold sirve productos; dividir por capas solo es útil si cada frontera tiene semántica de recuperación. Tests unitarios cubren funciones de DataFrame, mientras pruebas de integración crean datasets temporales y ejecutan updates. Lakeflow pipelines aporta la operación gestionada; el proyecto sigue necesitando control de versiones, revisión y promoción reproducible.

Definición pura

Función declarativa determinista que construye y devuelve un DataFrame sin ejecutar acciones ni producir efectos externos durante la evaluación.

Garantiza que el mismo código y configuración generen el mismo grafo en validación, despliegue y reintento.
Frontera de pipeline

Conjunto de datasets y flows que comparten actualización, configuración, permisos, observabilidad y estrategia de recuperación coordinada.

Equilibra aislamiento operativo con complejidad y evita agrupar dominios que no deberían fallar o desplegarse juntos.
Promoción reproducible

Proceso que despliega el mismo artefacto versionado en dev, test y prod cambiando únicamente configuración controlada y credenciales.

Reduce divergencias manuales y permite atribuir cada resultado a una versión concreta revisada y probada.
PythonConfiguración de entorno sin duplicar código
source_catalog = spark.conf.get("pipelines.source_catalog")
target_catalog = spark.conf.get("pipelines.target_catalog")
lateness = spark.conf.get("pipelines.orders_lateness", "15 minutes")

source_table = f"{source_catalog}.bronze.orders"
target_prefix = f"{target_catalog}.commerce"

assert source_catalog != target_catalog or target_catalog.endswith("_dev")

Aplica permisos y ownership en la configuración de despliegue; una aserción no sustituye las políticas del entorno.

Puntos clave

  • Configuración cambia por entorno; lógica y artefacto permanecen iguales.
  • Transformaciones puras son comprobables sin arrancar el pipeline completo.
  • Nombres, comentarios y propiedades de tabla forman parte del contrato gobernado.

Evita

  • Copiar el pipeline entero para prod y permitir que las versiones diverjan.
  • Introducir llamadas externas en funciones declarativas y crear resultados no deterministas durante reevaluación.

Recuerdo activo

¿Qué debe variar entre dev y prod?

Borrador privado · solo en este navegador
5 lecciones pendientes

Vista de lectura · sin ejecución

Learn Databricks

Learn Databricks · commit 08c378c

README.md