Saltar al contenido

lección 7

Leer y escribir Parquet a escala: Spark + Data Lake

Spark y Parquet son la pareja perfecta: columnar, compresión, predicate pushdown y particionado por carpetas.

50 min

### La pareja perfecta: Spark + Parquet

En la Skill 12 (Data Lakes) aprendiste que Parquet es un formato columnar, comprimido y con metadatos integrados. Ahora vas a ver por qué Spark y Parquet juntos son la combinación estándar de la industria para Data Lakes. No es coincidencia: Spark fue diseñado con Parquet en mente, y Parquet fue diseñado para sistemas como Spark.

La magia está en las optimizaciones automáticas: cuando Spark lee un Parquet, puede saltar columnas que no necesitas (projection pushdown), saltar bloques enteros que no cumplen tu filtro (predicate pushdown), y leer solo las particiones relevantes (partition pruning). Con CSV, Spark tiene que leer TODO el archivo siempre.

### CSV vs Parquet: la diferencia en números

1from pyspark.sql import SparkSession
2from pyspark.sql.functions import col
3import time
4
5spark = SparkSession.builder.appName("parquet-vs-csv").master("spark://spark-master:7077").getOrCreate()
6
7# Leer CSV
8start = time.time()
9ventas_csv = spark.read.option("header", "true").option("inferSchema", "true").csv("/home/jovyan/data/ventas.csv")
10ventas_csv.filter(col("ciudad") == "Madrid").select("producto", "precio_unitario").count()
11t_csv = time.time() - start
12
13# Primero, escribir como Parquet para poder comparar
14ventas_csv.write.mode("overwrite").parquet("/home/jovyan/data/ventas.parquet")
15
16# Leer Parquet
17start = time.time()
18ventas_pq = spark.read.parquet("/home/jovyan/data/ventas.parquet")
19ventas_pq.filter(col("ciudad") == "Madrid").select("producto", "precio_unitario").count()
20t_parquet = time.time() - start
21
22print(f"CSV: {t_csv:.2f}s")
23print(f"Parquet: {t_parquet:.2f}s")
24print(f"Parquet es {t_csv / t_parquet:.1f}x más rápido")
25
26# Comparar tamaño en disco
27import os
28import subprocess
29# CSV: ~60 MB para 1M filas
30# Parquet: ~15 MB (compresión Snappy por defecto, solo guarda lo necesario)

Parquet es más rápido Y más pequeño que CSV — sin esfuerzo extra

### Escribir Parquet con particionado por carpetas

La killer feature para Data Lakes: particionar los datos en carpetas por una o más columnas. Si particionas por "ciudad", Spark crea una carpeta por cada ciudad (ciudad=Madrid/, ciudad=Barcelona/...). Cuando filtras por ciudad="Madrid", Spark SOLO lee la carpeta de Madrid y ignora las demás. Con 1 TB de datos y 10 ciudades, leer una ciudad tarda 10x menos.

1# Escribir con particionado por ciudad
2(
3 ventas_csv
4 .write
5 .mode("overwrite")
6 .partitionBy("ciudad") # ← crea carpetas: ciudad=Madrid/, ciudad=Barcelona/, etc.
7 .parquet("/home/jovyan/output/ventas_particionadas/")
8)
9
10# Estructura en disco:
11# ventas_particionadas/
12# ciudad=Madrid/
13# part-00000.parquet
14# part-00001.parquet
15# ciudad=Barcelona/
16# part-00000.parquet
17# ciudad=Valencia/
18# part-00000.parquet
19# ...
20
21# Leer con partition pruning (solo lee Madrid)
22ventas_madrid = (
23 spark.read
24 .parquet("/home/jovyan/output/ventas_particionadas/")
25 .filter(col("ciudad") == "Madrid")
26 # Spark SOLO lee la carpeta ciudad=Madrid/ — ignora el resto
27)
28
29print(f"Filas Madrid: {ventas_madrid.count():,}")

Particionado por carpetas: Spark solo lee los datos que necesita

Consejo de senior: elige la columna de partición sabiamente. Ideal: columnas con baja cardinalidad (10-1000 valores) que uses frecuentemente en filtros. Fecha es la elección #1 (partición por año/mes o por día). NO particiones por una columna con millones de valores únicos — crearás millones de carpetas con archivos diminutos.

### El patrón Data Lake: zonas raw → silver → gold

En producción, un pipeline Spark típico sigue las zonas del Data Lake que viste en la Skill 12:

  1. 01.Raw (bronze): datos crudos tal cual llegan. CSV, JSON, logs. Se guardan sin modificar.
  2. 02.Silver: datos limpiados y tipados. Se convierten a Parquet con esquema validado.
  3. 03.Gold: datos agregados y modelados, listos para consumo por BI y analytics.
