Scaling up GPU acceleration in the Polars dataframe library
Large scale analysis of structured data underpins much of enterprise decision making. These analyses are often formulated using dataframe libraries, of which Polars is a popular example. It has a domain specific language for writing dataframe queries, exposed as a lazy API in Python, a query optimiser and multiple different execution engines.
In this talk, I will cover recent work developing the multi-GPU accelerated engine in Polars. I will give an overview of the computational patterns that appear in large scale data analytics, where they present challenges to efficient execution and how we address them. I'll show how this multi-engine offering allows seamless scaling of data analyses from laptop, to workstation, and beyond, providing an "interactive" experience even at terabyte scale.
If you've ever written dataframe code and wondered "what is actually going on here?", this talk might be for you.
The core technical contribution to Polars is a new GPU runtime that uses cooperative multitasking, built from dataframe primitives provided by cuDF. Multi-GPU and multi-node execution is enabled by tying together query plan fragments with collective communication. Our approach takes many cues from SPMD, domain-decomposed programming models that are the norm in HPC applications and adapts the ideas to dataframe analyses where a static data decomposition is not a good fit since the sizes of intermediates are not just data size, but also value dependent.
This challenge of value-dependent data size is similar to the one that SQL query optimisers must address. In that case, there is prior art showing that when only optimising with metadata even the best tuned heuristics can be hilariously wrong a significant portion of the time. We still rely on the Polars optimiser to give us a good query plan, but by moving away from static decomposition, we are able to adapt the execution to the data and do not need to commit to a particular algorithm up front. For example, we can switch from a broadcast join to a hash-partitioned shuffle join if we observe tables are too large at runtime.
To assess the performance of our approach, I will show scaling results for dataframe adaptations of the TPC-H and TPC-DS database benchmarks.
Lawrence Mitchell works at NVIDIA. His focus is on high-productivity, high-performance libraries for data analytics. He leads the technical design and implementation of the cuDF-accelerated Polars GPU engine. Prior to joining NVIDIA he was a lecturer in computer science and applied mathematics at the University of Durham with research interests in high performance simulation of continuum mechanics, structure-preserving numerical methods, and preconditioning techniques for coupled multiphysics problems. He was a founding co-lead and technical architect of the open source Firedrake project for finite element simulation.