TEL-420 Sistemas Paralelos P5d taller mapreduce spark · M5
M5 Guía de Práctica / Taller
Descargar .md

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:

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.