TEL-420 Sistemas Paralelos Módulo M3 · MPI y cústeres

Programación Paralela en Memoria Distribuida (MPI)

Módulo M3 de la asignatura. Diseñado conforme al apartado 20 del programa docente oficial.

EC3 Semanas 7-10 16 horas 28 % de la nota Contenido completo

Introducción del módulo

El módulo 2 abordó la paralelización en memoria compartida: múltiples hilos ejecutándose en los diferentes núcleos de un único procesador físico, compartiendo un espacio común de direcciones de memoria RAM. Ese modelo es cómodo y eficiente, pero tropieza con una barrera física insalvable: el número de núcleos que pueden conectarse a una misma placa base sin que el bus de memoria o el protocolo de coherencia de caché colapsen. Los servidores comerciales más potentes alcanzan típicamente entre 64 y 128 núcleos en dos o cuatro sockets. Resolver problemas científicos o de ingeniería de escala nacional o regional —como la simulación hidrodinámica del Río Pilcomayo o el procesamiento de imágenes satelitales multiespectrales del Gran Chaco— exige cientos o miles de núcleos trabajando en paralelo.

Para superar el límite de una sola máquina física, la computación de alto rendimiento (High Performance Computing, HPC) recurre a la memoria distribuida: interconectar múltiples computadoras independientes (llamadas nodos) a través de una red de comunicaciones de alta velocidad, formando lo que se conoce como un clúster. En este paradigma, ningún procesador puede acceder directamente a la memoria de otro nodo. Si el Proceso A en el Nodo 1 necesita datos calculados por el Proceso B en el Nodo 2, el dato debe empaquetarse explícitamente en un mensaje, enviarse a través de la tarjeta de red (NIC), viajar por el cableado y ser recibido y desempaquetado en la memoria local del receptor.

La herramienta estándar mundial para gobernar este paradigma es MPI (Message Passing Interface). MPI no es un compilador ni un lenguaje nuevo: es una especificación estándar de biblioteca portable, implementada fundamentalmente en C, C++ y Fortran. A diferencia de OpenMP, donde la paralelización se añade mediante directivas pragmáticas a un código secuencial preexistente, programar en MPI exige repensar desde la base la partición del problema, el balance de carga y la coreografía de comunicaciones. En memoria distribuida, la red es entre 1.000 y 10.000 veces más lenta que el acceso a la memoria RAM local; por consiguiente, un diseño ineficiente que comunique con excesiva frecuencia destruirá cualquier ganancia de aceleración.

Competencia que se desarrolla (EC3). Implementar soluciones de cómputo paralelo en memoria distribuida mediante el estándar MPI y herramientas de orquestación, optimizando la comunicación entre procesos y la escalabilidad del sistema.


3.1 Modelo de paso de mensajes y espacios de memoria privados

En el modelo de paso de mensajes (Message Passing Model), la computación se estructura como un conjunto de procesos autónomos, cada uno de los cuales posee su propio espacio de memoria virtual privado.

Modelo de memoria compartida vs distribuida

Figura 3.1. Comparación estructural entre el modelo de memoria compartida (OpenMP) y el modelo de memoria distribuida (MPI). En MPI, cada proceso se ejecuta de forma aislada con su propia memoria RAM física; cualquier intercambio de datos exige transferencias explícitas a través de la red de interconexión.

3.1.1 Arquitectura NORMA frente a UMA/NUMA

Las arquitecturas de memoria se clasifican según el acceso físico de los procesadores a la memoria:

CaracterísticaMemoria Compartida (UMA / NUMA)Memoria Distribuida (NORMA)
AcrónimoUniform / Non-Uniform Memory AccessNo-Remote Memory Access
Espacio de direccionesÚnico y global para todos los hilosMúltiples espacios privados disjuntos
Mecanismo de comunicaciónLectura/escritura directa en variablesEnvío y recepción explícita de mensajes
Coherencia de datosMantenida por hardware (bus snooping / directorio)Responsabilidad exclusiva del programador
Escalabilidad máximaLimitada (10^1 a 10^2 núcleos)Prácticamente ilimitada (10^3 a 10^6 núcleos)
Costo de hardwareElevado (placas multiprocesador especializadas)Moderado (clústeres de computadoras comerciales commodity)

Bajo una arquitectura NORMA, no existen variables globales compartidas. Si un proceso modifica el valor de una variable int total = 5;, dicha modificación ocurre únicamente en su memoria local. Los demás procesos continuarán viendo sus propios valores locales inalterados. Esta propiedad elimina de raíz las condiciones de carrera sobre memoria compartida que se estudiaron en el módulo 2. Sin embargo, introduce un desafío nuevo: el coste temporal de la red y la sincronización explícita entre procesos.

3.1.2 Anatomía de un mensaje

