Aller au contenu principal

Intégrer MLflow et Ray

MLflow est la plus grande plateforme d'ingénierie d'IA open source pour les agents, les LLM et les modèles de ML. MLflow permet aux équipes de toutes tailles de déboguer, d'évaluer, de surveiller et d'optimiser leurs applications d'IA tout en contrôlant les coûts et en gérant l'accès aux modèles et aux données. Avec plus de 30 millions de téléchargements mensuels, des milliers d'organisations s'appuient sur MLflow chaque jour pour déployer l'IA en production en toute confiance. La combinaison de Ray avec MLflow vous permet de distribuer des charges de travail avec Ray et de suivre les modèles, les métriques, les paramètres et les métadonnées générés pendant l'entraînement avec MLflow.

Cet article explique comment intégrer MLflow avec les composants Ray suivants :

  • Ray Core: applications distribuées à usage général non couvertes par Ray Tune et Ray Train

  • Ray Train: formation de modèles distribués

  • Ray Tune: Réglage distribué des hyperparamètres

  • Model Serving: Déploiement de modèles pour l'inférence en temps réel.

Intégrer Ray Core et MLflow

Ray Core fournit les blocs de construction fondamentaux pour les applications distribuées à usage général. Il vous permet de monter en charge les fonctions et les classes Python sur plusieurs nœuds.

Cette section décrit les schémas suivants pour intégrer Ray Core et MLflow :

  • Enregistrer les modèles MLflow à partir du processus Driver Ray
  • Consigner les modèles MLflow des exécutions enfants

Journalisez MLflow à partir du processus du Driver Ray

Il est généralement préférable d'enregistrer les modèles MLflow à partir du processus Driver plutôt qu'à partir des nœuds Worker. Cela est dû à la complexité accrue du transfert de références avec état aux Workers distants.

Par exemple, le code suivant échoue parce que le serveur MLflow Tracking n'est pas initialisé à l'aide de MLflow Client à partir des nœuds worker.

Python
import mlflow

@ray.remote
def example_logging_task(x):
# ...

# This method will fail
mlflow.log_metric("x", x)
return x

with mlflow.start_run() as run:
ray.get([example_logging_task.remote(x) for x in range(10)])

Au lieu de cela, renvoyez les métriques au nœud driver. Les métriques et les métadonnées sont généralement suffisamment petites pour être transférées vers le Driver sans causer de problèmes de mémoire.

Prenez l'exemple ci-dessus et mettez-le à jour pour Logs les métriques renvoyées par une tâche Ray :

Python
import mlflow

@ray.remote
def example_logging_task(x):
# ...
return x

with mlflow.start_run() as run:
results = ray.get([example_logging_task.remote(x) for x in range(10)])
for x in results:
mlflow.log_metric("x", x)

Pour les tâches qui nécessitent d'enregistrer des artefacts volumineux, tels qu'une grande table Pandas, des images, des tracés ou des modèles, Databricks vous recommande de conserver l'artefact en tant que fichier. Ensuite, rechargez l'artefact dans le contexte du Driver ou enregistrez directement l'objet avec MLflow en spécifiant le chemin d'accès au fichier enregistré.

Python
import mlflow

@ray.remote
def example_logging_task(x):
# ...
# Create a large object that needs to be stored
with open("/Volumes/<catalog>/<schema>/<volume>/myLargeFilePath.txt", "w") as f:
f.write(myLargeObject)
return x

with mlflow.start_run() as run:
results = ray.get([example_logging_task.remote(x) for x in range(10)])
for x in results:
mlflow.log_metric("x", x)
# Directly log the saved file by specifying the path
mlflow.log_artifact("/Volumes/<catalog>/<schema>/<volume>/myLargeFilePath.txt")

Enregistrer les tâches Ray en tant qu'exécutions enfants MLflow

