Puedes leer todo sin registrarte. Solo crearemos un perfil anónimo cuando decidas guardar tu progreso.
Professional
Proyecto de streaming con SLA
Entrega un flujo operable que soporte datos tardíos, recuperación e incidentes reproducibles.
- Cumplir SLA de frescura y completitud
- Recuperar sin duplicados
- Crear métricas y runbook
01Modelo mentalDefinición del SLA
Un SLA de streaming debe traducirse en indicadores medibles de frescura, completitud, corrección y disponibilidad, con ventanas y responsables explícitos.
+
Definición del SLA
Un SLA de streaming debe traducirse en indicadores medibles de frescura, completitud, corrección y disponibilidad, con ventanas y responsables explícitos.
- Objetivo
- Un SLA de streaming debe traducirse en indicadores medibles de frescura, completitud, corrección y disponibilidad, con ventanas y responsables explícitos.
- Duración estimada
- 30 min aprox.
- Dificultad
- Professional
- Prerrequisitos
- m16
Decir 'tiempo real' no permite diseñar ni operar. Un SLO útil puede exigir que el p95 de `published_at - event_ts` sea inferior a cinco minutos y que al menos el 99,5% de eventos válidos aparezcan antes de quince minutos. La completitud se reconcilia con una fuente autoritativa, no con que el stream siga activo.
Cada indicador necesita una fuente de medición, periodo y presupuesto de error. Un watermark de diez minutos no garantiza por sí mismo diez minutos de frescura: el backlog, el trigger y el sink también cuentan. La arquitectura se valida contra los SLO, no al revés.
Modelo mental
Un SLA de streaming no es la frase tiempo real; es un contrato cuantificado entre productor, plataforma y consumidor. Se descompone en indicadores: frescura del evento publicado, completitud respecto a la fuente, corrección de reglas, disponibilidad de lectura y tiempo de recuperación. Cada indicador necesita método, ventana, percentil, presupuesto de error y dueño. La frescura puede medirse como reloj menos máximo event time válido, mientras que la latencia de microbatch es solo un componente interno. Un p95 de cinco minutos permite colas ocasionales que un máximo estricto no permitiría. RPO expresa cuántos datos podrían perderse y RTO cuánto se tarda en restaurar el servicio; checkpoints y replay deben demostrar ambos. El contrato también define comportamiento degradado: si Kafka se retrasa, puede ser mejor servir datos marcados como stale que bloquear una tabla consistente y disponible.
Medida concreta del comportamiento observado, como frescura p95 o porcentaje de eventos reconciliados.
Convierte expectativas ambiguas en datos sobre los que se pueden alertar y mejorar.Objetivo interno para un SLI durante una ventana, normalmente más estricto que el límite contractual externo.
Crea margen operativo y guía decisiones de capacidad y fiabilidad antes de incumplir el SLA.Máxima pérdida de datos aceptable y máximo tiempo para recuperar el servicio tras un incidente.
Conecta checkpoints, retención, replay y runbooks con compromisos verificables de continuidad.WITH stream AS (
SELECT
date_trunc('hour', event_ts) AS hour,
percentile_approx(
unix_timestamp(published_at) - unix_timestamp(event_ts),
0.95
) AS freshness_p95_s,
count(DISTINCT event_id) AS published_events
FROM main.silver.clickstream
WHERE event_ts >= current_timestamp() - INTERVAL 24 HOURS
GROUP BY 1
)
SELECT * FROM stream ORDER BY hour DESC;La completitud necesita comparar `published_events` con un conteo independiente de la fuente o un manifiesto de productores.
Puntos clave
- Frescura compara tiempo de negocio publicado con el reloj; disponibilidad solo indica si el proceso está ejecutándose.
- Completitud requiere un denominador o reconciliación autoritativa.
- El presupuesto de error determina cuándo una desviación es incidente y qué cambios se priorizan.
Evita
- Medir solo duración de microbatch y llamarla frescura de negocio.
- Fijar un SLA sin definir zona horaria, percentil, ventana ni exclusiones de eventos inválidos.
Recuerdo activo
¿Por qué una consulta `ACTIVE` puede incumplir frescura?
Borrador privado · solo en este navegador02ImplementaciónDiseño de eventos y particiones
La arquitectura separa entrada auditable, transformación stateful, publicación idempotente y observabilidad para que cada frontera pueda recuperarse.
+
Diseño de eventos y particiones
La arquitectura separa entrada auditable, transformación stateful, publicación idempotente y observabilidad para que cada frontera pueda recuperarse.
- Objetivo
- La arquitectura separa entrada auditable, transformación stateful, publicación idempotente y observabilidad para que cada frontera pueda recuperarse.
- Duración estimada
- 30 min aprox.
- Dificultad
- Professional
- Prerrequisitos
- m16
Bronze conserva el evento original y coordenadas de origen. Silver valida esquema, deduplica y aplica reglas temporales. Gold sirve agregados con la latencia acordada. Cada consulta tiene un checkpoint exclusivo y cada tabla una clave o contrato que permite recomponerla desde una capa anterior.
La capacidad se diseña para la tasa sostenida más un margen de recuperación. Si tras una hora de caída llegan diez millones de eventos, el pipeline debe procesar más rápido que la tasa normal para volver al SLA. Separar ingestión de agregación evita que un operador stateful pesado bloquee la captura de entrada.
Modelo mental
Una arquitectura streaming recuperable separa fronteras de responsabilidad. La entrada bronze conserva el sobre original y una identidad reproducible; las transformaciones stateful usan checkpoints propios; la salida canónica se confirma idempotentemente; observabilidad registra tanto progreso técnico como verdad de negocio. Acoplar todas las etapas en una única consulta parece reducir latencia, pero amplía el dominio de fallo y hace que una API lenta bloquee ingestión. Separarlas con tablas Delta añade un commit y algo de latencia, a cambio de replay, aislamiento y evolución independiente. Cada consumidor tiene su checkpoint, por lo que reparar uno no rebobina los demás. El modelo medallion no es solo organización de calidad: sus commits son fronteras de recuperación. Los schemas, claves, secuencias y expectativas forman contratos versionados, y los efectos externos se derivan después de que exista una fuente canónica auditable.
Commit durable desde el que una etapa downstream puede reanudar o reconstruirse sin releer el sistema original.
Reduce radio de impacto y hace posible reparar una transformación sin afectar ingestión.Representación gobernada que constituye la verdad publicada para una etapa o dominio.
Desacopla efectos externos y consumidores de reintentos y formatos del transporte.Conjunto de etapas, datos y consumidores afectados cuando falla o cambia un componente.
Guía la decisión de separar consultas y checkpoints aunque aumente ligeramente la latencia.pipeline: clickstream
sources:
- main.bronze.web_events
targets:
silver: main.silver.web_events
gold: main.gold.sessions_5m
slo:
freshness_p95_seconds: 300
completeness_percent: 99.5
recovery:
rto_minutes: 30
replay_source: main.bronze.web_events
owner: data-platform-oncallEste contrato debe acompañarse de consultas que calculen cada indicador y enlaces al runbook.
Puntos clave
- Bronze inmutable proporciona replay y auditoría.
- Checkpoints independientes reducen el radio de impacto de un cambio o fallo.
- Capacidad de recuperación debe superar la tasa de llegada, no solo sostenerla.
Evita
- Acoplar ingestión, enrichments externos y agregación en una sola consulta sin punto de replay intermedio.
- Dimensionar únicamente para el promedio y no poder reducir backlog después de una interrupción.
Recuerdo activo
¿Qué condición permite que un pipeline recupere backlog?
Borrador privado · solo en este navegador03OperaciónEstado, CDC y calidad
Un plan de recuperación distingue reinicio, replay y rebuild, y nunca borra estado antes de capturar evidencia y delimitar el rango afectado.
+
Estado, CDC y calidad
Un plan de recuperación distingue reinicio, replay y rebuild, y nunca borra estado antes de capturar evidencia y delimitar el rango afectado.
- Objetivo
- Un plan de recuperación distingue reinicio, replay y rebuild, y nunca borra estado antes de capturar evidencia y delimitar el rango afectado.
- Duración estimada
- 30 min aprox.
- Dificultad
- Professional
- Prerrequisitos
- m16
Un fallo transitorio del executor suele resolverse reanudando desde el mismo checkpoint. Un cambio incompatible de estado requiere checkpoint nuevo y reconstrucción desde bronze. Un error lógico publicado exige identificar versiones o tiempos afectados, corregir código y escribir de forma idempotente en un destino aislado antes del cutover.
La retención de Kafka, CDF, archivos y Delta debe cubrir el RTO y el máximo tiempo de detección. Si el origen ya eliminó datos, restaurar compute no recupera completitud. Los runbooks incluyen criterios para pausar productores o consumidores, rutas alternativas y validación posterior.
Modelo mental
Recuperar no siempre significa reiniciar. Un reinicio restaura el mismo plan desde un checkpoint compatible tras un fallo transitorio. Un replay vuelve a procesar un intervalo conservando una fuente durable, normalmente hacia staging o una versión nueva. Un rebuild reconstruye estado y destino completos cuando cambió la semántica o se perdió compatibilidad. Elegir mal puede ocultar pérdida o duplicar efectos. Antes de actuar se inmoviliza evidencia: error, batch, offsets, versiones Delta, estado de sinks y métricas. Borrar checkpoint convierte recuperación en consulta nueva y elimina la frontera que permitía razonar. El runbook debe indicar condiciones de entrada, autoridad, RPO/RTO, validaciones y rollback. Un replay no se publica directamente sobre producción sin reconciliar claves, secuencias, conteos e invariantes; además debe evitar efectos externos hasta aprobar el resultado.
Continuación del mismo plan lógico desde el último checkpoint compatible después de un fallo operativo.
Es la opción de menor impacto cuando código, estado y datos siguen siendo válidos.Reprocesamiento deliberado de un intervalo histórico desde una fuente durable con límites explícitos.
Corrige resultados sin destruir el progreso de la consulta activa y ofrece comparación antes de publicar.Reconstrucción completa o sustancial del estado y salida con un checkpoint nuevo.
Es necesaria ante cambios incompatibles o corrupción, pero exige planificación de corte, coste y reconciliación.SELECT
min(event_ts) AS first_affected,
max(event_ts) AS last_affected,
count(*) AS affected_rows,
count(DISTINCT event_id) AS affected_events
FROM main.silver.web_events
WHERE pipeline_version = '2026.07.20-bad'
AND event_date BETWEEN DATE '2026-07-20' AND DATE '2026-07-21';Repara primero en una tabla sombra y compara claves/conteos antes de reemplazar o hacer `MERGE` en producción.
Puntos clave
- Reinicio conserva checkpoint; replay usa un rango y estado aislados; rebuild reconstruye una tabla completa.
- Captura `lastProgress`, offsets, versión de código y error antes de mutar estado.
- La retención de origen es un requisito de recuperación, no solo una decisión de coste.
Evita
- Borrar checkpoint como primer paso y perder el punto exacto del incidente.
- Reprocesar todo el histórico en el mismo destino sin aislar escrituras ni evitar duplicados.
Recuerdo activo
¿Cuándo es apropiado conservar el checkpoint durante la recuperación?
Borrador privado · solo en este navegador04DiagnósticoObservabilidad y alertas
La observabilidad combina progreso Spark, lag de fuente, calidad de datos y SLO de negocio para ofrecer alertas accionables.
+
Observabilidad y alertas
La observabilidad combina progreso Spark, lag de fuente, calidad de datos y SLO de negocio para ofrecer alertas accionables.
- Objetivo
- La observabilidad combina progreso Spark, lag de fuente, calidad de datos y SLO de negocio para ofrecer alertas accionables.
- Duración estimada
- 30 min aprox.
- Dificultad
- Professional
- Prerrequisitos
- m16
El progreso de la consulta muestra duración, tasas y estado; Kafka o Auto Loader aportan backlog; las tablas silver aportan máximo event time y descartes. Una alerta de frescura debe enlazar esas señales para diferenciar fuente parada, cuello de proceso, dato inválido o sink lento.
Las métricas se escriben en una tabla con `query_id`, `batch_id`, versión y timestamp. Un dashboard sin alertas ni propietario no reduce tiempo de recuperación. Cada alarma debe incluir umbral, periodo, severidad, enlace a evidencia y primera acción segura.
Modelo mental
La observabilidad útil enlaza cuatro planos: salud de la consulta, progreso de la fuente, comportamiento del estado y calidad del producto. Un dashboard que solo muestra cluster y estado `RUNNING` puede permanecer verde mientras se publican datos viejos o incompletos. Las métricas técnicas —offsets, batch duration, input rate, state rows, retries— responden cómo funciona el motor. Las métricas de negocio —máximo event time, pedidos por mercado, suma de importes, porcentaje válido— responden si entrega el producto prometido. Cada alerta debe combinar duración y severidad para evitar ruido durante microbatches vacíos o ráfagas normales, y debe señalar una primera comprobación. Logs estructurados incluyen run, query, batch, checkpoint version y correlación de la fuente; tablas de operaciones permiten tendencias y postmortems más allá de la retención de la interfaz.
Telemetría sobre ejecución, offsets, recursos, estado y fallos del motor streaming.
Localiza mecanismos operativos, pero necesita contexto de negocio para determinar impacto real.Indicadores sobre frescura, volumen, calidad y coherencia del producto de datos entregado.
Detecta pipelines técnicamente activos que producen resultados inútiles o incompletos.Regla con umbral sostenido, severidad, owner y primera hipótesis o acción segura asociada.
Reduce fatiga y transforma una señal en tiempo de diagnóstico y recuperación menor.import json
from pyspark.sql import Row
progress = query.lastProgress or {}
record = Row(
query_id=str(query.id),
batch_id=progress.get("batchId"),
observed_at=progress.get("timestamp"),
input_rps=progress.get("inputRowsPerSecond", 0.0),
processed_rps=progress.get("processedRowsPerSecond", 0.0),
progress_json=json.dumps(progress),
)
spark.createDataFrame([record]).write.mode("append").saveAsTable(
"main.ops.streaming_progress"
)Para una solución productiva usa un listener o monitor administrado y controla volumen/PII del JSON de progreso.
Puntos clave
- Una única métrica técnica rara vez explica un incumplimiento de negocio.
- Las alertas usan ventanas sostenidas para evitar ruido de un microbatch aislado.
- Versión de código y query ID permiten correlacionar regresiones con despliegues.
Evita
- Alertar por una tasa baja cuando la fuente no tiene eventos y el SLO sigue satisfecho.
- Guardar logs sin `batch_id`, versión ni owner, impidiendo correlación y respuesta.
Recuerdo activo
¿Qué señales distinguen una fuente parada de un consumidor lento?
Borrador privado · solo en este navegador05Decisión de diseñoGame day y recuperación
Un game day verifica con fallos controlados que checkpoints, capacidad, alertas y runbook cumplen el RTO sin duplicar ni perder datos.
+
Game day y recuperación
Un game day verifica con fallos controlados que checkpoints, capacidad, alertas y runbook cumplen el RTO sin duplicar ni perder datos.
- Objetivo
- Un game day verifica con fallos controlados que checkpoints, capacidad, alertas y runbook cumplen el RTO sin duplicar ni perder datos.
- Duración estimada
- 30 min aprox.
- Dificultad
- Professional
- Prerrequisitos
- m16
La prueba puede detener compute durante diez minutos, introducir un evento tardío y reiniciar desde el mismo checkpoint. Se mide tiempo hasta recuperar frescura, lag máximo, duplicados y completitud. Otra prueba despliega un cambio stateful incompatible en un entorno aislado para practicar rollback o rebuild.
El runbook se escribe como decisiones observables: si el checkpoint está íntegro y el código es compatible, reanudar; si faltan offsets, escalar pérdida y reconciliar; si el sink tiene efectos externos, verificar idempotencia. El resultado del ejercicio genera acciones con responsable y fecha.
Modelo mental
Un game day convierte supuestos de resiliencia en evidencia mediante fallos controlados. No busca demostrar que nada falla; comprueba que detección, recuperación, idempotencia y comunicación funcionan dentro de RTO/RPO. El experimento define hipótesis, alcance, guardas de seguridad, datos sintéticos o reversibles, responsables y criterio de aborto. Se eligen fallos representativos: matar el driver durante un commit, ralentizar un sink, detener una partición, enviar poison pills o llenar backlog. Antes se registra el estado esperado y después se reconcilia cada evento. Reiniciar con éxito no basta: hay que demostrar que no faltan ni sobran claves, que el state store recuperó, que alertas llegaron al owner y que el runbook no exigió conocimiento tribal. Los resultados alimentan capacidad, automatización y documentación; un fallo del ejercicio es aprendizaje antes de una incidencia real.
Afirmación medible sobre cómo responderá el sistema a un fallo concreto bajo condiciones definidas.
Permite declarar éxito o fracaso con evidencia en vez de aceptar que el Job volvió a verde.Límite técnico u operativo que contiene el impacto del experimento y activa aborto o rollback.
Hace posible probar escenarios realistas sin convertir el aprendizaje en un incidente incontrolado.Comparación de identidades, conteos, importes y efectos antes y después de recuperar.
Demuestra RPO e idempotencia, propiedades que una captura de estado RUNNING no puede probar.experiment: stop-streaming-compute
preconditions:
- bronze_backlog_is_zero
- baseline_freshness_p95_lt_300s
fault:
duration_minutes: 10
success:
- no_duplicate_event_ids
- completeness_percent_gte_99.5
- freshness_recovered_within_30m
rollback:
- resume_original_job
- preserve_checkpoint
owner: data-platform-oncallNo ejecutes un experimento de resiliencia en producción sin límites, observabilidad y autorización explícitos.
Puntos clave
- Define estado inicial y criterios de éxito antes de inyectar el fallo.
- Mide RTO y calidad final, no solo que el proceso vuelve a estado RUNNING.
- El game day debe ser reversible, aislado y aprobado por el owner del servicio.
Evita
- Declarar éxito cuando el query reinicia aunque el backlog y los duplicados sigan creciendo.
- Probar borrado de checkpoint productivo sin snapshot, aislamiento ni procedimiento de rollback.
Recuerdo activo