Un mensaje en MPI no es simplemente una secuencia de bytes planos en memoria; es un sobre formal que contiene metadatos indispensables para que el sistema operativo y el hardware de red puedan entregarlo con precisión. Todo mensaje consta de:

  1. Datos (Payload):
  2. Dirección inicial del buffer en memoria (void *buf).
  3. Número de elementos a transmitir (count).
  4. Tipo de dato elemental (datatype), por ejemplo MPI_INT o MPI_DOUBLE.
  5. Envoltorio (Envelope):
  6. Identificador del proceso origen (source).
  7. Identificador del proceso destino (dest).
  8. Etiqueta de mensaje (tag): entero que permite distinguir el propósito del mensaje (por ejemplo, tag 1 para datos climáticos, tag 2 para señales de sincronización).
  9. Contexto o Comunicador (communicator): ámbito cerrado de procesos dentro del cual el mensaje es válido.

3.2 Estándar MPI: inicialización, procesos y comunicadores

MPI nació en 1994 como un esfuerzo conjunto entre la academia y la industria (el MPI Forum) para unificar las múltiples bibliotecas propietarias incompatibles de paso de mensajes de la época (como PVM, P4 y NX). MPI es un estándar abierto, no un producto comercial cerrado. Sus especificaciones mayores son:

  • MPI-1.0 (1994): Comunicaciones punto a punto y colectivas básicas en topologías estáticas.
  • MPI-2.0 (1997): Creación dinámica de procesos, comunicación de un solo lado (Remote Memory Access, RMA) y E/S paralela paralela (MPI-IO).
  • MPI-3.0 (2012): Colectivas no bloqueantes, soporte mejorado de Fortran 2008 y colectivas sobre vecindarios.
  • MPI-4.0 (2021): Soporte para entornos resilientes a fallos, interfaces particionadas y operaciones para computación a exaescala.

En el laboratorio se emplean implementaciones maduras de software libre: OpenMPI (ampliamente extendido en clústeres Linux) y MPICH.

Comunicadores y Grupos en MPI

Figura 3.2. Estructura de comunicadores en MPI. Al arrancar, todos los procesos pertenecen al comunicador predeterminado MPI_COMM_WORLD. Mediante MPI_Comm_split, los procesos pueden subdividirse en comunicadores especializados.

3.2.1 El ciclo de vida de un programa MPI

Todo programa MPI sigue una estructura canónica obligatoria en C:

#include <mpi.h>
#include <stdio.h>

int main(int argc, char *argv[]) {
    int rank, size;

    // 1. Inicialización del entorno paralelo
    MPI_Init(&argc, &argv);

    // 2. Consulta del tamaño del comunicador (cuántos procesos hay en total)
    MPI_Comm_size(MPI_COMM_WORLD, &size);

    // 3. Consulta del rango del proceso (quién soy yo)
    MPI_Comm_rank(MPI_COMM_WORLD, &rank);

    printf("Hola desde el proceso %d de un total de %d procesos\n", rank, size);

    // 4. Finalización ordenada del entorno paralelo
    MPI_Finalize();
    return 0;
}

Reglas estrictas de ejecución:

  • MPI_Init: Debe llamarse antes de cualquier otra función de la biblioteca MPI. Configura los búferes internos, establece los enlaces de red y asigna las estructuras de control.
  • MPI_Finalize: Limpia los recursos asignados, libera los canales de red y sincroniza la terminación. Ninguna función MPI puede invocarse después de MPI_Finalize.
  • Rango (Rank): Es un identificador entero continuo que va desde 0 hasta \text{size} - 1. El rango es único dentro de cada comunicador. Por convención, el proceso con rank == 0 suele actuar como proceso coordinador o maestro (Master).

3.2.2 Compilación y ejecución

Un programa MPI no se compila con gcc directamente, sino utilizando el envoltorio (wrapper) proporcionado por la instalación:

# Compilación
mpicc -O3 -Wall mpi_hola.c -o mpi_hola

# Ejecución de 4 procesos en la máquina local
mpirun -np 4 ./mpi_hola

El lanzador mpirun (o mpiexec) se encarga de crear las 4 instancias del binario, propagar las variables de entorno, conectar los canales de comunicación por sockets TCP o memoria compartida interproceso (IPC) y orquestar el inicio sincronizado.

3.2.3 Subcomunicadores con MPI_Comm_split

Cuando un algoritmo requiere dividir tareas (por ejemplo, un grupo de procesos dedicado al cálculo matricial y otro grupo a escribir los resultados en disco), no es eficiente filtrar mediante constantes condicionales if (rank < 4). Se divide formalmente el comunicador:

MPI_Comm comm_grupo;
int color = (rank < 3) ? 0 : 1; // Dos grupos: grupo 0 y grupo 1
MPI_Comm_split(MPI_COMM_WORLD, color, rank, &comm_grupo);

