Aller au contenu principal

HorovodRunner : deep learning distribué avec Horovod

important

Horovod et HorovodRunner sont désormais obsolètes. Les versions après 15.4 LTS ML n’auront pas ce package préinstallé. Pour le deep learning distribué, Databricks recommande d'utiliser TorchDistributor pour l'entraînement distribué avec PyTorch ou l'API tf.distribute.Strategy pour l'entraînement distribué avec TensorFlow.

Découvrez comment effectuer l'entraînement distribué de modèles de machine learning à l'aide de HorovodRunner pour lancer des Jobs d'entraînement Horovod en tant que Jobs Spark sur Databricks.

Qu'est-ce que HorovodRunner ?

HorovodRunner est une API générale pour exécuter des charges de travail de deep learning distribuées sur Databricks à l'aide du framework Horovod. En intégrant Horovod au mode barrière de Spark, Databricks est en mesure d'offrir une plus grande stabilité pour les jobs d'entraînement de deep learning de longue durée sur Spark. HorovodRunner utilise une méthode Python qui contient du code d'entraînement de deep learning avec des hooks Horovod. HorovodRunner sérialise la méthode sur le driver et la distribue aux workers Spark. Un job MPI Horovod est intégré en tant que job Spark à l'aide du mode d'exécution de barrière. Le premier exécuteur collecte les adresses IP de tous les exécuteurs de tâches à l'aide de BarrierTaskContext et déclenche un job Horovod à l'aide de mpirun. Chaque processus MPI Python charge le programme utilisateur sérialisé, le désérialise et l'exécute.

HorovodRunner

Entraînement distribué avec HorovodRunner

HorovodRunner vous permet de lancer des jobs d'entraînement Horovod en tant que jobs Spark. L'API HorovodRunner prend en charge les méthodes présentées dans le tableau. Pour plus de détails, consultez la documentation de l'API HorovodRunner.

Méthode et signature.

Description

init(self, np)

Créez une instance de HorovodRunner.

**run(self, main, **kwargs)**

Exécutez un job d'entraînement Horovod en invoquant main(**kwargs). La fonction principale et les arguments de mot-clé sont sérialisés à l'aide de cloudpickle et distribués aux Worker du cluster.

Méthode et signature.

Description

init(self, np)

Créez une instance de HorovodRunner.

**run(self, main, **kwargs)**

Exécutez un job d'entraînement Horovod en invoquant main(**kwargs). La fonction principale et les arguments de mot-clé sont sérialisés à l'aide de cloudpickle et distribués aux Worker du cluster.

L'approche générale pour développer un programme d'entraînement distribué à l'aide de HorovodRunner est :

  1. Créez une instance HorovodRunner initialisée avec le nombre de nœuds.
  2. Définissez une méthode d'entraînement Horovod selon les méthodes décrites dans l'utilisation de Horovod, en veillant à ajouter toutes les déclarations d'importation à l'intérieur de la méthode.
  3. Transmettez la méthode d'entraînement à l'instance HorovodRunner.

Par exemple :

Python
hr = HorovodRunner(np=2)

def train():
import tensorflow as tf
hvd.init()

hr.run(train)

Pour exécuter HorovodRunner sur le Driver uniquement avec n sous-processus, utilisez hr = HorovodRunner(np=-n). Par exemple, s'il y a 4 GPU sur le nœud du driver, vous pouvez choisir n jusqu'à 4. Pour plus de détails sur le paramètre np, consultez la documentation de l'API HorovodRunner. Pour plus de détails sur la manière de pin un GPU par sous-processus, veuillez consulter le guide d'utilisation de Horovod.

Une erreur courante est que les objets TensorFlow sont introuvables ou ne peuvent pas être sérialisés. Cela se produit lorsque les instructions d'importation de la bibliothèque ne sont pas distribuées aux autres exécuteurs. Pour éviter ce problème, incluez toutes les instructions d'importation (par exemple,)import tensorflow as tf à la fois en haut de la méthode d'entraînement Horovod et à l'intérieur de toute autre fonction définie par l'utilisateur appelée dans la méthode d'entraînement Horovod.

Enregistrer l'entraînement Horovod avec Horovod Timeline

Horovod a la capacité d'enregistrer la chronologie de son activité, appelée Horovod Timeline.

important

Horovod Timeline a un impact significatif sur les performances. Le throughput d'Inception3 peut diminuer d'environ 40 % lorsque Horovod Timeline est activé. Pour accélérer les Jobs HorovodRunner, n'utilisez pas la chronologie Horovod.

Vous ne pouvez pas visualiser la chronologie Horovod pendant que la formation est en cours.

Pour enregistrer une chronologie Horovod, définissez la variable d'environnement HOROVOD_TIMELINE à l'emplacement où vous souhaitez enregistrer le fichier de chronologie. Databricks recommande d'utiliser un emplacement sur un stockage partagé afin que le fichier de chronologie puisse être facilement récupéré. Par exemple, vous pouvez utiliser les API de fichiers locaux DBFS, comme illustré :

Python
timeline_dir = "/dbfs/ml/horovod-timeline/%s" % uuid.uuid4()
os.makedirs(timeline_dir)
os.environ['HOROVOD_TIMELINE'] = timeline_dir + "/horovod_timeline.json"
hr = HorovodRunner(np=4)
hr.run(run_training_horovod, params=params)

