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

Backpressure, lag y capacidad

Contenido abierto

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

Lección 5 de 5

Backpressure, lag y capacidad

Particiones, límites de offsets y métricas de lag permiten regular throughput sin confundir una protección temporal con una solución de capacidad.

Duración
17 min aprox.
Objetivo
Particiones, límites de offsets y métricas de lag permiten regular throughput sin confundir una protección temporal con una solución de capacidad.
Siguiente paso
Continuar con el laboratorio
Ver detalles del módulo

Kafka, buses de eventos y garantías de entrega

Integra sistemas de eventos sin confundir offsets, claves, orden y garantías end-to-end.

Al terminar podrás
  • Configurar lectura Kafka con seguridad
  • Interpretar particiones y offsets
  • Diseñar idempotencia entre source y sink
Ver fuentes y revisión

Metadatos editoriales

Última revisión
21 jul 2026
Nivel
Professional
Ruta relacionada
streaming
Dominios blueprint
Message buses · Streaming ingestion
Estado
Revisión editorial interna
Fuentes principales
Kafka connector for Structured Streaming · Databricks · Spark API options reference · Databricks
Reportar un error
05
Decisión de diseño

Backpressure, lag y capacidad

Particiones, límites de offsets y métricas de lag permiten regular throughput sin confundir una protección temporal con una solución de capacidad.

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

El paralelismo máximo de lectura está condicionado por las particiones Kafka, aunque Spark pueda dividir rangos grandes en más tareas según capacidades de la fuente. `maxOffsetsPerTrigger` limita el volumen de un microbatch y protege al sink durante recuperación, pero si queda por debajo de la tasa de llegada el lag crecerá indefinidamente.

Las métricas `avgOffsetsBehindLatest`, `maxOffsetsBehindLatest` y bytes estimados muestran atraso por fuente. Deben correlacionarse con duración del trigger, distribución de particiones y tasas de procesamiento. Una sola partición caliente puede dominar el SLA aunque la media parezca saludable.

Modelo mental

El throughput de Kafka está acotado inicialmente por sus particiones: una partición proporciona un flujo ordenado que una tarea consume por rango de offsets, mientras que varias particiones permiten paralelismo. Más particiones no garantizan equilibrio si la key concentra tráfico, y más workers que particiones no crean lectores útiles. Structured Streaming puede limitar offsets por trigger para proteger memoria, estado o sink durante picos. Ese límite regula admisión; no aumenta capacidad sostenida. Si la llegada media supera el proceso medio, el lag seguirá creciendo aunque los lotes sean pequeños. El objetivo operativo es mantener headroom, detectar skew por partición y estimar tiempo de vaciado del backlog. Cambiar particionado también afecta orden por entidad y puede requerir coordinación con productores, no es solo una opción del consumidor.

Consumer lag

Diferencia por partición entre el último offset disponible en Kafka y el offset confirmado por la consulta.

Mide backlog real y permite estimar si la frescura se recupera o se deteriora.
Partición caliente

Partición que recibe o procesa mucha más carga que las demás por distribución sesgada de keys.

Limita el lote completo y no se resuelve simplemente añadiendo workers o usando un promedio global.
Control de admisión

Límite deliberado sobre cuántos offsets entran en un microbatch.

Protege downstream durante ráfagas, pero debe distinguirse de una mejora de capacidad permanente.
PythonExtracción de lag desde el progreso
progress = query.lastProgress or {}
for source in progress.get("sources", []):
    metrics = source.get("metrics", {})
    print({
        "description": source.get("description"),
        "avg_lag": metrics.get("avgOffsetsBehindLatest"),
        "max_lag": metrics.get("maxOffsetsBehindLatest"),
        "bytes_behind": metrics.get("estimatedTotalBytesBehindLatest"),
    })

Alerta por tendencia y tiempo estimado de recuperación, no por un valor aislado durante un pico esperado.

Puntos clave

  • `maxOffsetsPerTrigger` controla lote, no aumenta capacidad.
  • La key del productor puede crear skew persistente entre particiones.
  • Mide lag máximo y por partición, además de la media.

Evita

  • Reducir el lote hasta estabilizar la duración mientras el backlog crece silenciosamente.
  • Añadir workers cuando el topic tiene una sola partición caliente y no ofrece paralelismo útil.

Recuerdo activo

¿Por qué un `maxOffsetsPerTrigger` demasiado bajo puede incumplir el SLA aunque cada lote termine rápido?

Borrador privado · solo en este navegador
5 lecciones pendientes

Fuente revisada · vista externa

Structured Streaming with Event Hubs or Kafka

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