int nuevo_rank, nuevo_size;
MPI_Comm_rank(comm_grupo, &nuevo_rank);
MPI_Comm_size(comm_grupo, &nuevo_size);

Cada nuevo subcomunicador tiene su propia numeración de rangos empezando desde 0.


3.3 Comunicaciones punto a punto: bloqueantes y no bloqueantes

Las comunicaciones punto a punto representan el intercambio elemental de información entre dos procesos específicos: un emisor y un receptor.

3.3.1 Primitivas bloqueantes (MPI_Send y MPI_Recv)

La firma estándar en C es:

int MPI_Send(const void *buf, int count, MPI_Datatype datatype, int dest,
             int tag, MPI_Comm comm);

int MPI_Recv(void *buf, int count, MPI_Datatype datatype, int source,
             int tag, MPI_Comm comm, MPI_Status *status);

El término bloqueante (blocking) significa que la función no retorna el control al programa hasta que el búfer de memoria asociado pueda ser reutilizado con seguridad:

  • En MPI_Send: Retorna cuando los datos han salido del búfer del usuario (bien porque ya se transmitieron por la red, o bien porque el subsistema de MPI los copió a un búfer temporal interno del sistema).
  • En MPI_Recv: Retorna únicamente cuando el mensaje ha llegado por completo y ha sido depositado en el búfer de destino en memoria.

Protocolos internos: Eager vs. Rendezvous

Las implementaciones de MPI eligen automáticamente entre dos protocolos según el tamaño del mensaje:

  1. **Protocolo Eager (mensajes pequeños, típicamente < 64\text{ KB}):** El emisor envía los datos inmediatamente a la red asumiendo que el receptor tendrá espacio en sus búferes internos para alojarlo. Si el receptor aún no ejecutó MPI_Recv, el mensaje queda retenido en la memoria del sistema. Es rápido pero consume memoria de sistema.
  2. Protocolo Rendezvous (mensajes grandes): El emisor envía primero un mensaje corto de solicitud de envío (Handshake). El emisor se detiene y no transfiere los datos masivos hasta que el receptor confirme con un mensaje de respuesta que ya reservó el búfer final. Evita desbordamiento de memoria, pero impone una penalización de dos viajes de red (Round Trip Time).

El peligro de interbloqueo (Deadlock)

Un error común en principiantes ocurre cuando dos procesos intentan enviarse mensajes mutuamente de forma cruzada:

// PELIGRO DE INTERBLOQUEO (DEADLOCK)
if (rank == 0) {
    MPI_Send(datos_a, N, MPI_DOUBLE, 1, 0, MPI_COMM_WORLD);
    MPI_Recv(datos_b, N, MPI_DOUBLE, 1, 0, MPI_COMM_WORLD, MPI_STATUS_IGNORE);
} else if (rank == 1) {
    MPI_Send(datos_b, N, MPI_DOUBLE, 0, 0, MPI_COMM_WORLD);
    MPI_Recv(datos_a, N, MPI_DOUBLE, 0, 0, MPI_COMM_WORLD, MPI_STATUS_IGNORE);
}

Si el tamaño N supera el umbral del protocolo Eager, ambos procesos se bloquean en MPI_Send esperando que el otro emita un MPI_Recv. Como ninguno puede avanzar, la aplicación queda congelada indefinidamente. Para evitarlo se utiliza alternancia de órdenes, MPI_Sendrecv (que realiza el envío y recepción atómicamente), o comunicaciones no bloqueantes.

Punto a punto bloqueante vs no bloqueante

Figura 3.3. Comparativa temporal. En A (bloqueante), los procesadores pierden ciclos valiosos esperando el handshake de red. En B (no bloqueante), el cálculo local se ejecuta simultáneamente mientras la tarjeta de red transmite los datos mediante DMA sin intervención de la CPU.

3.3.2 Comunicaciones no bloqueantes (MPI_Isend e MPI_Irecv)

Las funciones que inician transferencias no bloqueantes llevan el prefijo I (Immediate). Inician la operación y retornan inmediatamente, devolviendo un descriptor de petición (request):

MPI_Request req_envio, req_recepcion;

// Inicia el envío sin esperar
MPI_Isend(buffer_salida, N, MPI_DOUBLE, destino, tag, MPI_COMM_WORLD, &req_envio);

// Inicia la recepción sin esperar
MPI_Irecv(buffer_entrada, N, MPI_DOUBLE, origen, tag, MPI_COMM_WORLD, &req_recepcion);

// COMPUTAR AQUÍ: Código de cálculo independiente de buffer_salida y buffer_entrada
calcular_zonas_interiores();

