Aller au contenu principal

start a Ray cluster on Databricks

Databricks simplifie le processus de démarrage d'un cluster Ray en gérant la configuration des clusters et des tâches de la même manière qu'il le fait pour tout Job Apache Spark. C'est parce que le cluster Ray est effectivement démarré au-dessus du cluster Apache Spark géré.

Exécuter Ray sur Databricks

Python
from ray.util.spark import setup_ray_cluster
import ray

# If the cluster has four workers with 8 CPUs each as an example
setup_ray_cluster(num_worker_nodes=4, num_cpus_per_worker=8)

# Pass any custom configuration to ray.init
ray.init(ignore_reinit_error=True)

Cette approche fonctionne à n'importe quelle échelle de cluster, de quelques nœuds à des centaines. Les clusters Ray sur Databricks prennent également en charge la mise à l'échelle automatique.

Après avoir créé le cluster Ray, vous pouvez exécuter n'importe quel code d'application Ray dans un Notebook Databricks.

important

Databricks recommande d'installer toutes les bibliothèques nécessaires pour votre application avec %pip install <your-library-dependency> afin de garantir qu'elles soient disponibles pour votre cluster Ray et votre application en conséquence. La spécification des dépendances dans l'appel de fonction d'initialisation Ray installe les dépendances dans un emplacement inaccessible aux nœuds Worker Apache Spark, ce qui entraîne des incompatibilités de version et des erreurs d'importation.

Par exemple, vous pouvez exécuter une application Ray simple dans un Notebook Databricks, comme suit :

Python
import ray
import random
import time
from fractions import Fraction

ray.init()

@ray.remote
def pi4_sample(sample_count):
"""pi4_sample runs sample_count experiments, and returns the
fraction of time it was inside the circle.
"""
in_count = 0
for i in range(sample_count):
x = random.random()
y = random.random()
if x*x + y*y <= 1:
in_count += 1
return Fraction(in_count, sample_count)

SAMPLE_COUNT = 1000 * 1000
start = time.time()
future = pi4_sample.remote(sample_count=SAMPLE_COUNT)
pi4 = ray.get(future)
end = time.time()
dur = end - start
print(f'Running {SAMPLE_COUNT} tests took {dur} seconds')

pi = pi4 * 4
print(float(pi))

Arrêter un cluster Ray

Les clusters Ray s'arrêtent automatiquement dans les circonstances suivantes :

  • Vous détachez votre Notebook interactif de votre cluster Databricks.
  • Votre Job Databricks est terminé.
  • Votre cluster Databricks est redémarré ou arrêté.
  • Il n'y a aucune activité pour le temps d'inactivité spécifié.

Pour arrêter un cluster Ray exécuté sur Databricks, vous pouvez appeler l'API ray.utils.spark.shutdown_ray_cluster.

Python
from ray.utils.spark import shutdown_ray_cluster
import ray

shutdown_ray_cluster()
ray.shutdown()

Étapes suivantes