Any suggestions, clues or guide how to best optimi...
# general
a
Any suggestions, clues or guide how to best optimize the partitionning and/or other parameters (max connections, num_cpu, concurrency etc.) to not kill memory when; 1. downloading + decoding images from MinIO a. Applying some chained computations on those images 2. wiriting deltalake into that same minio store? 3. (not using ray - yet) Specs:32 cores, 64Gb RAM Objects: ~2.5mil images (mix of RGB, Grayscale, RGBA,...) of varying resolutions (800x600 - 2000x1800)
j
I think it’s easiest if you try the new native executor @Colin Ho can guide you here The problem with memory is usually you need to know your partitioning upfront (total number of partitions). The new executor does streaming execution which lets you set morsel sizes instead.
a
Is that using `daft.context.set_runner_native()`right? This works with all versions >=0.3.11 right? I've tried using it but didn't seem to fix the OOM issues - or do I need to go through
daft.set_execution_config
j
I think you might still need to reduce the sizes of your morsels using set_execution_config, especially on a 32 core machine (we’re probably running too aggressively).
👍 1
a
Would appreciate if there are any guidelines or rule of thumbs to help set the best size ? 😄 Otherwise I'll go through trial and error 🙂
j
You can think of the total number of rows running in parallel as: num_cpus * morsel size
So given you have 32 CPUs, maybe try something like 256? Really depends how large your images are 😝 We’re working on more automated mechanisms to do dynamic sizing/parallelism based on memory pressure but right now it’s a little bit manual!
a
As I said, the worst case scenario (AFAIK so far) is images of size ~2000x1800x4
And yeah looking forward to that more automated stuff ;D - quite liking daft in general so far thank you!
(and still learning a lot on the side)
Just to get this straight, you suggest to set the following to 256? (and the default is 131072)
• default_morsel_size – Default size of morsels used for the new local executor. Defaults to 131072 rows.
Also note, again, I'm not keeping the images in the rows! Just downloading them processing, computing and hoping that they get flushed after getting the results
c
Yes, try 256, likely the lower the better if you're doing url-downloads and decodes.
👍 1
a
Can we use the morsel settings while using ray clusters? Seems to not work correctly?
Also, for single node, thanks that worked great! 😄
c
no the morsel size only works when using the native runner, i..e
set_runner_native()
a
What would be the best approach when using ray then? Currently, just partitionning, but it doesn't seem efficient / optimal at all?
c
For ray, using
into_partitions
, will split up your data into smaller chunks. Is that what you're currently using today?
a
Yes, I compute chunk size in worst case scenario based, but of course that's not optimal most of the time