// Esperar a que la transferencia de red culmine antes de usar los datos
MPI_Wait(&req_envio, MPI_STATUS_IGNORE);
MPI_Wait(&req_recepcion, MPI_STATUS_IGNORE);

Esta técnica es la piedra angular del solapamiento cómputo-comunicación (computation-communication overlap): mientras los datos viajan a través de los cables de red física mediante controladores DMA (Direct Memory Access), la CPU ejecuta operaciones aritméticas a plena velocidad. La latencia de la red queda virtualmente oculta.


3.4 Comunicaciones colectivas: broadcast, scatter, gather y reduce

Una comunicación colectiva involucra a todos los procesos pertenecientes a un comunicador determinado. Si un solo proceso del comunicador no invoca la rutina colectiva, la aplicación sufrirá un interbloqueo irrecuperable.

Colectivas principales en MPI

Figura 3.4. Las cuatro primitivas colectivas fundamentales en MPI: Broadcast (uno a todos), Scatter (distribución disjunta), Gather (recolección) y Reduce (agregación con operador asociativo).

3.4.1 Difusión (Broadcast): MPI_Bcast

Un proceso distinguido (denominado raíz o root) transmite una copia idéntica del mismo dato a todos los procesos del grupo:

int MPI_Bcast(void *buffer, int count, MPI_Datatype datatype, int root, MPI_Comm comm);

En lugar de que el proceso raíz envíe p - 1 mensajes secuenciales individuales (lo que tendría una complejidad temporal lineal \mathcal{O}(p)), las implementaciones modernas de MPI construyen un árbol binomial de difusión. El proceso 0 envía al 1; luego 0 y 1 envían en paralelo al 2 y al 3; y así sucesivamente. La complejidad de tiempo de un broadcast en árbol binomial es:

$$T_{bcast} \approx \lceil \log_2 p \rceil \cdot (\alpha + \beta \cdot m)$$

donde p es el número de procesos, \alpha la latencia de red, \beta el tiempo por byte y m el tamaño del mensaje.

3.4.2 Distribución (Scatter) y Recolección (Gather)

Permiten particionar y consolidar arreglos de datos disjuntos:

  • MPI_Scatter: El proceso raíz toma un arreglo de tamaño N = p \times k y reparte un bloque contiguo de tamaño k a cada uno de los p procesos: ``c MPI_Scatter(sendbuf, k, MPI_DOUBLE, recvbuf, k, MPI_DOUBLE, root, MPI_COMM_WORLD); ``
  • MPI_Gather: Operación inversa. Cada proceso aporta un bloque local de tamaño k, y el proceso raíz los concatena en orden de rango en su búfer receptor: ``c MPI_Gather(sendbuf, k, MPI_DOUBLE, recvbuf, k, MPI_DOUBLE, root, MPI_COMM_WORLD); ``

3.4.3 Reducción global: MPI_Reduce y MPI_Allreduce

Toman valores calculados localmente por cada proceso y los combinan utilizando una operación matemática asociativa y conmutativa:

Operador MPIOperación matemática equivalente
MPI_SUMSuma aritmética (\sum x_i)
MPI_PRODProducto (\prod x_i)
MPI_MAX / MPI_MINValor máximo / mínimo global
MPI_LAND / MPI_LOROperaciones lógicas AND / OR
MPI_MAXLOC / MPI_MINLOCValor extremo y el rango del proceso que lo posee

Sintaxis en C:

double local_pi_sum, global_pi_sum;
// Reduce: El resultado final solo queda en el proceso root (0)
MPI_Reduce(&local_pi_sum, &global_pi_sum, 1, MPI_DOUBLE, MPI_SUM, 0, MPI_COMM_WORLD);

// Allreduce: El resultado final queda disponible en TODOS los procesos
MPI_Allreduce(&local_pi_sum, &global_pi_sum, 1, MPI_DOUBLE, MPI_SUM, MPI_COMM_WORLD);

El costo de una reducción global sobre un árbol de procesamiento es:

$$T_{reduce} \approx \lceil \log_2 p \rceil \cdot (\alpha + \beta \cdot m + \gamma \cdot m)$$

donde \gamma representa el tiempo de la CPU en ejecutar la operación aritmética por elemento.


3.5 Patrones de comunicación en sistemas distribuidos

En la mayoría de los problemas científicos de simulación (geofísica, dinámica de fluidos, modelos epidemiológicos o climáticos), los procesos no se comunican con todos los demás de manera aleatoria. Se comunican casi exclusivamente con sus vecinos espaciales inmediatos.

Patrón de intercambio de halos en 2D

Figura 3.5. Patrón de intercambio de celdas fantasma (halo exchange) en una malla bidimensional. Cada proceso calcula sus celdas interiores mientras transmite y recibe los bordes Norte, Sur, Este y Oeste de sus procesos vecinos.

