Repository navigation
Upserting large table extremely slow #2159
Description
Activity
Haha I was just looking at this last week, it is especially slow for time series data with tens of millions of rows (that is a huuuuuge filter) -- I wonder if it would make sense for a user to be able to supply their own filter into the util? That was what enabled me to squeeze more performance out after hard coding it in
Being able to provide a "hint" seems like a decent workaround, but then we have to rely on the user providing the correct filter, otherwise the upsert won't work properly.
Ultimately I think upsert should be able to make use of the table's
PartitionSpecto determine what files to read given the table to be upserted. (This would only work when the row identifiers are also being used to partition but I believe that is a pretty natural use pattern).Reacted by Matt CorleyThis aligns well with the discussion here: #2138 (comment)
While there have been improvements to
upsert- like reducing memory pressure and avoiding recursion increate_match_filter- performance is still suboptimal. This is largely becauseupsertis currently built on top of operations likedeleteandoverwrite.For example,
overwriteinternally performs adeletefollowed by anappend. Thedeletestep may even rewrite entire data files by applying a filter and preserving rows that don’t match, then rewriting the result using the current schema, partition spec, and sort order.Interestingly, this "replace" behavior in
deleteis already very close to what anupsertneeds to do - especially when new rows are appended immediately after. Here's a simplified view of how files are rewritten:df = ArrowScan(...).to_table(...) filtered_df = df.filter(preserve_row_filter) if len(filtered_df) == 0: replaced_files.append((original_file.file, [])) elif len(df) != len(filtered_df): replaced_files.append((original_file.file, [...]))Because
upsertreuses these higher-level constructs, we lose the opportunity to optimize the operation at a lower level. Treatingupsertas a first-class primitive, likedeleteoroverwrite, would allow us to optimize each step more precisely and avoid unnecessary rewrites, which is especially important for large tables.cc @Fokko
Hey @koenvo thanks for raising this discussion. Nothing is set in stone, so there are always possibilities to optimize, and I agree, we started with rough building blocks.
The nice thing of the current approach is, when there is nothing to delete, it will only create an append operation. Also creating two snapshots (
DELETE+APPEND) instead of just aOVERWRITE, makes it more transparent to other clients to what happend to the table. Although this should not be at every expense. Currently, we actually produce:DELETE,APPENDfor the delete operationAPPENDfor the insert operation
Also, potentially generating quite a bit of data and metadata, so there is definitely room for improvement.
Here's a simplified view of how files are rewritten:
I don't think that's correct, because there we pull in all the data in memory, but we filter on only relevant files based on the Iceberg metadata.
I think @Anton-Tarazi's original point -- creating a bunch of (Python object) filter expressions for every row in a large dataframe is going to be slow, and we do that before any real i/o happens -- this is done so that, as Fokko says,
we filter on only relevant files based on the Iceberg metadata.
My guess is for tables of a certain size/layout you'd do better to just skip this metadata-based file pruning given the expensive python object creations.
Honestly, I think it would be a better use of community resources to invest more in the iceberg-rust/datafusion path so that the bulk of this logic can be moved out of pure python, than to spend a ton of effort trying to optimize the Python implementation.
Reacted by Jayce Slesar, Andre Luis Anastacio, Kevin Liu and Dan HomolaHonestly, I think it would be a better use of community resources to invest more in the iceberg-rust/datafusion path so that the bulk of this logic can be moved out of pure python, than to spend a ton of effort trying to optimize the Python implementation.
Agreed, I think if we spent more effort rolling wheels from the rust version and targeting a tighter relationship between pyiceberg and the rust implementation everyone wins
Reacted by Fokko Driesprong, Yingjian Wu, Andre Luis Anastacio and Lukas ValatkaTotally agree. Lets start exploring the iceberg-rust codebase
Reacted by Lukas Valatka@Fokko @kevinjqliu do you think its worth setting up a roadmap for what should be candidates for rolling wheels from rust? Would really help focus efforts on lacking parts of rust if we want parity!
Reacted by Matthias QI've created #2396 as I think that would already speed up things quite a bit.
Posted this on Slack as well but during my initial test with one month of data.
I read Parquet source files in batches of 128 MB (based on file size) and used .upsert() to write to iceberg. I tried with the table identifier as a single column (xxhash 128) and multiple columns (the actual business columns).
initial upsert (on empty table) 962,000 rows in 75 seconds
batch 2 upsert 967,000 rows in 256 seconds
batch 3 upsert 728,000 rows in 328 secondsThen I tried to process 1 day of data and I ran into an OOM error (with delta-rs this same test took 3 seconds to completely successfully).
Each batch took longer than the previous batch and everything took way too long. I am going to test if I can use partial overwrites to build a more specific predicate/expression.
Reacted by Lukas ValatkaUp! we have a peculiar case that upsert (of this package here) basically silently exits when rows are matched. Any ideas why? Probably OOM, but not certain.
Context: 27k rows. That's a miserable amount imho (we're talking big datalakehouse tooling here).
Hi, i’m coming from #2138 where I first wrote about this issue. I like the idea of leveraging rust-based plugins but I just wanted to throw a thought out there.
Is the Java version of Iceberg facing a similar problem ? I haven’t had the time to investigate, so maybe if anyone already knows about it and can provide some insights.
It feels like at least, with the presence of metadata, it should be able to do some physical pushdown predicates and locate relevant partitions in the case of a simple equality query (WHERE col=val).
Nope, Java version i.e. the one PySpark just flies... We fully pivoted to PySpark for our ETL data sync jobs into Lakehouse. We benchmarked Spark's MERGE INTO vs. equivalent PyIceberg upsert.
@astronautas I realize I might have been unclear about whom this message was addressed to. I was actually referring to the members who mentioned Rust as the go‑to solution. While I agree that would be the ideal scenario, I wanted to check whether the Java implementation of Iceberg i.e. https://github.com/apache/iceberg could serve as a useful source of inspiration, given that it seems to work well there. After spending a bit more time looking into it, it appears that @EnyMan has implemented exactly that in #2943. Let’s see how it goes!
Reacted by Lukas Valatka@astronautas I realize I might have been unclear about whom this message was addressed to. I was actually referring to the members who mentioned Rust as the go‑to solution. While I agree that would be the ideal scenario, I wanted to check whether the Java implementation of Iceberg i.e. https://github.com/apache/iceberg could serve as a useful source of inspiration, given that it seems to work well there. After spending a bit more time looking into it, it appears that @EnyMan has implemented exactly that in #2943. Let’s see how it goes!
Sorry to be blunt, but what's the point of getting inspiration? Python native implementations can never catch-up with Java / Rust / whatnot implementations in terms of performance. It means essentially rebuilding some query engine part, including sorts / joins / windows for selecting the latest value per key. Also that PR's been built by AI, so I am doubtful.
imho we should strive to integrate with them, and the Rust one seems most easy to interop - heck, tons of Python libs moved to Rust extensions for compute-intense parts.
@astronautas Please feel free to review:
Features in the Python implementation often follow the trends and logic established in the Java implementation.
I agree with the Rust interop direction. My point is simply that we shouldn’t discard interim efforts such as #2943, which has already been prepared.
The java was 10x faster on the same benchmark dataset 30s vs 3s just the upsert call (times differ from PR because of other optimizations i did especially to partitions and snapshot rewrite). The next slowest parts are the overwrite and append parts, i haven't looked at them but they could also be potentially improved getting us close to nice speeds.
On the other hand i dont think this needs to be light speed fast since we store data so reasonably fast should suffice. Don't get me wrong i am all for Rust where it makes sense. I mentioned the motivation behind why this was done in pure Python briefly in the PR.
Reacted by Mohamad Alhalabi and Lukas ValatkaI'm also very disappointed in
upsert()performance.. as it seems it doesn't use partition_filter at all! Am I missing something?
I have a large table (100m+ rows) partitioned bycompany_id,day(created_at)and each row then has a uniqueid..I would expect upserting with
join_cols=company_id,day(created_at),idto actually use partition filtering so it doesn't scan all files, and then be super fast.But it doesn't (am I missing something here?) and believe @EnyMan's #2943 doesn't address that?
@psantus isn't row_filter in the DataScan basically that? I think the issue you are encountering is the explosion of files across so many partitions and the complexity of the filter applied to that many files. This also applies to batches (row groups in the parquet files), that's why I tried to make the filters smaller and thus faster to execute on so many files.
Is there a recommended way to scale upserts from Python for large tables, i.e. via Rust bindings?
@ayushk7102, not really, you can try using my fork, but no guarantees. There is also this discussion that might shed some light: https://github.com/apache/iceberg-python/discussions/3118
Reacted by Ayush Kumar- added a commit that references this issue
on May 26, 2026 This issue has been automatically marked as stale because it has been open for 180 days with no activity. It will be closed in next 14 days if no further activity occurs. To permanently prevent this issue from being considered stale, add the label 'not-stale', but commenting on the issue is preferred when possible.
This issue has been closed because it has not received any activity in the last 14 days since being marked as 'stale'
Feature Request / Improvement
Feature Request / Improvement
Upserting large dataframes (tens of millions of rows) in un-usably slow due to creating a massive
BooleanExpressioninupsert_util.create_match_filter. This is before any IO is even started. It would be nice if upserts could support large tables.I'd be happy to work on this issue.