Puedes leer todo sin registrarte. Solo crearemos un perfil anónimo cuando decidas guardar tu progreso.
Professional
Kafka, buses de eventos y garantías de entrega
Integra sistemas de eventos sin confundir offsets, claves, orden y garantías end-to-end.
- Configurar lectura Kafka con seguridad
- Interpretar particiones y offsets
- Diseñar idempotencia entre source y sink
01Modelo mentalKafka topics, partitions y consumer offsets
Kafka entrega registros binarios con metadatos de topic, partición y offset; deserializar el payload sin perder esa trazabilidad es la primera responsabilidad del consumidor.
+
Kafka topics, partitions y consumer offsets
Kafka entrega registros binarios con metadatos de topic, partición y offset; deserializar el payload sin perder esa trazabilidad es la primera responsabilidad del consumidor.
- Objetivo
- Kafka entrega registros binarios con metadatos de topic, partición y offset; deserializar el payload sin perder esa trazabilidad es la primera responsabilidad del consumidor.
- Duración estimada
- 17 min aprox.
- Dificultad
- Professional
- Prerrequisitos
- m14
El DataFrame de Kafka expone `key`, `value`, `topic`, `partition`, `offset`, `timestamp` y headers. `key` y `value` llegan como bytes, por lo que hacer solo `CAST(value AS STRING)` no valida el contrato. Un esquema explícito con `from_json` permite distinguir una fila corrupta de una evolución compatible.
Conservar `topic`, `partition` y `offset` en bronze permite investigar pérdidas, reconstruir un rango y demostrar qué mensaje originó una fila. La clave Kafka también importa: el productor la usa para particionar y preservar orden dentro de una partición, pero Kafka no ofrece orden global entre particiones.
Modelo mental
Kafka no entrega objetos de negocio; entrega registros ordenados dentro de particiones, identificados por topic, partición y offset, con key, value, headers y timestamps en representación binaria. El consumidor debe preservar primero ese sobre técnico y deserializar después el payload con un contrato versionado. El offset no es un identificador global ni representa el orden entre particiones: solo avanza dentro de una partición concreta. La key influye en el particionado del productor y, por tanto, en qué entidades comparten orden. Structured Streaming proyecta esos metadatos como columnas y conserva su progreso en el checkpoint. Eliminarlos antes de validar dificulta auditar duplicados, reconstruir un rango o localizar un productor defectuoso. Una arquitectura robusta mantiene una capa bronze con bytes y metadatos inmutables, añade resultado de parsing y errores, y solo promueve a silver los eventos que satisfacen el contrato.
Coordenada única de un registro dentro de la retención de un topic Kafka.
Permite auditoría, replay preciso y deduplicación técnica sin asumir un offset global inexistente.Bytes usados normalmente por el productor para elegir partición y agrupar el orden de entidades relacionadas.
Una key estable preserva orden por entidad y evita particiones calientes por estrategias defectuosas.Conjunto de metadatos técnicos y payload crudo que rodea al evento de negocio.
Conservarlo en bronze hace diagnosticables los errores de contrato y posibilita volver a decodificar sin releer Kafka.kafka_raw = (
spark.readStream.format("kafka")
.option("kafka.bootstrap.servers", bootstrap_servers)
.option("subscribe", "orders.v1")
.load()
)
bronze = kafka_raw.selectExpr(
"CAST(key AS STRING) AS message_key",
"CAST(value AS STRING) AS payload",
"topic", "partition", "offset",
"timestamp AS kafka_timestamp"
)No registres credenciales en opciones ni notebooks; obtén endpoints y secretos desde configuración gobernada.
Puntos clave
- El offset solo es único dentro de un par topic-partición.
- La key dirige particionado y orden local; no sustituye la clave de negocio del payload.
- Bronze debe conservar metadatos Kafka y el payload original o una referencia auditable.
Evita
- Tratar `(offset)` como identificador global y provocar colisiones entre particiones.
- Descartar metadatos inmediatamente, haciendo imposible demostrar qué mensaje produjo un error.
Recuerdo activo
¿Qué combinación identifica de forma inequívoca una posición Kafka?
Borrador privado · solo en este navegador02Implementaciónkey, value, headers y timestamps
Las opciones de suscripción y offsets solo fijan el inicio de una consulta nueva; al reanudar, el checkpoint gobierna la posición.
+
key, value, headers y timestamps
Las opciones de suscripción y offsets solo fijan el inicio de una consulta nueva; al reanudar, el checkpoint gobierna la posición.
- Objetivo
- Las opciones de suscripción y offsets solo fijan el inicio de una consulta nueva; al reanudar, el checkpoint gobierna la posición.
- Duración estimada
- 17 min aprox.
- Dificultad
- Professional
- Prerrequisitos
- m14
`subscribe` elige topics concretos, `subscribePattern` usa una expresión regular y `assign` fija particiones. Debe configurarse exactamente uno. En streaming, `startingOffsets` es `latest` por defecto y solo se consulta si no existe progreso previo; cambiarlo después no rebobina una consulta con checkpoint.
Para un backfill se usa un checkpoint y destino separados o una lectura batch con rangos explícitos, no se modifica a ciegas una consulta productiva. `failOnDataLoss=false` permite continuar cuando offsets ya no existen, pero acepta una posible pérdida y debe acompañarse de reconciliación; no es una solución genérica para errores.
Modelo mental
La posición inicial de Kafka se decide una vez, cuando nace una consulta sin checkpoint. Después, el checkpoint es la autoridad sobre offsets; cambiar `startingOffsets` no rebobina una consulta existente. Esta distinción evita dos errores frecuentes: creer que `latest` salta datos en cada reinicio o intentar un backfill modificando opciones mientras se reutiliza el mismo estado. `subscribe` sigue topics explícitos, `assign` fija particiones concretas y `subscribePattern` descubre topics que coinciden con un patrón; cada opción cambia la topología y debe gobernarse. `earliest` procesa lo aún retenido, no una historia ilimitada. Si Kafka ha eliminado offsets que el checkpoint necesita, la decisión entre fallar, saltar o reconstruir afecta completitud y no debe ocultarse. Un replay fiable usa una consulta y destino separados con límites de offsets, dejando intacta la continuidad de producción.
Posición usada para inicializar cada partición únicamente cuando no existe progreso restaurable en un checkpoint.
Aclara por qué cambiar `earliest` o `latest` no modifica una consulta ya iniciada.Regla que determina qué topics y particiones forman la fuente, mediante subscribe, patrón o asignación explícita.
Forma parte de la identidad de la consulta y condiciona descubrimiento, permisos y compatibilidad de recuperación.Política por la que el broker elimina segmentos antiguos independientemente del progreso del consumidor.
Define la ventana máxima para recuperar backlog o hacer replay directamente desde Kafka.orders = (
spark.readStream.format("kafka")
.option("kafka.bootstrap.servers", bootstrap_servers)
.option("subscribe", "orders.v1")
.option("startingOffsets", "earliest")
.option("failOnDataLoss", "true")
.option("maxOffsetsPerTrigger", 500000)
.load()
)Tras el primer commit, `startingOffsets` deja de decidir la posición; el checkpoint continúa desde los offsets confirmados.
Puntos clave
- Una consulta reanudada toma offsets del checkpoint, no de `startingOffsets`.
- `earliest` en una consulta nueva puede consumir toda la retención y generar un backlog considerable.
- `failOnDataLoss=false` cambia una garantía de integridad y exige una decisión operativa explícita.
Evita
- Cambiar `startingOffsets` a `earliest` esperando que un stream existente relea su historia.
- Desactivar `failOnDataLoss` para silenciar una retención insuficiente sin medir el hueco perdido.
Recuerdo activo
¿Qué ocurre si se cambia `startingOffsets` en una consulta que conserva el mismo checkpoint?
Borrador privado · solo en este navegador03OperaciónAutenticación y secretos
Un contrato de evento separa deserialización, validación y evolución para que los mensajes incompatibles no detengan ni contaminen el flujo válido.
+
Autenticación y secretos
Un contrato de evento separa deserialización, validación y evolución para que los mensajes incompatibles no detengan ni contaminen el flujo válido.
- Objetivo
- Un contrato de evento separa deserialización, validación y evolución para que los mensajes incompatibles no detengan ni contaminen el flujo válido.
- Duración estimada
- 17 min aprox.
- Dificultad
- Professional
- Prerrequisitos
- m14
Con JSON, un `StructType` explícito evita inferencia por microbatch y permite detectar `_corrupt_record` o campos obligatorios nulos. Avro o Protobuf con un schema registry ofrecen contratos más fuertes, pero siguen necesitando una política de compatibilidad entre productor y consumidor.
Un patrón robusto publica el mensaje crudo en bronze, enruta fallos a cuarentena con el motivo y solo expone columnas tipadas en silver. Añadir un campo opcional suele ser compatible; renombrar o cambiar tipo requiere una migración coordinada o lectura multi-versión mediante un campo `schema_version`.
Modelo mental
Un contrato de evento define cómo convertir bytes en un hecho confiable y cómo evolucionará esa conversión. Incluye formato, versión, campos requeridos, tipos, semántica, compatibilidad y tratamiento de datos desconocidos. El parsing no debe mezclarse con reglas de negocio: primero se determina si el payload puede interpretarse; después se valida si el hecho es aceptable, y finalmente se transforma. Esta separación evita que una nueva columna opcional se confunda con corrupción o que un importe negativo derribe el decoder. En esquemas evolucionados, añadir un campo opcional con valor predeterminado suele ser compatible; renombrar, cambiar tipo o reinterpretar unidades puede romper consumidores aunque el JSON siga siendo válido. Una capa bronze guarda el original y la versión; silver materializa un esquema canónico y una cuarentena explicable. El objetivo no es aceptar todo, sino hacer explícita y reversible cada decisión de compatibilidad.
Capacidad de una versión nueva del lector para interpretar datos escritos con versiones anteriores del contrato.
Permite desplegar consumidores antes o después sin bloquear el historial retenido.Registro cuya forma o contenido provoca de manera determinista el fallo repetido del consumidor.
Debe aislarse con evidencia para evitar que retries conviertan un error de datos en indisponibilidad del pipeline.Representación interna estable a la que se normalizan versiones externas compatibles.
Desacopla la evolución del productor de todas las transformaciones downstream y simplifica tests.from pyspark.sql import functions as F, types as T
order_schema = T.StructType([
T.StructField("order_id", T.StringType(), False),
T.StructField("event_ts", T.TimestampType(), False),
T.StructField("amount", T.DecimalType(18, 2), False),
T.StructField("schema_version", T.IntegerType(), False),
])
decoded = bronze.withColumn(
"event", F.from_json("payload", order_schema)
)Mantén `payload` hasta que la fila haya superado validación para poder diagnosticar y reintentar.
Puntos clave
- El esquema se valida antes de aplicar reglas de negocio.
- La cuarentena conserva payload y coordenadas Kafka para remediación.
- La evolución se negocia entre productor y consumidor; no se resuelve concediendo tipos `string` a todo.
Evita
- Inferir el esquema en cada ejecución y aceptar cambios accidentales del productor.
- Descartar mensajes inválidos sin contador, payload ni coordenadas de origen.
Recuerdo activo
¿Qué datos mínimos necesita una cuarentena Kafka para reprocesar una fila?
Borrador privado · solo en este navegador04DiagnósticoAt-least-once y exactly-once práctico
Kafka, checkpoint y Delta pueden ofrecer procesamiento exactamente una vez, pero cualquier sink externo vuelve a exigir idempotencia end-to-end.
+
At-least-once y exactly-once práctico
Kafka, checkpoint y Delta pueden ofrecer procesamiento exactamente una vez, pero cualquier sink externo vuelve a exigir idempotencia end-to-end.
- Objetivo
- Kafka, checkpoint y Delta pueden ofrecer procesamiento exactamente una vez, pero cualquier sink externo vuelve a exigir idempotencia end-to-end.
- Duración estimada
- 17 min aprox.
- Dificultad
- Professional
- Prerrequisitos
- m14
Spark registra rangos de offsets en el checkpoint y Delta confirma cada microbatch de manera transaccional. Un fallo puede hacer que el motor vuelva a calcular el lote, pero el protocolo del sink evita materializarlo dos veces. Esto no deduplica dos mensajes distintos que el productor publicó con el mismo `order_id`.
Cuando el destino es una API, una base no transaccional o varios sistemas, `foreachBatch` ofrece semántica at-least-once a menos que la función sea idempotente. Un outbox del productor, identificadores de evento estables y `MERGE` en silver resuelven capas distintas del problema.
Modelo mental
El recorrido Kafka–Spark–Delta puede acercarse a exactamente una vez porque cada componente ofrece una identidad durable: offsets por partición, commits por microbatch y transacciones Delta. Sin embargo, esa composición depende de no introducir una frontera que desconozca el protocolo. El checkpoint hace que Spark vuelva a leer un rango cuando no quedó confirmado; Delta puede hacer que la repetición converja al mismo resultado si se usa el sink integrado o un `MERGE` determinista. Una API externa, otra base o dos tablas escritas secuencialmente pueden observar intentos parciales. Por eso se diseña una salida canónica única y se derivan efectos posteriores con consumidores independientes. También se distingue duplicado de origen —dos registros Kafka con el mismo evento— de reintento técnico —el mismo rango ejecutado otra vez—: el primero requiere clave de negocio, el segundo coordinación o idempotencia del sink.
Coordenada del transporte, como topic-partition-offset o batch id, que identifica un intento dentro de una ejecución.
Permite rastrear reintentos, pero no sustituye una identidad de negocio entre reconstrucciones.Tabla transaccional de efectos pendientes escrita junto con el estado canónico y consumida de forma independiente.
Evita intentar una transacción distribuida con APIs externas y hace reparables las publicaciones parciales.Propiedad por la que reejecutar datos conduce al mismo estado final pese a intentos repetidos o desordenados.
Es una formulación práctica y comprobable de corrección para pipelines recuperables.MERGE INTO main.silver.orders AS target
USING staged_orders AS source
ON target.order_id = source.order_id
WHEN MATCHED AND source.event_ts > target.event_ts THEN
UPDATE SET *
WHEN NOT MATCHED THEN
INSERT *;Antes del `MERGE`, conserva una sola fila ganadora por `order_id` dentro del microbatch para evitar múltiples matches.
Puntos clave
- Exactly-once de procesamiento no elimina duplicados creados por el productor.
- Un checkpoint no puede deshacer un efecto externo ya confirmado fuera de Spark.
- La clave `event_id` permite deduplicación de negocio además de coordinación técnica de offsets.
Evita
- Prometer exactly-once porque Kafka usa offsets, aunque el sink sea una API sin clave idempotente.
- Confundir reejecución del mismo offset con dos mensajes distintos enviados por el productor.
Recuerdo activo
¿Puede el checkpoint eliminar un cargo duplicado ya creado en una API externa?
Borrador privado · solo en este navegador05Decisión de diseñoBackpressure, 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.
+
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.
- 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.
- Duración estimada
- 17 min aprox.
- Dificultad
- Professional
- Prerrequisitos
- m14
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.
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 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.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.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