3.5.1 Descomposición de dominio y celdas fantasma (Ghost Cells)

Consideremos la simulación hidrológica en 2D del Río Pilcomayo. El terreno se modela como una grilla de N_x \times N_y celdas. Para calcular el estado del agua en la posición (i, j) en el instante de tiempo t + 1, el método de diferencias finitas requiere los valores en (i, j), (i \pm 1, j) e (i, j \pm 1) en el instante t:

$$U_{i, j}^{t+1} = \Phi(U_{i, j}^t, U_{i+1, j}^t, U_{i-1, j}^t, U_{i, j+1}^t, U_{i, j-1}^t)$$

Al particionar la grilla entre p procesos, las celdas situadas en el borde de un bloque local necesitan los valores de la fila o columna del bloque asignado al proceso vecino. Para resolver esto sin acceder a memoria remota, cada proceso reserva un anillo perimetral de memoria adicional denominado celdas fantasma (ghost cells o halos).

En cada paso de tiempo, el algoritmo sigue el siguiente patrón estricto:

  1. Iniciar el envío de los bordes locales a los vecinos con MPI_Isend.
  2. Iniciar la recepción de los halos de los vecinos con MPI_Irecv.
  3. Calcular los puntos interiores: Las celdas que no dependen de los bordes se calculan inmediatamente.
  4. Ejecutar MPI_Waitall para asegurar que los halos vecinos llegaron.
  5. Calcular los puntos de frontera: Actualizar las celdas del perímetro local usando los datos recién recibidos.

3.5.2 Topologías virtuales cartesianas

MPI provee soporte nativo para organizar los procesos lógicamente en una matriz 2D o 3D:

int dims[2] = {2, 2};     // Malla de 2x2 procesos (4 procesos en total)
int periods[2] = {0, 0};  // Condiciones de frontera no periódicas
MPI_Comm cart_comm;

MPI_Cart_create(MPI_COMM_WORLD, 2, dims, periods, 1, &cart_comm);

int rank_norte, rank_sur, rank_oeste, rank_este;
// Encuentra automáticamente los vecinos en el eje Y (desplazamiento 1)
MPI_Cart_shift(cart_comm, 0, 1, &rank_norte, &rank_sur);
// Encuentra los vecinos en el eje X (desplazamiento 1)
MPI_Cart_shift(cart_comm, 1, 1, &rank_oeste, &rank_este);

Si un proceso está en el límite del dominio físico, MPI_Cart_shift asigna el valor especial MPI_PROC_NULL a ese vecino. Las funciones de envío y recepción ignoran automáticamente MPI_PROC_NULL, simplificando enormemente el código al evitar ramas condicionales complejas.


3.6 Arquitectura de clústeres modernos

Un clúster HPC es un sistema de cómputo distribuido integrado por múltiples computadoras independientes (nodos) que cooperan a través de una red dedicada de interconexión para actuar como una sola entidad computacional de alta potencia.

Arquitectura de clúster en laboratorio

Figura 3.6. Topología física y lógica del clúster de laboratorio sobre estaciones Dell OptiPlex 7010 (32 GB RAM). El nodo maestro orquesta los servicios de control y almacenamiento compartido, mientras los nodos de cómputo ejecutan los procesos paralelos de cómputo.

3.6.1 Componentes de un clúster HPC

  1. Nodo Maestro o Cabecera (Head Node / Login Node):
  2. Es el punto de entrada para los usuarios y administradores.
  3. Aloja el planificador de trabajos (scheduler, ej. Slurmctld).
  4. Servidor de almacenamiento compartido (NFS / Lustre).
  5. No debe utilizarse para ejecutar cómputos pesados directamente, a fin de no saturar la gestión.
  6. Nodos de Cómputo (Compute / Worker Nodes):
  7. Máquinas dedicadas exclusivamente a procesar trabajos en segundo plano sin interfaces gráficas ni servicios innecesarios.
  8. Ejecutan el demonio de control del gestor de recursos (slurmd).
  9. Redes de Interconexión:
  10. Red de Administración / Gestión: Típicamente Ethernet 1 Gbps para acceso SSH, monitoreo IPMI/BMC y comandos del sistema operativo.
  11. Red de Datos / HPC: Red de ultra-baja latencia y alto ancho de banda dedicada exclusivamente al tráfico MPI y acceso a almacenamiento.
  12. Almacenamiento Compartido Distribuido:
  13. Un directorio común (como /shared/hpc o /home/usuario) que se encuentra físicamente en un nodo o servidor SAN/NAS y se monta mediante NFS (Network File System) en todos los nodos con la misma ruta absoluta.
  14. Esto garantiza que el binario compilado y los archivos de entrada estén disponibles de forma idéntica en cualquier nodo del clúster sin necesidad de copiarlos manualmente a cada máquina.

