At my company (healthcare analytics) we have very ...
# general
c
At my company (healthcare analytics) we have very large clients and smaller/medium size clients. Our default for working with data is using Databricks (and PySpark). Some clients require us to use Databricks and others don't really care. We'd like to develop more of an engine agnostic process and are investigating Daft. If we went with Daft solely as our DataFrame library, would that work as well on Databricks? And just to be clear, we currently do mostly "serverless". Not sure if that makes a difference. How does Daft "distribute" in this scenario? Would Ray still need to be included on Databricks? If we could do this, we would probably run smaller/cheaper clusters using Ray for smaller clients or clients that don't require Databricks.
c
Daft has two modes, local and distributed. The local mode (which is the default) should work fine on Databricks. The distributed mode requires ray, and Databricks does seem to provide a way to run ray: https://docs.databricks.com/en/machine-learning/ray/start-ray.html, but I personally have not tried it out. @jay have you played around with it before?
c
Thanks. I noticed this after I sent my question: > Limitations > • Ray on Apache Spark is supported for single user (assigned) access mode, dedicated access mode, no isolation shared access mode, and jobs clusters only. A Ray cluster cannot be initiated on clusters using serverless-based runtimes. See Access modes.
Follow up. What does it mean to say "local mode on Databricks"? Is that mean Daft Dataframes are confined to a single node cluster (but I can take advantage of multiple CPUs)? PySpark obviously will run on multiple worker nodes but I'm guessing running "local mode on Databricks" won't take advantage of multiple workers. Is that correct?
c
Yes that is correct. Daft will not take advantage of multiple workers, but will take advantage of multiple cpus on the driver node
j
You will be pleasantly surprised at how much you can do with Daft without distributed mode actually. We've seen folks see upwards of 7x speedups for some workloads vs running Spark, just because of how much overhead there is with running a cluster. (do your own benchmarking of course!)
c
Thanks. I'll do some investigation using a Job Cluster with autoscaling and Ray (I believe that is supported). That will be kind of close to serverless. Most of my jobs are daily batch jobs anyway. For certain clients I might be able to get away with a single node.
🙌 2
j
Yes, I think up to 100G or so single-node would be sufficient for most queries. You'll want to start running something distributed once you're reading a lot more data from cloud storage, to take advantage of increased aggregate network bandwidth
c
@jay - when you say "without distributed mode actually" are you referring to running on a large VM or something (vs Spark)? Or are you referring to a Databricks Single Node Cluster?
j
running on a large VM
👍 1
I have indeed tried using Ray on Databricks before. I found the integration to be a little finicky tbh
c
In the context of Daft + Ray and accessing/using Unity Catalog or just Ray in general?
j
(Just Ray-on-Databricks-Spark in general)
c
Thanks for all of your help.
🙌 1
j
Definitely! Would love to keep hearing feedback as you go along
c
Started doing some testing of PySpark vs Daft on Databricks. I'm getting this error when trying to access the Unity Catalog. I'm running a single-node cluster with Python 3.11.
image.png
I can't seem to find any documentation on that error. Any help would be appreciated.
nm. I had to pin "httpx==0.27.2" on the node. thx
🙌 1