Apache Arrow and the “Things I Hate About Pandas”
wesmckinney.com
wesmckinney.com
All that being said, I'd stress pretty clearly that I never let a single line of pandas into production. There are a few reasons that I've long wanted to summarise, but just real quick: 1) It's a heavy dependency and things can go wrong, 2) It can act in unexpected ways - throw in an empty value in a list of integers and you suddenly get floats (I know why, but still), or increase the number of rows beyond a certain threshold and type inference works differently. 3) It can be very slow, especially if your workflow is write heavy (at the same time it's blazing fast for reads and joins in most cases, thanks to its columnar data structure). 4) The API evolves and breaking changes are not infrequent - that's a great thing for exploratory work, but not when you want to update libs on your production.
pandas is an amazing library, the best at exploratory work, bar none. But I would not let it power some unsupervised service.
Instead, I wrote a small wrapper around numpy to provide a data frame like object (850 lines of code by sloccount). So far, this has worked well for us.
2) Depends a bit on your background, but to me this is not really unexpected. Integers don't have a well-defined "missing" value while Floats do, so Pandas is trying to help you by not using python objects and instead converting to the "most useful" array type. It only does so if it can convert the integers without loss of precision.
3) This one I totally get, I wrote a custom, msgpack-based serialisation due to that for our usage (before Arrow was around, seriously considering that for data exchange now).
4) Apart from the changes to `resample` all of those breaking changes had a prior `DeprecationWarning`, IIRC.
In the end, I implemented a format that is essentially a list of columns, each having a specific type and some other metadata. The format for each column depends on the stored type, fallback is a msgpack sequence of strings/ints/dates mixed with nils, but array data is stored as-is and datetimes are stored either as an array of int64 with the unit attached or as begin, end, frequency.
Then I found plasma and that has blown my mind.
@wesm, how hard have you pushed plasma?
For everyone else: http://arrow.apache.org/blog/2017/08/08/plasma-in-memory-obj...
Plasma, an in-memory object store that is being developed as part of Apache Arrow. Plasma holds immutable objects in shared memory so that they can be accessed efficiently by many clients across process boundaries.
...
One of the goals of Apache Arrow is to serve as a common data layer enabling zero-copy data exchange between multiple frameworks. A key component of this vision is the use of off-heap memory management (via Plasma) for storing and sharing Arrow-serialized objects between applications.
I will add ... "in python." That is definitely true, but to call it "the best at exploratory work" is not accurate. I might be opening up a completely separate debate, but for down and dirty exploratory work nothing beats R's dplyr and ggplot.
With that being said, I now do most of my work in python because of putting models into production. I also haven't had any issues with pandas in production; maybe because I'm not doing high throughput operations and our ML application is relatively lightweight.
1. the weird handling of types and null values (#4) 2. the verbosity of filtering like `dataframe[dataframe.column == x]` and transformations like `dataframe.col_a - dataframe.col_b`, compared to `dplyr` in R 3. warts on the indexing system (including MultiIndex, which is very powerful but confusing)
For those of us who use Pandas as an alternative to R, these usability shortcomings matter way more than memory efficiency.
https://github.com/deeplearning4j/datavec
https://deeplearning4j.org/datavec
It vectorizes/tensorizes most major data types to put them in shape for machine learning. It also lets you save the data pipeline as a reusable object.
There is dplython, but it doesn't quite work the same so I don't use it much. https://github.com/dodger487/dplython
df %>% select(var1, var2) %>% rbind(df2) %>% na.omit()
etc? That was the big benefit I saw from using dplyr.`
df >> call(pd.dropna, axis=1)
Strings are a killer -- indeed any variable-length object makes array programming tricky when it's nested, so a sound strategy is to intern your strings first (in some way; KDB has enumerations, but language support isn't necessary: hash the strings and save an inverted index works good enough for a lot of applications). Interning strings means you see integers in your data operations, which is about as un-fun to program in as it sounds. People want to be able to write something like:
….str.extract('([ab])(\d)', expand=False)
and then get disappointed that it's slow. Everything is slow when you do it a few trillion times, but slow things are really slow when you do them a few trillion times.If we think about how we build our tables, we can store these as a single-byte column (or even a bitmask) and an int (or long) column, then we get fast again.
However it is clear "fast" and "let's use JSON" are incompatible, and a good middleware or storage system isn't going to make me trade.
As far as I understand you want to handle nested arrays of strings in your data. Ok
The "right" way is to build an index of the strings we are storing and then store the index values (hashes of some kind) in the arrays as longs
This way our arrays are doing numbers and we handwavy search for or use strings through some wrapper
Is this right?
And I am guessing the middleware you want does this transparently? Maybe storing the index alongside the data in some fashion
I don't think the rule is that firm:
You don't have to use longs because if there are only 256 unique values, why waste so many bytes?
Meanwhile if there are so many unique values, is longs enough?
> And I am guessing the middleware you want does this transparently?
Maybe.
I personally think explicit is fine, as long as it's the most easiest and obvious way to do it.
But I get that there's a lot of "data scientists" (probably the gross majority) that really struggle with (what I think are) basic data structures -- they'll rarely produce an efficient solution, and we'll see "yet another" article about how awk+sort is faster than a 160 node hadoop cluster...
> Maybe storing the index alongside the data in some fashion
Maybe not. What works for 30k strings falls over with 300m unique strings.
The real trick with nested objects is to invert the query -- to convert the select statements you want to run, into an insert/upsert statement you run when you're loading the data.
More SQL-style indexing would be a lot more intuitive at least for me.
However I do prefer the R data.table model, which is what you descibe. You can set an index on one or more columns in the table, and that's that.
I use Oracle APEX because it has a killer "interactive report" feature (ie a data grid on steroids), which enables non-programers to easily filter, aggregate,export,report on, etc, the data. However, although APEX is a free option that comes with the DB, it ties you to Oracle.
It would be great if there was a similar, database independent, low-code tool like APEX out there, so am curious what you have seen to work well.
In industry, does Pandas tend to power the application layer, or does it find more use as an exploratory data tool?
If the latter, do people prefer to push computation down into OLAP databases for performance reasons?
And if so, what impact will the convergence of libraries and database functionality have on product development? These features strike me as things that you'd find in a database, e.g. query optimizer. I know in the past couple years there have been a couple commercial acquisitions of in-memory execution engines, e.g. Hyper by Tableau.
E: that's not to say pandas isn't good. It's really good. Thanks for the software, Wes!
My experience echos yours, Pandas from my observation, is more like a post-modeling tool, that people use to further process data that they digest from certain DB query or Spark jobs.
After reading through Arrow homepage, I am left somewhat baffled about where it seats. If my reading is correct, it is a client-side protocol that abstracts away the underlying data storage implementations? If so, isn't it still limited by how much data the client machine can handle? Or the benefits is about the unified interface of accessing different storage system? No matter what, it seems pretty ambitious. Looking forward to see how it goes.
Data processing systems need runtime memory formats. Arrow is an efficient one for analytical data processing. It has the additional benefit of zero-copy data interchange layer for sharing memory between processes written in any language.
The existence of other database systems that perform equivalent tasks isn't useful if they are not accessible to Python programmers with a convenient API.
I pull the relevant data out of a production database, clean it, add relevant columns, filter out trash data, use seaborn to produce some simple plots to see what my data approximately looks/structured like then off to sklearn.
At the moment the general work flow is:
* Internal library based over Pandas which abstracts our mess of internal databases
* Application specific model code that utilises the internal library to pull data in. This is then fed into a trained scikit-learn model and then further processed by Pandas.
* Internal monitoring tools (dashboards based upon Ploty and Flask as well as an alerting system) are built using the internal library and Pandas as the glue.
From a design decision we focused upon Pandas as the root source of all data. Everything is a DataFrame throughout the entire application.
Painpoints:
* Writing to a database is pretty painful (SQL Server here as Windows shop).
* Minor API changes can be irritating.
* Pandas MultiIndexing is both very painful and mind bending at the same time trying to get the slice syntax to work.
Overall though, Pandas is a huge value add and we've gradually rolled out from 2 people to approximately 9-10 people who hadn't used python in anger before.
Almost all reporting functionality is being migrated into Pandas instead of SQL stored procs, excel, tableau etc for the additional flexibility it provides.
What stopped you from contributing improvements to pandas? Have you taken alternate routes to open source your work?
And I get that I am being a little harsher than reality dictates. However, the testing and "build" process that surrounds most of "notebooks" is laughably like what we specifically avoided in software when we said your build should be standardized in an external file. And not scripting the main IDE that you happen to be using.
Indeed, I am perplexed by folks that don't know how to move between IDEs or who won't bother to understand how they are pulling dependencies into their system. Notebooks, though, seem to embrace that.
Which, as I've indicated elsewhere, is great for interactive use, but seems a major step backwards for serious solutions.
I don't know why that bothers me, but it definitely does.
It bothers me, and I can tell you why! Between pandas and matplotlib, the royal road to liberation from routine analysis tasks is paved with Python scripts. Jupyter has an important ancillary role in aiding discovery. But this whole notion of "executable notebooks" seems designed to keep people in bondage to fragile workflows based on capturing and replaying user input. It caters to the least common denominator, to the one person on the team who can't be trusted to read things. I'm infuriated on behalf of anybody subjected to such foolishness.
http://jupyter-dashboards-layout.readthedocs.io/en/latest/us...
> Alice is a Jupyter Notebook user. Alice prototypes data access, modeling, plotting, interactivity, etc. in a notebook. Now Alice needs to deliver a dynamic dashboard for non-notebook users. Today, Alice must step outside Jupyter Notebook and build a separate web application. Alice cannot directly transform her notebook into a secure, standalone dashboard application.
I find this pretty unconvincing. The gap between "stuff I did in a notebook" and a secure, let alone correct, application is nontrivial. There's no way for Alice to do this without learning to write computer programs for real. And if she does that, she'll find that it's a lot easier when you don't pull in a huge dependency like Jupyter.
import seaborn as sns
sns.violinplot(data=dataframe)
Set %matplotlib inline and you don't need more commands in your notebook.
https://github.com/jupyter/enhancement-proposals/blob/master...
Also, there's a world of difference between this and the way I'd recommend people use Jupyter. Jupyter is great for exploratory data analysis. It caputres every step you take along the way, especially if you're disciplined about not reusing cells. At the end, you have something you can paste into Powerpoint. If the boss asks you for that same plot the next day, you don't reuse the notebook -- that means repeating all of your mistakes. You pull out the parts of the analysis you want to keep into a Python script, and you run that in the future. In no way is it a good idea to try to use the notebook operationally.
The last time I was trying to use pandas, it was the hackernews data dump. It wasn't big. However when pandas started using the memory, my 32GB was just too little.
I just ended to convert the data within postgres, much faster, with sensible memory usage.
https://stackoverflow.com/questions/25962114/how-to-read-a-6...
I have loaded that csv file to postgres, the database had similar size. With indexing it was 15GB on disk. All queries were quite fast.
So instead of loading the data to pandas, and making searches, I just wrote some SQL, and got the same results. However much faster, and with much smaller memory usage.
I recently moved some data processing from Python/pandas into a database, and with SQL, the processing time went from several minutes to a couple seconds, (and that on a tiny VM).
I understand that not everyone is familiar with databases and SQL, and so default to the toolset they know. But, the performance gains can make learning databases and SQL highly worthwhile. (And, much can be learned in a just few days, especially for those already familiar with working with data.)
From my point of view, the SQL database can store the huge dataset, with its changes, and I can iterate through the results, and make lots of nice queries getting only the data I need.
And not all operations can be done incrementally by iterating through the results.
It doesn't. There's a binary version of the protocol. The output conversion for that is near trivial (transformation to big endian).
Wow. Didn't knew that.. strings is the one thing I always minimize in my datasets, due to speed and memory considerations..
BTW, I wanted to thank you for your 2012 slides on how you used hashes to group and join data. It led me to learn more about categoricals and I ended up implementing a Factor() object [1] in the other tool I use (Stata) that ended up being a life saver. In fact, once you have a powerful and fast categorical type, with a set of key functions, you can do anything from group the data, to count distinct categories, to run fixed effect regressions in no time.
[1] http://fmwww.bc.edu/repec/scon2017/Baltimore17_Correia.pdf
factor(x = character(), levels, labels = levels, exclude = NA, ordered = is.ordered(x), nmax = NA)
[...]
levels
: an optional vector of the values (as character strings) that x might
have taken. The default is the unique set of values taken by
as.character(x), sorted into increasing order of x. Note that this set
can be specified as smaller than sort(unique(x)).
labels
: either an optional character vector of (unique) labels for the
levels (in the same order as levels after removing those in exclude), or
a character string of length 1.
This way, you can do something like that: > x <- 1:3
> factor(x, levels = 1:2, labels = c("foo", "bar"))
[1] foo bar <NA>
Levels: foo bar
But this actually is: > factor(as.character(x), levels = c("1", "2"), labels = c("foo", "bar"))
[1] foo bar <NA>
Levels: foo barMy impression is that factors in R are borderline-deprecated, especially in the tidyverse, in favor of just using the equivalent non-factor vector.
Could you please provide a real world example where this actually is a problem.
Point 10: R has lazy evaluation, which means here that a function will not be evaluated when you define it, it will be evaluated when you call it (maybe not quite the same as some other language's lazy eval). I'm not aware of any built in feature for query planning, if you ask for nrow(some_func(myframe)), it will evaluate the some_func(myframe) function and then count up the rows. You could always write your own query planning function I suppose.
Point 11: R has several multicore/cluster libraries, and some are actually decent. If you are like most R users, you use StackOverflow a lot, and you'll end up with one algorithm that uses the snow package and another algorithm that uses multicore, one that uses parallel, and so on. A few very well written packages have hooks that make going to multiple cores easy, but most do not and you typically have to roll your own.
Out of SAS, R, and Python or even C++ and Java. Data frame are native.
So personally, the syntax especially for data frame is beautiful compares to python.
Python don't even have missing value built in or subsetting data frame.
It seems unlikely you meant this as stated. How is it possible to "evaluate a function" when you define it? Certainly you need to give it arguments, and that can only happen when you call it.
I can pass an argument but R won’t try to evaluate unless it needs it. This can be beneficial when only some of a function’s branches need the argument. You can pass a solver for the traveling salesman problem but R won’t waste CPU cycles until it reaches a point where it has to solve the TSP to get an answer. Maybe the first branch of the function is a feasibility check, and the TSP will be skipped for something else.
This is a little harder to explain for data frames, but you can create R functions that act a little like generators in python. This can help with memory management where instead of a gigantic matrix you have a function that generates the part of the matrix that you need.
Welcome to R!
If you look at attempts to do this stuff in python---e.g. patsy, which emulates R's formula DSL, and there's another project that emulates dplyr I don't recall the name of---you see they have to resort to parsing and eval'ing strings instead of working on expressions (language objects that represent ASTs), which is not nearly as nice or safe.
Edit: But just to emphasize your surprise -- yes, you can definitely be surprised by delayed evaluation in many contexts if you're used to more traditional languages.
> y <- 10
> wat <- function(x=10*y) { y = -y; x }
> wat()
[1] -1000
But (1) good library writers don't play these kinds of tricks, so it doesn't come up too often in practice; and (2) when writing/debugging my own code, I've not found it too hard to reason about, anticipate, and avoid these effects. > y <- 10
> less_wat <- function(x=10*y) { force(x); y = -y; x }
> less_wat()
[1] 1000The data.table package (https://github.com/Rdatatable/data.table/wiki) does make progress on some of these - I'd say #1, #3, maybe #7, #8. Dplyr has a query planner too, fwiw.
Just like Feather is built on top of Arrow, ONNX can be based on top of Arrow.
Well it's good to see that open source works and competitors can benefit from each other's work.
> A multicore schedular for parallel evaluation of operator graphs
Does anything like this already exist somewhere?
Ibis (also by Wes McKinney) does the first part, but it offloads scheduling and execution to the underlying database you are using.
Is there a list of major projects that are leveraging Apache Arrow?
There are probably folks using it for exploration but probably not a production system.
If someone is using Pandas, Spark, or whatever for an important product, it's probably best for them to maintain whatever underlying data layer until the Arrow devs (I guess that means you) are willing to commit to a somewhat stable API. A stable API and a relatively bug-free experience is what typically marks a 1.0 release.
There are plenty of smaller projects that should be perfectly happy to use the 0.7 release and grow/evolve as Arrow does. Especially when using Pandas+Arrow, since it's probably not a production environment and I can spare a few hours to fix a confusing bug.
In new projects it is common for a method or class to change direction or role as it develops, which may bring with it a refactor of identifier names, parameters, and such. After the class has been used for some time in different contexts these changes happen less and less, hence we call the class stable. A 1.0 release implies this stability.
Third party dependencies are still free to mutate their APIs, but when maintaining a production system you don't want to be mole-whacking API changes every point release.
For example, Apache Drill and the follow-on startup Dremio is in use in various places, and I believe they use the Arrow runtime. In contrast, when we worked on Graphistry/NodeJS <> MapD bindings, we did zero-copy data interop by agreeing on just the file format.
For people used to building frameworks like the above, Arrow is at a fine point. We hoped to use the runtime, but it wasn't necessary so far. More importantly, as framework builders, it got us past the typical decision of roll-your-own format via ~flatbuffers vs. dealing with orc/parquet. As a user, you'd be leveraging Pandas, Spark, etc., and only calling Arrow when occasionally talking between systems using your framework's internally supported interop layer. Data engineers will eventually be more exposed to this as they focus more on say streaming arrows vs non-streaming parquets, but that seems early for now.
Ex: We recently released https://github.com/apache/arrow/tree/master/js and are generally using Arrow as a way to compose interactive-time columnar CPU/GPU visual analytics technologies, including our own: https://devblogs.nvidia.com/parallelforall/goai-open-gpu-acc... . I'm hoping we'll have the cycles to start describing how the piece fit together for how we're rethinking visual analytics web apps (and interactive-time ETL in general), but you can start to guess one aspect of it based on the above.
What good reason? Why on earth would every single person on HN who happened to click on that link know what `pandas` is? How would they even know whose blog it is on?
Unless your blog is specifically, only for people who are already familiar with your `thing` (in which case, a mailing list might be better), then it simply makes sense to always have a header with a tagline explaining what your `thing` is and a link back to the main project website. Just in case it gets featured on HN or something.
You're leveling the criticism that since his blog deals with a niche that he should drop to a mailing list. Seriously? HN has been developing a bad reputation, and nitpicky comments presented without decorum create an incentive for people to stay away.
It's perfectly reasonably for the author of a blog post to assume that people familiar with his blog, his usual readers, will know what's up.
It's not too unreasonable to suggest that the author provide a summary if it's not clear what the piece or blog is about, except in this case it is described on the blog's About page on the list of open source projects. Redefining all of these things in every post would get annoying very fast.
It's reasonable to say that people working with statistical and data computing are likely as aware of Pandas as a web developer is aware of several of the largest JavaScript libraries.
No, I'm not. I'm offering a suggestion. You're being absurdly thin-skinned on someone else's behalf.
> Redefining all of these things in every post would get annoying very fast.
So put it in a header so it automatically appears at the top of every post. As I suggested originally.
This is simply good practise for any project where you want to attract people who haven't heard of your software before. If you don't care about these people, then a mailing list is a better idea so you can hide it away from everyone. Again, just a suggestion...
> No, I'm not. I'm offering a suggestion. You're being absurdly thin-skinned on someone else's behalf
No, we're being annoyed at the crazy level of bike shedding that hacker news is starting to see.
It's the author of Pandas's blog. He didn't submit it to HN, someone else did. It's totally reasonable to expect that if people are reading his, the author of Pandas, blog, then they already know what Pandas is.
What difference does that make? It's still here, right? So people are still reading it from here, right?
>It's totally reasonable to expect that if people are reading his, the author of Pandas, blog, then they already know what Pandas
No, not when it's linked from elsewhere. Which is what happens with blogs. Like it was right now.