Class ConvertEqualityDeletes
java.lang.Object
org.apache.iceberg.flink.maintenance.api.ConvertEqualityDeletes
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:
- Planner (p=1): scans staging branch, emits file-level ReadCommands with phase timestamps
- Reader (p=N): reads files, emits row-level IndexCommands
- PKIndex (p=N): maintains PK index shards, resolves equality deletes to DV positions
- DVWriter (p=N, keyed by data file path): buffers positions per file, writes Puffin DVs inline
- 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.
-
Nested Class Summary
Nested Classes -
Method Summary
-
Method Details
-
builder
-