apache / apache/iceberg-python

Upsert with 1M rows extremely slow due to `create_match_filter` and `txn.delete()` performance

Ouverte
#3,129 1 commentaire 0 réactions 0 personnes assignées Voir sur GitHub
Langage dominant
Python
Étoiles
1.1k
Forks
581
Merge moyen
1 j 17 h
PR mergées (30 j)
77

Description

### Apache Iceberg version

0.11.0

### Please describe the bug 🐞

CC @goutamvenkat-anyscale @koenvo @Fokko

Hello! We are implementing distributed writes from Ray Data to Iceberg. As part of upserts, we:

1. Write data files in parallel across Ray workers (each worker writes its share of Parquet files directly to storage and returns `DataFile` metadata + the upsert key columns back to the driver)
2. On the driver, concatenate all upsert keys collected from workers, call `create_match_filter` to build a delete predicate, then call `txn.delete()` followed by an append to commit

Upserting 1M rows (383 MiB) into an Iceberg table takes **~17.5 minutes**, almost entirely in the delete step:

```
create_match_filter (1M keys → In filter): 10.26s
txn.delete(): 1054.35s
append + commit: 1.14s
─────────────────────────────────────────────────────
Total upsert commit: 1065.75s
```

PyIceberg version `0.11.0`

This matches what's reported in #2159 and #2138.

The bottlenecks are:
1. **`create_match_filter`** — constructs a Python `BooleanExpression` node per row, which is expensive at 1M+ keys
2. **`txn.delete()`** — evaluates the resulting giant `In` expression against the table's data files with no partition pruning, effectively doing a full table scan

We have a few questions:

1. **Merge-on-read upserts** — is this on the roadmap, and if so, roughly when? MoR would let us avoid the expensive delete + rewrite cycle entirely for large upserts.
2. **Optimizing `create_match_filter` or `txn.delete()`** — is there a recommended way to speed these up today? For example, batching the `In` filter, or passing a partition-level hint to constrain the file scan?
3. **Partition-aware deletes** — if the upsert key columns overlap with partition columns, is there a supported way to restrict `txn.delete()` to only the relevant partitions, rather than scanning the full table?

## Related

- #2159 — Upserting large table extremely slow
- #2138 — Upsertion memory usage grows exponentially as table size grows
- #2943 — Optimize upsert performance for large datasets

### Willingness to contribute

- [ ] I can contribute a fix for this bug independently
- [x] I would be willing to contribute a fix for this bug with guidance from the Iceberg community
- [ ] I cannot contribute a fix for this bug at this time

Guide de contribution

Aucun guide de contribution indexé pour ce dépôt

Piste de recherche

Commencez par les points d’entrée create_match_filter et txn.delete() décrits dans le rapport, puis reproduisez le benchmark d’upsert de 1M lignes à partir des mesures fournies. Consultez les issues associées #2159, #2138 et #2943 pour prendre connaissance du contexte existant. Le travail est terminé lorsqu’une amélioration mesurée ou une méthode documentée et prise en charge pour éviter le coût signalé de la suppression de toute la table est disponible.

Rédigé par le modèle d'indexation à partir du texte de l'issue.

Évaluation

Stack technique
python
Domaine
data-engineering, databases
Type d'issue
Bug
Difficulté
5/5
Temps estimé
Plus d'une semaine
Activité
Calme
Clarté
À clarifier
Accessibilité débutants
35/100

Recevez les nouvelles issues par e-mail

Un résumé court des issues GitHub adaptées aux débutants.