---
titulo: "Práctica P5d: Taller de procesamiento distribuido a gran escala con Apache Spark y MapReduce"
modulo: M5
ec: EC5
contenido: "5.4"
semanas: "16"
horas: 4
tipo: "Taller en laboratorio"
asignatura: "Sistemas Paralelos"
sigla: TEL-420
docente: "Ing. Elias Cassal Baldiviezo"
institucion: "Universidad Autónoma Juan Misael Saracho — Facultad de Ciencias Integradas de Yacuiba / F.I.R.N.T."
---

# P5d · Taller de procesamiento masivo de datos con Apache Spark y MapReduce

**Contenido del programa:** §20.5 · contenido 5.4  
**Horas de laboratorio:** 4 · **Modalidad:** individual o parejas  

---

## 1. Propósito

Implementar algoritmos paralelos sobre conjuntos masivos de datos (*Big Data*) utilizando el modelo funcional **MapReduce** y el motor en memoria **Apache Spark (PySpark)**, comprendiendo el concepto de colecciones distribuidas tolerantes a fallos (RDD / DataFrames), transformaciones perezosas (*Lazy Transformations*) y acciones.

---

## 2. Fundamentación técnica

A diferencia de MPI, donde el programador gestiona manualmente el envío y recepción de mensajes y el código se detiene si un solo nodo falla, los frameworks de Big Data como Apache Spark ofrecen:
- **Tolerancia a fallos automática:** Los datos se particionan y se registra su linaje (*Lineage Graph*); si un nodo colapsa, Spark recalcula únicamente la partición perdida en otro nodo trabajador.
- **Procesamiento en memoria:** Mantiene los datos en la memoria RAM distribuida entre transformaciones sucesivas, siendo hasta $100\text{x}$ más rápido que el almacenamiento intermedio en disco de Hadoop clásico.

---

## 3. Actividades guiadas

### Ejercicio 1: Pipeline MapReduce en PySpark
Procesar un dataset de registros climáticos o focos de calor del Gran Chaco:

```python
from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("AnalisisFocosCalorChaco") \
    .master("local[*]") \
    .getOrCreate()

sc = spark.sparkContext

# Lectura distribuida
lineas = sc.textFile("data/focos_calor_chaco.csv")

# Pipeline MapReduce: Conteo por municipio
resultado = lineas \
    .filter(lambda l: not l.startswith("fecha")) \
    .map(lambda l: l.split(",")) \
    .map(lambda cols: (cols[2].strip(), 1)) \
    .reduceByKey(lambda a, b: a + b) \
    .sortBy(lambda par: par[1], ascending=False)

for municipio, total in resultado.take(10):
    print(f"Municipio: {municipio:<20} | Focos: {total}")

spark.stop()
```

### Ejercicio 2: Monitoreo en Spark Web UI
1. Lanzar el trabajo con `spark-submit`.
2. Acceder a la interfaz web de Spark (`http://localhost:4040`).
3. Inspeccionar el grafo acíclico dirigido (DAG), las etapas de computación (*Stages*) y las fases de redistribución de datos por red (*Shuffle Read / Shuffle Write*).

---

## 4. Entregables
Código fuente del script Spark y reporte de métricas analizando cómo influye la fase de *Shuffle* en el tiempo total de procesamiento.
