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

Modelo incremental de Structured Streaming

Contenido abierto

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

Lección 1 de 5

Modelo incremental de Structured Streaming

Structured Streaming ejecuta una consulta incremental como una secuencia de microbatches y conserva el progreso necesario para continuar tras un fallo.

Duración
17 min aprox.
Objetivo
Structured Streaming ejecuta una consulta incremental como una secuencia de microbatches y conserva el progreso necesario para continuar tras un fallo.
Siguiente paso
Continuar con la siguiente lección
Ver detalles del módulo

Structured Streaming, triggers y checkpoints

Construye consultas incrementales con estado durable y semántica de recuperación clara.

Al terminar podrás
  • Explicar microbatches y progreso
  • Configurar triggers y checkpoints
  • Diseñar sinks idempotentes
Ver fuentes y revisión

Metadatos editoriales

Última revisión
21 jul 2026
Nivel
Professional
Ruta relacionada
streaming
Dominios blueprint
Data Ingestion and Acquisition · Structured Streaming
Estado
Revisión editorial interna
Fuentes principales
Structured Streaming checkpoints · Databricks · Monitor Structured Streaming queries · Databricks
Reportar un error
01
Modelo mental

Modelo incremental de Structured Streaming

Structured Streaming ejecuta una consulta incremental como una secuencia de microbatches y conserva el progreso necesario para continuar tras un fallo.

Mostrar prerrequisitos
Dificultad
Professional
Prerrequisitos
m12
Reportar un error en esta lección

Una lectura con `readStream` describe un flujo que todavía no está en ejecución. Spark construye el plan, descubre qué datos nuevos están disponibles en cada trigger y solo entonces materializa un microbatch. La misma API de DataFrames permite reutilizar transformaciones batch, pero una operación válida en batch puede ser inviable en streaming si exige conservar estado sin límite.

En un pipeline de pedidos conviene separar tres decisiones: cómo se descubre la entrada, qué transformaciones son incrementales y cómo se confirma la salida. El SLA no lo decide `readStream`; lo determinan el trigger, el volumen por lote, la capacidad y el comportamiento del sink ante reintentos.

Modelo mental

Imagina Structured Streaming como un motor de consultas incrementales que avanza una frontera de datos confirmados, no como un bucle que vuelve a ejecutar una consulta batch completa. `readStream` declara una fuente cuyo final se mueve; las transformaciones construyen un plan lógico y `writeStream` crea una consulta activa. En cada microbatch, Spark determina un rango nuevo de entrada, ejecuta solo ese delta y coordina el resultado con el checkpoint. La API parece idéntica a DataFrames batch porque comparte álgebra, pero la viabilidad cambia: una proyección es stateless, mientras que contar por cliente obliga a recordar información entre lotes. Por eso una solución de producción se diseña alrededor de tres fronteras: qué progreso ofrece la fuente, cuánto estado requiere el plan y cómo confirma el sink. El trigger regula cuándo se intenta avanzar; no garantiza por sí mismo latencia, capacidad ni exactamente una vez.

Consulta incremental

Plan que procesa únicamente el rango de entrada aún no confirmado y conserva continuidad entre ejecuciones.

Evita razonar como si cada trigger recalculara toda la fuente y permite estimar coste y recuperación correctamente.
Microbatch

Unidad transaccional de progreso que agrupa un rango de offsets, su cómputo y el intento de escritura al sink.

Es la frontera práctica para reintentos, métricas, idempotencia y diagnóstico de latencia.
Operador stateful

Transformación que mantiene información de lotes anteriores, como una agregación, deduplicación o join stream-stream.

Su estado condiciona memoria, checkpoint, compatibilidad de cambios y necesidad de watermarks.
PySparkConsulta incremental mínima con destino Delta
orders = spark.readStream.table("main.bronze.orders")

query = (
    orders.where("order_id IS NOT NULL")
      .writeStream
      .option("checkpointLocation", "/Volumes/main/ops/checkpoints/orders_silver")
      .trigger(processingTime="1 minute")
      .toTable("main.silver.orders")
)

El objeto `query` permite consultar estado y progreso; el checkpoint debe ser exclusivo de esta consulta.

Puntos clave

  • `readStream` y `writeStream` definen una consulta continua; una acción batch no la inicia.
  • Cada microbatch procesa un rango identificable de entrada y lo confirma en el checkpoint.
  • Una transformación stateful requiere límites temporales y una estrategia de recuperación explícita.

Evita

  • Confundir la latencia del trigger con el tiempo real de proceso: si un lote tarda dos minutos, un trigger de diez segundos no crea capacidad adicional.
  • Usar el mismo checkpoint para dos consultas o destinos distintos, mezclando offsets y commits incompatibles.

Recuerdo activo

¿Qué parte del código inicia realmente la consulta?

Borrador privado · solo en este navegador
5 lecciones pendientes

Fuente revisada · vista externa

Azure Databricks Hands-on

Azure Databricks Hands-on · commit a91650b

HandsOn.dbc

Archivo importable

Este notebook se abre desde su fuente revisada

El repositorio no permite republicar su contenido dentro de Lakehouse Lab. Conservamos la misma experiencia lateral, la ruta exacta y el commit auditado, y dejamos la lectura en GitHub para respetar la autoría.

Autor
Tsuyoshi Matsuzaki
Licencia
No verificada
Formato
dbc
Abrir / descargar .dbc