Pandas 2.0 and the Arrow revolution
datapythonista.me
datapythonista.me
pandas DataFrames are persisted in memory. The rule of thumb was for RAM capacity / dataset size in memory to be 5-10x as of 2017. Let's assume that pandas has made improvements and it's more like 2x now.
That means you can process datasets that take up 8GB of RAM in memory on a 16GB machine. But 8GB of RAM in memory is a lot different than what you'd expect with pandas.
pandas historically persisted string columns as objects, which was wildly inefficient. The new string[pyarrow] column type is around 3.5 times more efficient from what I've seen.
Let's say a pandas user can only process a string dataset that has 2GB of data on disk (8GB in memory) on their 16GB machine for a particular analysis. If their dataset grows to 3GB, then the analysis errors out with an out of memory exception.
Perhaps this user can now start processing string datasets up to 7GB (3.5 times bigger) with this more efficient string column type. This is a big deal for a lot of pandas users.
I am glad Dask exists, but it is not a seamless solution where you can just swap in a Dask dataframe.
In my case I run analysis on an Apple M1 with 16GB of RAM and my files on disk are thousands of parquet files that add up to hundreds of gigs.
Apart from that: Being able to go from duckdb to pandas very quickly to make operations that make more sense on the other end and come back while not having to change the format is super powerful. The author talks about this with Polars as an example.
I can't stress enough how much I think this is truly transformative. It's generally nice as a working pattern, but much more importantly it lets the scale of a problem that a tool needs to solve shrink dramatically. Pandas doesn't need to do everything, nor does DuckDB, nor does some niche tool designed to perform very specific forecasting - any can be slotted into an in memory set of processes with no overhead. This lowers the barrier to entry for new tools, so they should be quicker and easier to write for people with detailed knowledge just in their area.
It extends beyond this too, as you can then also get free data serialisation. I can read data from a file with duckdb, make a remote gRPC call to a flight endpoint written in a few lines of python that performs whatever on arrow data, returns arrow data that gets fed into something else... in a very easy fashion. I'm sure there's bumps and leaky abstractions if you do deep work here but I've absolutely got remote querying of dataframes & files working with a few lines of code, calling DuckDB on my local machine through ngrok from a colab instance.
Another way to solve this problem is using a Lakehouse storage format like Delta Lake so you only read in a fraction of the data to the pandas DataFrame (disclosure: I am on the Delta Lake team). I've blogged about this and think predicate pushdown filtering / Z ORDERING data is more straightforward that adding an entire new tech like Polars / DuckDB to the stack.
If you're using Polars of course, it's probably best to just keep on using it rather than switching to pandas. I suppose there are some instances when you need to switch to pandas (perhaps to access a library), but think it'll be better to just stick with the more efficient tech normally.
If there are efforts to help this be even better, I heartily welcome them.
This should be a step in the right direction, but it will probably still require manually specifying types for CSVs.
If Dang or any other mods see this, please correct the link.
<link href="/blog/" rel="canonical" />
They're probably also screwing up their SEO this way.Just beware that polars is not as mature, so take this into consideration if choosing it for your next project. It also currently lacks some of the more advanced data operations, but you can always convert back and forth to pandas for anything special (of course paying a price for the conversion).
As the blog mentions, once Pandas 2.0 is released, if you use Arrows types, converting between Pandas and Polars will thankfully be almost immediate (requiring only metadata adjustments).
I would also be curious about Numpy, since I know you can transparently map data to a Numpy array, but that's just "raw" fixed-width binary data and not something more structured like Arrow.
capacity_df - outage_df
prices.loc['2023-01'] *= 1
There’s many workflows and models that do thousands of these types of operations. prices.loc['2023-01'] *= 1
You can always do df.to_pandas() ... prices.loc['2023-01'] *= 1 ... from_pandas() :)More seriously, you are right, this is a tough one to do in polars. Polars seems to want to work with whole columns at a time, it doesn't give you write access to row sets.
There's also all types of annoying workarounds you have to do while tuples as indexes resulting from it converting to a MultiIndex. For example
srs = pd.Series({('a'):1,('b','c'):2})
is a len(2) Series. srs.loc[('b','c')] throws an error while srs.loc[('a')] and srs.loc[[('b','c')]] do not. Not to vent my frustrations, but this maybe gives an idea of why this change is important and I very much look forward to improvements in the area!
some_pandas_object.values
to get to the raw numpy, because often dealing with raw np buffer is more efficient or ergonomic. Hopefully losing numpy foundations will not affect (the efficiency of) code which does this.https://pandas.pydata.org/docs/dev/reference/api/pandas.Data...
If I can't rely on them in my packages or long-running projects, then I don't see the point in learning them to use in interactive work.
You are entitled to your opinion (which I see in every thread that discusses the tidyverse), but in my opinion this is a considerable and outdated exaggeration which will mislead the less experienced. Let's balance it out with some different perspective.
1. The core of the tidyverse, the data manipulation package dplyr, reached version 1.0 in May 2020 and made a promise to keep the API stable from that point on. To the best of my knowledge, they have done so, and anyone who wants to verify this can look at the changelogs. That's nearly 3 years of stability.
2. For several years, functions in the tidyverse have had a lifecycle indicator that appears in the documentation. It tells you if the function has reached a mature state or is still experimental. To the best of my knowledge, they have kept the promises from the lifecycle indicator.
3. I have been a full-time R and tidyverse user since dplyr was first released in 2014, and my personal experience is consistent with the two observations above. I agree with the parent commenter that the tidyverse API used to be unstable, but this was mainly in 2019 or earlier, before dplyr went to 1.0. And even back then, they were always honest about when the API might change. So now that the tidyverse maintainers are saying dplyr and other tidyverse packages are stable, I see no rational basis to doubt them.
4. Finally, even during the early unstable API period of the tidyverse, I personally did not find it such a great burden to upgrade my code as the tidyverse improved. It was actually quite thrilling to watch Hadley's vision develop and incrementally learn and use the concepts as he built them out. To use the tidyverse is to be part of something greater, part of the future, part of a new way of thinking that makes you a better data analyst.
IMHO, the functionality and ergonomics of the tidyverse are light-years ahead of any other* data frame package, and anyone who doesn't try it because of some negative anecdotes is missing out.
*No argument from me if you prefer data.table. It's got some performance advantages on large data and a different philosophy that may appeal more to some. Financial time series folks often prefer it. YMMV.
Wow.
I think it's totally fine to use base R or data.table or whatever else you like. There is no one right way, and I have used all of these and more in different contexts. But if people are getting impressions of pros and cons from the discussion here on HN, they should be aware that claims of API instability are several years out of date. It would be a shame if people were scared about instability that isn't there.
This never happens. Angular does not become React, React gets created instead. CoffeScript does not become TypeScript, etc.
There is too much baggage and backward compatibility that prevents such radical transformations.
Instead a new thing is created.
https://pandas.pydata.org/docs/dev/whatsnew/v2.0.0.html#copy...
This is doomed then. Pandas API is already extremely bloated
I admit that the API has issues (if/else? being the most glaring to me), notwithstanding Pandas has mass adoption because the benefits outweigh the warts.
(I happen to wish that 2.0 deprecated some of the API, but Python 3 burned a deep scar that many don't wish to relive.)
dataframe.column
vs dataframe['column']
as one example comes to mind but there is surely much moreI am of the philosophy of 'The Zen of Python'
There should be one-- and preferably only one --obvious way to do it.
Pandas is a powerful library, but when I have to use it in a workplace it usually gives me a feeling of dread, knowing I am soon to face various hacks and dataframes full of NaNs without them being handled properly, etc.I would get rid of the .column accessor, but you will see a lot of pushback. Notably from the R camp.
I think, in general, that Python has a lot of surprising stuff going on, but my standard for "simple" is Scheme.
Some other key differentiatiors:
- multi-threaded: almost all operations are multi-threaded and share a single threadpool that has low contention (not multiprocessing!). Polars often is able to completely saturate all your CPU cores with useful work.
- out-of-core: polars can process datasets much larger than RAM.
- lazy: polars optimizes your queries and materializes much less data.
- completely written in rust: polars controls every performance critical operation and doesn't have to defer to third parties, this allows it to have tight control over performance and memory.
- zero-required dependencies. This greatly reduces latency. A pandas import takes >500ms, a polars import ~70/80ms.
- declarative and strict API: polars doesn't adhere to the pandas API because we think it is suboptimal for a performant OLAP library.
Polars will remain a much faster and more memory efficient alternative.
However I am curious on how Arrow beats NumPy on regular ints and floats.
For the last 10 years I've been under impression that int and float columns in Pandas are basically NumPy ndarrays with extra methods.
Then NumPy ndarrays are basically C arrays with well defined vector operations which are often trivially parallelizable.
So how does Arrow beat Numpy when calculating
mean (int64) 2.03 ms 1.11 ms 1.8x
mean (float64) 3.56 ms 1.73 ms 2.1x
What is the trick?Maybe some optimization on multithreading
Still I am reasonably sure Numpy utilizes multithreading already. https://in.nau.edu/arc/parallelism-in-python
And here we assume there are no NaNs.
So something else.
Context: https://lists.apache.org/thread/sdxr8b0lj82zd0ql7zhk9472opq3...
Also see ADBC, which aims to provide a unified API on top of various Arrow-native and non-Arrow-native database APIs: https://arrow.apache.org/blog/2023/01/05/introducing-arrow-a...
(I am an Arrow contributor.)
pandas is just for 2D columnar stuff; it's sugar on numpy