Daft: Distributed DataFrame for Python
github.com
github.com
We do have a dependency on the Arrow2 crate like Polars does, but that has been deprecated recently so both projects are having to deal with that right now.
[0] -- https://github.com/search?q=repo%3AEventual-Inc%2FDaft%20pol...
https://www.getdaft.io/projects/docs/en/latest/faq/dataframe...
Daft is distributed and vectorized with plenty of benefits from lazy evaluation.
[1]: https://www.tpc.org/tpch/results/tpch_advanced_sort5.asp?PRI...
For hardware, we were using AWS i3.2xlarge machines in a distributed cluster. And on the storage side we are reading Parquet files over the network from AWS S3. This is most representative of how users run query engines like Daft.
The TPC-H benchmarks are usually performed on databases which have pre-ingested the data into a single-node server-grade machine that’s running the database.
Note that Daft isn’t really a “database”, because we don’t have proprietary storage. Part of the appeal of using query engines like Daft and Spark is to able to read data “at rest” (as Parquet, CSV, JSON etc). However this will definitely be slower than a database which has pre-ingested the data into indexed storage and proprietary formats!
Hope that helps explain the discrepancies!
> The performance metric reported by TPC-H is called the TPC-H Composite Query-per-Hour Performance Metric (QphH@Size), and reflects multiple aspects of the capability of the system to process queries.
so it's not directly comparable to seconds/query reported by Daft.
One thing i haven't seen yet and I'd be really interested in knowing if column level type hinting is on the plan at all for this (or any similar tools)?
I'm a data engineer and my team use mypy pretty heavily, but that can only really tell us if something is a "dataframe" type. It always strikes me that 90% of data issues during development (incorrect joins, column misspellings or non-type appropriate operations, not factoring in nulls) would be eliminated if mypy had information of the columns and their data types.
1. Construct a dataframe (performs schema inference)
2. Access (now well-typed) columns and operations on those columns in the dataframe, with associated validations.
Unfortunately step (1) can only happen at runtime and not at type-checking-time since it requires running some schema inference logic, and step (2) relies on step (1) because the expressions of computation are "resolved" against those inferred types.
However, if we can fix (1) to happen at type-checking time using user-provided type-hints in place of the schema inference, we can maybe figure out a way to propagate this information through to mypy.
Would love to continue the discussion further as an Issue/Discussion on our Github!
Regex for strings and more detail on how partitioning works would be interesting
And thanks for the feedback! We’ll add more capabilities for regex, as well as flesh out our documentation for partitioning.
Edit: added a new issue for regex support :) https://github.com/Eventual-Inc/Daft/issues/1962
The network indeed becomes the bottleneck. In 2 main ways:
1. Reading data from cloud storage is very expensive. Here’s a blogpost where we talk about some of the optimizations we’ve done in that area: https://blog.getdaft.io/p/announcing-daft-02-10x-faster-io
2. During a global shuffle stage (e.g. sorts, joins, aggregations) network transfer of data between nodes becomes the bottleneck.
This is why the advice is often to stick with a local solution such as DuckDB, Polars or Pandas if you can keep vertically scaling!
However, horizontally scaling does have some advantages:
- Higher aggregate network bandwidth for performing I/O with storage
- Auto-scaling to your workload’s resource requirements
- Scaling to large workloads which may not fit on a single machine. This is more common in Daft usage because we also work with multimodal data such as images, tensors and more for ML data modalities.
Hope this helps!
Depends on your workload size, if the compute time is less than time it takes to transfer the data back and forth, then it might be a bad idea to use it indeed.
Typically these solutions should only be tested after maxing out vertical scaling, before applying for horizontal scaling.
Network is one hell of a destroyer when it comes to advantages gained from distributed computing.
We actually already have read support. Check out the pyiceberg docs' Daft section: https://py.iceberg.apache.org/api/#daft
It's also very easy to use from Daft itself: `daft.read_iceberg(pyiceberg_table)`. Give it a shot and let us know how it works for you!
> To disable this behavior, set the following environment variable: DAFT_ANALYTICS_ENABLED=0
> [0] In short, we collect the following:
> On import, we track system information such as the runner being used, version of Daft, OS, Python version, etc.
> On calls of public methods on the DataFrame object, we track metadata about the execution: the name of the method, the walltime for execution and the class of error raised (if any). Function parameters and stacktraces are not logged, ensuring that user data remains private.
Is this telemetry really necessary?
[0] https://www.getdaft.io/projects/docs/en/latest/faq/telemetry...