Skip to content

Support Native Iceberg MOR write #6240

Description

@unikdahal

What is the problem?

Comet supports native Iceberg scans and data-file writes, but Iceberg V2 merge-on-read DELETE, UPDATE, and MERGE still execute the position-delta write through Iceberg's JVM WriteDelta path.

The goal of this issue is to make the executor-side Iceberg position-delta write native while keeping Iceberg planning, validation, and the final table commit on the Spark/Iceberg JVM side.

MergeRowsExec / CometMergeRowsExec
        |
        | operation + _file + _pos + _spec_id + _partition
        v
CometIcebergDeltaWriteExec
        |
        +-- data rows -----> iceberg-rust data writer
        |
        +-- deletes -------> iceberg-rust position-delete writer
        |
        v
Iceberg WriterCommitMessage
        |
        v
existing Iceberg JVM commit

The native path should support:

  • Iceberg format v2
  • DELETE, UPDATE, and MERGE
  • Parquet position deletes
  • PARTITION and FILE delete granularity
  • partition evolution / historical specs
  • replacement of existing file-scoped position deletes
  • JVM fallback whenever the native path cannot safely reproduce Iceberg behavior

Prerequisites

iceberg-rust

Tracked under apache/iceberg-rust#3287.

These keep generic Iceberg delete-file loading/rewrite behavior in iceberg-rust rather than reimplementing it in Comet.

Comet writer/runtime foundations

General native Iceberg writer productionization remains tracked by #5649.

Comet implementation

Proposed incremental split:

  • WriteDelta planning and JVM fallback

    • recognize Iceberg WriteDelta
    • extract operation/data/row-id/metadata projections
    • preserve existing Iceberg JVM execution as the fallback path
  • Native position-delta engine

    • route data rows and position deletes
    • load/rewrite existing file-scoped deletes
    • handle historical partition specs
    • produce native data/delete file metadata
  • Native WriteDelta integration

    • serde / task payload
    • CometIcebergDeltaWriteExec
    • reconstruct Iceberg writer commit messages
    • preserve JVM-side table commit
  • Spark-version integration

    • integrate with native MergeRowsExec where its Spark contract can be preserved
    • preserve JVM fallback for unsupported row-level execution shapes

MERGE compatibility

#5318 provides native MergeRowsExec independently of native Iceberg writes.

Its current compatibility boundary is intentionally:

  • Spark 3.5.x / 4.0.x — native MergeRowsExec
  • Spark 4.1+ — JVM fallback until action metrics and MergeSummary / V2 writer commit semantics are preserved end-to-end
  • Spark 4.2 InsertOnlyMergeExec — separate operator/follow-up

DELETE and UPDATE do not depend on native MergeRowsExec and can use the native delta writer independently.

Scope

This issue is limited to Iceberg V2 merge-on-read position-delta writes.

Out of scope:

Related

Activity

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

Metadata

Metadata

Assignees

Labels

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions