What you will be able to do
- Build a single-machine Hyperopt fmin() workflow from an objective function, a search space, a search algorithm and max_evals
- Distribute that workflow across a Databricks cluster by passing a SparkTrials object to fmin()
- Configure SparkTrials parallelism and timeout, and explain the trade-off between parallelism and adaptivity
- Choose between SparkTrials and the default Trials class depending on whether the model is single-machine or distributed
- Describe how SparkTrials logs parent and child runs to MLflow, and recognise when distributing trials will not speed things up
Key concept
SparkTrials — A Databricks-developed Trials object that you pass to Hyperopt's fmin(). The driver keeps proposing hyperparameter settings, and each trial (one model fit) runs as a Spark task on a worker, so many single-machine models such as scikit-learn estimators train at the same time.
1.Start from a single-machine fmin() workflow
You can't parallelize a tuning run until you have one, so start with the plain single-machine Hyperopt workflow. Databricks describes four steps. First, define a function to minimize. Second, define a search space over the hyperparameters. Third, select a search algorithm. Fourth, run the tuning algorithm with Hyperopt fmin(). The Databricks example notebook tunes a scikit-learn support vector classifier on the Iris dataset, and the only hyperparameter it searches is the regularization parameter C.
Most of the code is in the objective function. It receives a candidate value of C, trains a model, scores it, and returns a loss. Hyperopt always *minimizes*. If your score is one where higher is better, such as accuracy, return its negative. The notebook scores each candidate with cross_val_score, but this objective only needs to return a loss. How you validate the model inside it is up to you.
def objective(C):
# Create a support vector classifier model
clf = SVC(C=C)
# Use the cross-validation accuracy to compare the models' performance
accuracy = cross_val_score(clf, X, y).mean()
# Hyperopt tries to minimize the objective function. A higher accuracy value means a better model, so you must return the negative accuracy.
return {'loss': -accuracy, 'status': STATUS_OK}The search space in the example is hp.lognormal('C', 0, 1.0). For the algorithm, the notebook names two main choices. hyperopt.tpe.suggest is Tree of Parzen Estimators, a Bayesian method that picks new settings based on past results. hyperopt.rand.suggest is random search, which samples the space without adapting. Keep this difference in mind, because it matters again when you turn up parallelism. The table below lists the core fmin() arguments.
| Argument | What it controls |
|---|---|
| fn | The objective function. Hyperopt calls it with values drawn from the search space, and it returns the loss |
| space | The hyperparameter space to search: categorical options or distributions such as uniform and log |
| algo | The search algorithm, most commonly hyperopt.rand.suggest (random search) or hyperopt.tpe.suggest (TPE) |
| max_evals | The number of hyperparameter settings to try, which is the number of models to fit |
Checkpoint 1 of 8· Fill the gap
This single-machine call should fit and evaluate at most 16 models. Which argument completes it?
argmin = fmin(
fn=objective,
space=search_space,
algo=algo,
? =16)max_evals is the maximum number of points in the hyperparameter space to test, which is the number of models fitted. parallelism is a SparkTrials argument, max_queue_len sets how many settings are generated ahead of time, and timeout is measured in seconds.
Source: docs.databricks.com2.Distribute the same run by adding SparkTrials
Only one thing changes. To distribute tuning, you add one more argument to fmin(): a Trials class called SparkTrials. Databricks developed SparkTrials so that a scikit-learn model, which only knows how to train on one machine, can be tuned faster on a cluster. Each model still trains on a single machine, but many models train at once on different workers. The objective, the search space, the algorithm and max_evals all stay the same.
from hyperopt import SparkTrials
# To display the API documentation for the SparkTrials class, uncomment the following line.
# help(SparkTrials)
spark_trials = SparkTrials()
with mlflow.start_run():
argmin = fmin(
fn=objective,
space=search_space,
algo=algo,
max_evals=16,
trials=spark_trials)
# Print the best value found for C
print("Best value found: ", argmin)Note two details in this sample. SparkTrials() is created with no arguments, so it uses the defaults for both of its optional arguments, which the next section covers. The call to fmin() also sits inside with mlflow.start_run():. Automated MLflow tracking is on by default, and this wrapper is what Databricks recommends. Section 5 explains why.
Checkpoint 2 of 8· Exam question
A machine learning engineer is training a `scikit-learn` `RandomForestClassifier` on a pandas DataFrame that fits comfortably in the memory of the Databricks driver node. The model itself trains as a single-node process, but the engineer wants to test 50 hyperparameter combinations by distributing the trials across the worker nodes of an all-purpose cluster using Hyperopt's `fmin`. Which approach correctly parallelizes this hyperparameter search?
Correct answer: A — Pass a `SparkTrials` object as the `trials` argument to `fmin`, since it distributes each single-node scikit-learn model evaluation to a separate Spark task on the cluster's workers.
- A. SparkTrials is the Hyperopt trials class built for exactly this case: it takes a single-node model's objective function and schedules each trial as its own Spark task, so evaluations run concurrently across the cluster's workers instead of one at a time on the driver.
- B. The base Trials class evaluates every trial sequentially on the driver process. It has no cluster-awareness of its own, so swapping in Trials would leave the scikit-learn search running one combination at a time regardless of how many workers the cluster has.
- C. Repartitioning a Spark DataFrame changes how Spark's own data processing is parallelized, but Hyperopt's trial scheduling is governed entirely by which trials object is passed to fmin, not by how the training data happens to be partitioned.
- D. max_evals only sets how many hyperparameter combinations fmin will evaluate in total; it says nothing about how many of those evaluations run at the same time. Concurrency is controlled by the parallelism setting on a SparkTrials object, not by the size of the evaluation budget.
3.How SparkTrials schedules trials: parallelism, timeout and the cluster
In Hyperopt, a trial generally means fitting one model on one setting of hyperparameters. With SparkTrials, the work is split between the driver and the workers. The driver generates new trials, and the workers evaluate them. Each trial is generated with a Spark job that has one task, and that task runs on a worker. If the cluster allows several tasks per worker, one worker can evaluate several trials at the same time.
SparkTrials takes two optional arguments.
**parallelism** is the maximum number of trials to evaluate at the same time. A higher value lets you test more hyperparameter settings at once. Because Hyperopt proposes new trials based on past results, there is a trade-off between parallelism and adaptivity. The concepts page gives the default as the number of Spark executors available and the maximum as 128. If you ask for more parallelism than the cluster configuration allows in concurrent tasks, SparkTrials lowers it to that limit.
**timeout** is the maximum number of seconds an fmin() call can take. When it is exceeded, all runs are terminated and fmin() exits. Information about completed runs is saved, so the trials that finished before the timeout are kept.
Two fmin() arguments interact with these. max_queue_len sets how many settings Hyperopt generates ahead of time. TPE generation can be slow, so raising it above the default of 1 can help, but it should generally be no larger than the parallelism setting. early_stop_fn can end the run before max_evals is reached. Under SparkTrials this function is polled, so it isn't guaranteed to run after every trial.
Checkpoint 3 of 8· Match them up
Match each setting to what it controls in a distributed Hyperopt run
Tap a term, then the definition that fits it.
parallelism and timeout are the two optional SparkTrials arguments. max_queue_len and early_stop_fn are fmin() arguments that interact with them.
“it can be helpful to increase this beyond the default value of 1, but generally no larger than the SparkTrials setting parallelism.”Source: docs.databricks.com
Checkpoint 4 of 8· Exam question
A team is tuning hyperparameters for an MLlib `RandomForestClassifier` that Spark already trains as a distributed algorithm across the cluster, spreading each training job over multiple executors. They ask whether wrapping the tuning loop in `SparkTrials` will speed up the search the same way it did for their earlier scikit-learn model. What should they be told?
Correct answer: A — Use a plain `Trials` object instead of `SparkTrials`, because MLlib already parallelizes the training of each individual model, and layering `SparkTrials` on top would compete for the same executor slots.
- A. MLlib already spreads a single model's training across the cluster's executors on its own. Adding SparkTrials on top would try to run multiple already-distributed training jobs concurrently, contending for the same executor resources, so the plain sequential Trials class is the correct fit here.
- B. The two-level distribution that helped the scikit-learn workload worked because that model trained on one machine and had spare cluster capacity to exploit. An MLlib model already consumes cluster resources to train itself, so stacking a higher parallelism value on top does not give the same benefit and instead risks resource contention.
- C. MLlib does not require a fixed ratio of hyperparameter combinations to executors, and max_evals is unrelated to worker count. This describes a constraint that does not exist in either Hyperopt or MLlib's training model.
- D. Converting the data to pandas would defeat the purpose of using MLlib, which is designed to train on distributed Spark DataFrames, and Hyperopt is not restricted to pandas-based objective functions in the first place.
4.SparkTrials or Trials: match the class to the model
SparkTrials is only for models that are *not* already distributed. When the objective function trains a model with a distributed algorithm such as Spark MLlib or Horovod, that training already uses the whole cluster. In that case, don't pass SparkTrials. Use Hyperopt's default Trials class, or leave out the trials argument altogether. Hyperopt then evaluates each trial on the driver, and from there the algorithm can start distributed training across the cluster's resources.
| What the objective function trains | trials argument | Where each trial runs | Automatic MLflow logging |
|---|---|---|---|
| Single-machine models such as scikit-learn | SparkTrials | In a single Spark task on a worker; many trials run at once | Yes, as parent and child runs |
| Distributed algorithms such as Spark MLlib or Horovod | Trials (the default), or omit trials | Launched from the driver, which gives the trial access to the full cluster | No, you must call MLflow manually |
Checkpoint 5 of 8· Check yourself
Your objective function trains a Spark MLlib model on a large DataFrame. A colleague suggests passing SparkTrials(parallelism=8) to fmin() to speed up tuning. What is correct?
SparkTrials is for algorithms that are not distributed themselves. For MLlib or Horovod, use Trials so that each trial runs from the driver and the algorithm handles the distribution.
“Use Trials when you call distributed training algorithms such as MLlib methods or Horovod in the objective function.”Source: docs.databricks.com
Checkpoint 6 of 8· Exam question
A data scientist is tuning a `scikit-learn` `LogisticRegression` model on a dataset with 2,000 rows. Each trial trains and scores in under half a second. After wrapping the search in `SparkTrials` with `parallelism=8`, the engineer notices the tuning job takes longer in wall-clock time than running the same search sequentially with a plain `Trials` object on the driver. What most likely explains this?
Correct answer: A — The objective function is so cheap that the overhead of launching and scheduling Spark tasks for each trial outweighs the time saved by running trials concurrently.
- A. For a fast-running objective function, the fixed cost of starting a Spark job and scheduling a task for each trial can exceed the trial's own training time. SparkTrials pays off on expensive objective functions; on cheap ones like sub-second logistic regression fits, the distribution overhead itself becomes the bottleneck.
- B. max_evals sets the total number of trials for both trials classes identically; SparkTrials does not silently evaluate more combinations than requested. The slowdown here comes from per-trial overhead, not from an inflated evaluation count.
- C. Neither trials class re-reads data from disk on every trial by default; how the objective function loads its data is up to the code the data scientist wrote, not a behavior SparkTrials imposes. This does not explain the wall-clock difference described.
- D. The parallelism setting controls how many trials run concurrently; it does not reset or discard the search algorithm's state between trials. TPE continues incorporating completed trial results regardless of the parallelism value chosen.
5.What SparkTrials logs to MLflow
Running many trials in parallel would be hard to review without tracking. Databricks Runtime ML supports logging to MLflow from workers, and SparkTrials records the tuning as nested runs:
- Main (parent) run: the fmin() call itself. If a run is already active, SparkTrials logs to it and leaves it open when fmin() returns. If no run is active, SparkTrials creates a run, logs to it, and ends it before fmin() returns.
- Child runs: each hyperparameter setting tested gets its own child run under the main run. Anything logged to MLflow from the workers is stored under the matching child run.
Databricks recommends wrapping each fmin() call in with mlflow.start_run():. That gives every tuning call its own main run and makes it easier to add extra tags, parameters or metrics to it. Inside the objective function, you don't manage runs yourself. A call such as mlflow.log_param there logs to the child run of that trial.
All three calls are logged to the same main run. When logged parameters or tags have conflicting names, MLflow appends a UUID to them. To get one main run per tuning call, give each fmin() its own mlflow.start_run() block.
Checkpoint 7 of 8· Check yourself
After an fmin() run with SparkTrials and max_evals=16, wrapped in mlflow.start_run(), how are the results organized in MLflow?
SparkTrials logs the fmin() call as the main run and each trial as a child run under it. Workers don't need to manage runs themselves.
“is logged as a child run under the main run.”Source: docs.databricks.com
To see the effect of C in the notebook example, select the resulting runs in the MLflow UI, click Compare, and plot C on the X-axis against loss on the Y-axis in the scatter plot.
Sources2
6.When distributing trials pays off, and Hyperopt's status
Distributing trials doesn't speed up every job. Each trial costs a Spark job plus Hyperopt bookkeeping. In the Iris example, the objective runs so quickly that the cost of starting Spark jobs dominates the calculation time, and the distributed version takes *longer* than the single-machine one. The best-practices page draws the general lesson: when trials take only a few tens of seconds, the overhead from Hyperopt and Spark can dominate, and the speedup may be small or even zero. For typical real-world problems, the objective function does more work, and distributing it with SparkTrials is faster than tuning on one machine.
Checkpoint 8 of 8· Check yourself
Each of your trials finishes in a few seconds. What should you expect from distributing them with SparkTrials?
Per-trial Spark and Hyperopt overhead is fixed, so for very short trials it can dominate and the speedup may be small or even zero.
“Both Hyperopt and Spark incur overhead that can dominate the trial duration for short trial runs (low tens of seconds).”Source: docs.databricks.com
Two more points help when you read results. The search algorithms are stochastic, so the loss usually doesn't fall steadily from one run to the next. A reported loss of NaN usually means the objective returned NaN for that setting. It doesn't affect other runs, and adjusting the search space or the objective can prevent it. Because SparkTrials runs trials on worker nodes, the data the objective function trains on has to reach those workers too.
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.SparkTrials is the right choice for any model you want to tune on a cluster, including Spark MLlib or Horovod models.Why is that wrong?
SparkTrials is for single-machine models such as scikit-learn. Distributed algorithms already use the whole cluster, so use the default Trials class, which runs each trial from the driver.
Covered in SparkTrials or Trials: match the class to the model
2.Setting parallelism as high as possible always gives the best hyperparameters, just faster.Why is that wrong?
For a fixed max_evals, higher parallelism means each new trial is proposed with fewer completed results, so an adaptive algorithm like TPE can produce worse results.
Covered in How SparkTrials schedules trials: parallelism, timeout and the cluster
3.Running SparkTrials on an autoscaling cluster lets the tuning job grow its parallelism as workers are added.Why is that wrong?
Hyperopt picks the parallelism value when execution begins and can't use workers the cluster adds later. Don't use SparkTrials on autoscaling clusters.
Covered in How SparkTrials schedules trials: parallelism, timeout and the cluster
4.Tuning with the default Trials class logs every trial to MLflow automatically, just as SparkTrials does.Why is that wrong?
Automatic MLflow logging comes with SparkTrials. With the Trials class, you must log trials to MLflow yourself.
Covered in SparkTrials or Trials: match the class to the model
5.Distributing trials with SparkTrials always finishes faster than single-machine tuning.Why is that wrong?
For short trials, Spark and Hyperopt overhead can dominate, and the speedup can be small, zero or even negative.
Covered in When distributing trials pays off, and Hyperopt's status
Practise it for real
Turn the single-machine SVC tuning run from the Databricks notebook into a distributed SparkTrials run and inspect its MLflow runs (needs Databricks Runtime ML 16.4 LTS or earlier, where Hyperopt is still included).
1.In a notebook on a non-autoscaling cluster, load the Iris data and define the objective(C) function that returns {'loss': -accuracy, 'status': STATUS_OK}.
Why: Hyperopt minimizes the loss, so a higher-is-better accuracy has to be negated.
You should see: Calling objective(1.0) directly returns a dictionary with a negative loss.
2.Set search_space = hp.lognormal('C', 0, 1.0) and algo=tpe.suggest, then run fmin(fn=objective, space=search_space, algo=algo, max_evals=16) without a trials argument.
Why: This gives you the single-machine baseline before you distribute anything.
You should see: fmin() prints a best value for C after 16 evaluations.
3.Create spark_trials = SparkTrials() and rerun the same fmin() call with trials=spark_trials inside a with mlflow.start_run(): block.
Why: SparkTrials evaluates the trials on workers, and the explicit run gives this fmin() call its own main MLflow run.
You should see: The 16 trials run as Spark jobs and fmin() returns a best value for C.
4.Open the notebook's Experiment panel, select the child runs, click Compare, and plot C on the X-axis against loss on the Y-axis.
Why: Each trial was logged as a child run under the main run, so you can compare hyperparameter settings side by side.
You should see: One parent run with 16 child runs, and a scatter plot of loss against C.
5.Compare the wall-clock time of the single-machine run with the SparkTrials run.
Why: The Iris objective is very fast, so Spark job overhead should dominate.
You should see: The distributed run is likely no faster, and may be slower, than the single-machine run.
Stuck? Get a nudge
If the distributed run is slower, that matches the notebook's own result for this objective. Distribution pays off when each trial does substantial work.
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.https://docs.databricks.com/aws/en/machine-learning/automl-hyperparam-tuning/hyperopt-spark-mlflow-integrationOfficial docs
“A higher accuracy value means a better model, so you must return the negative accuracy.”
↩︎ Start from a single-machine fmin() workflow“hyperopt.tpe.suggest: Tree of Parzen Estimators, a Bayesian approach which iteratively and adaptively selects new hyperparameter settings to explore based on past results”
↩︎ Start from a single-machine fmin() workflow“To distribute tuning, add one more argument to fmin(): a Trials class called SparkTrials.”
↩︎ Distribute the same run by adding SparkTrials“the overhead of starting the Spark jobs dominates the calculation time, so the calculations for the distributed case take more time.”
↩︎ When distributing trials pays off, and Hyperopt's status - 2.https://docs.databricks.com/aws/en/machine-learning/automl-hyperparam-tuning/hyperopt-conceptsOfficial docs
“Number of hyperparameter settings to try (the number of models to fit).”
↩︎ Start from a single-machine fmin() workflow“allows you to distribute a Hyperopt run without making other changes to your Hyperopt code”
↩︎ Distribute the same run by adding SparkTrials“Each trial is generated with a Spark job which has one task, and is evaluated in the task on a worker machine.”
↩︎ How SparkTrials schedules trials: parallelism, timeout and the cluster“Default: Number of Spark executors available. Maximum: 128.”
↩︎ How SparkTrials schedules trials: parallelism, timeout and the cluster“If the value is greater than the number of concurrent tasks allowed by the cluster configuration, SparkTrials reduces parallelism to this value.”
↩︎ How SparkTrials schedules trials: parallelism, timeout and the cluster“When this number is exceeded, all runs are terminated and fmin() exits. Information about completed runs is saved.”
↩︎ How SparkTrials schedules trials: parallelism, timeout and the cluster“When using SparkTrials, the early stopping function is not guaranteed to run after every trial, and is instead polled.”
↩︎ How SparkTrials schedules trials: parallelism, timeout and the cluster“Use SparkTrials when you call single-machine algorithms such as scikit-learn methods in the objective function.”
↩︎ SparkTrials or Trials: match the class to the model“If there is no active run, SparkTrials creates a new run, logs to it, and ends the run before fmin() returns.”
↩︎ What SparkTrials logs to MLflow“MLflow log records from workers are also stored under the corresponding child runs.”
↩︎ What SparkTrials logs to MLflow“wrap the call to fmin() inside a with mlflow.start_run(): statement.”
↩︎ What SparkTrials logs to MLflow“When you call fmin() multiple times within the same active MLflow run, MLflow logs those calls to the same main run.”
↩︎ What SparkTrials logs to MLflow“you do not need to manage runs explicitly in the objective function.”
↩︎ What SparkTrials logs to MLflow“Hyperopt is not included in Databricks Runtime for Machine Learning after 16.4 LTS ML.”
↩︎ When distributing trials pays off, and Hyperopt's status“Databricks recommends using either Optuna for single-node optimization or RayTune”
↩︎ When distributing trials pays off, and Hyperopt's status“SparkTrials accelerates single-machine tuning by distributing trials to Spark workers.”
↩︎ Key concept“For models created with distributed ML algorithms such as MLlib or Horovod, do not use SparkTrials.”
↩︎ Exam trap 1“Because Hyperopt proposes new trials based on past results, there is a trade-off between parallelism and adaptivity.”
↩︎ Exam trap 2“greater parallelism speeds up calculations, but lower parallelism may lead to better results since each iteration has access to more past results.”
↩︎ Prediction“it can be helpful to increase this beyond the default value of 1, but generally no larger than the SparkTrials setting parallelism.”
↩︎ Checkpoint“Use Trials when you call distributed training algorithms such as MLlib methods or Horovod in the objective function.”
↩︎ Checkpoint“is logged as a child run under the main run.”
↩︎ Checkpoint - 3.https://docs.databricks.com/aws/en/machine-learning/automl-hyperparam-tuning/hyperopt-best-practicesOfficial docs
“Do not use SparkTrials on autoscaling clusters.”
↩︎ How SparkTrials schedules trials: parallelism, timeout and the cluster“Both Hyperopt and Spark incur overhead that can dominate the trial duration for short trial runs (low tens of seconds).”
↩︎ When distributing trials pays off, and Hyperopt's status“Because Hyperopt uses stochastic search algorithms, the loss usually does not decrease monotonically with each run.”
↩︎ When distributing trials pays off, and Hyperopt's status“Hyperopt selects the parallelism value when execution begins.”
↩︎ Exam trap 3“The speedup you observe may be small or even zero.”
↩︎ Exam trap 5 - 4.https://docs.databricks.com/aws/en/machine-learning/automl-hyperparam-tuning/hyperopt-distributed-mlOfficial docs
“Each trial is executed from the driver node, giving it access to the full cluster resources.”
↩︎ SparkTrials or Trials: match the class to the model“Databricks does not support automatic logging to MLflow with the Trials class.”
↩︎ SparkTrials or Trials: match the class to the model“Databricks does not support automatic logging to MLflow with the Trials class.”
↩︎ Exam trap 4