Mini pipeline distribuido en casa: Proxmox + Spark + Grafana (KPIs de vivienda)

🏠⚡️ En esta entrada monto un pipeline de datos ligero con Proxmox (LXC), Apache Spark y Grafana para calcular y visualizar KPIs de vivienda a partir de un CSV (1460 filas, 81 columnas). La idea: levantar un master + worker de Spark en contenedores, ejecutar un job que compute precio medio global, precio por barrio, por área y por calidad, y publicar los resultados en Grafana con ficheros CSV servidos desde Nextcloud.

✅ Entorno real 100% on-prem: Proxmox (LXC) · Spark Standalone · Python · Nextcloud (URLs de descarga) · Grafana (datasource CSV).


Proxmox con LXC y Spark UI mostrando el worker ALIVE

Arquitectura

  • Proxmox con dos LXC Debian: spark-master y spark-worker (red interna 192.168.1.0/24).
  • Spark Standalone en el master: spark://192.168.1.69:7077; un worker conectado.
  • Nextcloud sirve los CSV finales con URL de descarga pública (solo lectura).
  • Grafana lee esos CSV (datasource CSV → URL) y renderiza paneles.

1) Preparar Spark en el master

# En el master (LXC)
export PYSPARK_PYTHON=/usr/bin/python3
export SPARK_LOCAL_IP=192.168.1.69

# Estructura de trabajo
mkdir -p /opt/spark/jobs /root/warehouse_out

2) Script del job: KPIs de vivienda

Este script lee train.csv (distribuido vía --files), calcula los KPIs y deja salidas en /root/warehouse_out tanto en formato directorio como en CSV único (más cómodo para Grafana).
# /opt/spark/jobs/kpis_real_estate.py

from pyspark.sql import SparkSession, functions as F
from pyspark import SparkFiles

"""
KPIs generados:
- avg_price_zone.csv      → Neighborhood, avg_price (orden desc)
- avg_price_area.csv      → LotAreaRange, avg_price (binning)
- avg_price_quality.csv   → OverallQual, avg_price (1..10)
Además, CSV único por KPI para Grafana.
"""

spark = (SparkSession.builder
         .appName("KPIs Real Estate")
         .config("spark.sql.shuffle.partitions", "2")
         .getOrCreate())

INPUT  = SparkFiles.get("train.csv")
OUTDIR = "/root/warehouse_out"

# Carga
df = (spark.read
      .option("header", True)
      .option("inferSchema", True)
      .csv(INPUT))

# 1) Precio medio por barrio
zone_avg = (df.groupBy("Neighborhood")
              .agg(F.round(F.avg("SalePrice"), 2).alias("avg_price"))
              .orderBy(F.desc("avg_price")))
zone_avg.coalesce(1).write.mode("overwrite").option("header", True)\
       .csv(f"{OUTDIR}/avg_price_zone")
zone_avg.coalesce(1).write.mode("overwrite").option("header", True)\
       .csv(f"{OUTDIR}/avg_price_zone.csv")

# 2) Precio por rangos de área
area_binned = df.withColumn("LotAreaRange",
  F.when(F.col("LotArea") <= 5000, "0–5000") .when((F.col("LotArea") > 5000)  & (F.col("LotArea") <= 10000), "5000–10000") .when((F.col("LotArea") > 10000) & (F.col("LotArea") <= 15000), "10000–15000") .when((F.col("LotArea") > 15000) & (F.col("LotArea") <= 20000), "15000–20000") .otherwise(">20000"))

area_avg = (area_binned.groupBy("LotAreaRange")
            .agg(F.round(F.avg("SalePrice"), 2).alias("avg_price")))
area_avg.coalesce(1).write.mode("overwrite").option("header", True)\
        .csv(f"{OUTDIR}/avg_price_area.csv")

# 3) Precio por calidad
qual_avg = (df.groupBy("OverallQual")
            .agg(F.round(F.avg("SalePrice"), 2).alias("avg_price"))
            .orderBy(F.asc("OverallQual")))
qual_avg.coalesce(1).write.mode("overwrite").option("header", True)\
         .csv(f"{OUTDIR}/avg_price_quality.csv")

spark.stop()

3) Lanzar el job

# CSV de entrada en el master: /root/train.csv

/opt/spark/bin/spark-submit \
  --master spark://192.168.1.69:7077 \
  --deploy-mode client \
  --conf spark.driver.host=192.168.1.69 \
  --driver-memory 512m \
  --executor-memory 512m \
  --conf spark.executor.cores=1 \
  --total-executor-cores=1 \
  --files /root/train.csv \
  /opt/spark/jobs/kpis_real_estate.py
Spark UI con la app FINISHED y un worker ALIVE

4) Salidas generadas

/root/warehouse_out/
  ├── avg_price_zone/           # part-*.csv + _SUCCESS
  ├── avg_price_zone.csv/       # CSV único para Grafana
  ├── avg_price_area.csv/
  └── avg_price_quality.csv/

5) Grafana: paneles y transformaciones

Datasources: CSVSource: URL → pega la URL de descarga de Nextcloud de cada CSV (shared > Download).
  • Gráficos de barras (zona / área / calidad):
    • Parsing options → marcar Skip empty lines y Relax column count.
    • Format: Table. Visualización: Bar chart.
    • Eje X = columna categórica (Neighborhood, LotAreaRange, OverallQual).
    • Orden: desc por avg_price (añade transformación Sort by si lo necesitas).
    • Mostrar valores: Always.
  • Precio medio global (tarjeta Stat):
    • Transform → Convert field type: avg_price → Number.
    • Transform → Reduce fields (Mean).
    • Transform → Organize fields para ocultar columnas no numéricas.
    • Unidad: Euro (€); Display name: “Precio Medio Global”.
  • Extremos (barrio más caro y más barato):
    • Duplicas el panel de zona (Stat o Bar) y aplicas Sort by avg_price asc/desc + Limit 1 (o Reduce Max/Min).
    • Usa dos paneles: “Barrio” (texto) y “Precio” (valor) o un Stat con Name & Value.
Grafana: precio medio por barrio

Resultados

  • Precio medio global€184k.
  • Más caro: NoRidge (~€335k) vs más barato: MeadowV (~€99k).
  • A mayor LotArea, mayor precio medio por tramos (binning).
  • OverallQual (1–10) es un gran predictor del precio.

Solución de problemas

  • FileNotFound con SparkFiles → garantiza --files /root/train.csv y en el script usa SparkFiles.get("train.csv").
  • Grafana “No numeric fields” → añade Convert field type y luego Reduce fields. Si ves texto extra, Organize fields para ocultarlo.
  • CSV multi-part → para Grafana genero también carpeta *.csv/ con un único part-*.csv.
Conclusión: con LXC en Proxmox y Spark Standalone puedes prototipar un flujo distribuido en casa con consumo mínimo. Los CSV resultantes publicados en Nextcloud permiten levantar paneles bonitos en Grafana en minutos. Próximos pasos: guardar en Parquet (TrueNAS), orquestar con cron/Airflow mini y añadir alertas en Grafana.
Todos los derechos reservados © | Política de Privacidad | Política de Cookies
Scroll al inicio