Skip to main content

Command Palette

Search for a command to run...

So your Databricks notebook is painfully slow?

Updated
•10 min read•View as Markdown
So your Databricks notebook is painfully slow?
J

I am Data Engineer with a big passion for learning as much as he can. I enjoy the outdoors, mountain biking, finding cool ways to solve new coding problems, and teaching others to code.

So you've been working on a Databricks notebook, and you've built most of your transformations, but now when you try to run the notebook it takes FOREVER to run?

Well, you're not alone in this experience. As someone who's worked mostly in standard OLAP or OLTP databases, the learning curve for building a performant Databricks notebook was HUGE. I had always heard about the speed and power of Spark, namely when it came to large datasets; clusters and all that. Much to my surprise, it wasn't quite what I expected.

Jobs that would typically take seconds on a standard SQL server, would take minutes. Even more complex jobs would take 30+ minutes (or even longer). I had mostly just accepted the performance, I would either have to accept the execution times or throw more compute at it.

Not until I had to dive deeper into Databricks, did I realize I should have spent more time learning different optimization techniques to improve my notebook performance.

Not Your Typical SQL Database

The long short is that Databricks isn't your typical SQL database. Meaning, a lot of the optimizations that are innate to most standard SQL servers, that we often take for granted, aren't present in Spark. It's because Spark isn't... a database. Seems pretty obvious when you say it out loud, doesn't it?

Sure, there are some internal optimizations with Databricks and Spark, and you can always throw more compute at something, but to get performant code you really have to put some thought into optimizing it.

On top of that, the initial training I had as well as the documentation I read did not provide much in terms of guidance for writing performant Databricks notebooks. There were some more general recommendations, but nothing particularly explicit. As well, I made the assumption Spark would work just like other databases, and all the typical techniques for optimizations would work.

Optimizations

So, what are some things I would've known when I started to help with optimizing my pipelines? While this is nowhere near an exhaustive list, below are some ways to optimize your notebook performance.

Choose the Right Compute

This one is probably the most integral to performance. If you don't choose the right compute, your notebook will run slow. Seems obvious right?

Unfortunately, choosing the right compute isn't always intuitive. So, I recommend trying multiple configurations to see what works best for your particular notebook.

As well, you should be aware that the bigger the compute or the more nodes in the cluster, the more expensive they become. However, the bigger compute or the more nodes, generally the faster things will run. So, it's just important to maintain a balance between cost, speed, and overall performance.

There are a variety of things to consider when choosing different compute. However, probably the biggest determinant is the type of work.

Type of Work

Generally, the type of work the cluster will do will impact which one you choose. You can refer to Azure Databricks Cluster Guidance for more general guidance.

Basic ETL

With basic ETL, about any cluster configuration will work, as there aren't generally a lot of joins or transformations. Just consider the amount of data when choosing the size of your compute.

Complex ETL or Analytics

With complex ETL or analytics jobs (which are probably the most common type of job), there are generally a lot of transformations, joins, or aggregates. With these types of jobs, ironically enough single-node setups (i.e., a driver only) or fewer-node clusters are generally the best.

You might be saying "Well, Spark is all about clusters?" Well, in some ways you're right. However, there's an innate issue with clusters, and it's something called "shuffling".

Shuffling is what happens when data is moved across nodes in a cluster. Shuffling can be incredibly expensive because data has to move from cluster to cluster. This results in both compute and network time, which can result in slow execution. As you can imagine, shuffling can be a real headache and slow your job down.

The reason a single-node is best to deal with shuffling, is well, because there's only one machine. So the data doesn't get shuffled between nodes. It's for this reason, for most kinds of work, a single-node or fewer-node cluster is likely going to be the most ideal setup.

If you do go with a single-node setup, just be sure to consider the following:

  • Memory Intensive Jobs

Large joins, transformations, window functions, or other memory-intensive operations can be problematic for a single-node. You can use memory-optimized compute (i.e., more ram), however, you might run into driver failure issues if you are maximizing resources on the node. So if you have regular failures, either size up your compute, or create a cluster.

  • Data Size

Goes without saying, but the more data, the more issues a single-node setup will have. When data size grows this is where the real power of Spark and clusters comes to bear. So, if you've reached a max in compute size, create a cluster. Just be mindful the more nodes, the more places data can be shuffled.

Multi-Use Clusters

This isn't as much a type of work, but more about how many people are utilizing the cluster. If multiple people will be utilizing the same compute, create a cluster. This helps with the typical high concurrency needs of multiple people using the same cluster.

Consider what kind of notebooks they are writing and how much data they will be analyzing, and scale up the compute in the cluster as required. Also, when creating a Databricks cluster, there is an option for session isolation. This is important as it keeps sessions separate if the same people are utilizing the same cluster.

Other Compute Optimizations

There are other optimization options when choosing your compute.

Cache Accelerated

I will get into caching in the next section, however, Databricks has a cache optimized compute choice. These machines are optimized for reading and writing to the cache, which can speed up your code tremendously.

Photon

Databricks has an optimized runtime called Photon. Long short, it speeds up operations tremendously. It does increase the cost of execution, but in turn, it will speed up execution - so there is a balance to be had between execution time and cost.