1from pyspark.sql.functions import col, to_date, when, lit, current_timestamp
2
3# === ZONA RAW → SILVER: limpiar y tipar ===
4raw = spark.read.option("header", "true").csv("/home/jovyan/data/ventas.csv")
5
6silver = (
7 raw
8 # Tipar columnas
9 .withColumn("id", col("id").cast("integer"))
10 .withColumn("fecha", to_date(col("fecha"), "yyyy-MM-dd"))
11 .withColumn("cantidad", col("cantidad").cast("integer"))
12 .withColumn("precio_unitario", col("precio_unitario").cast("double"))
13 .withColumn("cliente_id", col("cliente_id").cast("integer"))
14 # Limpiar: eliminar filas sin fecha o con cantidad <= 0
15 .filter(col("fecha").isNotNull())
16 .filter(col("cantidad") > 0)
17 # Añadir metadatos
18 .withColumn("procesado_en", current_timestamp())
19)
20
21# Escribir silver particionado por mes
22(
23 silver
24 .withColumn("mes", col("fecha").substr(1, 7)) # "2024-01"
25 .write.mode("overwrite")
26 .partitionBy("mes")
27 .parquet("/home/jovyan/output/lake/silver/ventas/")
28)
29
30# === ZONA SILVER → GOLD: agregar para negocio ===
31gold_ventas_diarias = (
32 spark.read.parquet("/home/jovyan/output/lake/silver/ventas/")
33 .groupBy("fecha", "ciudad", "categoria")
34 .agg(
35 spark_sum("cantidad").alias("unidades_vendidas"),
36 spark_sum(col("cantidad") * col("precio_unitario")).alias("facturacion"),
37 countDistinct("cliente_id").alias("clientes_unicos"),
38 )
39)
40
41gold_ventas_diarias.write.mode("overwrite").partitionBy("fecha").parquet("/home/jovyan/output/lake/gold/ventas_diarias/")
42print("✓ Pipeline raw → silver → gold completado")

Pipeline completo: raw CSV → silver Parquet limpio → gold Parquet agregado

Cuidado con los "small files problem": si particionas por día Y ciudad Y categoría, puedes generar miles de archivos diminutos (1 KB cada uno). Los archivos pequeños son ineficientes para lectura — Spark tiene overhead por cada archivo que abre. Regla: cada archivo Parquet debería pesar al menos 100 MB. Usa coalesce antes de escribir si es necesario.

### Modos de escritura

  • overwrite: borra todo lo anterior y escribe de nuevo. Simple pero peligroso.
  • append: añade archivos nuevos sin borrar los existentes. Para ingesta incremental.
  • ignore: no escribe si la ruta ya existe. Para idempotencia simple.
  • error/errorifexists (default): falla si la ruta ya existe.
1# Modos de escritura
2df.write.mode("overwrite").parquet("/ruta/") # Borra y reescribe
3df.write.mode("append").parquet("/ruta/") # Añade sin borrar
4df.write.mode("ignore").parquet("/ruta/") # No escribe si existe
5df.write.mode("error").parquet("/ruta/") # Falla si existe (default)
6
7# Controlar compresión
8df.write.option("compression", "snappy").parquet("/ruta/") # default, rápido
9df.write.option("compression", "gzip").parquet("/ruta/") # más comprimido, más lento
10df.write.option("compression", "zstd").parquet("/ruta/") # mejor ratio, moderno

Modos de escritura y opciones de compresión

## ejercicios

[01]

Convertir CSV a Parquet optimizado

Lee el CSV de ventas, conviértelo a Parquet particionado por categoría, con compresión snappy y exactamente 3 archivos por partición. Verifica que la lectura filtrada es más rápida.

Cargando editor...
[02]

Pipeline completo: raw → silver → gold

Implementa un pipeline de 3 capas. Raw: lee el CSV. Silver: limpia (elimina nulos, tipa, añade columna "ingreso" = cantidad * precio). Gold: facturación diaria por ciudad.

Cargando editor...
[03]

Validar esquema antes de escribir

Lee el Parquet de silver y verifica que tiene el esquema esperado antes de construir gold. Si falta alguna columna, lanza un error claro.

Cargando editor...

Regístrate para guardar tu progreso.

## comentarios

Reporta erratas, ayuda a otros o comparte tu opinión. Sé constructivo.

Inicia sesión para comentar y responder.

cargando comentarios...