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

Modelo declarativo y grafo

Contenido abierto

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

Lección 2 de 5

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.

Duración
17 min aprox.
Objetivo
Una streaming table procesa nuevas filas de una fuente streaming y conserva semántica incremental apropiada para bronze y silver append/CDC.
Siguiente paso
Continuar con la siguiente lección
Ver detalles del módulo

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.

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

Mostrar prerrequisitos
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
5 lecciones pendientes

Vista de lectura · sin ejecución

Learn Databricks

Learn Databricks · commit 08c378c

README.md