Autoscaling

Autoscaling is an option when configuring a cluster. It allows nodes to be added or removed as the load shifts when running a notebook or job. It allows the cluster to size up when transformations become more demanding, and drop them off when they aren't needed. This is also beneficial for cost control.

Caching

If I had any one thing I wish I knew about sooner, it would have been caching. It's fundamentally changed how I coded in Databricks. Caching is a mechanism for storing data locally on the clusters.

If you know about Databricks and Spark, you know why this is valuable. Network time (i.e., reading data from a distributed file system) can kill speed, and caching can reduce network time. It becomes more valuable the more you utilize a dataset in your notebook. As well, you can hash parts (columns or rows) of tables for smaller transformations.

You can cache a dataframe, a view, a table, or a subset (select) of a table.

#creating a dataframe from a table
df = spark.sql("select * from schema.table")
#caching the dataframe
df.cache()

#creating a view
df.createOrReplaceTempView("some_view")
spark.sql("cache table some_view")

#caching a select
spark.sql("cache select column1, column2 from schema.table where some_column > 1")

As well, you can enable caching on Databricks. This will cache data automatically when it is first read. spark.conf.set("spark.databricks.io.cache.enabled","true").

Reduce Column or Row Count

It goes without saying, but the less data you have to work with, the faster the process will be. This goes for both the length and the width of the data you are working with. So if you only need two columns, I would recommend limiting the data to only the two columns you are utilizing.

Namely, limiting the number of columns is where you will see the biggest improvements. Delta tables (and parquet files) are columnar stores of data, and they are optimized for reading lots of rows of data, quickly. That being said, the more columns that are added, performance takes a hit.

This also works incredibly well with caching. If you only have to use two or three columns, try caching that data set to the cluster for subsequent transformations.

Avoid UDFs

UDF's, or better known as user-defined functions, are ways to perform custom transformations on your data. The unfortunate thing, is they are generally not natively optimized.

There are a variety of ways to optimize UDF's, such as higher-order functions or pandas user-defined functions. However, I would recommend attempting that transformation utilizing native functionality, as any built in function or transformation will generally be optimized from the get-go.

That being said, Photon vectorizes UDF's. So, if you have to utilize a UDF, either try to optimize it or utilize Photon. Optimization work will only pay off, so don't just rely on Photon if there's a way to develop the UDF in a better way.

"Hard Write"

By "hard write" I simply mean writing the table out to a distributed file store of your choice. It's generally helpful if you have a lot of transformations to perform, as it allows you to write out intermediate steps. As well, writing out the table periodically while developing can help you recover and start over if a session is lost, or you need to step away.

This is one you should consider using more infrequently, as it requires both compute time and network time, so it can be very expensive. That being said, it can be very helpful in certain situations.

Sorting

Sorting your data can be tremendously beneficial. It may seem arbitrary, but if you end up joining the data later, sorting can speed up searches for matches.

As well, when you perform window or analytical functions such as row_number(), dense_rank(), or any other partitioned function, sorting the data based on your partition fields can speed up execution time significantly.

This works very well when caching, so if you have to calculate a window function utilizing only a handful of columns, pulling those into memory and sorting them, then executing the window function can speed up execution time tremendously.

Other Optimizations

These aren't explicitly technical, but these are some other recommendations for optimizing notebooks.

Split Your Code into Reasonable Steps

Always split your code into reasonable steps. This helps with both readability and performance. If you have complex calculations you need to perform, cache the required columns and data, perform the calculation, then join/add it back to the final table you are working on. This will be way more performant than attempting to calculate everything at once.

Avoid Iterative Views

Spark is based on views. However, just like normal SQL views, they can become problematic when they are iteratively built on top of each other. Especially when they contain joins or complex transformations. Try to limit building views on top of views. Try caching iterative tables or views to pull them into memory to speed up later transformations or joins. You can also perform a "hard write" on a view if that makes sense given your use case.

Optimizing Joins

Joining utilizing certain field types is more performant than others. For example, joining on an integer field versus string field will generally be more performant.

So, if you have a long string you are attempting to utilize in a join, try finding a way to optimize that join such as hashing into an integer field utilizing hash() or generating a row_number() to provide a unique identifier per record or key. However, take note that with hash(), collisions can happen, so you will still likely have to join on the existing string field.

Experiment

I think this is probably the most important thing I can recommend when trying to optimize anything. Try different ways of approaching the same problem, and see if you can speed things up. Do research, google, creep on Stack Overflow, and/or ask ChatGPT; you might hit a wall for a while, but you'll have a breakthrough eventually. It's also tremendously helpful as you will learn the best ways to do something natively.

Conclusion

Databricks and Spark are powerful tools for working with big data. However, building performant notebooks isn't as simple as it seems. There is no one size fits all of optimization, but hopefully, you can utilize these techniques to help improve the runtime of your notebooks.

There are also plenty of other techniques not covered here, so feel free to start with Databricks documentation for optimization.

Thanks for reading! I hope you learned something!

More from this blog

J

JaggedArray

16 posts

I'm a constantly curious Data Engineer, who loves nothing more than to learn new things, and help others learn new things by making coding approachable.