8 ms·
Apache Arrow Datafusion 5.0.0 release
- xiaodai 5y agoI want to understand something, why use Spark and DataFusion? As the way to interact with them is SQL. Why not just use traditional DBMS like PostgreSQL? Are there explainer articles on this topic? See Quora: https://www.quora.com/unanswered/Whats-the-advantage-of-doing-SQL-in-Spark-and-Apache-Arrow-DataFusion-vs-doing-SQL-on-a-traditional-DBMS-like-MySQL-or-PostgreSQL https://www.quora.com/unanswered/Whats-the-advantage-of-doin...
- simonw 5y agoBecause if you have many TBs of data it's cheaper to run something like Spark across a bunch of smaller machines than it is to try and set up a many-TB PostgreSQL instance. A trick that many data warehousing tools use these days is to farm out computing to where the data is stored. You might have a PB of data spread across 100 different instances. When a SQL query comes in you break that up into a query plan that can be run in parallel against the subset of data on each of those instances, then aggregate together the results. It's cheaper to send the computation out to run next to the data than it is to copy the data back to the nodes that are executing the computation. It's all variants of the classic map/reduce technique. As a result, a data warehouse may be able to run a dumb SQL 'like' query against everything it is storing in a reasonable amount of time - since it gets to run in parallel. The trade-off is that you don't have consistency - ACID etc - or real-time results against data changes - generally your data warehouse will be repopulated on a schedule, but it won't be great at answering questions about changes that just happened a few seconds ago.
- tomnipotent 5y ago> Why not just use traditional DBMS like PostgreSQL Serialization. It's usually a non-trivial percentage of time spent in a lot of distributed systems, and for some workloads it can be the bulk of time spent. If I want to grab 3GB of data from a remote host and process it locally, we have to agree on how that data is going to be transferred so I can use it. Could very well be SQL, so we have some sort of network-based tabular data stream. Maybe it's Parquet files so we're using NFS/S3 to copy the files to local disks before reading into a completely separate in-memory data structure. At the end of this workload, I have 1GB of data I now want to write it back. Maybe the data is stored in-memory as an array of mixed-type structs, but I can't just send those bytes as-is to SQL server or mmap to the filesystem and expect Parquet to know what it means. Apache Arrow and DataFusion aims to eliminate all that work in rewriting bytes between hosts. Imagine being able to create a cost-based optimized query plan on Host A, send it to Hosts M-P for processing, and even have that query plan trickle down to Parquet predicates when reading files from disk, before returning data to Host A which can simply be copied from the network into local memory and you can start working with it right away.
- FridgeSeal 5y agoData formats like Parquet/Arrow and DataFusion are optimised for high speed read/write (and the processing) of large amounts of data, which is generally what you’re going to be using them for. Additionally, as others have mentioned, clustered processing for larger-than-machine/ram datasets is a bit easier to manage compared to setting up a database cluster. Another benefit is ephemeral-compute: we have Kubernetes cluster, and a particular message in a Kafka topic can kickstart a a spark job across several machines in the cluster (possibly causing auto scaling) which processes the x-TB’s of data it needs to, writes the results out and then finishes. Faster, cheaper and more suited than keeping a multi-node db cluster going. Also lets us run non-SQL stages with less bottlenecks: bulk ML scoring, bulk data enrichment, etc.
- legg0myegg0 5y agoHow would you compare the goals, vision, and current status of DataFusion with DuckDB? (www.duckdb.org) Could DuckDB be an execution engine for Ballista?
- hashjoiner 5y agoDataFusion committee here. DuckDB can work on Arrow data, so I think it could coexist with DataFusion / Ballista quite well. I am not sure whether having different execution engines for Ballista is on the roadmap, but it's certainly a possibility!
- etareduce 5y agoCongrats to the launch! The python binding to datafusion is also soon to be released. Exciting time!
- houqp 5y agoOne of the Arrow Datafusion committers here. Happy to help answer any question.
- pella 5y agoWhat is the "DataFusion"? - not in the FAQ ( https://arrow.apache.org/faq/ https://arrow.apache.org/faq/ ) - not in the Release page.
- pella 5y agoOK: I have found: https://github.com/apache/arrow-datafusion https://github.com/apache/arrow-datafusion "DataFusion is an extensible query execution framework, written in Rust, that uses Apache Arrow as its in-memory format.DataFusion supports both an SQL and a DataFrame API for building logical query plans as well as a query optimizer and execution engine capable of parallel execution against partitioned data sources (CSV and Parquet) using threads. DataFusion also supports distributed query execution via the Ballista crate." "Use Cases: DataFusion is used to create modern, fast and efficient data pipelines, ETL processes, and database systems, which need the performance of Rust and Apache Arrow and want to provide their users the convenience of an SQL interface or a DataFrame API."
- houqp 5y agoYou beat me to it, was about to post the github link :) Readme is a good starting place to learn more about the project.
- troelsSteegin 5y agoIs this the best current view of a roadmap for Datafusion? https://www.mail-archive.com/dev@arrow.apache.org/msg23576.html https://www.mail-archive.com/dev@arrow.apache.org/msg23576.h...
- alamb 5y ago<DataFusion committer here> I do think that is the best current view of a RoadMap and Vision -- it would be great to flesh it out a bit more. In fact, I'll make a note to try and add some more higher level context into the project on our goals.
- pauldix 5y agoThe new core we're building for InfluxDB (named InfluxDB IOx) uses Datafusion for query execution. We have multiple team members contributing to this release and we're super excited to be involved with it. I think it's a really exciting time for new OLAP systems because of Arrow, Rust, and the rise of object store + ephemeral compute for analytical and time series data.
- houqp 5y agoIndeed, big shout out to the InfluxDB team!
- rektide 5y ago> the rise of object store + ephemeral compute Question, does Datafusion itself integrate with Arrow's Plasma object store? Or is it more agnostic?
- andygrove 5y agoThere is no support for Plasma. There is a TableProvider API for custom file formats and there is built-in support for CSV, Parquet, and JSON.
- liminal 5y agoWhat's the relationship between Datafusion and Ballista? They seem to have been merged into a single repo. Do they share a release schedule? Are they a single product or still separate?
- andygrove 5y agoBallista started out as a separate project and was donated in April 2021. They currently share a release schedule (but have different versioning) and this was the first release of DataFusion to include the Ballista crate. My hope is that Ballista and DataFusion become more integrated over time but remain separate, with DataFusion being an embedded / single-process query engine and Ballista providing distributed execution.