Introducción Práctica a Angel Parameter Server (PS) para Machine Learning a Gran Escala

Angel es un sistema de Parameter Server (PS) de alto rendimiento diseñado para manejar modelos de aprendizaje automático masivos que superan la capacidad de memoria de una sola máquina. Su arquitectura permite desacoplar el almacenamiento de parámetros del cálculo, optimizando la comunicación y la escalabilidad en entornos distribuidos.

Arquitectura de Angel PS

El sistema se fundamenta en cuatro componentes esenciales que interactúan para gestionar el entrenamiento distribuido:

  • PServer (Parameter Server): Gestiona el almacenamiento de los parámetros del modelo. Soporta particionamiento automático y actualizaciones asíncronas.
  • Worker: Nodo encargado de la computación local. Solicita parámetros (pull), calcula gradientes y envía actualizaciones (push).
  • Master: Coordina los recursos, supervisa el estado de las tareas y gestiona la tolerancia a fallos.
  • Client: Punto de entrada para la configuración del trabajo y el control del ciclo de vida del entrenamiento.

Uso de PSModel para la Gestión de Parámetros

La abstracción PSModel es la interfaz principal para interactuar con los servidores de parámetros. A continuación se detallan las operaciones fundamentales.

1. Definición e Inicialización

// Definición de un modelo de pesos densos
val dimensionFeatures = conf.getInt(MLConf.ML_FEATURE_INDEX_RANGE, 20000)
val modelWeights = PSModel("lr_model_weights", 1, dimensionFeatures)
 .setRowType(RowType.T_DOUBLE_DENSE)
 .setAverage(true) // Promedio de gradientes activado
 .setHogwild(true) // Optimización para concurrencia

2. Recuperación de Datos (Pull)

// Obtención síncrona de una fila de parámetros
val currentWeights = modelWeights.getRow(0)

// Recuperación de múltiples filas simultáneamente
val indicesFilas = Array(0, 1, 5)
val multipleVectors = modelWeights.getRows(indicesFilas)

// Pull selectivo mediante índices (útil para modelos dispersos)
val featureIndices = Array(10, 50, 100)
val sparseData = modelWeights.getRowWithIndex(0, featureIndices)

3. Actualización de Parámetros (Push)

// Actualización incremental mediante gradientes
modelWeights.increment(deltaGradient)

// Sincronización de reloj para confirmar la iteración
modelWeights.syncClock()

Implementación de Regresión Logística Distribuida

El sigueinte ejemplo muestra cómo estructurar un algoritmo de Regresión Logística (LR) utilizando el framework Angel.

Definición del Modelo

class CustomLRModel(conf: Configuration, _ctx: TaskContext = null) 
 extends MLModel(conf, _ctx) {
 
 val featRange = conf.getInt(MLConf.ML_FEATURE_INDEX_RANGE, 10000)
 
 // Inicializamos el contenedor de parámetros en el PS
 val weightsPS = PSModel("lr.params", 1, featRange)
   .setRowType(RowType.T_DOUBLE_DENSE)
 
 addPSModel(weightsPS)
 
 setSavePath(conf)
 setLoadPath(conf)

 override def predict(data: DataBlock[LabeledData]): DataBlock[PredictResult] = {
   // Lógica de inferencia omitida por brevedad
   null
 }
}

Lógica de Entrenamiento

class LRTrainProcess extends TrainTask[LabeledData] {
 
 override def parse(key: LongWritable, value: Text): LabeledData = {
   DataParser.parseVector(key, value, feaNum, "dummy", negY = true)
 }

 override def train(ctx: TaskContext): Unit = {
   val lrModel = new CustomLRModel(conf, ctx)
   val psRef = lrModel.weightsPS
   val alpha = conf.getDouble(MLConf.ML_LEARN_RATE, 0.01)
   val maxIter = conf.getInt(MLConf.ML_EPOCH_NUM, 50)

   for (epoch <- 0 until maxIter) {
     // 1. Obtener pesos actuales del PS
     val localWeights = psRef.getRow(0)
     
     // 2. Calcular gradiente localmente
     val localGradient = calculateGrad(localWeights, ctx.getDataBlock)
     
     // 3. Aplicar tasa de aprendizaje y enviar al PS
     localGradient.timesBy(-1.0 * alpha)
     psRef.increment(localGradient)
     
     // 4. Sincronizar con otros workers
     psRef.clock().get()
     ctx.incIteration()
   }
 }

 private def calculateGrad(w: TVector, data: DataBlock[LabeledData]): TVector = {
   val gradAccumulator = new DenseDoubleVector(w.getDimension)
   val reader = data.readingIterator()
   
   while (reader.hasNext) {
     val item = reader.next()
     val x = item.getX
     val y = item.getY
     val score = 1.0 / (1.0 + Math.exp(-w.dot(x)))
     gradAccumulator.plusBy(x.timesBy(score - y))
   }
   gradAccumulator
 }
}

Configuración y Despliegue en Yarn

Para ejecutar el entrenamiento en un clúster Hadoop, se utiliza el script de envío de Angel especificando los recursos para Workers y PServers.

./bin/angel-submit \
 --action.type train \
 --angel.app.submit.class com.example.CustomLRRunner \
 --angel.train.data.path /data/input \
 --angel.save.model.path /models/output \
 --ml.epoch.num 50 \
 --ml.feature.num 10000 \
 --angel.workergroup.number 4 \
 --angel.worker.memory.mb 8192 \
 --angel.ps.number 2 \
 --angel.ps.memory.mb 4096

Estrategias de Optimización de Rendimiento

Para maximizar la eficiencia en modelos de gran escala, considere las siguientes técnicas:

Gestión de Memoria y Almacenamiento

  • Uso de Tipos Dispersos: Si los datos tienen muchas dimensiones vacías, utilice RowType.T_DOUBLE_SPARSE para reducir drásticamente el consumo de memoria en el PServer.
  • Hogwild: Habilite esta opción para permitir actualizaciones asíncronas de hilos sin bloqueos, aumentando el rendimiento a costa de una ligera inconsistencia temporal.

Protocolos de Sincronización

Modo Descripción Ventaja
BSP (Bulk Synchronous Parallel) Sincronización estricta en cada iteración. Garantiza convergencia estable.
SSP (Stale Synchronous Parallel) Permite que algunos workers se adelanten N pasos. Reduce el impacto de workers lentos (stragglers).
Asíncrono Sin barreras de sincronización. Máximo rendimiento de hardware.

Resolución de Problemas Comunes

Si encuentra errores de OutOfMemory (OOM), verifique la configuración de memoria en el PServer (--angel.ps.memory.mb). Un PServer almacena fragmentos de la matriz de parámetros; si el modelo es muy grande, aumente el número de PServers para distribuir mejor la carga.

Para problemas de convergencia lenta, ajuste el MLConf.ML_LEARN_RATE o verifique la magnitud de los gradientes antes de enviarlos al PS mediante registros de log controlados en el TaskContext.

Etiquetas: angel-ps distributed-computing machine-learning scala parameter-server

Publicado el 8-1 12:00