3.6.2 Comparativa de tecnologías de interconexión

TecnologíaAncho de Banda TípicoLatencia de Mensaje Corto (\alpha)Mecanismo de TransferenciaCosto Relativo
Gigabit Ethernet (1 GbE)1\text{ Gbps} (125\text{ MB/s})40 - 80\ \mu\text{s}Pila TCP/IP del KernelMuy bajo (estándar comercial)
10 Gigabit Ethernet (10 GbE)10\text{ Gbps} (1.25\text{ GB/s})10 - 20\ \mu\text{s}TCP/IP o RoCEv2Bajo / Medio
InfiniBand HDR200\text{ Gbps} (25\text{ GB/s})0.6 - 1.2\ \mu\text{s}RDMA nativo por hardwareAlto (HPC profesional)
InfiniBand NDR400\text{ Gbps} (50\text{ GB/s})0.4 - 0.8\ \mu\text{s}RDMA nativo por hardwareMuy alto (Centros Nacionales)

Nota técnica para el laboratorio en Yacuiba: Las estaciones Dell OptiPlex 7010 cuentan con interfaces Gigabit Ethernet integradas. La latencia observada en el laboratorio ronda los $50\ \mu\text{s}$, en contraste con los $1\ \mu\text{s}$ típicos de InfiniBand. Esto impone una lección de ingeniería fundamental para el estudiante: para que un programa sea escalable en la infraestructura local, la granularidad del cómputo debe ser suficientemente gruesa para que el tiempo de cálculo opaque la latencia del enlace Ethernet.


3.7 Gestores de recursos: Slurm y orquestación

En un entorno multi-usuario o de clúster de producción, los programadores no ejecutan directamente mpirun en las terminales de los nodos. Si dos usuarios lanzan procesos MPI pesados sobre la misma máquina al mismo tiempo, los núcleos físicos se sobrecargan, los procesos compiten por la memoria RAM hasta provocar un Out-Of-Memory (OOM Killer) y las mediciones de rendimiento carecen de validez científica.

La solución universal en supercómputo es un Gestor de Cargas de Trabajo y Recursos (Workload Manager), y el estándar de la industria es Slurm (Simple Linux Utility for Resource Management).

Arquitectura de Slurm

Figura 3.7. Arquitectura y ciclo de vida de un trabajo en Slurm. El usuario remite un script de lotes con sbatch. El demonio central slurmctld gestiona la cola y asigna los nodos libres, donde slurmd ejecuta los procesos paralelos a través de srun.

3.7.1 Componentes y demonios de Slurm

  • slurmctld (Slurm Controller Daemon): Reside en el nodo maestro. Mantiene la lista de nodos disponibles, vigila su estado de salud mediante latidos periódicos (heartbeats), gestiona las particiones (colas) y decide qué trabajo pasa a ejecutarse según prioridades y recursos libres.
  • slurmd (Slurm Daemon): Se ejecuta como servicio del sistema en cada uno de los nodos de cómputo. Recibe las órdenes del controlador, reserva los núcleos de CPU y memoria requeridos mediante cgroups de Linux, lanza las tareas y supervisa su finalización.
  • slurmdbd (Slurm Database Daemon): Registra el uso histórico de recursos y contabilidad en una base de datos MySQL/MariaDB.

3.7.2 Comandos esenciales de usuario

  1. sinfo: Inspecciona el estado de los nodos y particiones del clúster:
  2. idle: Nodo libre listo para aceptar trabajos.
  3. alloc: Nodo ocupado procesando tareas asignadas.
  4. down / drain: Nodo fuera de servicio por mantenimiento o error de red.
  5. squeue: Consulta los trabajos en cola de espera (PD, pendiente) o en ejecución (R, running).
  6. scancel <JOBID>: Cancela y retira un trabajo propio del sistema.
  7. srun: Lanza una tarea paralela de forma interactiva en tiempo real.
  8. sbatch: Envía un script de procesamiento por lotes (batch script) para ejecución desatendida.

3.7.3 Elaboración de un script sbatch para MPI

Un script de Slurm es un script de shell Bash ordinario que comienza con directivas #SBATCH interpretadas por el planificador antes de lanzar el trabajo:

#!/bin/bash
#SBATCH --job-name=pilcomayo_mpi      # Nombre del trabajo
#SBATCH --output=salida_%j.log        # Archivo de salida estándar (%j = ID del trabajo)
#SBATCH --error=error_%j.log          # Archivo de registro de errores
#SBATCH --partition=normal            # Partición o cola solicitada
#SBATCH --nodes=3                     # Solicitar exactamente 3 nodos físicos/virtuales
#SBATCH --ntasks-per-node=4           # 4 tareas (procesos MPI) por cada nodo (Total: 12 tareas)
#SBATCH --cpus-per-task=1             # 1 núcleo físico de CPU por cada tarea MPI
#SBATCH --mem=4G                      # Límite estricto de memoria RAM por nodo
#SBATCH --time=00:15:00               # Tiempo máximo de ejecución: 15 minutos

echo "=== Inicio del Trabajo Slurm $SLURM_JOB_ID ==="
echo "Ejecutándose en los nodos: $SLURM_JOB_NODELIST"
echo "Total de tareas MPI asignadas: $SLURM_NTASKS"

# Cargar módulos de software o variables de entorno
module load openmpi/4.1.5 2>/dev/null || true

# Ejecución paralela controlada por Slurm
srun ./montecarlo_mpi 100000000

echo "=== Trabajo finalizado exitosamente ==="

El comando sbatch job_mpi.sh retorna de inmediato indicando el ID asignado (por ejemplo, Submitted batch job 1042). El usuario puede desconectarse de la terminal; Slurm asignará los nodos cuando estén disponibles, ejecutará el código y depositará los resultados en salida_1042.log.


3.8 Optimización de rendimiento en sistemas distribuidos

En sistemas de memoria distribuida, el rendimiento no depende únicamente del número de instrucciones por segundo (MFLOPS/GFLOPS) que procesan las CPUs. En la práctica, el factor determinante es la relación cómputo-comunicación.

Modelo LogP de comunicación

*Figura 3.8. Parámetros del modelo formal LogP. La latencia de red (L), la sobrecarga del procesador (o) y la brecha temporal entre inyecciones consecutivas (g) determinan el coste de transferir mensajes en clúster.*

3.8.1 El modelo analítico de comunicación

El tiempo necesario para transmitir un mensaje continuo de m bytes a través de un canal de red se modela matemáticamente como una función afín:

$$T_{comm}(m) = \alpha + \beta \cdot m$$

donde:

  • \alpha (Latencia de inicio o startup time): Tiempo fijo en segundos consumido por el software y el hardware para preparar el mensaje, empaquetarlo, negociar el protocolo y poner el primer bit en el cable físico. Es independiente del tamaño del mensaje.
  • \beta (Tiempo de transferencia por byte): Inverso del ancho de banda efectivo de la red (\beta = \frac{1}{B}). Si el ancho de banda es B = 125\text{ MB/s} (Gigabit Ethernet), \beta = 8 \times 10^{-9}\text{ s/byte}.
  • m: Tamaño del mensaje en bytes.

Consecuencia algorítmica: Empaquetamiento de datos

Transmitir 1.000 números flotantes de 8 bytes en 1.000 mensajes individuales de 8 bytes cuesta:

$$T_{1000\_mensajes} = 1000 \cdot (\alpha + 8\beta) = 1000\alpha + 8000\beta$$

Transmitir los mismos 1.000 números agrupados en un solo arreglo contiguo de 8.000 bytes cuesta:

$$T_{1\_mensaje} = \alpha + 8000\beta$$

Dado que en Gigabit Ethernet \alpha \approx 50\ \mu\text{s}, enviar mil mensajes sueltos consume 1000 \times 50\ \mu\text{s} = 50.000\ \mu\text{s} = 50\text{ ms} solo en latencia de arranque. Enviar el paquete consolidado consume apenas 50\ \mu\text{s} de latencia: ¡una diferencia de mil veces en sobrecarga de red! En MPI siempre deben agruparse los datos antes de transmitirlos.

3.8.2 Granularidad y relación cómputo-comunicación

La granularidad de un algoritmo paralelo se define formalmente como la relación entre el tiempo dedicado al cálculo aritmético útil y el tiempo invertido en comunicación:

$$q = \frac{T_{computo}}{T_{comunicacion}}$$

  • **Granularidad fina (q \ll 1):** El proceso pasa más tiempo enviando y esperando mensajes que calculando. La aplicación no escalará; añadir más nodos reducirá el rendimiento global en lugar de mejorarlo.
  • **Granularidad gruesa (q \gg 1):** Los procesos computan durante períodos prolongados entre ráfagas breves de comunicación. Es la condición necesaria para obtener alta eficiencia en clústeres basados en redes estándar Ethernet.

3.8.3 Escalabilidad Fuerte frente a Escalabilidad Débil

Escalabilidad Fuerte (Strong Scaling — Ley de Amdahl)

El tamaño total del problema se mantiene fijo (N = \text{cte}) mientras se incrementa el número de procesadores p:

$$S(p) = \frac{1}{(1 - f) + \frac{f}{p} + \frac{T_{comm}(p)}{T_1}}$$

