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.
- Explicar microbatches y progreso
- Configurar triggers y checkpoints
- Diseñar sinks idempotentes
01Modelo mentalModelo 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.
+
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.
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.
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.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.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.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