Also, how is memory managed in such cases? ```df_i...
# general
a
Also, how is memory managed in such cases?
Copy code
df_img = df_img.with_column(
    'temp_results',
    df_img['path']
    .url.download(max_connections=8, io_config=io_config)
    .image.decode()
    .apply(computed_dict, return_dtype=daft.DataType.python())
)
will the images / binaries be flushed after the computation, or after each partition, or when can I assume the memory to be freed?
c
In the case of your example, memory will be incrementally freed upon completion of the expressions, per-partition. For example, the binaries from
url.download
will be flushed once
image.decode
is completed, and the images from
image.decode
will be flushed once
.apply
is completed.
a
Thank you! That's what I assumed. But, I'm seeing some drastically increasing memory load over time, I'm still unsure where it's coming from, it's most probably not the whole binaries/image. Trying to work on a ~5Tb dataset, I keep crashing. Currently still trying to investigate / profile what's happening.
c
If you're not already using it, I suggest https://github.com/bloomberg/memray to profile memory. It outputs binary files that can be processed into flamegraphs, tables, etc, of which you could send to us and we can also take a look at it if you find anything interesting!
🙌 1
n
Does this patten also follow across transformations? ie., does daft perform the same cleanup in this situation? Or do these behave differently?
Copy code
df_img = df_img.with_column(
    'temp_results',
    df_img['path']
    .url.download(max_connections=8, io_config=io_config)
    .image.decode())
df_img = df_img.with_column('path', df_img['path'].apply(computed_dict, return_dtype=daft.DataType.python())
)
c
It looks to me that it will be the same. More generally though, you can inspect the physical plan to get a more accurate sense of what will happen with a query.
df.explain(True)
Daft consolidates all transformations into a plan, and you can see each step when doing
.explain(True)
, example below. In a project step for example, daft will execute some number of expressions, i.e.
image_decode(download(col(path))) as temp_results
per partition. You can think of the execution of these expressions as like doing a depth first search. Retrieve the column, 'path', do the download, then decode.
In general, once a computed expression like
url_download
is not used or reference any more in any other expression, it should be dropped.
n
Cool - is there more information on how to debug what daft is doing while it's doing it? I've been trying to understand when it decides to map out the work and when it collects it all back together.
c
Currently the best bet for live debugging would be the progress bar, 🥲, which we understand is not very ergonomic and are working on something better. Or if you're running on Ray, checking the job in the ray dashboard to see the submitted tasks.