Introducción a las Variables Copmartidas
En el entorno de procesamiento distribuido de Spark, el manejo de estado global reuqiere mecanismos específicos para evitar inconsistencias. Para facilitar la estadística y administración de información común, el framework proporciona dos tipos de variables compartidas: las variables Broadcast (solo lectura) y los Accumulators (acumuladores). Mientras que las primeras distribuyen datos a los nodos, los acumuladores permiten agregar información desde los ejecutores hacia el controlador principal.
El ciclo de vida de un Accumulator está centralizado en el Driver. Aunque cada Task mantiene una copia local en su Executor para realizar operaciones de suma parciales, el valor global solo se consolida en el Driver una vez finalizadas las tareas. Es crucial entender dos características fundamentales:
- Naturaleza Aditiva: Solo permiten operaciones de suma o fusión; no son variables de estado general modificables arbitrariamente.
- Evaluación Lazy: Al igual que las transformaciones de RDD, los acumuladores no se actualizan hasta que se ejecuta una acción (Action). Consultar su valor antes de una acción puede devolver el estado inicial.
Tipologías de Accumulators en Spark
En versiones recientes (Spark 2.x en adelante), se distinguen principalmente dos categorías:
1. Tipos Nativos
Spark incluye implementaciones listas para usar para casos comunes:
- LongAccumulator: Para conteos enteros.
- DoubleAccumulator: Para sumas de valores flotantes.
- CollectionAccumulator: Para agregar elementos a una colección distribuida.
La instanciación se realiza directamente desde el contexto de Spark, asignando un nombre visible en la interfaz web:
LongAccumulator contadorErrores = sc.longAccumulator("contador_errores");
2. Accumulators Personalizados
Para lógicas de agregación complejas, es necesario implementar la clase abstracta AccumulatorV2. Esto requiere definir cómo se inicializa, copia, resetea, agrega un valor, fusiona con otro acumulador y cómo se devuelve el valor final.
A continuación, se presenta una implementación refactorizada que utiliza un mapa interno para contar categorías dinámicas, en lugar de manipular cadenas de texto concatenadas:
package com.example.spark.core;
import org.apache.spark.util.AccumulatorV2;
import java.util.HashMap;
import java.util.Map;
import java.util.Objects;
public class MetricAggregatorAccumulator extends AccumulatorV2<String, Map<String, Long>> {
private Map<String, Long> metricsMap;
private final Map<String, Long> initialState = new HashMap<>();
public MetricAggregatorAccumulator() {
reset();
}
@Override
public boolean isZero() {
return metricsMap.equals(initialState);
}
@Override
public AccumulatorV2<String, Map<String, Long>> copy() {
MetricAggregatorAccumulator copy = new MetricAggregatorAccumulator();
copy.metricsMap = new HashMap<>(this.metricsMap);
return copy;
}
@Override
public void reset() {
this.metricsMap = new HashMap<>();
}
@Override
public void add(String category) {
if (category == null) return;
metricsMap.put(category, metricsMap.getOrDefault(category, 0L) + 1);
}
@Override
public void merge(AccumulatorV2<String, Map<String, Long>> other) {
if (!(other instanceof MetricAggregatorAccumulator)) {
throw new UnsupportedOperationException("Solo se pueden fusionar instancias del mismo tipo");
}
MetricAggregatorAccumulator otherAcc = (MetricAggregatorAccumulator) other;
for (Map.Entry<String, Long> entry : otherAcc.metricsMap.entrySet()) {
metricsMap.put(entry.getKey(),
metricsMap.getOrDefault(entry.getKey(), 0L) + entry.getValue());
}
}
@Override
public Map<String, Long> value() {
return new HashMap<>(metricsMap);
}
}
Mecánica de Ejecución y Ciclo de Vida
El funcionamiento interno de los acumuladores sigue un flujo estricto entre el Driver y los Executors:
Registro en el Driver
La definición e inicialización ocurren exclusivamente en el programa principal. El acumulador debe registrarse en el SparkContext para que el framwork pueda rastrearlo. Este registro permite que la variable sea serializada y enviada junto con las tareas a los nodos del clúster. El Driver actúa como el punto de convergencia donde los resultados parciales se combinan conforme las tareas se completan.
Procesamiento en el Executor
Cuando un Executor recibe una tarea, deserializa tanto el RDD como las funciones asociadas, incluyendo el objeto acumulador. Durante la ejecución del Task, las llamadas al método add modifican la copia local del acumulador. Al finalizar la tarea, el Executor reporta los cambios de estado al Driver, quien ejecuta la lógica de merge para actualizar el valor global manteniendo la consistencia de los datos agregados.