apache / apache/iceberg-rust

[Epic] Add MERGE INTO support for DataFusion integration

Open
#2,201 4 comments 3 reactions 0 assignees View on GitHub
epic
Dominant language
Rust
Stars
1.4k
Forks
567
Avg merge
2d 2h
Merged PRs (30d)
93

Description

### What's the feature are you trying to implement?

Add support for SQL `MERGE INTO` (UPSERT) operations in the iceberg-datafusion integration. This enables atomic row-level updates and inserts based on join conditions, essential for CDC pipelines, incremental updates, and data synchronization. I already have a [PoC branch](https://github.com/wirybeaver/iceberg-rust/commits/feature/merge-into/).

Refer to the [insert_into support](https://github.com/apache/iceberg-rust/issues/1540)

The Spark **SPJ** (Storage Partition Join) style is the key optimization I wanted to introduce. The Datafusion currently doesn't support `merge_into` sql parsing and logic plan yet. I am contributing the "MERGE INTO" in datafusion as well: https://github.com/apache/datafusion/issues/20746.

**SQL Example:**
```sql
MERGE INTO target_table t
USING source_table s
ON t.id = s.id
WHEN MATCHED THEN
UPDATE SET t.value = s.value
WHEN NOT MATCHED THEN
INSERT (id, value) VALUES (s.id, s.value)
```

---

### Query Execution Plan (CoW Mode)

#### Baseline Plan (Unpartitioned Table)

```
┌─────────────────────────────┐
│ IcebergMergeCommitExec │ Commits via RowDelta transaction
│ (add + remove data files) │ Outputs: record count
└──────────────┬──────────────┘

┌──────────────▼──────────────┐
│ CoalescePartitionsExec │ Merges all partitions into
│ │ single stream for atomic commit
└──────────────┬──────────────┘

┌──────────────▼──────────────┐
│ IcebergMergeWriteExec │ Writes merged rows to new Parquet
│ │ files via TaskWriter; tracks
│ │ _file values as deleted_files
└──────────────┬──────────────┘

┌──────────────▼──────────────┐
│ IcebergMergeExec │ FULL OUTER JOIN (HashJoinExec)
│ │ Classifies rows:
│ │ MATCHED → apply UPDATE exprs
│ │ NOT MATCHED → apply INSERT exprs
└──────┬──────────────┬───────┘
│ │
┌──────▼──────┐ ┌─────▼───────┐
│ IcebergTable│ │ Source Plan │
│ Scan │ │ (any exec) │
│ (target, │ │ │
│ with _file)│ │ │
└─────────────┘ └─────────────┘
```

#### SPJ-Optimized Plan (Partitioned Table)

When **all** partition columns appear in the join keys and use hash-compatible
transforms (Identity or Bucket), the optimizer wraps both sides with
repartitioning to eliminate cross-partition shuffles:

```
┌─────────────────────────────┐
│ IcebergMergeCommitExec │
└──────────────┬──────────────┘

┌──────────────▼──────────────┐
│ CoalescePartitionsExec │
└──────────────┬──────────────┘

┌──────────────▼──────────────┐
│ IcebergMergeWriteExec │
└──────────────┬──────────────┘

┌──────────────▼──────────────┐
│ IcebergMergeExec │
└──────┬──────────────┬───────┘
│ │
┌──────▼──────┐ ┌─────▼───────┐
│Repartition │ │Repartition │
│Exec (Hash │ │Exec (Hash │
│on _partition)│ │on _partition)│
└──────┬──────┘ └──────┬──────┘
│ │
┌──────▼──────┐ ┌──────▼──────┐
│Projection │ │Projection │
│Exec (adds │ │Exec (adds │
│_partition) │ │_partition) │
└──────┬──────┘ └──────┬──────┘
│ │
┌──────▼──────┐ ┌──────▼──────┐
│ IcebergTable│ │ Source Plan │
│ Scan │ │ (any exec) │
│(target, │ │ │
│ with _file) │ │ │
└─────────────┘ └─────────────┘
```

---

The following tasks are already completed on the PoC branch. Will raise formal PRs one after another as the fork repo doesn't support stacking PRs.
- [x] https://github.com/apache/iceberg-rust/pull/2203
- [ ] Add IcebergMergeExec with HashJoinExec integration and row classification
- [ ] Add IcebergMergeWriteExec and IcebergMergeCommitExec nodes
- [ ] Implement full MERGE execution logic with file tracking
- [ ] Integrate MERGE INTO into IcebergTableProvider
- [ ] Add comprehensive MERGE INTO integration tests
- [ ] Add partition-aware merge optimization (spark storage partition join style)

### Willingness to contribute

I would be willing to contribute to this feature with guidance from the Iceberg Rust community

Contributor guide

Open the contributing guide

Research direction

Start by reviewing the PoC branch and completed PR #2203, then trace the listed IcebergMergeExec, IcebergMergeWriteExec, IcebergMergeCommitExec, and IcebergTableProvider entry points. Use the planned MERGE INTO integration tests to verify row classification, file tracking, commits, and partition-aware optimization across the remaining checklist.

Written by the indexing model from the issue text.

Assessment

Tech stack
rust, sql
Domain
backend, databases
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Quiet
Clarity
Mostly clear
Newbie friendliness
25/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.