Data Science

Understanding Lazy Evaluation in Polars

Understand what is eager and lazy execution and how you can use lazy execution to optimize your queries

Wei-Meng Lee
July 15, 20227 min read
Photo by Hans-Jurgen Mager on Unsplash
Photo by Hans-Jurgen Mager on Unsplash

In my previous article on Polars, I introduced you to the Polars DataFrame library that is much more efficient than the Pandas DataFrame.

Getting Started with the Polars DataFrame Library

In this article, I am going to dive deeper into what makes Polar so fast - lazy evaluation. You will learn the difference between eager execution and lazy evaluation/execution.

Implicit Lazy Evaluation

To understand the effectiveness of lazy evaluation, it is useful to compare with how things are done in Pandas.

For this exercise, I am going to use the flights.csv file located at https://www.kaggle.com/datasets/usdot/flight-delays. This dataset contains the flight delays and cancellation details for flights in the US in 2015. It was collected and published by the DOT's Bureau of Transportation Statistics.

Licensing - CC0: Public Domain (https://creativecommons.org/publicdomain/zero/1.0/).

We shall use Pandas to load the flights.csv file, which contains 5.8 million rows and 31 columns. Usually, if you load this on a machine with limited memory, Pandas will take a long time to load it into a dataframe (if it loads at all). The following code does the following:

  • Load the flights.csv file into a Pandas DataFrame

  • Filters the dataframe to look for those flights in December and whose origin airport is SEA and destination airport is DFW

  • Measure the amount of time needed to load the CSV file into a dataframe and the filtering

text
import pandas as pdimport time
text
start = time.time()
text
df = pd.read_csv('flights.csv')df = df[(df['MONTH'] == 12) &         (df['ORIGIN_AIRPORT'] == 'SEA') &        (df['DESTINATION_AIRPORT'] == 'DFW')]
text
end = time.time()print(end - start)df

On my M1 Mac with 8GB RAM, the above code snippet took about about 7.74 seconds to load and display the following result:

Image by author
Image by author

The main issue with Pandas is that you have to load all the rows of the dataset into the dataframe before you can do any filtering to remove all the unwanted rows.

While you can load the first or last n rows of a dataset into a Pandas DataFrame, to load specific rows (based on certain conditions) into a dataframe requires you to load the entire dataset before you can perform the necessary filtering.

Let's now use Polars and see if the loading time can be reduced. The following code snippet does the following:

  • Load the CSV file using the read_csv() method of the Polars library

  • Perform a filter using the filter() method and specify the conditions for retaining the rows that we want

text
import polars as plimport time
text
start = time.time()
text
df = pl.read_csv('flights.csv').filter(        (pl.col('MONTH') == 12) &         (pl.col('ORIGIN_AIRPORT') == 'SEA') &                                                  (pl.col('DESTINATION_AIRPORT') == 'DFW'))
text
end = time.time()print(end - start)display(df)

On my computer, the above code snippet took about 3.21 seconds, a vast improvements over Pandas:

Image by author
Image by author

Now, it is important to know that the read_csv() method uses eager execution mode, which means that it will straight-away load the entire dataset into the dataframe before it performs any filtering. In this aspect, this block of code that uses Polars is similar to that of that using Pandas. But you can already see that Polars is much faster than Pandas.

Notice here that the filter() method works on a Polars DataFrame object

The next improvement is to replace the read_csv() method with one that uses lazy execution - scan_csv(). The scan_csv() method delays execution until the collect() method is called. It analyzes all the queries right up till the collect() method and tries to optimize the operation. The following code snippet shows how to use the scan_csv() method together with the collect() method:

text
import polars as plimport time
text
start = time.time()df = pl.scan_csv('flights.csv').filter(        (pl.col('MONTH') == 12) &         (pl.col('ORIGIN_AIRPORT') == 'SEA') &        (pl.col('DESTINATION_AIRPORT') == 'DFW')).collect()end = time.time()print(end - start)display(df)

The scan_csv() method is known as an implicit lazy method, since by default it uses lazy evaluation. It is important to remember that the scan_csv() method does not return a DataFrame - it returns a LazyFrame instead.

For the above code snippet, instead of loading all the rows into the dataframe, Polars optimizes the query and loads only those rows satisfying the conditions in the filter() method. On my computer, the above code snippet took about 2.67 seconds, a further reduction in processing time compared to the previous code snippet.

Notice here that the filter() method works on a Polars LazyFrame object

Explicit Lazy Evaluation

Remember earlier on I mentioned that the read_csv() method uses eager execution mode? What if you want to use lazy execution mode on all its subsequent queries? Well, you can simply call the lazy() method on it and then end the entire expression using the collect() method, like this:

text
import polars as plimport time
text
start = time.time()df = pl.read_csv('flights.csv')       .lazy()       .filter(         (pl.col('MONTH') == 12) &          (pl.col('ORIGIN_AIRPORT') == 'SEA') &         (pl.col('DESTINATION_AIRPORT') == 'DFW')).collect()end = time.time()
text
print(end - start)display(df)

By using the lazy() method, you are instructing Polars to hold on the execution for subsequent queries and instead optimize all the queries right up to the collect() method. The collect() method starts the execution and collects the result into a dataframe. Essentially, this method instructs Polars to eagerly execute the query.

Understanding the LazyFrame object

Let's now break down a query and see how Polars actually works. First, let's use the scan_csv() method and see what it returns:

text
pl.scan_csv('titanic_train.csv')

Source of Data: The data source for this article is from https://www.kaggle.com/datasets/tedllh/titanic-train.

Licensing - Database Contents License (DbCL) v1.0 https://opendatacommons.org/licenses/dbcl/1-0/

The above statement returns a "polars.internals.lazy_frame.LazyFrame" object. In Jupyter Notebook, it will show the following execution graph (I will talk more about this as we go along):

Image by author
Image by author

The execution graph shows the sequence in which Polars will execute your query.

The LazyFrame object that is returned represents the query that you have formulated, but not yet executed. To execute the query, you need to use the collect() method:

text
pl.scan_csv('titanic_train.csv').collect()

You can also enclose the query using a pair of parentheses and assign it to a variable. To execute the query, you simply call the collect() method of the query, like this:

text
q = (    pl.scan_csv('titanic_train.csv')    )q.collect()

The advantage of enclosing your queries in a pair of parentheses is that it allows you to chain multiple queries and put them in separate lines, thereby greatly enhancing readability.

The above code snippet shows the following output:

Image by author
Image by author

For debugging purposes, sometimes it is useful to just return a few rows to examine the output, and so you can use the fetch() method to return the first n rows:

text
q.fetch(5)

The above statement returns the first five rows of the result:

You can chain the various methods in a single query:

text
q = (    pl.scan_csv('titanic_train.csv')    .select(['Survived','Age'])    .filter(        pl.col('Age') > 18    ))

The show_graph() method displays the execution graph that you have seen earlier, with a parameter to indicate if you want to see the optimized graph:

text
q.show_graph(optimized=True)

The above statement shows the following execution graph. You can see that the filtering based on the Age column is done together during the loading of the CSV file:

Image by author
Image by author

In contrast, let's see how the execution flow will look like if the queries are executed in eager mode (i.e. non-optimized):

text
q.show_graph(optimized=False)

As you can see from the output below, the CSV file is first loaded, followed by the selection of the two columns, and finally the filtering is performed:

Image by author
Image by author

To execute the query, call the collect() method:

text
q.collect()

The following output will be shown:

If you only want the first five rows, call the fetch() method:

text
q.fetch(5)
Image by author
Image by author

Join Medium with my referral link - Wei-Meng Lee

I will be running a workshop on Polars in the upcoming ML Conference (22–24 Nov 2022) in Singapore. If you want a jumpstart on the Polars DataFrame, register for my workshop at https://mlconference.ai/machine-learning-advanced-development/using-polars-for-data-analytics-workshop/.

Summary

I hope you now have a better idea of how lazy execution works in Polars and how to enable it even for queries that supports only eager execution. Displaying the execution graph makes it easier for you to understand how your queries are being optimized. In the next few articles, I will continue my discussion on the Polars DataFrame and the various ways to manipulate them. If there is any particular topic that you want me to focus on, leave me a comment!

Related Articles