Aller au contenu principal

Entraînement distribué de modèles XGBoost à l’aide de xgboost.spark

info

Aperçu

Cette fonctionnalité est en aperçu public.

Le package Python xgboost>=1.7 contient un nouveau module xgboost.spark. Ce module inclut les estimateurs PySpark xgboost xgboost.spark.SparkXGBRegressor, xgboost.spark.SparkXGBClassifier et xgboost.spark.SparkXGBRanker. Ces nouvelles classes prennent en charge l'inclusion d'estimateurs XGBoost dans les Pipelines SparkML. Pour plus de détails sur l'API, consultez la documentation de l'API Python Spark API XGBoost.

Exigences

Databricks Runtime 12,0 ML et versions supérieures.

xgboost.spark paramètres

Les estimateurs définis dans le module xgboost.spark prennent en charge la plupart des paramètres et arguments utilisés dans XGBoost standard.

  • Les paramètres pour le constructeur de classe, la méthode fit et la méthode predict sont en grande partie identiques à ceux du module xgboost.sklearn.
  • La dénomination, les valeurs et les defaults sont en grande partie identiques à ceux décrits dans les parameters XGBoost.
  • Les exceptions sont quelques paramètres non pris en charge (tels que gpu_id, nthread, sample_weight, eval_set), et les paramètres spécifiques à l'estimateur pyspark qui ont été ajoutés (tels que featuresCol, labelCol, use_gpu, validationIndicatorCol). Pour plus de détails, consultez la documentation de l'XGBoost Python Spark API.

Formation distribuée

Les estimateurs PySpark définis dans le module xgboost.spark prennent en charge la formation distribuée XGBoost en utilisant le paramètre num_workers. Pour utiliser l'entraînement distribué, créez un classifieur ou un régresseur et définissez num_workers sur le nombre de tâches Spark exécutées simultanément pendant l'entraînement distribué. Pour utiliser tous les emplacements de tâches Spark, définissez num_workers=sc.defaultParallelism.

Par exemple :

Python
from xgboost.spark import SparkXGBClassifier
classifier = SparkXGBClassifier(num_workers=sc.defaultParallelism)
remarque
  • Vous ne pouvez pas utiliser mlflow.xgboost.autolog avec XGBoost distribué. Pour consigner un modèle xgboost Spark à l'aide de MLflow, utilisez mlflow.spark.log_model(spark_xgb_model, artifact_path).
  • Vous ne pouvez pas utiliser XGBoost distribué sur un cluster sur lequel l'autoscaling est activé. Les nouveaux nœuds Worker qui start dans ce paradigme de mise à l'échelle élastique ne peuvent pas recevoir de nouveaux ensembles de tâches et restent inactifs. Pour obtenir des instructions sur la désactivation de l'autoscaling, voir Activer l'autoscaling.

Activer l'optimisation pour l'entraînement sur des dataset à fonctionnalités éparses

Les estimateurs PySpark définis dans le module xgboost.spark prennent en charge l'optimisation pour l'entraînement sur des jeux de données avec des caractéristiques clairsemées. Pour activer l'optimisation des ensembles de fonctionnalités rares, vous devez fournir un dataset à la méthode fit qui contient une colonne de fonctionnalités composée de valeurs de type pyspark.ml.linalg.SparseVector et définir le parameter d'estimateur enable_sparse_data_optim sur True. De plus, vous devez définir le paramètre missing sur 0.0.

Par exemple :

Python
from xgboost.spark import SparkXGBClassifier
classifier = SparkXGBClassifier(enable_sparse_data_optim=True, missing=0.0)
classifier.fit(dataset_with_sparse_features_col)

Entraînement GPU

Les estimateurs PySpark définis dans le module xgboost.spark prennent en charge l'entraînement sur GPU. Définissez le parameter use_gpu sur True pour activer l'entraînement GPU.

remarque

Pour chaque tâche Spark utilisée dans l'entraînement distribué XGBoost, un seul GPU est utilisé pour l'entraînement lorsque l'argument use_gpu est défini sur True. Databricks recommande d’utiliser la valeur default de 1 pour la configuration du cluster Spark spark.task.resource.gpu.amount. Sinon, les GPU supplémentaires alloués à cette tâche Spark sont inactifs.

Par exemple :

Python
from xgboost.spark import SparkXGBClassifier
classifier = SparkXGBClassifier(num_workers=sc.defaultParallelism, use_gpu=True)

Dépannage

Pendant la formation multi-nœuds, si vous rencontrez un message NCCL failure: remote process exited or there was a network error, cela indique généralement un problème de communication réseau entre les GPU. Ce problème survient lorsque NCCL (NVIDIA Collective Communications Library) ne peut pas utiliser certaines interfaces réseau pour la communication GPU.

Pour résoudre, définissez la sparkConf du cluster pour spark.executorEnv.NCCL_SOCKET_IFNAME à eth. Cela définit essentiellement la variable d'environnement NCCL_SOCKET_IFNAME sur eth pour tous les Worker d'un nœud.

