Fast analysis with DuckDB and Pyarrow
tech.gerardbentley.com
tech.gerardbentley.com
This post also links to another discussion about the Parquet data format (https://pythonspeed.com/articles/pandas-read-csv-fast/), also supported by Arrow, which is also extremely useful but I never see anyone talking about it. Granted, Parquet data can't natively be imported into Excel which is likely the main cause.
The future is very promising, I am personally very excited about DuckDB. But it’s too soon to be griping about old tutorials.
DuckDB is a convenient drop-in for OLAP in the way that SQLite is for OLTP, it's a great library. Python's death on the vine as a data analytics platform cannot come fast enough. There's a reason people are abandoning pandas and spark/databricks in droves and fleeing to useful SQL-like tooling such as Snowflake. Once this tech bubble bursts, and the tens-of-thousands of subsidized Python-packing data scientists/engineers/analysts get laid off, the language will return to what it should be: another scripting language that generates SQL for talking to a real data engine.
I find it's "limitations" generally enforce structure that makes rewriting and distilling bad queries easier. Well written and formatted queries can be extremely readable and easy to navigate/grep in a way that I have never found the others to be.
The flexibility of dropping in and out of a full fledged programming language and a query api, I find it tends to grow unnecessarily complex in the hands of many practitioners.
Although my preference is a SQL cursor or equivalent interface in my chosen language. For some reason I find the strict separation between SQL (the declarative expression of business logic and relations) and the language (imperative or functional control) very helpful.
SQL also has the benefits of portability. Almost every query computation engine supports SQL these days. While, obviously you will need to rewrite the parts that rely on unique extensions, the migration path is greatly simplified.
The one thing I think pandas has going for it that I desperately wish was picked up as a new standard in SQL are aggregations for seemlessy moving between different time series frequencies. Pandas as problematic as it is, I have yet to find anything else that makes time frequency conversions as convenient and predictable.
Is Parquet better/faster/stronger than keeping everything in schemaless CSVs? 100%, but it has historically meant that I have to make trade-offs to benefit the tooling rather than how I want to interactively approach my data.
Most people do not have "big data" problems where the performance differences of Parquet vs csv matter.
I might submit a few feature requests, but one that immediately comes to mind: csv -> parquet. Perhaps out of scope for the original vision, but having a single utility that could roundtrip data would be fantastically useful.
Arrow has been truly revolutionary in this regard, providing a solid in-memory data format (with performant APIs in many languages) for interchange between different engines and even formats.
You can go from ORC to Parset to CSV on a local FS or S3.
With DuckDB, it’s like you can build your own AWS Athena at likely a fraction of the cost. Now if only someone would integrate vaex with DuckDB, it will make your powerful Apple Silicon machines a compelling alternative to running a full fledged Spark/Hadoop cluster.
Isn't the whole purpose of Athena to scale to large amounts of data that don't fit into memory? How does duckdb fit in here? I thought it's an in-memory database?
df = pd.read_csv("large.csv", engine="pyarrow")
https://pythonspeed.com/articles/pandas-read-csv-fast/Getting non-technical people to do all that (a project manager wanting to see some simple stats for example) becomes difficult. Having a single binary that's easier to distribute, faster to execute and can be plugged into an existing query tool (like DataGrip) is a huge win for actual usability.
With that said, the resident memory size of stuff-in-parquet/arrow/duckdb is a lot less in practice than stuff-in-pandas (i know) and stuff-in-sqlite (i believe), so it still enables more in-memory workloads than you might otherwise be able to do.
The key ingredient is the on-disk format: hdf5 and more recently also arrow.
Looks like DuckDb would force me to use SQL, which is not what I typically want. That being said, if this offered better data management than a binary blob like hdf5, it might be worthwhile. If the use case requires processing historical data for example, a real DB would be better than juggling 10 x 40GB files or having one big 400GB file that keeps growing.
https://github.com/duckdb/duckdb/blob/master/examples/python...
Plus, if you are working in Python, you can use DuckDB as the engine underneath Ibis, Fugue, Siuba, or anything that works with SQLAlchemy (using the DuckDB-engine driver)! In R, you can use dplyr or dbplyr.
DuckDB's file format is one way to persist data (it uses a single file), but you can also write out to Parquet, or write out to Apache Arrow and then parquet (in a partitioned format I believe).
Disclaimer - I write docs for DuckDB!
"Reading the full CSV without datetime parsing is in line in terms of speed though." This sentence was a bit ambiguous, but is important: if you read this file in pandas with engine='pyarrow' but don't convert the date/time column to a pandas datetime dtype you get the same ~100 ms read time as calling PyArrow directly. So basically the entire time is spent converting the strings to dates. If this datetime thing isn't an issue for you then you can just use the engine='pyarrow' argument to read CSVs with Pandas.
In my own tests with various datesets Polars has always been much faster than duckdb/pyarrow. For this relatively small dataset it's about 2x faster, which is about the smallest margin I've found. Polars is also much easier to write, as its query optimization is so effective - you don't need to know all the tricks that Pandas requires (and is still 3x-10x faster than Pandas even when you apply all the tricks).
I've started making videos to address the need for more guidance in Polars - see this new on reading CSVs: https://www.youtube.com/watch?v=nGritAo-71o
I'm curious how it loads so fast initially.
This post prompted me to give pyarrow a spin and WOW did it load the CSVs fast compared to what I'm used to. That said, I'm not entirely sure what to do with them from there. When I tried table["col"] it gave back a list of lists (maybe because of the multi-threaded loading?). And while the article talked about "querying" the arrow files with duckdb, how easy is it to work with them otherwise, in terms of updating them? And what about non-SQL style updates? Sometimes you just have to do an itertuples or apply and do random non-declarative modifications to the data.
And everyone uses it because it's what you do when your boss tells you "we're transforming the analytics team, you're all to become data scientists because everyone has data scientists now". You just grab whatever had the biggest mindshare on SO and in random yt tutorials. Can't blame them.
But hooo boy does pd get on my nerves.
Care to share what propriety stuff you were using?
Honestly I still feel like I'm missing some sort of larger story about the semantics of Pandas (like the "functions" explanation above), so if anyone knows of anything that made Pandas click, please let me know.
One big issue I see with Pandas is that it assumes that I want multidimensional keys represented via axes labels on some sort of grid (or hypercube). This is evident from the __getitem__ interface, which is frankly confusing. In reality, selection is the most important detail, and managing columns is just a nice-to-have. Pandas's MVP is the Series object, not DataFrame.
In article #1 in your link, "SettingWithCopy" is emblematic of the issue that happens when you allow mutability, so mutability delenda est. This is why FridgeSeal is confused below. One day I'll write a blog post about ways Pandas could be better, but in the meantime I will probably monkey patch a lot of the nonsense out.
It's fast! The API is nice. But the documentation is not great. While using pandas I can just open the docs and find all functions and examples of how to use, etc. Polars felt really bare.
It's something I'd like to use in personal projects. But other data scientists would not accept using polars as it is today.
We hope to have such a clear API that once the expression language clicks, it is less needed to have the docs open constantly as it should feel natural.
I used it a month or two ago, so I don't remember specifics, but. The docs weren't unclear, they were bare compared to pandas tough. Once I found an example of how to do something I managed to do it. But a couple of times I opened a page and there was WIP message.
There's definitely a bit of unfairness in the comparison, as pandas as you say has been the standard for python for years. But that means I already know all the syntax and how the API works, so even just a method name or required params will get me a long way.
When I was fumbling around with polars I would have liked more examples of how to use each method/function. I mean, once the language clicks I wouldn't need it anymore. But at the start I really need that hand-holding. I'd have to reread other pages to remember simple syntax.
I used on an analysis that was taking long and it went way quicker.