-
Notifications
You must be signed in to change notification settings - Fork 413
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
feat: no longer load full table into ram in write by using concurrent write #2289
base: main
Are you sure you want to change the base?
Conversation
@aersam I think we should max out the concurrent streams for python users. In most use cases we are passing a recordBatchReader where the recordBatches are already in memory before constructing the reader, in that case you won't see any memory difference. And it wouldn't be different than the prior behavior since the reader was always collected. I also have one suggestion on the python side, I think it's better if we simplify it and just provide a parameter called |
How about |
@aersam that also works! :) |
I finally had the time to update this branch with the new parallel parameter in python. Hope it's looking good now! |
@aersam btw, did you have any profiling numbers on speed ups/memory trade offs when parallel is True. Would be nice to share those in the release notes later on |
I only did some manual test on my own data, but could probably write some benchmark in python, using duckdb or polars as source. Would it make sense to add this to the code somehow? |
@aersam here you could add it, and even maybe reuse some of the benchmarks there: https://github.com/delta-io/delta-rs/tree/main/crates/benchmarks |
I did some very basic benchmarking, but the results were not as I hoped :) While RAM consumption is significantly lower, the speed is not good enough yet. I think maybe the channel must be bigger, I'll do some more testing I did my test quick and dirty using python, I can share the code if you want. Basically it's this: import duckdb
from deltalake.writer import write_deltalake
from uuid import uuid4
with duckdb.connect() as con: # get your 42.parquet here: https://duckdb.org/2024/03/26/42-parquet-a-zip-bomb-for-the-big-data-age.html
con.execute("select b, random() as a from read_parquet('42.parquet') limit 300000000")
reader = con.fetch_record_batch()
write_deltalake(f"_test/{uuid4()}", reader, schema=reader.schema, mode="overwrite", engine="rust") |
Pretty sure the non-async write causes issues. But object_store 0.10 will change a lot there, so maybe better to wait for that |
Yes let's see how effective these changes are with new upload trait |
@aersam fyi, ObjectStore just got bumped in the repo to 0.10 |
@aersam hey, do you think you have time to resolve the merge conflicts? |
I'm sorry, I had a bit a shift in priorities, so it will take time to do so. Especially since there were quite some changes in the writer as I see |
No worries! Just ping me once it's ready for another review round |
Description
This is a followup of #2265
It additionally uses streams/channels to concurrently write at the cost of more memory consumption. Default is keeping one recordbatch in RAM only, so it's opt-in.
I tested this with a local file and it went from 700s to 200s if I work with 10 concurrent streams. Of course memory consumption goes up, but given that we currently load the whole table in RAM, it's OK :)
This adds a depenency on async-channel as I need a multi-consumer channel.
Related Issue(s)
Fixes #2255