Puedes leer todo sin registrarte. Solo crearemos un perfil anónimo cuando decidas guardar tu progreso.
Lección 1 de 5
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.
- Duración
- 17 min aprox.
- 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.
- Siguiente paso
- Continuar con la siguiente lección
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.
- 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.
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