Databricks Runtime 5.2 ML Features Multi-GPU Workflow, Pregel API, and Performant GraphFrames

We are excited to announce the release of Databricks Runtime 5.2 for Machine Learning. This release includes several new features and performance improvements to help developers easily use machine learning on the Databricks Unified Analytics Platform.

Continuing our efforts to make developers’ lives easy to build deep learning applications, this release includes the following features and improvements:

  • HorovodRunner includes a simplified workflow for multi-GPU machines and support for a return value.
  • GraphFrames introduces a Pregel-like API for bulk-synchronous message-passing using DataFrame operations, with performance optimizations on Databricks.
  • Clusters now start faster.

Using HorovodRunner for Distributed Training

In Databricks Runtime 5.0 ML, we introduced HorovodRunner, a new API for distributed deep learning training. In this release, we introduce two new features.

First, HorovodRunner provides built-in support for using nodes that each have multiple GPUs. On a GPU cluster, each Horovod process maps to one GPU on the cluster, and those processes are placed as groups on compute nodes. For example, if you run a job with np=7 processes on a GPU cluster with 4 GPUs on each node, then you will have 4 processes on the first node and 3 processes on the second node. This simplifies job setup while reducing the inter-task communication costs.

Second, the call can return the value from MPI process 0. This makes it easier for data scientists to fetch helpful results, such as training metrics or the trained model, as in the following code snippet.

def train():
  “““ Method for running training on each Horovod worker ”””

  model = get_model()

  # NEW: We use the return value for evaluation metrics.
  eval_results = model.evaluate(...)
  return eval_results

hr = HorovodRunner(np=8)
eval_results =

To learn about how to run distributed deep learning training on Databricks Runtime 5.2 ML, see the docs at Azure Databricks and AWS.

Pregel API in GraphFrames

GraphFrames is the open-source graph processing library built on top of Apache Spark DataFrames. In the latest release, GraphFrames exposes the Pregel API, which is a bulk-synchronous message-passing API for implementing iterative graph algorithms. For example, the snippet below runs PageRank.

val ranks = graph.pregel
  .withVertexColumn("rank", lit(1.0 / numVertices),
    coalesce(Pregel.msg, lit(0.0)) * (1.0 - alpha) + alpha / numVertices)
  .sendMsgToDst(Pregel.src("rank") / Pregel.src("outDegree"))

For more details, check the Scala API and Python API.

On Databricks Runtime 5.2 ML, we further improved the speed of Pregel implementation from open-source GraphFrames by up to 10x.

Performance improvements

Including PyTorch in Databricks Runtime 5.1 ML Beta increased cluster start times. In this release, we removed some duplicate libraries that helped lead to 25% faster start times.

Other Package Updates

We updated the following packages:

  • Horovod 0.15.0 to 0.15.2
  • TensorBoard 1.12.0 to 1.12.2


  • Read more about Databricks Runtime 5.2 ML Beta for Azure Databricks and AWS.
  • Try the example notebooks for distributed deep learning training for Azure Databricks and AWS on Databricks Runtime 5.2 ML Beta.