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).
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 leetrain.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
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: CSV → Source: 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”.
- Transform → Convert field type:
- 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.
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.csvy en el script usaSparkFiles.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 únicopart-*.csv.