Exemple de Notebook

Ce Notebook montre l'utilisation du package Python xgboost.spark avec Spark MLlib.

notebook PySpark-XGBoost

Guide de migration pour le module sparkdl.xgboost obsolète

  • Remplacer from sparkdl.xgboost import XgboostRegressor par from xgboost.spark import SparkXGBRegressor et remplacer from sparkdl.xgboost import XgboostClassifier par from xgboost.spark import SparkXGBClassifier.
  • Modifiez tous les noms de parameter dans le constructeur d'estimateur, du style camelCase au style snake_case. Par exemple, changer XgboostRegressor(featuresCol=XXX) en SparkXGBRegressor(features_col=XXX).
  • Les paramètres use_external_storage et external_storage_precision ont été supprimés. Les estimateurs xgboost.spark utilisent l’API d’itération de données DMatrix pour utiliser la mémoire plus efficacement. Il n'est plus nécessaire d'utiliser le mode de stockage externe inefficace. Pour les très grands datasets, Databricks vous recommande d’augmenter le paramètre num_workers, ce qui permet à chaque tâche d’entraînement de partitionner les données en partitions de données plus petites et plus gérables. Envisagez de définir num_workers = sc.defaultParallelism, qui définit num_workers sur le nombre total d’emplacements de tâches Spark dans le cluster.
  • Pour les estimateurs définis dans xgboost.spark, le paramètre num_workers=1 exécute l'entraînement du modèle à l'aide d'une seule tâche Spark. Ceci utilise le nombre de cœurs de processeur spécifié par le paramètre de configuration du cluster Spark spark.task.cpus, qui est de 1 par default. Pour utiliser plus de cœurs de processeur pour entraîner le modèle, augmentez num_workers ou spark.task.cpus. Vous ne pouvez pas définir le paramètre nthread ou n_jobs pour les estimateurs définis dans xgboost.spark. Ce comportement est différent du comportement précédent des estimateurs définis dans le package sparkdl.xgboost obsolète.

Convertir le modèle sparkdl.xgboost en modèle xgboost.spark

sparkdl.xgboost les modèles sont enregistrés dans un format différent de celui des modèles xgboost.spark et ont des paramètres différents. Utilisez la fonction utilitaire suivante pour convertir le modèle :

Python
def convert_sparkdl_model_to_xgboost_spark_model(
xgboost_spark_estimator_cls,
sparkdl_xgboost_model,
):
"""
:param xgboost_spark_estimator_cls:
`xgboost.spark` estimator class, e.g. `xgboost.spark.SparkXGBRegressor`
:param sparkdl_xgboost_model:
`sparkdl.xgboost` model instance e.g. the instance of
`sparkdl.xgboost.XgboostRegressorModel` type.

:return
A `xgboost.spark` model instance
"""

def convert_param_key(key):
from xgboost.spark.core import _inverse_pyspark_param_alias_map
if key == "baseMarginCol":
return "base_margin_col"
if key in _inverse_pyspark_param_alias_map:
return _inverse_pyspark_param_alias_map[key]
if key in ['use_external_storage', 'external_storage_precision', 'nthread', 'n_jobs', 'base_margin_eval_set']:
return None
return key

xgboost_spark_params_dict = {}
for param in sparkdl_xgboost_model.params:
if param.name == "arbitraryParamsDict":
continue
if sparkdl_xgboost_model.isDefined(param):
xgboost_spark_params_dict[param.name] = sparkdl_xgboost_model.getOrDefault(param)

xgboost_spark_params_dict.update(sparkdl_xgboost_model.getOrDefault("arbitraryParamsDict"))

xgboost_spark_params_dict = {
convert_param_key(k): v
for k, v in xgboost_spark_params_dict.items()
if convert_param_key(k) is not None
}

booster = sparkdl_xgboost_model.get_booster()
booster_bytes = booster.save_raw("json")
booster_config = booster.save_config()
estimator = xgboost_spark_estimator_cls(**xgboost_spark_params_dict)
sklearn_model = estimator._convert_to_sklearn_model(booster_bytes, booster_config)
return estimator._copyValues(estimator._create_pyspark_model(sklearn_model))

# Example
from xgboost.spark import SparkXGBRegressor

new_model = convert_sparkdl_model_to_xgboost_spark_model(
xgboost_spark_estimator_cls=SparkXGBRegressor,
sparkdl_xgboost_model=model,
)

Si vous avez un modèle pyspark.ml.PipelineModel contenant un modèle sparkdl.xgboost comme dernière étape, vous pouvez remplacer l’étape du modèle sparkdl.xgboost par le modèle xgboost.spark converti.

Python
pipeline_model.stages[-1] = convert_sparkdl_model_to_xgboost_spark_model(
xgboost_spark_estimator_cls=SparkXGBRegressor,
sparkdl_xgboost_model=pipeline_model.stages[-1],
)