Vous pouvez intégrer Ray Core à MLflow en utilisant des exécutions enfants. Cela implique les étapes suivantes :

  1. Créez une exécution parente : Initialisez une exécution parente dans le processus du driver. Cette exécution sert de conteneur hiérarchique pour toutes les exécutions enfant suivantes.
  2. Créer des exécutions enfants : au sein de chaque tâche Ray, lancez une exécution enfant sous l’exécution parente. Chaque exécution enfant peut indépendamment enregistrer ses propres métriques.

Pour implémenter cette approche, assurez-vous que chaque tâche Ray reçoit les identifiants client nécessaires et le parent run_id. Cette configuration établit la relation hiérarchique parent-enfant entre les exécutions. L’extrait de code suivant montre comment récupérer les informations d’identification et les transmettre au parent run_id:

Python
from mlflow.utils.databricks_utils import get_databricks_env_vars
mlflow_db_creds = get_databricks_env_vars("databricks")

username = "" # Username path
experiment_name = f"/Users/{username}/mlflow_test"

mlflow.set_experiment(experiment_name)

@ray.remote
def ray_task(x, run_id):
import os
# Set the MLflow credentials within the Ray task
os.environ.update(mlflow_db_creds)
# Set the active MLflow experiment within each Ray task
mlflow.set_experiment(experiment_name)
# Create nested child runs associated with the parent run_id
with mlflow.start_run(run_id=run_id, nested=True):
# Log metrics to the child run within the Ray task
mlflow.log_metric("x", x)

return x

# Start parent run on the main driver process
with mlflow.start_run() as run:
# Pass the parent run's run_id to each Ray task
results = ray.get([ray_task.remote(x, run.info.run_id) for x in range(10)])

Ray Train et MLflow

Le moyen le plus simple d'enregistrer les modèles Ray Train dans MLflow est d'utiliser le point de contrôle généré par l'exécution de l'entraînement. Une fois l'exécution de l'entraînement terminée, rechargez le modèle dans son framework de deep learning natif (tel que PyTorch ou TensorFlow), puis enregistrez-le avec le code MLflow correspondant.

Cette approche garantit que le modèle est stocké correctement et est prêt pour l'évaluation ou le déploiement.

Le code suivant recharge un modèle à partir d'un point de contrôle Ray Train et le logs dans MLflow :

Python
result = trainer.fit()

checkpoint = result.checkpoint
with checkpoint.as_directory() as checkpoint_dir:
# Change as needed for different DL frameworks
checkpoint_path = f"{checkpoint_dir}/checkpoint.ckpt"
# Load the model from the checkpoint
model = MyModel.load_from_checkpoint(checkpoint_path)

with mlflow.start_run() as run:
# Change the MLflow flavor as needed
mlflow.pytorch.log_model(model, "model")

Bien qu'il soit généralement préférable de renvoyer les objets au nœud Driver, avec Ray Train, sauvegarder les résultats finaux est plus facile que l'historique complet de l'entraînement du processus Worker.

Pour stocker plusieurs modèles à partir d'une exécution d'entraînement, spécifiez le nombre de points de contrôle à conserver dans le ray.train.CheckpointConfig. Les modèles peuvent ensuite être lus et enregistrés de la même manière que le stockage d'un modèle unique.

remarque

MLflow n'est pas responsable de la gestion de la tolérance aux pannes pendant l'entraînement du modèle, mais plutôt du suivi du cycle de vie du modèle. La tolérance aux pannes est plutôt gérée par Ray Train lui-même.

Pour stocker les métriques d'entraînement spécifiées par Ray Train, récupérez-les de l'objet de résultat et stockez-les à l'aide de MLflow.

Python
result = trainer.fit()

with mlflow.start_run() as run:
mlflow.log_metrics(result.metrics_dataframe.to_dict(orient='dict'))

# Change the MLflow flavor as needed
mlflow.pytorch.log_model(model, "model")

