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

Kafka topics, partitions y consumer offsets

Contenido abierto

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.

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
01
Modelo mental

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.

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

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.

Topic-partition-offset

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.
Record key

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.
Envelope

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.
PySparkLectura con metadatos Kafka preservados
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 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