Is there a way to cache to a path on disk instead ...
# daft-dev
k
Is there a way to cache to a path on disk instead (basically doing a write_parquet and then straightaway read_parquet)? I'm wondering if that could help alleviate some of the object store usage
d
Hmm could you elaborate a little? If you're writing to cloud storage, presumably there's a dataframe that's being written, so why not reuse that dataframe?
k
For example, in this function from dedupe, there's multiple collect() statements. I'm wondering if writing them to disk instead would help the process to not OOM. Once I have gotten the duplicate pairs, I would write the results out. However, I'm not sure why an antijoin to my original DF works if I run it separately but it doesn't work if I continue to run it as another statement sequentially after the pairs were saved.
Copy code
def components(df: DataFrame) -> DataFrame:
    b = df.select(daft.to_struct("u", "v").alias("e")).collect()
    while True:
        a = (b
             .select(large_star_map(col("e.u"), col("e.v")).alias("e")) 
             .explode("e")
             .select("e.*")
             .groupby("u").agg_list("v")
             .select(large_star_reduce(col("u"), col("v")).alias("e"))
             .explode("e")
             .where(~col("e").is_null())
             .distinct()
             .collect()
        )
        b = (a
             .select(small_star_map(col("e.u"), col("e.v")).alias("e"))
             .select("e.*")
             .groupby("u").agg_list("v")
             .select(small_star_reduce(col("u"), col("v")).alias("e"))
             .explode("e")
             .where(~col("e").is_null())
             .distinct()
             .collect()
        )
        # check convergence
        a_hash = a.select(col("e").hash().alias("hash")).sum("hash").to_pydict()["hash"][0]
        b_hash = b.select(col("e").hash().alias("hash")).sum("hash").to_pydict()["hash"][0]
        if a_hash == b_hash:
            return b.select("e.*")
j
What about running a .collect() right after the read? That would materialize it fully (spilling if necessary).
k
Trying to see if there's a way to reduce the spilling and also materialize less even if it's slower
But i guess that doesn't work out well because of the groupby which will need the whole thing to be materialized anyway?
j
Unfortunately if you’re running .collect() that materializes everything and will spill if your cluster doesn’t have enough memory (by default 30% of cluster memory) to hold everything. The groupby does perform a shuffle, which does indeed perform a full materialization of the data in the process.
Sometimes you just need more memory 😝 (Or smaller partitions)
plus one 1
😅 1
k
Hmm.. okay haha will have to figure out how to get more memory For the other question on the antijoin succeeding if run separately but failing if it's run sequentially after the dedupe, is there perhaps a way to clear out the object store?
d
Everything in the object store is reference counted. So if you get rid of the references, it should garbage collect itself
k
Would there be a difference between
Copy code
df = read_parquet()
df = ...
df.write_parquet(path)
df.read_parquet(path)
df2 = other_processes(df)
vs
Copy code
df = read_parquet()
df = ...
df2 = other_processes(df)
vs
Copy code
df = read_parquet()
df = ...
df.write_parquet(path)
df2 = other_processes(df)
d
In your first example, did you mean to say
Copy code
df = read_parquet()
df = ...
df.write_parquet(path)
df = read_parquet(path)
df2 = other_processes(df)
?
k
I meant the difference among all three variants
d
Right I'm just double checking that in the first variant you're reassigning to the same variable name
k
The first one writes and reads, the second one doesn't save at all, and the third one writes but continues using the original df
Oh yes I did mean that
d
I want to add a fourth variant here which would look something like
Copy code
def f():
  df = read_parquet()
  df = ...
  df.write_parquet(path)
 
f()
df = read_parquet(path)
df2 = other_process(df)
It's a little subtle, but if you're using cpython (the reference implementation of Python), then I believe when you do
Copy code
A = object1
A = object2
The actual steps under the hood are: 1. Object1 will be created and assigned to dict['A'] 2. Object2 is created 3. Dict['A'] is reassigned to point to object2 4. Object1 is dropped So both objects (in our case both dataframes) will exist in memory during the reassignment
But if we wrap the operations in the closure f(), then the references are dropped when f exits
k
I see, so if that's the case then for variant 1 there would be a short period of time whereby they both coexist and cause a memory explosion?
d
That's my understanding, I can do some experiments to benchmark. But I believe Ray advises wrapping work in functions for this reason
k
I see
How about the other two variants?
d
It looks like they'll all have the same impact as using the closure f()
I could be wrong, just riffing
k
For variants 2 and 3 wouldn't df and df2 coexist throughout because df is used in df2's definition and they would exist until df2 is materialized and then dropped?
d
Yeah I think variants 2, 3, 4 face that issue
Oh hm but then again df is never materialized, it's passed into otherprocess
k
True haha
1,3,4 materializes df when it's being written but I guess there's no caching there because when I run a count_rows on df after writing to disk I think it just does the whole process again..
I try to run variant 1 as if it's a manually implemented .persist() so that the row count can be faster as well without having to pressurize the nodes to hold the entire dataset within memory with .collect() but then sometimes i get errors that the parquets are missing footers so I'm not sure if it's possible that a process is marked as completed but the write has not yet really completed.. or that there is some kind of lag in the file system