Pour configurer correctement vos clusters Spark et Ray et éviter les problèmes d'allocation de ressources, vous devez ajuster le paramètre resources_per_worker. Plus précisément, définissez le nombre de CPU pour chaque worker Ray à un de moins que le nombre total de CPU disponibles sur un nœud worker Ray. Cet ajustement est crucial, car si le formateur réserve tous les cœurs disponibles pour les acteurs Ray, cela peut entraîner des erreurs de contention de ressources.

Ray Tune et MLflow

L’intégration de Ray Tune à MLflow vous permet de suivre et de consigner efficacement les expérimentations de réglage des hyperparamètres dans Databricks. Cette intégration tire parti des capacités de suivi d'expérimentation de MLflow pour enregistrer les métriques et les résultats directement à partir des tâches Ray.

Approche d'exécution enfant pour la journalisation

Similaire à la journalisation à partir des tâches Ray Core, les applications Ray Tune peuvent utiliser une approche d'exécution enfant pour journaliser les métriques de chaque essai ou itération de réglage. Utilisez les étapes suivantes pour mettre en œuvre une approche d’exécution enfant :

  1. Créez une exécution parente : Initialisez une exécution parente dans le processus du driver. Cette exécution sert de conteneur principal pour toutes les exécutions enfant suivantes.
  2. Log child runs: Chaque tâche Ray Tune crée une exécution enfant sous l'exécution parente, en maintenant une hiérarchie claire des résultats d'expérimentation.

L'exemple suivant montre comment s'authentifier et consigner des informations à partir des tâches Ray Tune à l'aide de MLflow.

Python
import os
import tempfile
import time

import mlflow
from mlflow.utils.databricks_utils import get_databricks_env_vars

from ray import train, tune
from ray.air.integrations.mlflow import MLflowLoggerCallback, setup_mlflow

mlflow_db_creds = get_databricks_env_vars("databricks")

EXPERIMENT_NAME = "/Users/<WORKSPACE_USERNAME>/setup_mlflow_example"
mlflow.set_experiment(EXPERIMENT_NAME)

def evaluation_fn(step, width, height):
return (0.1 + width * step / 100) ** (-1) + height * 0.1

def train_function_mlflow(config, run_id):
os.environ.update(mlflow_db_creds)
mlflow.set_experiment(EXPERIMENT_NAME)

# Hyperparameters
width = config["width"]
height = config["height"]

with mlflow.start_run(run_id=run_id, nested=True):
for step in range(config.get("steps", 100)):
# Iterative training function - can be any arbitrary training procedure
intermediate_score = evaluation_fn(step, width, height)
# Log the metrics to MLflow
mlflow.log_metrics({"iterations": step, "mean_loss": intermediate_score})
# Feed the score back to Tune.
train.report({"iterations": step, "mean_loss": intermediate_score})
time.sleep(0.1)

def tune_with_setup(run_id, finish_fast=True):
os.environ.update(mlflow_db_creds)
# Set the experiment or create a new one if it does not exist.
mlflow.set_experiment(experiment_name=EXPERIMENT_NAME)

tuner = tune.Tuner(
tune.with_parameter(train_function_mlflow, run_id),
tune_config=tune.TuneConfig(num_samples=5),
run_config=train.RunConfig(
name="mlflow",
),
param_space={
&quot;width&quot;: tune.randint(10, 100),
&quot;height&quot;: tune.randint(0, 100),
&quot;steps&quot;: 20 if finish_fast else 100,
},
)
results = tuner.fit()

with mlflow.start_run() as run:
mlflow_tracking_uri = mlflow.get_tracking_uri()
tune_with_setup(run.info.run_id)

Model Serving

L'utilisation de Ray Serve sur les clusters Databricks pour l'inférence en temps réel pose des défis en raison des limitations de sécurité réseau et de connectivité lors de l'interaction avec des applications externes.

Databricks recommande d'utiliser Model Serving pour déployer des Modèles de machine learning en production sur un Endpoint d'API REST. Pour plus d'informations, consultez Vue d'ensemble des modèles personnalisés.