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.
- Crear tablas streaming y materialized views
- Comparar declarativo con Structured Streaming
- Distinguir Spark Declarative Pipelines de la oferta gestionada Lakeflow
02ImplementaciónModelo 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.
+
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.
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.
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 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.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.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