Package org.apache.iceberg.flink.maintenance.operator
package org.apache.iceberg.flink.maintenance.operator
-
ClassDescriptionCommits the rewrite changes using
RewriteDataFilesCommitManager.Plans the rewrite groups using theBinPackRewriteFilePlanner.Executes a rewrite for a singleDataFileRewritePlanner.PlannedGroup.Delete the files using theFileIOwhich implementsSupportsBulkOperations.A deletion vector position emitted byEqualityConvertPKIndexwhen resolving an equality delete.Result from theEqualityConvertDVWritercontaining new and rewritten DV files written by a single resolver task.Commits data files and DVs to the target branch.Keyed parallel resolver that buffersDVPositions per data-file path, then writes Puffin DV files directly viaBaseDVFileWriter.Parallel operator that maintains a shard of the primary key index.Result of equality convert planning.Planner for the equality delete conversion pipeline.Parallel reader that processesReadCommands from the planner and emitsIndexCommands.Calls theExpireSnapshotsto remove the old snapshots and emits the filenames which could be removed in theExpireSnapshotsProcessor.DELETE_STREAMside output.A specialized reader implementation that extracts file names from Iceberg table rows.A key selector implementation that extracts a normalized file path from a file URI string.Command from theEqualityConvertPlannerto theEqualityConvertPKIndex.Recursively lists the files in the `location` directory.Lists the metadata files referenced by the table.Event sent from TriggerManagerOperator to TriggerManagerCoordinator to register a lock release handler.Event sent from LockRemoverOperator to LockRemoverCoordinator to notify that a lock has been released.Manages locks and collectMetricfor the Maintenance Tasks.Plans the splits to read a metadata table content.Monitors an Iceberg table for changesA specialized co-process function that performs an anti-join between two streams of file URIs.Envelope from theEqualityConvertPlannerto theEqualityConvertReader, wrapping an IcebergContentScanTaskplus the metadata the reader needs to process it.Serialized primary key used as a Flink keyed state key.Implementation of the Source V2 API which uses an iterator to read the elements, and uses a single thread to do so.Skip file deletion processing when an error is encountered.Event describing changes in an Iceberg tableAggregates results of the operators for a given maintenance task.TriggerManager starts the Maintenance Tasks by emittingTriggermessages which are calculated based on the incomingTableChangemessages.