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 SparkSession2from pyspark.sql.functions import col3import time45spark = SparkSession.builder.appName("parquet-vs-csv").master("spark://spark-master:7077").getOrCreate()67# Leer CSV8start = 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() - start1213# Primero, escribir como Parquet para poder comparar14ventas_csv.write.mode("overwrite").parquet("/home/jovyan/data/ventas.parquet")1516# Leer Parquet17start = 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() - start2122print(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")2526# Comparar tamaño en disco27import os28import subprocess29# CSV: ~60 MB para 1M filas30# 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 ciudad2(3 ventas_csv4 .write5 .mode("overwrite")6 .partitionBy("ciudad") # ← crea carpetas: ciudad=Madrid/, ciudad=Barcelona/, etc.7 .parquet("/home/jovyan/output/ventas_particionadas/")8)910# Estructura en disco:11# ventas_particionadas/12# ciudad=Madrid/13# part-00000.parquet14# part-00001.parquet15# ciudad=Barcelona/16# part-00000.parquet17# ciudad=Valencia/18# part-00000.parquet19# ...2021# Leer con partition pruning (solo lee Madrid)22ventas_madrid = (23 spark.read24 .parquet("/home/jovyan/output/ventas_particionadas/")25 .filter(col("ciudad") == "Madrid")26 # Spark SOLO lee la carpeta ciudad=Madrid/ — ignora el resto27)2829print(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:
- 01.Raw (bronze): datos crudos tal cual llegan. CSV, JSON, logs. Se guardan sin modificar.
- 02.Silver: datos limpiados y tipados. Se convierten a Parquet con esquema validado.
- 03.Gold: datos agregados y modelados, listos para consumo por BI y analytics.
1from pyspark.sql.functions import col, to_date, when, lit, current_timestamp23# === ZONA RAW → SILVER: limpiar y tipar ===4raw = spark.read.option("header", "true").csv("/home/jovyan/data/ventas.csv")56silver = (7 raw8 # Tipar columnas9 .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 <= 015 .filter(col("fecha").isNotNull())16 .filter(col("cantidad") > 0)17 # Añadir metadatos18 .withColumn("procesado_en", current_timestamp())19)2021# Escribir silver particionado por mes22(23 silver24 .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)2930# === 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)4041gold_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 escritura2df.write.mode("overwrite").parquet("/ruta/") # Borra y reescribe3df.write.mode("append").parquet("/ruta/") # Añade sin borrar4df.write.mode("ignore").parquet("/ruta/") # No escribe si existe5df.write.mode("error").parquet("/ruta/") # Falla si existe (default)67# Controlar compresión8df.write.option("compression", "snappy").parquet("/ruta/") # default, rápido9df.write.option("compression", "gzip").parquet("/ruta/") # más comprimido, más lento10df.write.option("compression", "zstd").parquet("/ruta/") # mejor ratio, moderno
Modos de escritura y opciones de compresión
## ejercicios
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.
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.
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.
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...