Logo GCP con Eduardo

GCP con Eduardo

Descarga el código de la lección

Arquitectura Medallion en Spark: JSON a Parquet, Datasets Tipados y Window Functions en Scala
✦ Guía Técnica & Arquitectura

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.

Laboratorio técnico impartido por Eduardo Martínez Agrelo, AI & Data Architect.

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.

Por qué JSON es una mala elección como formato analítico y por qué Parquet puede reducir radicalmente el almacenamiento y el I/O.
Cómo implementar Datasets tipados con case class en Scala para hacer explícito el contrato de datos.
Cómo utilizar Window Functions para obtener el Top 3 de canciones por país y día sin recurrir a lógica distribuida artesanal.
Cómo encajar Bronze, Silver y Gold dentro de una arquitectura Medallion orientada a producción.

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.

Engine Apache Spark Language Scala SQL Spark SQL Format JSON Format Parquet API Dataset[T] Analytics Window Functions Architecture Medallion

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.

Scala · Apache Spark · Medallion Pipeline
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.

1

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.

2

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.

3

Diseña particiones según consultas

Elige las columnas de particionamiento basándote en selectividad, cardinalidad y patrones de acceso. Mide antes de asumir.

4

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.

5

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.

6

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.

7

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

¿El JSON está limitado a la zona de ingestión cuando realmente es necesario?
¿Silver está persistido en un formato columnar adecuado para las consultas?
¿El schema de los datos está versionado y validado?
¿Las particiones están justificadas por patrones de consulta?
¿Se monitorizan small files y skew?
¿La Window Function utiliza la semántica de ranking correcta?
¿El job puede reejecutarse sin generar duplicados?
¿Se mide coste total y no únicamente tamaño de almacenamiento?

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.

Sobre el Autor: Eduardo Martínez Agrelo

AI & Data Architect

Eduardo Martínez Agrelo trabaja en la intersección entre inteligencia artificial, ingeniería de datos y arquitectura tecnológica. Su enfoque combina fundamentos técnicos, diseño de plataformas y decisiones pragmáticas orientadas a producción, con especial atención a rendimiento, escalabilidad, calidad y coste.

En estos laboratorios, el objetivo es ir más allá de la sintaxis de Spark: entender por qué una decisión de almacenamiento cambia el coste, cómo un sistema de tipos puede mejorar un pipeline y cuándo una Window Function es la abstracción correcta para resolver un problema analítico distribuido.

© Eduardo Martínez Agrelo · AI & Data Architect · Guía técnica de Apache Spark y Scala