I've started a comparison of Daft to Databricks (P...
# general
c
I've started a comparison of Daft to Databricks (PySpark). I've run into a performance issue that I can't explain. Some context: 1. I'm running the equivalent PySpark and Daft on identical Databricks SINGLE-node cluster (32 GB, 4 cores) 2. The table/file is 135 million rows of lab results (delta table - approx 70 parquet files) 3. For PySpark I'm using standard Unity catalog table read and applying filter; for Daft I'm connecting to the data location (ADLS) and using delta read. I'm applying equivalent filters. For PySpark it takes usually 2s to 8s. For Daft it takes usually 50s to 60s. I'll attach more details in thread.
๐Ÿ‘‹ 3
daft_code
spark_code
daft explain plan
spark explain plan
This GIF shows Daft executing on Databricks. It shows some additional output as it's executing. Not sure how helpful that is.
๐Ÿ™Œ 2
j
Oo loving the screen recording. Super helpful.
d
Cool, thanks for running this comparison! Instead of a
.count()
I'm curious to compare the performance on a
.collect()
for daft vs spark. I feel like the performance difference is either in this aggregation or in data skipping on the scans, and
.collect()
will tell us which it is
a .min()
.sum()
on a single column will also be a pretty interesting aggregation to look at
c
For:
collect
- about the same for each. to be be clear the same as their respective previous runs.
sum
- same for spark and for daft it improved quite a bit. instead of 55s it was 28s (about hafl)
j
@Desmond Cheong any thoughts here? Iโ€™m pretty sure this is because Photon is leveraging statistics written specially into DeltaLake by the databricks platform, if these files were written through databricks Spark, which Daft doesnโ€™t have access to. (E.g. sum of a given column per-file) On the count side, I actually think this might not be too bad as a new feature we can add to do a
count(lit(1))
since our engine should technically support micropartitions with no columns but that have a row size now โ€” @Desmond Cheong perhaps we could try to knock that out quickly?
Iโ€™d be super interested to see if databricks Delta provides some additional stats we could maybe access from Daft for speeding up these operations
d
Yeah I'm doing some experiments on my side. Delta provides num rows, null count, min, and max stats for the first 32 columns of a table. Using num rows would allow us to do a pure delta metadata read for some of these aggregations. But what struck me as odd is that we're still slower for collect workloads, so my guess is that our data skipping is not as aggressive as dbr's
j
Iโ€™m also pretty sure that Databricks I/O is going through their dbfs, which is going to be cached โ€” vs Daft reads are going straight to S3 every time because thatโ€™s what deltalake provides us
๐Ÿ’ฏ 1
Too many variables ๐Ÿ˜ญ
They definitely have the home ground advantage here
c
In looking at some of the Databricks/Spark runtime stats, it appears like there are a lot of parallel operations. It appears to be spawn 67 tasks which correspond to each of the partitions. I don't think it's "skipping" any of the parquet files. Is it possible that Daft and Databricks are doing the "same" execution plan but Daft is not parallelizing? Is there some specific setting/configuration that I need to do to "turn on" parallelization?
j
Daft's single node engine is very different from Spark's -- it's a streaming execution model with parallelism determined by the number of available CPUs on your machine. cc @Colin Ho for thoughts as well, but you can actually get much more indepth information about what Daft is doing by turning on the environment variable:
DAFT_DEV_ENABLE_EXPLAIN_ANALYZE=1
This will dump a file that contains a little graph that we can then visualize!
c
j
Huh. That looks like potentially the wrong execution. Is there another file available that has the count aggregations and stuff?
c
That was for a filter with a projection of one column that is collected.
๐Ÿ‘ 1
I'll run one with count
count_explain.txt
c
CPU Time = 46.96ms for the count, which means bulk of time must be from the scan
j
Gotcha. Yeah I'm guessing the difference in times with the fact that Databricks is reading from their optimized
dbfs
filesystem, but Daft is reading directly from S3. The actual compute in this query is really fast (just a few ms) We'll have to do some more experiments on our end to see if we can maybe special-case I/O operations when run from inside a databricks environment to hit dbfs.
Also the other thing we can do here is for a count we really shouldn't need to read the first column (which looks like is 6.78GB of data here). We should have a count operation translate into
sum(lit(1))
I think, which should work.
c
I'm using
read_deltalake
. Is there an issue with delta-rs not scanning/reading in parallel? I could see reading/scanning each parquet in sequence would be slow.
j
No, we skip delta-rs for the actual data reads. The way that reads work in Daft: 1. We ask delta-rs for file locations and other metadata about each file 2. We then use those file locations to perform reads in parallel The problem here is likely that delta-rs is returning us non-optimal file locations (likely pointers to the data in S3 directly), but the databricks spark engine is instead accessing some optimized storage layer (i.e. dbfs). @Desmond Cheong any chance you could do some quick investigations here on a databricks instance?
๐Ÿ‘ 1
c
I see. Point of clarification. I'm on Azure. Files on ADLS v2. Not sure that makes a difference, but wanted to point out not using S3.
I found out something today you will probably find interesting. I was using an app registration/service principal (with client secret) to access ADLS. When I switched to using
access_key
in AzureConfig the performance increased significantly. I went from consistently 1m+/- runtimes to closer to 15s. Still slower than PySpark but huge improvement.
Also, if I explicitly turn off caching on Databricks,
collect()
takes about 12s consistently for PySpark (vs 15s for Daft using access_key). @jay - I think your theory on caching (at least for subsequent runs has some merit).
๐Ÿ˜ฎ 1
d
That's interesting! Was this the config you disabled?
spark.databricks.io.cache.enabled
c
Yes.
spark.conf.set("<http://spark.databricks.io|spark.databricks.io>.cache.enabled", False)
spark.catalog.clearCache()
Not sure the clearCache is strictly required.
In an ideal world I would like to switch completely from PySpark to Daft/Ray. However, with the performance differences that I'm seeing, I don't think that is possible (at least for now). The work that I'm doing is a quasi-DSL that will run on Databricks or VMs. Since the DSL is an abstraction anyway, I'll have to write separate implementations (1 in Daft and 1 in PysSpark). The good news is there is a lot of overlap of the PySpark API and the Daft API so even though there will be some duplicative effort I think there can be some shared code too (specifically around filtering using
col()
). If you have suggestions on how to improve performance on Databricks, I'm happy to give it a shot. Thanks.
d
If you have suggestions on how to improve performance on Databricks
One thing that we might be able to do on our side is take advantage of databrick's DBIO cache as well. Need to prototype it a little but I'll keep you in the loop!
๐Ÿ‘ 1
j
Thanks for giving us a spin on databricks @CR โค๏ธ Keep giving us feedback and the team will work to make it better. One project we're actively working on which might interest you is an integration with Spark to give Daft a PySpark-compatible API... It's currently making progress and we should have some announcements about it soon ๐Ÿ™‚
c
integration with Spark to give Daft a PySpark-compatible API
That's interesting. Is there more info on this? A github issue or ?
j
Docs about this beta feature coming soon! cc @ChanChan Mao @Cory Grinstead
๐Ÿ‘€ 1
c
Thanks. I look forward to hearing more.
c
That's interesting. Is there more info on this? A github issue or ?
we have an issue label for spark related work. https://github.com/Eventual-Inc/Daft/issues?q=is%3Aissue%20state%3Aopen%20label%3Adaft-connect feel free to look at open/closed issues to get an idea of the current state of "daft-connect" (our spark connect implementation)
If you want to start playing around with it, here's a quick snippet to get you started! Note that many operations are not yet implemented, the SQL api should be the most complete.
โค๏ธ 3