Class ConvertEqualityDeletes

java.lang.Object
org.apache.iceberg.flink.maintenance.api.ConvertEqualityDeletes

@Experimental public class ConvertEqualityDeletes extends Object
Creates the equality delete to DV conversion data stream. Runs a single iteration of the conversion for every Trigger event.

The pipeline reads equality delete files from a staging branch, converts them to deletion vectors (DVs) using a primary key index stored in Flink state, and commits the data files and DVs to the target branch.

The conversion is split into parallel stages:

  1. Planner (p=1): scans staging branch, emits file-level ReadCommands with phase timestamps
  2. Reader (p=N): reads files, emits row-level IndexCommands
  3. PKIndex (p=N): maintains PK index shards, resolves equality deletes to DV positions
  4. DVWriter (p=N, keyed by data file path): buffers positions per file, writes Puffin DVs inline
  5. Committer (p=1): commits data files and DVs to the target branch

Mutual exclusion with concurrent maintenance tasks (e.g. compaction) is enforced by the Flink maintenance framework lock.