Should you ditch Spark for DuckDB or Polars?
milescole.dev
milescole.dev
DuckDB on AWS EC2's price performance rate is 10x that of Databricks and Snowflake with its native file format, so it's a better deal if you're not processing petabyte-level data. That's unsurprising, given that DuckDB operates in a single node (no need for distributed shuffles) and works primarily with NVME (no use of object stores such as S3 for intermediate data). Thus, it can optimize the workloads much better than the other data warehouses.
If you use SQL, another gotcha is that DuckDB doesn't have advanced catalog features in cloud data warehouses. Still, it's possible to combine DuckDB compute and Snowflake Horizon / Databricks Unity Catalog thanks to Apache Iceberg, which enables multi-engine support in the same catalog. I'm experimenting this multi-stack idea with DuckDB <> Snowflake, and it works well so far: https://github.com/buremba/universql
Your data catalog point remains, it doesn’t offer anything other basic SQL describe functionality.
I have rarely seen data people familiar with K8s as they mostly use managed services, but feel free to prove me wrong!
That's why I started Universql in the first place; use their interface/protocol as an open-source tech and reduce the compute cost (often ~90% of the total cost of DWHs) by using DuckDB. You still get all the catalog features without paying the premium price.
Using popularity as justification of anything is the definition of slippery slope.
Disclaimer: I work at Altinity.
What I'm a little curious about with these "single node" solutions - is redundancy not a concern with setups like this? Is it assumed that you can just rebuild your data warehouse from some form of "cold" storage if you lose your nvme data?
You can use spot instances if you can afford more latency or might prefer on-demand instances and keep them warm if you need low latency. Databricks (compute) and Snowflake (warehouse) do that automatically for you for the premium price.
If a node fails when running a process (e.g. for an external reason not related to your own code or data: like your spot EC2 instance terminating due to high demand), you just run it again. When you're done running your processes, normally the processing node is completely terminated.
tl;dr: you treat them like cattle with a very short lifecycle around data processes. The specifics of resource/process scheduling being dependent on your data needs.
So if you want to use the duckdb native format, and you have a lot of data... what do you do? How do you keep your duckdb file up to date with incoming data to your data lake? Maybe you just... don't? Have a daily rebuild of your ephemeral-ish duckdb node?
It seems like it would be kind of a pain to manage once you get to moderate scale (say, TB+ datasets).
The author focuses on read/write performance on Delta (makes sense for the scope of the comparison). I think if an engineer is considering switching from spark to duckdb/polars for their data warehouse, they would likely be open to data formats other than Delta, which is tightly coupled to the spark (and even more so to the closed-source Databricks implementation). In my use case, we saw enough speed wins and cost savings that it made sense to fully migrate our data warehouse to a self managed duckdb warehouse using duckdb's native file format.
I'm thinking the same wrt dropping Parquet.
I don't need concurrent writes, which seems to me about the only true caveat DuckDB would have.
Two other questions I am asking myself:
1) Is there a an upper file size limit in duckdb's native format where performance might degrade?
2) Are there significant performance degradations/ hard limits if I want to consolidate from X DuckDB's into a single one by programmatically attaching them all and pulling data in via a large `UNION ALL` query?
Or would you use something like Polars to query over N DuckDB's?
1) I haven't personally run into upper size limits to the point of non linear performance degradation. However, some caveats to that are (a) most my files are in the range of 2-10gb with a few topping out near 100gb. (b) I am running a single r6gd metal as the primary interface with this which has 512 gb of ram. So, essentially, any one of my files can fit into ram.
Even given that setup, I will mention that I find myself hand tuning queries a lot more than I was with Spark. Since duckdb is meant to be more lightweight the query optimization engine is less robust.
2) I can't speak too much towards this use case. I haven't had any occasion to query across duckdb files. However, I experimented on top of delta lake between duckdb and polars and never really found a true performance case for polars in my (again atypical use case) set of test. But definitely worth doing your own benchmarking on your specific use case :)
Databricks have also bought into Iceberg and will probably lead with that or unify the two in future.
That aside, I was more pointing out that Delta, particularly via a commercial offering, is a data format biased towards Spark in terms of performance, since it is being developed primarily by Databricks as a part of the spark ecosystem. If you are plan to use Delta regardless of your compute engine, it makes perfect sense as a benchmark. However, for certain circumstances, the performance wins could be (in my case was) worth it to switch data formats.
Something under appreciated about polars is how easy it is to build a plugin. I recently took a rust crate that reimplemented the h3 geospatial coordinate system, exposed it at as a polars plugin and achieved performance 5X faster than the DuckDB version.
With knowing 0 rust and some help from AI it only took me 2ish days - I can’t imagine doing this in C++ (DuckDB).
Pandas' use of the dataframe concepts and APIs were informed by R and a desire to provide something familiar and accessible to R users (i.e. ease of user adoption).
Likewise, when the Spark development community somewhere around the version 0.11 days began implementing the dataframe abstraction over its original native RDD abstractions, it understood the need to provide a robust Python API similar to the Pandas APIs for accessibility (i.e. ease of user adoption).
At some point those familiar APIs also became a burden, or were not-great to begin with, in several ways and we see new tools emerge like DuckDB and Polars.
However, we now have a non-unique issue where people are learning and applying specific tools versus general problem-solving skills and tradecraft in the related domain (i.e. the common pattern of people with hammers seeing everything as nails). Note all of the "learn these -n- tools/packages to become a great ____ engineer and make xyz dollars" type tutorials and starter-packs on the internet today.
But from my experience in almost all cases it is misguided requirements e.g. we want to support 100x data requirements in 5 years that drive in hindsight bad choices. Not resume driven development.
And at least in enterprise space having a vendor who can support the technology is just as important as the merits of the technology itself. And vendors tend to spring up from popular, trendy technologies.
The problem with resume-driven technology choices are a kind of tech debt that typically costs a lot of cash to operate and perhaps worse, delivers terrible opportunity costs, both which really do sink businesses.
Premature scaling doesn’t. The challenge is not leaving it too late. Even doing that though only has led to a very few high-profile business failures.
All I think when I read this is, standing up new environments, observability, dev/QA training, change control, data migration, mitigating risks to business continuity, integrating with data sources and sinks, and on and on...
I've got enough headaches already without another one of those projects.
There are also some interesting points in the following podcast about ease of use and transactional capabilities of duckdb which are easy to overlook (you can skip the first 10 mins): https://open.spotify.com/episode/7zBdJurLfWBilCi6DQ2eYb
Of course, if you have truly massive data, you probably still need spark
I've also experimented with duckdb whilst on a databricks project, and did also think "we could do this whole thing with duckdb and a large EC2 instance spun up for an few hours a week".
But of course duckdb was new then, and you can't re-architect on a hunch. Thanks for the aricle.
I didn't think that people used polars a lot for ELT. I've usually seen it used for aggregations with small outputs (which, as you called out, it does a great job at).
I'm currently thinking of ditching Parquet all together, and going all in DuckDB files.
I don't need concurrent writes, my data would rarely exceed 1TB and if it were, I could still offload to Parquet.
Conceptually I can't see a reason for this not working, but given the novelty of the tech I'm wondering if it'll hold up.
I'd be interested in hearing about experiences of using duckdb files though, i can see instances where it could be useful to us
Example:
df = (
df.mutate(new_column=df.old_column.dosomething())
.alias('temp_table')
.sql('SELECT db_only_function(new_column) AS newer_column from temp_table')
.mutate(other_new_column = newer_column.do_other_stuff())
)It's super flexible and duckdb makes it very performant. The general vice i experience creating overly complex transforms but otherwise it's super useful and really easy to mix dataframes and SQL. Finally it supports pretty much every backend including pyspark and polars
I might have missed it, but the integration of duckdb and the arrow library makes mixing and matching dataframes and sql syntax fairly seamless.
I’m convinced the simplicity of duckdb is worth a performance penalty compared to spark for most workloads. Ime, people struggle with fully utilizing spark.
I'd recommend checking out their architecture whitepaper: https://docs.google.com/document/d/1tBw9A4j62ruI5omIJbMxly-l...
Imagine spark without the JVM baggage and with no need to spill to disk / serialize between steps unless it's necessary.
I think they’ve diluted the “brand” a bit with this approach and would be better off sticking with “Ray” for the distributed computing and spinning up the others as something completely separate, but that’s just me.
> Both persistent and in-memory databases use spilling to disk to facilitate larger-than-memory workloads (i.e., out-of-core-processing).
I don’t have personal experience with it though.
JFYI. I think the article itself is pretty unbiased but I feel like its worth putting this disclaimer for the author.
And I found the entire article very informative with little room for bias.
"Before writing this blog post, honestly, I couldn’t have answered with anything besides a gut feeling largely based on having a confirmation bias towards Spark."
That is - all these code assistants are going to be 10x as useful on spark/pandas as they would be on duckdb/polars, due to the age of the former and the continued rate of change in the latter.
Makes me wonder if each new technology project will need to make an effort to sythensize data to allow for "indexing" by LLMs. It will be like a new form of SEO
DuckDB is fantastic, though. I’ve never really built big data streaming situations so I can really accomplish anything I’ve needed with DuckDB. I’m not sure about building full data pipelines with it. Any time I’ve tried, it feels a little “duck-tapey” but the ecosystem has matured tremendously in the past couple years.
Polars never got a lot of love from me, though I love what they’re doing. I used to do a lot of work in pandas and python but I kind of moved onto greener pastures. I really just prefer any kind of ETL work in SQL.
Compute was kind of always secondary to developer experience for me. It kills me to say this but my favorite tool for data exploration is still PowerBI. If I have a strange CSV it’s tough to beat dragging it into a BI tool and exploring it that way. I’d love something like DuckDB Harlequin but for BI / Data Visualization. I don’t really love all the SaaS BI platforms I’ve explored. I did really like Plotly.
Totally open to hearing other folks’ experiences or suggestions. Nothing here is an indictment of any particular tool, just my own ADD and the pressure of needing to ship.
This pushed me to finally investigate the DuckDB hype, and it’s hard to believe I can just write SQL and it just works.
Thanks for the feedback on marketing! Daft is indeed distributed using Ray, but to do so involves Daft being architected very carefully for distributed computing (e.g. using map/reduce paradigms).
Ray fulfills almost a Kubernetes-like role for us in terms of orchestration/scheduling (admittedly it does quite a bit more as well especially in the area of data movement). But yes the technologies are very complementary!
- auto loader/cloud files. Can attach to a blob storage of e.g csv or json, and give them as batches. As new files come in you get batches containing only the new files.
-structured streaming and it's checkpoints. It keeps tracks across runs of how far in the source it has read (including cloud files sources), and it's easy to either continue the job with only the new data, or delete the checkpoint and rebuild everything.
How can you do something similar with duckdb? If you have e.g a growing blob store of csv/avro/json files? Just read everything every day? Create some homegrown setup?
I guess what I describe above is independent of the actual compute library, you could use any transformation library to do the actual batches (and with foreachbatch you can actually use duckdb in spark like this).
This will enable you to run incremental batches or do structured backfilling.
Another option would be to use a data loading framework like dlt, which will also do the same but even more light weight (without the orchestration part, you write your batch flow logic in it).
This blog post suggests that it has been supported since 2021 and matches my experience.
The link you shared shows how DuckDB can run SQL queries on a pandas dataframe (e.g. `duckdb.query("<SQL query>")`. The SQL query in this case is a string. A dataframe API would allow you to write it completely in Python. An example for this would be polars dataframes (`df.select(pl.col("...").alias("...")).filter(pl.col("...") > x)`).
Dataframe APIs benefit from autocompletion, error handling, syntax highlighting, etc. that the SQL strings wouldn't. Please let me know if I missed something from the blog post you linked!
However it's somewhat weakened by the possibility that some parts of the SQL string are resolved by the surrounding python context.
There may be another benefit: a lot of LLMs are getting good at how do I do X in Duckdb.
Progress is being tracked on Github Discussions[3].
[1]: https://duckdb.org/docs/api/python/spark_api.html
The author ran Spark in Fabric, which has V-Order write enabled by default. DuckDB and Polars don't have this, as it's an MS proprietary algorithm. V-Order adds about 15% overhead to write, so it does change the result a bit.
The data sizes were bit on a large size, at least for the data amounts I see daily. There definitely are tables in the 10GB, 100GB, and even in 1TB size range, but most tables traveling through data pipelines are much smaller.
Yeah I was debating whether to share all of the source code. I may share a portion of it soon.
realistically means keeping in mind that the processing itself also requires memory as well as prerequisites like indexes which also need to be kept in memory.
maximum memory at AWS would be 1.5TB using r8g.metal-48xl. so, assuming 50% usable for the raw data means about 750GB are realistic.
I guess the scale of data here ~100GB is manageable with something like DuckDB but once data gets past a certain scale, wouldn't single machine performance have no way of matching a distributed spark cluster?
We normally still would only need to process say, the last 10 days of user data to get decent recommendations, but occasionally it would make sense for processes running over the entire dataset.
Also this isn't that large when you consider binary artifacts (say, healthcare imaging) being stored in a database, which pretty sure that's what a lot of electronic healthcare record systems do.
A random company I bumped into has a 40TB OLTP database to this effect.
If you can afford it you should be using the host NVME drives or FSX and relying on Spark to handle outages. The difference in performance will be orders of magnitude different.
And in this case you won't have the ability to store 64TB. The average max is 2TB.
Microsoft’s synapse is their «ripoff product» which tries to compete directly by offering the same thing (only worse) with MS branding and better azure integration.
I’ve yet to see Spark being used outside of these products but would be happy to hear of such use.
They happen all the time if you work for banks, large finance companies or government.
It’s not just the 2TB databases - it’s the 100 analysts all doing it their own thing with that data at the same time.
This is my problem with databricks. It seems like in the course of selling their product they have taken the received wisdom of "do not run expensive and complex compute clusters unless absolutely necessary" and turned it into "it's fun and easy to run distributed compute clusters - everyone's doing it and you should too" regardless of how contextually appropriate it is.
When it's more common to be manipulating lots of smaller datasets together in a way where you need to have more than 100GB of disk space. And in this situation you really need Spark.
Recent podcast https://talkpython.fm/episodes/show/488/multimodal-data-with...