Ensuite, ajoutez le code spécifique à la chronologie au début et à la fin de la fonction d'entraînement. L'exemple de notebook suivant inclut un exemple de code que vous pouvez utiliser comme solution de contournement pour afficher la progression de l'entraînement.

Exemple de Notebook de chronologie Horovod

Pour download le fichier de chronologie, utilisez la CLI Databricks, puis utilisez la fonction chrome://tracing du navigateur Chrome pour l'afficher. Par exemple :

Chronologie Horovod

Workflow de développement

Voici les étapes générales pour migrer le code de deep learning à nœud unique vers l'entraînement distribué. Les exemples : Migrer vers le deep learning distribué avec HorovodRunner de cette section illustrent ces étapes.

  1. Préparer le code à nœud unique : Préparez et testez le code à nœud unique avec TensorFlow, Keras ou PyTorch.

  2. Migrer vers Horovod : suivez les instructions de la page Utilisation de Horovod pour migrer le code avec Horovod et le tester sur le driver :

    1. Ajoutez hvd.init() pour initialiser Horovod.
    2. pin un GPU serveur à utiliser par ce processus à l'aide de config.gpu_options.visible_device_list. Avec la configuration habituelle d'un GPU par processus, cela peut être défini sur le rang local. Dans ce cas, le premier processus sur le serveur se verra attribuer le premier GPU, le deuxième processus se verra attribuer le deuxième GPU, et ainsi de suite.
    3. Incluez une partition du dataset. Cet opérateur de dataset est très utile lors de l'exécution d'un entraînement distribué, car il permet à chaque worker de lire un sous-ensemble unique.
    4. Montez en charge le taux d'apprentissage par le nombre de Worker. La taille effective du batch dans l'entraînement distribué synchrone est mise à l'échelle par le nombre de Worker. L'augmentation du taux d'apprentissage compense l'augmentation de la taille du batch.
    5. Enveloppez l'optimiseur dans hvd.DistributedOptimizer. L’optimiseur distribué délègue le calcul du gradient à l’optimiseur d’origine, effectue la moyenne des gradients à l’aide d’allreduce ou d’allgather, puis applique les gradients moyennés.
    6. Ajoutez hvd.BroadcastGlobalVariablesHook(0) pour diffuser les états initiaux des variables du rang 0 à tous les autres processus. Ceci est nécessaire pour assurer une initialisation cohérente de tous les Worker lorsque l'entraînement est démarré avec des poids aléatoires ou restauré à partir d'un point de contrôle. Alternativement, si vous n'utilisez pas MonitoredTrainingSession, vous pouvez exécuter l'opération hvd.broadcast_global_variables après que les variables globales ont été initialisées.
    7. Modifiez votre code pour enregistrer les points de contrôle uniquement sur le Worker 0 afin d'empêcher d'autres Workers de les corrompre.
  3. Migrer vers HorovodRunner : HorovodRunner exécute le Job d'entraînement Horovod en appelant une fonction Python. Vous devez encapsuler la procédure d'entraînement principale dans une seule fonction Python. Vous pouvez ensuite tester HorovodRunner en mode local et en mode distribué.

Mettre à jour les bibliothèques de deep learning

Si vous mettez à niveau ou à niveau inférieur TensorFlow, Keras ou PyTorch, vous devez réinstaller Horovod afin qu'il soit compilé par rapport à la bibliothèque nouvellement installée. Par exemple, si vous souhaitez mettre à niveau TensorFlow, Databricks recommande d'utiliser le script d'initialisation des instructions d'installation de TensorFlow et d'y ajouter le code d'installation de Horovod spécifique à TensorFlow. Veuillez consulter les instructions d'installation de Horovod pour utiliser différentes combinaisons, telles que la mise à niveau ou la rétrogradation de PyTorch et d'autres bibliothèques.

Bash
add-apt-repository -y ppa:ubuntu-toolchain-r/test
apt update

# Using the same compiler that TensorFlow was built to compile Horovod
apt install g++-7 -y
update-alternatives --install /usr/bin/gcc gcc /usr/bin/gcc-7 60

HOROVOD_GPU_ALLREDUCE=NCCL HOROVOD_CUDA_HOME=/usr/local/cuda pip install horovod==0.18.1 --force-reinstall --no-deps --no-cache-dir

Exemples : Migrer vers l'apprentissage profond distribué avec HorovodRunner

Les exemples suivants, basés sur le dataset MNIST, montrent comment migrer un programme de deep learning à nœud unique vers le deep learning distribué avec HorovodRunner.

Limitations

  • Lorsque vous travaillez avec des fichiers de workspace, HorovodRunner ne fonctionnera pas si np est supérieur à 1 et si le notebook importe à partir d’autres fichiers relatifs. Envisagez d’utiliser horovod.spark au lieu de HorovodRunner.
  • Si vous rencontrez des erreurs comme WARNING: Open MPI accepted a TCP connection from what appears to be a another Open MPI process but cannot find a corresponding process entry for that peer, cela indique un problème de communication réseau entre les nœuds de votre cluster. Pour résoudre cette erreur, ajoutez l'extrait suivant dans votre code d'entraînement pour utiliser l'interface réseau principale.
Python
import os
os.environ["OMPI_MCA_btl_tcp_if_include"]="eth0"
os.environ["NCCL_SOCKET_IFNAME"]="eth0"