Puedes leer todo sin registrarte. Solo crearemos un perfil anónimo cuando decidas guardar tu progreso.
Lección 4 de 5
Offsets, commits y checkpoints
La garantía extremo a extremo depende tanto del seguimiento de offsets como de que el sink confirme lotes de forma idempotente.
- Duración
- 17 min aprox.
- Objetivo
- La garantía extremo a extremo depende tanto del seguimiento de offsets como de que el sink confirme lotes de forma idempotente.
- 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
04DiagnósticoOffsets, commits y checkpoints
La garantía extremo a extremo depende tanto del seguimiento de offsets como de que el sink confirme lotes de forma idempotente.
+
Offsets, commits y checkpoints
La garantía extremo a extremo depende tanto del seguimiento de offsets como de que el sink confirme lotes de forma idempotente.
Los sinks Delta integrados coordinan los commits de cada microbatch con el checkpoint. Si una tarea falla después de escribir pero antes de confirmar, el reinicio puede volver a presentar el mismo lote y el sink debe reconocerlo. La garantía de la fuente por sí sola no impide duplicados en un servicio externo.
`foreachBatch` abre la puerta a `MERGE`, múltiples destinos o APIs, pero transfiere responsabilidad al código. El `batch_id` y una clave de negocio estable permiten registrar el lote o hacer upsert determinista. Enviar eventos a un endpoint sin clave idempotente conserva, como mucho, semántica at-least-once.
Modelo mental
Exactly-once no es una propiedad que la fuente Kafka, el checkpoint o Delta puedan declarar aisladamente; es un resultado extremo a extremo. Spark puede volver a presentar un microbatch cuando no sabe si el intento anterior quedó publicado. Un sink transaccional integrado puede reconocer el lote y coordinarlo con el progreso, mientras que una API REST, un correo o dos destinos independientes no participan automáticamente en ese protocolo. Por eso el modelo mental correcto separa entrega de procesamiento: at-least-once admite reintentos y posible repetición; idempotencia hace que repetir produzca el mismo estado final. `foreachBatch` ofrece toda la expresividad batch, incluido `MERGE`, pero traslada al autor la responsabilidad de claves, orden y atomicidad. Un `batch_id` identifica un intento lógico de la consulta; una clave de negocio identifica el hecho y suele sobrevivir mejor a reconstrucciones con un checkpoint nuevo.
Garantía de que un registro no se pierde, aunque un fallo pueda provocar uno o más intentos de entrega.
Obliga a diseñar consumidores que toleren repetición en vez de inferir unicidad por la existencia del checkpoint.Propiedad por la que aplicar varias veces la misma operación deja el mismo estado que aplicarla una vez.
Convierte reintentos inevitables en recuperación segura para `MERGE`, APIs y efectos externos.Identificador estable del hecho o entidad, independiente del lote y del intento técnico que lo transporta.
Permite deduplicar incluso después de reconstruir una consulta con otro checkpoint o redistribuir eventos.from delta.tables import DeltaTable
def upsert_orders(batch_df, batch_id):
target = DeltaTable.forName(spark, "main.silver.orders_current")
(target.alias("t")
.merge(batch_df.alias("s"), "t.order_id = s.order_id")
.whenMatchedUpdateAll()
.whenNotMatchedInsertAll()
.execute())
(orders.writeStream
.foreachBatch(upsert_orders)
.option("checkpointLocation", "/Volumes/main/ops/checkpoints/orders_upsert")
.start())El `MERGE` debe resolver múltiples cambios de la misma clave dentro del lote antes del upsert.
Puntos clave
- Exactly-once es una propiedad del recorrido fuente–estado–sink, no una etiqueta aislada.
- `foreachBatch` permite lógica batch por microbatch y debe tolerar la repetición del mismo lote.
- Una clave de negocio suele ser más útil que confiar únicamente en el número de lote.
Evita
- Hacer `append` en `foreachBatch` y asumir que el checkpoint evita cualquier repetición posterior al fallo.
- Ejecutar varias acciones sobre `batch_df` sin persistirlo cuando el coste de recomputación sea relevante.
Recuerdo activo