A medida que p crece, la fracción de cómputo local por proceso disminuye (\frac{N}{p} \to 0), pero el número de comunicaciones y la latencia aumentan. Inevitablemente, a partir de cierto umbral crítico de nodos, el tiempo de comunicación domina por completo y el Speedup decae bruscamente.

Escalabilidad Débil (Weak Scaling — Ley de Gustafson)

El tamaño del problema crece proporcionalmente con el número de procesadores (N \propto p), de modo que la cantidad de trabajo asignada a cada proceso individual permanece constante:

$$S_{Gustafson}(p) = p - \alpha_{serie} \cdot (p - 1)$$

Bajo escalabilidad débil, si un solo nodo procesa un millón de celdas de la cuenca del Pilcomayo en 10 segundos, 16 nodos procesarán 16 millones de celdas en aproximadamente los mismos 10 segundos (más una pequeña penalización de comunicación de fronteras). Es el modelo preferido en HPC para justificar inversiones en clústeres de gran envergadura.

Actividades de Aprendizaje Autónomo — Módulo 3 (EC3)

Asignatura: TEL-420 · Sistemas Paralelos Módulo 3: Programación Paralela en Memoria Distribuida (MPI y Cústeres) Docente: Ing. Elias Cassal Baldiviezo Horas de dedicación autónoma: 8 horas Semanas de ejecución: 7 a 10


1. Justificación y propósito pedagógico

Conforme al apartado 20.3 del Programa Docente del Proyecto Formativo, el aprendizaje autónomo en el módulo 3 afianza las competencias de diseño y orquestación de sistemas distribuidos, desafiando al estudiante a coordinar procesos aislados con memoria privada y a evaluar científicamente el impacto de la red de interconexión física sobre la escalabilidad.


2. Bloque de Actividades Autónomas Obligatorias

Actividad A1: Desarrollo progresivo de aplicaciones MPI (Punto a Punto y Colectivas)

  • Momento de entrega: Fin de la Semana 9.
  • Instrumento de evaluación: file:///home/eliasdev/sistemas_paralelos/10-modulos/M3-mpi-custeres/rubrica-codigo-mpi.md (40% de EC3).
  • Consigna de trabajo:
  • Completar la implementación de las cuatro aplicaciones MPI solicitadas:
  • Anillo de procesos bloqueante libre de interbloqueos cíclicos.
  • Anillo no bloqueante con solapamiento comprobado de cómputo aritmético mediante MPI_Isend / MPI_Irecv.
  • Aproximación de Monte Carlo para \pi con MPI_Bcast y MPI_Reduce utilizando generadores congruenciales lineales (LCG) con semillas disjuntas por rango.
  • Multiplicación matricial distribuida por bloques de filas con MPI_Scatter y MPI_Gather.
  • Incluir pruebas automatizadas de verificación que aseguren que para p \in \{2, 4, 8\} el resultado numérico coincide con la versión secuencial con un error menor a 10^{-9}.

Actividad A2: Diseño, configuración y documentación de clúster virtualizado con Slurm

  • Momento de entrega: Fin de la Semana 9.
  • Instrumento de evaluación: file:///home/eliasdev/sistemas_paralelos/10-modulos/M3-mpi-custeres/rubrica-informe-cluster.md (35% de EC3).
  • Consigna de trabajo:
  • Construir la receta de Docker Compose o máquinas virtuales KVM/Multipass que despliegue el clúster multi-nodo en las estaciones Dell OptiPlex (Core i7, 32 GB RAM).
  • Configurar la red virtual, acceso SSH sin contraseñas entre nodos, exportación del sistema de archivos NFS /shared/hpc y los demonios slurmctld y slurmd.
  • Redactar el informe técnico en PDF documentando la topología de red, el direccionamiento IP, el archivo slurm.conf y las pruebas de validación con sinfo y squeue.

Actividad A3: Experimentación de escalabilidad distribuida y validación del modelo LogP

  • Momento de entrega: Fin de la Semana 10.
  • Instrumento de evaluación: file:///home/eliasdev/sistemas_paralelos/10-modulos/M3-mpi-custeres/rubrica-escalabilidad.md (25% de EC3).
  • Consigna de trabajo:
  • Automatizar la ejecución de pruebas con un script de benchmarking sobre Slurm para p \in \{1, 2, 4, 8, 12\} y 3 tamaños de problema.
  • Calcular medias aritméticas y desviaciones estándar descartando la primera corrida de calentamiento.
  • Estimar los parámetros del modelo de comunicación T_{comm} = \alpha + \beta \cdot m en la red Ethernet del laboratorio.
  • Graficar y contrastar las curvas empíricas de Speedup frente a la Ley de Amdahl (escalabilidad fuerte) y Ley de Gustafson (escalabilidad débil).
M3

Diagnóstico

Examen HTML autocontenido interactivo.

M3

Sumativo

Examen HTML autocontenido interactivo.