Puedes leer todo sin registrarte. Solo crearemos un perfil anónimo cuando decidas guardar tu progreso.
Lección 3 de 5
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.
- Duración
- 17 min aprox.
- 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.
- 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
03Operació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.
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