Arquitectura Medallion en Spark: De JSON a Bronze, Silver y Gold con Scala
Tres laboratorios que concentran algunas de las decisiones más importantes al diseñar pipelines analíticos con Apache Spark: abandonar formatos ineficientes, imponer contratos de datos tipados y resolver analítica avanzada con Window Functions. El objetivo no es únicamente conseguir que el job funcione, sino construir una plataforma que escale en rendimiento, coste, mantenibilidad y gobernanza.
Lo que aprenderás
Los Laboratorios 2, 3 y 4 forman una progresión natural: primero optimizamos la representación física de los datos, después establecemos un contrato semántico con tipos y finalmente explotamos esos datos para resolver una consulta analítica real.
Stack tecnológico
El laboratorio está diseñado alrededor de Spark SQL y la API de Spark en Scala, utilizando formatos y abstracciones que son especialmente relevantes en plataformas Data Lake y Lakehouse.
Matriz de decisión: JSON, Parquet, Dataset y Window Functions
No todos estos elementos compiten entre sí. Cada uno resuelve una capa distinta del problema: representación física, contrato de datos y lógica analítica.
| Elemento | Función | Rendimiento | Coste / I/O | Cuándo utilizarlo |
|---|---|---|---|---|
| JSON | Formato textual y flexible para intercambio o ingestión. | Bajo para analítica a escala. | Alto volumen de bytes y parsing. | Bronze, APIs, eventos o sistemas fuente que ya producen JSON. |
| Parquet | Formato columnar optimizado para analítica. | Alto gracias a compresión, estadísticas y column pruning. | Mucho menor I/O cuando el acceso es selectivo. | Silver/Gold y almacenamiento analítico persistente. |
| Dataset[T] | Representación tipada de registros en Scala. | Compatible con el optimizador de Spark. | No es un formato físico; define el contrato de aplicación. | Transformaciones donde seguridad de tipos y dominio importan. |
| Window Function | Analítica relativa dentro de particiones. | Distribuida, pero puede requerir shuffle y ordenación. | Coste dependiente de particionamiento y cardinalidad. | Rankings, Top-N, acumulados, deduplicación y series temporales. |
La arquitectura Medallion como hilo conductor
La arquitectura Medallion separa la ingesta de la normalización y del consumo analítico. Esto permite que cada capa tenga una responsabilidad clara y evita que los consumidores tengan que interpretar directamente los formatos de los sistemas fuente.
Bronze — preservar el origen
La capa Bronze conserva los datos lo más cerca posible de la fuente. Si el sistema entrega JSON, puede tener sentido aterrizar inicialmente ese JSON para mantener trazabilidad y capacidad de replay.
Principio: raw first, transform later.
Silver — imponer calidad y estructura
Aquí ocurre la transformación crítica: parsing, normalización, casting, validación y escritura en Parquet. El objetivo es que el formato de consumo interno sea eficiente y tenga un esquema coherente.
Principio: schema + quality + efficient storage.
Gold — resolver el negocio
Gold no debería ser simplemente una copia de Silver. Debe contener datasets preparados para casos de uso concretos: KPIs, rankings, agregados y tablas dimensionales o de serving.
Principio: optimize for consumption.
El error conceptual
Medallion no significa necesariamente tres carpetas sin más. Es una separación lógica de responsabilidades. En producción puede implementarse con tablas, catálogos, particiones, políticas de retención y contratos de datos.
Principio: architecture, not folder naming.
Antipatrón crítico: utilizar JSON como almacenamiento analítico
⚠ El problema no es que JSON sea “malo” en términos absolutos
El problema aparece cuando un formato diseñado principalmente para intercambio y legibilidad humana se utiliza como almacenamiento principal de grandes volúmenes analíticos.
JSON es textual y normalmente exige que Spark lea y parsee registros completos. En cambio, Parquet organiza los datos por columnas y permite que el motor aproveche column pruning, compresión y estadísticas para reducir la cantidad de información que debe leer.
Por eso, afirmar que “Parquet ahorra un 85%” puede ser perfectamente posible en un laboratorio concreto, pero no debe convertirse en una constante arquitectónica universal. El ratio real depende de cardinalidad, tipos, compresibilidad, contenido, codec, distribución y estructura de los datos.
La métrica correcta no es solamente GB almacenados: hay que observar bytes leídos, duración del job, shuffle, CPU, número de archivos y coste total de la consulta.
Implementación práctica en Scala
El siguiente pipeline muestra el patrón completo: leer JSON, normalizarlo mediante un Dataset tipado, persistirlo como Parquet y finalmente utilizar una Window Function para calcular el Top 3.
import org.apache.spark.sql.{Dataset, SparkSession}
import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._
case class SongPlay(
country: String,
day: String,
song: String,
artist: String,
plays: Long
)
object MusicMedallionPipeline {
def main(args: Array[String]): Unit = {
val spark = SparkSession.builder()
.appName("MusicMedallionPipeline")
.getOrCreate()
import spark.implicits._
// -------------------------------------------------------------------------
// BRONZE
// Datos de origen: JSON.
// En producción, esta capa debería conservar trazabilidad y replay.
// -------------------------------------------------------------------------
val bronze = spark.read
.option("multiLine", "false")
.json("data/bronze/music.json")
// -------------------------------------------------------------------------
// SILVER
// Convertimos los registros a un Dataset tipado.
// El cast explícito evita depender de inferencias ambiguas del origen.
// -------------------------------------------------------------------------
val silver: Dataset[SongPlay] = bronze
.select(
col("country").cast("string"),
col("day").cast("string"),
col("song").cast("string"),
col("artist").cast("string"),
col("plays").cast("long")
)
.as[SongPlay]
.filter(
$"country".isNotNull &&
$"day".isNotNull &&
$"song".isNotNull &&
$"plays".isNotNull
)
// Parquet pasa a ser el formato persistente de la capa Silver.
silver.write
.mode("overwrite")
.partitionBy("country", "day")
.parquet("data/silver/music")
// -------------------------------------------------------------------------
// GOLD
// Ranking diario por país.
//
// row_number() garantiza exactamente N filas por partición.
// Si se desean empates con la misma posición, valorar dense_rank().
// -------------------------------------------------------------------------
val rankingWindow = Window
.partitionBy($"country", $"day")
.orderBy(
$"plays".desc,
$"song".asc
)
val top3 = silver
.withColumn("position", row_number().over(rankingWindow))
.filter($"position" <= 3)
// Dataset Gold listo para consumo analítico.
top3.write
.mode("overwrite")
.partitionBy("country", "day")
.parquet("data/gold/top3_songs")
spark.stop()
}
}
Datasets tipados: Scala frente a DataFrames y PySpark
La principal ventaja de una case class no es hacer que Spark “procese más
rápido”. Su valor está en el contrato que establece entre el dominio de la aplicación
y los registros que manipula.
Scala Dataset[T]
Permite modelar explícitamente entidades como SongPlay. El compilador
puede detectar errores de tipos en el código Scala antes de ejecutar el job.
Es especialmente útil cuando el pipeline contiene lógica de dominio significativa, transformaciones reutilizables y equipos que quieren aprovechar el sistema de tipos de Scala.
DataFrame / PySpark
Un DataFrame ofrece una API estructurada y es una opción excelente para pipelines declarativos, SQL, exploración y transformaciones basadas en expresiones.
La ausencia de tipado estático del lenguaje no significa ausencia de esquema: el DataFrame sigue teniendo un schema de Spark que debe gobernarse y validarse.
Window Functions: Top 3 de canciones por país y día
Un ranking como “las tres canciones más escuchadas de cada país cada día” no se
resuelve correctamente con un simple groupBy. Primero necesitamos
establecer el universo de comparación mediante una ventana.
1. Definir la partición
partitionBy(country, day) crea un ranking independiente para cada
combinación de país y día.
2. Definir el orden
orderBy(plays.desc) coloca las canciones con mayor número de
reproducciones primero. Añadir una segunda clave como song.asc proporciona
determinismo cuando dos canciones tienen las mismas reproducciones.
3. Elegir la función de ranking
row_number() devuelve una posición única por fila. Es apropiado cuando el requisito es exactamente Top 3.
dense_rank() conserva empates y puede producir más de tres filas. La elección depende de la semántica del negocio, no de una preferencia sintáctica.
4. Filtrar después de calcular la ventana
Finalmente, filtramos position <= 3. Este patrón es una herramienta
fundamental para problemas de Top-N, deduplicación, latest-record y análisis
secuencial.
Decisiones de diseño que importan en producción
Particionar no significa “particionar todo”
Particionar Parquet por país y día puede acelerar lecturas selectivas, pero una cardinalidad excesiva puede generar demasiados directorios y archivos pequeños. El diseño debe responder a los patrones reales de consulta.
Evitar small files
Un pipeline que escribe miles de archivos diminutos puede degradar el rendimiento aunque el formato sea Parquet. El tamaño y número de archivos debe formar parte de la estrategia de operación del Data Lake.
El shuffle es el enemigo silencioso
Las Window Functions requieren normalmente redistribuir y ordenar datos dentro de las particiones. En datasets grandes, hay que inspeccionar el plan físico y vigilar skew, particiones desbalanceadas y presión de memoria.
Coste = almacenamiento + compute + I/O
El ahorro de Parquet no debe medirse únicamente por el tamaño final. Un diseño eficiente también reduce bytes leídos y trabajo de CPU, lo que puede tener un impacto mayor que el almacenamiento bruto.
Schema evolution bajo control
La flexibilidad de JSON resulta útil durante la ingesta, pero trasladar esa flexibilidad sin gobierno a Silver crea deuda técnica. Los cambios de schema deben detectarse, versionarse y tratarse explícitamente.
Idempotencia y replay
Una arquitectura Medallion empresarial debe permitir reejecutar una partición o ventana temporal sin duplicar datos. La estrategia de escritura y las claves de negocio deben diseñarse pensando en este requisito.
Framework de implementación paso a paso
Una implementación empresarial puede estructurarse como una secuencia de decisiones verificables.
Define el contrato de Bronze
Identifica qué datos deben conservarse del origen, qué metadatos de ingesta se necesitan y cómo se podrá reproducir un procesamiento fallido.
Normaliza hacia Silver
Aplica tipos explícitos, validación de nulls, normalización de nombres y reglas de calidad. Persistir en Parquet debe ser una decisión consciente, no una simple conversión de extensión.
Diseña particiones según consultas
Elige las columnas de particionamiento basándote en selectividad, cardinalidad y patrones de acceso. Mide antes de asumir.
Introduce tipos donde aporten valor
Utiliza case class y Dataset[T] para modelos de dominio
donde el tipado estático mejore la seguridad y mantenibilidad del pipeline.
Implementa Top-N con Window Functions
Define primero la partición, después el orden y finalmente la función de ranking.
Decide conscientemente entre row_number, rank y
dense_rank.
Observa el plan físico
No optimices únicamente el código fuente. Examina exchanges, sorts, scans, particiones y volumen de datos. El plan físico explica dónde está realmente el coste.
Gobierna y opera
Añade métricas, validaciones de calidad, control de schema, alertas, políticas de retención y mecanismos de replay. El pipeline termina cuando puede operarse, no cuando compila.
Checklist de arquitectura
Preguntas frecuentes sobre Spark, Scala y Medallion
¿Por qué convertir JSON a Parquet en una arquitectura Spark?
JSON es un formato textual y flexible, pero resulta poco eficiente como almacenamiento analítico a gran escala. Parquet es columnar y permite a Spark aplicar técnicas como column pruning y aprovechar compresión y estadísticas. Por eso, la conversión de Bronze a Silver suele reducir I/O y coste de procesamiento. El porcentaje exacto de ahorro, sin embargo, depende de los datos y debe medirse.
¿Qué ventaja aporta un Dataset tipado con case class frente a un DataFrame?
En Scala, un Dataset[T] permite representar los registros mediante un
tipo de dominio como case class SongPlay. Esto proporciona comprobaciones
de tipos en las partes de la API tipada y hace que el contrato del modelo sea más
explícito. No sustituye la validación del schema de entrada ni significa que un
Dataset sea siempre superior a un DataFrame.
¿Cómo obtener el Top 3 de canciones por país y día?
Hay que crear una ventana con partitionBy(country, day), ordenar las
canciones por reproducciones descendentes y aplicar row_number() o una
función de ranking alternativa. Después se filtran las filas cuyo ranking sea menor
o igual que 3. Si los empates deben conservarse, dense_rank() puede ser
más apropiado.
