Skip to content

Upserting large table extremely slow #2159

Description

@Anton-Tarazi

Feature Request / Improvement

Feature Request / Improvement

Upserting large dataframes (tens of millions of rows) in un-usably slow due to creating a massive BooleanExpression in upsert_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.

Activity

  1. jayceslesar commented on Jun 28, 2025

    @jayceslesar
    Contributor

    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

  2. Anton-Tarazi commented on Jun 28, 2025

    @Anton-Tarazi
    ContributorAuthor

    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 PartitionSpec to 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).

  3. koenvo commented on Jul 2, 2025

    @koenvo
    Contributor

    This aligns well with the discussion here: #2138 (comment)

    While there have been improvements to upsert - like reducing memory pressure and avoiding recursion in create_match_filter - performance is still suboptimal. This is largely because upsert is currently built on top of operations like delete and overwrite.

    For example, overwrite internally performs a delete followed by an append. The delete step 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 delete is already very close to what an upsert needs 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 upsert reuses these higher-level constructs, we lose the opportunity to optimize the operation at a lower level. Treating upsert as a first-class primitive, like delete or overwrite, would allow us to optimize each step more precisely and avoid unnecessary rewrites, which is especially important for large tables.

    cc @Fokko

  4. Fokko commented on Jul 3, 2025

    @Fokko
    Contributor

    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 a OVERWRITE, 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, APPEND for the delete operation
    • APPEND for 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.

  5. corleyma commented on Jul 3, 2025

    @corleyma

    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.

  6. jayceslesar commented on Jul 3, 2025

    @jayceslesar
    Contributor

    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.

    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

  7. koenvo commented on Jul 3, 2025

    @koenvo
    Contributor

    Totally agree. Lets start exploring the iceberg-rust codebase

  8. jayceslesar commented on Jul 3, 2025

    @jayceslesar
    Contributor

    @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!

  9. Fokko commented on Aug 28, 2025

    @Fokko
    Contributor

    I've created #2396 as I think that would already speed up things quite a bit.

  10. ldacey commented on Oct 1, 2025

    @ldacey

    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 seconds

    Then 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.

  11. astronautas commented on Jan 23, 2026

    @astronautas

    Up! 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).

  12. mwa28 commented on Jan 28, 2026

    @mwa28

    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).

    @Fokko @jayceslesar @koenvo

  13. astronautas commented on Jan 29, 2026

    @astronautas

    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.

  14. mwa28 commented on Jan 29, 2026

    @mwa28

    @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!

  15. astronautas commented on Jan 29, 2026

    @astronautas

    @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.

  16. mwa28 commented on Jan 29, 2026

    @mwa28

    @astronautas Please feel free to review:

    #2166 (comment)

    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.

  17. EnyMan commented on Jan 29, 2026

    @EnyMan

    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.

  18. psantus commented on Feb 12, 2026

    @psantus

    I'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 by company_id,day(created_at) and each row then has a unique id..

    I would expect upserting with join_cols=company_id,day(created_at),id to 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?

  19. EnyMan commented on Feb 13, 2026

    @EnyMan

    @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.

  20. ayushk7102 commented on Mar 7, 2026

    @ayushk7102

    Is there a recommended way to scale upserts from Python for large tables, i.e. via Rust bindings?

  21. EnyMan commented on Mar 12, 2026

    @EnyMan

    @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

  22. github-actions commented on Sep 9, 2026

    @github-actions

    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.

  23. github-actions commented on Sep 24, 2026

    @github-actions

    This issue has been closed because it has not received any activity in the last 14 days since being marked as 'stale'

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions