|
iceberg-cpp
|
API for overwriting files in a table by partition. More...
#include <replace_partitions.h>
Public Member Functions | |
| ReplacePartitions & | AddFile (const std::shared_ptr< DataFile > &file) |
| Add a data file to the table. | |
| ReplacePartitions & | ValidateAppendOnly () |
| Validate that no partitions will be replaced and the operation is append-only. | |
| ReplacePartitions & | ValidateFromSnapshot (int64_t snapshot_id) |
| Set the snapshot ID used in validations for this operation. | |
| ReplacePartitions & | ValidateNoConflictingData () |
| Enable validation that data added concurrently does not conflict with this operation. | |
| ReplacePartitions & | ValidateNoConflictingDeletes () |
| Enable validation that deletes that happened concurrently do not conflict with this operation. | |
| std::string | operation () override |
| A string that describes the action that produced the new snapshot. | |
Public Member Functions inherited from iceberg::SnapshotUpdate | |
| Kind | kind () const override |
| Return the kind of this pending update. | |
| bool | IsRetryable () const override |
| Whether this update can be retried after a commit conflict. | |
| Status | Commit () override |
| Apply the pending changes and commit. | |
| auto & | ReportWith (this auto &self, std::shared_ptr< MetricsReporter > reporter) |
| Set the metrics reporter for this snapshot update. | |
| auto & | DeleteWith (this auto &self, std::function< Status(const std::string &)> delete_func) |
| Set a callback to delete files instead of the table's default. | |
| auto & | StageOnly (this auto &self) |
| Stage a snapshot in table metadata, but do not make it current. | |
| auto & | ScanManifestsWith (this auto &self, Executor &executor) |
| Configure an executor for manifest planning work. | |
| auto & | ToBranch (this auto &self, const std::string &branch) |
| Perform operations on a particular branch. | |
| auto & | Set (this auto &self, const std::string &property, const std::string &value) |
| Set a summary property in the snapshot produced by this update. | |
| auto & | WriteManifestsWith (this auto &self, Executor &executor, int32_t parallelism) |
| Configure an executor and max writer count for writing new manifests. | |
| Result< ApplyResult > | Apply () |
| Apply the update's changes to create a new snapshot. | |
| Status | Finalize (Result< const TableMetadata * > commit_result) override |
| Finalize the snapshot update, cleaning up any uncommitted files. | |
Public Member Functions inherited from iceberg::PendingUpdate | |
| PendingUpdate (const PendingUpdate &)=delete | |
| PendingUpdate & | operator= (const PendingUpdate &)=delete |
| PendingUpdate (PendingUpdate &&) noexcept=default | |
| PendingUpdate & | operator= (PendingUpdate &&) noexcept=default |
Public Member Functions inherited from iceberg::ErrorCollector | |
| ErrorCollector (ErrorCollector &&)=default | |
| ErrorCollector & | operator= (ErrorCollector &&)=default |
| ErrorCollector (const ErrorCollector &)=default | |
| ErrorCollector & | operator= (const ErrorCollector &)=default |
| template<typename... Args> | |
| auto & | AddError (this auto &self, ErrorKind kind, const std::format_string< Args... > fmt, Args &&... args) |
| Add a specific error and return reference to derived class. | |
| auto & | AddError (this auto &self, Error err) |
| Add an existing error object and return reference to derived class. | |
| auto & | AddError (this auto &self, std::unexpected< Error > err) |
| Add an unexpected result's error and return reference to derived class. | |
| bool | has_errors () const |
| Check if any errors have been collected. | |
| size_t | error_count () const |
| Get the number of errors collected. | |
| Status | CheckErrors () const |
| Check for accumulated errors and return them if any exist. | |
| void | ClearErrors () |
| Clear all accumulated errors. | |
| const std::vector< Error > & | errors () const |
| Get read-only access to all collected errors. | |
Static Public Member Functions | |
| static Result< std::unique_ptr< ReplacePartitions > > | Make (std::string table_name, std::shared_ptr< TransactionContext > ctx) |
| Create a new ReplacePartitions instance. | |
Protected Member Functions | |
| Status | Validate (const TableMetadata ¤t_metadata, const std::shared_ptr< Snapshot > &snapshot) override |
| Validate the current metadata. | |
| Result< std::vector< ManifestFile > > | Apply (const TableMetadata ¤t_metadata, const std::shared_ptr< Snapshot > &snapshot) override |
| Apply changes, translating the fail-any-delete error. | |
Protected Member Functions inherited from iceberg::MergingSnapshotUpdate | |
| Status | CleanUncommitted (const std::unordered_set< std::string > &committed) override |
| Clean up any uncommitted manifests that were created. | |
| std::unordered_map< std::string, std::string > | Summary () override |
| Get the summary map for this operation. | |
| MergingSnapshotUpdate (std::string table_name, std::shared_ptr< TransactionContext > ctx) | |
| Constructor; reads merge configuration from table properties. | |
| Status | AddDataFile (std::shared_ptr< DataFile > file) |
| Stage a data file to be added to the table. | |
| Status | AddDeleteFile (std::shared_ptr< DataFile > file) |
| Stage a delete file to be added to the table. | |
| Status | ValidateNewDeleteFile (const TableMetadata &metadata, const DataFile &file) |
| Validate a delete file against the table format version rules. | |
| Status | AddDeleteFile (std::shared_ptr< DataFile > file, int64_t data_sequence_number) |
| Stage a delete file with an explicit data sequence number. | |
| Status | AddManifest (ManifestFile manifest) |
| Add all files in a pre-existing data manifest to the new snapshot. | |
| Status | DeleteDataFile (std::shared_ptr< DataFile > file) |
| Register a data file (by object) to be deleted from the table. | |
| Status | DeleteDeleteFile (std::shared_ptr< DataFile > file) |
| Register a delete file (by object) to be removed from the table. | |
| Status | DeleteByPath (std::string_view path) |
| Register a data file path to be deleted from the table. | |
| Status | DeleteByRowFilter (std::shared_ptr< Expression > expr) |
| Register an expression to delete matching rows. | |
| Status | DropPartition (int32_t spec_id, PartitionValues partition) |
| Register a partition to be dropped. | |
| void | FailMissingDeletePaths () |
| Fail if any registered delete path is not found in any manifest. | |
| void | FailAnyDelete () |
| Fail if any manifest entry matches a delete condition. | |
| void | SetNewDataFilesDataSequenceNumber (int64_t sequence_number) |
| Override the data sequence number assigned to all newly-added data files. | |
| bool | HasDataSequenceNumber () const |
| Returns true if SetNewDataFilesDataSequenceNumber was called. | |
| void | CaseSensitive (bool case_sensitive) |
| Set case sensitivity for row filter and expression evaluation. | |
| bool | IsCaseSensitive () const |
| Returns true if case-sensitive matching is enabled (default: true). | |
| bool | AddsDataFiles () const |
| Returns true if any data files have been staged for addition. | |
| bool | AddsDeleteFiles () const |
| Returns true if any delete files have been staged for addition. | |
| bool | DeletesDataFiles () const |
| Returns true if any data files have been registered for deletion. | |
| bool | DeletesDeleteFiles () const |
| Returns true if any delete files have been registered for removal. | |
| const std::shared_ptr< Expression > & | RowFilter () const |
| Returns the row-filter expression set via DeleteByRowFilter, or nullptr. | |
| Result< std::shared_ptr< PartitionSpec > > | DataSpec () const |
| Returns the single partition spec for all staged data files. | |
| std::vector< std::shared_ptr< DataFile > > | AddedDataFiles () const |
| Returns all data files staged for addition. | |
| Status | ValidateNoNewDeletesForDataFiles (const TableMetadata &metadata, std::optional< int64_t > starting_snapshot_id, const DataFileSet &replaced_files, const std::shared_ptr< Snapshot > &parent, std::shared_ptr< FileIO > io) const |
| Return an error if any snapshot after starting_snapshot_id, or from the beginning if unset, added a delete file that covers a file in replaced_files. | |
| Status | ValidateAddedDVs (const TableMetadata &metadata, std::optional< int64_t > starting_snapshot_id, std::shared_ptr< Expression > conflict_filter, const std::shared_ptr< Snapshot > &parent, std::shared_ptr< FileIO > io) const |
| Return an error if a staged deletion vector conflicts with a deletion vector added since starting_snapshot_id. | |
Protected Member Functions inherited from iceberg::SnapshotUpdate | |
| SnapshotUpdate (std::shared_ptr< TransactionContext > ctx) | |
| Result< std::vector< ManifestFile > > | WriteDataManifests (std::span< const std::shared_ptr< DataFile > > files, const std::shared_ptr< PartitionSpec > &spec, std::optional< int64_t > data_sequence_number=std::nullopt) |
| Write data manifests for the given data files. | |
| Result< std::vector< ManifestFile > > | WriteDeleteManifests (std::span< const ContentFileWithSequenceNumber > files, const std::shared_ptr< PartitionSpec > &spec) |
| const std::string & | target_branch () const |
| bool | can_inherit_snapshot_id () const |
| const std::string & | commit_uuid () const |
| int32_t | manifest_count () const |
| int32_t | attempt () const |
| int64_t | target_manifest_size_bytes () const |
| virtual bool | CleanupAfterCommit () const |
| Check if cleanup should happen after commit. | |
| int64_t | SnapshotId () |
| Get or generate the snapshot ID for the new snapshot. | |
| Status | DeleteFile (const std::string &path) |
| Delete a file at the given path. | |
| std::string | ManifestPath () |
| std::string | ManifestListPath () |
| SnapshotSummaryBuilder & | summary_builder () |
| SnapshotSummaryBuilder | BuildManifestCountSummary (std::span< const ManifestFile > manifests, int32_t replaced_manifests_count) |
Protected Member Functions inherited from iceberg::PendingUpdate | |
| PendingUpdate (std::shared_ptr< TransactionContext > ctx) | |
| const TableMetadata & | base () const |
Additional Inherited Members | |
Public Types inherited from iceberg::PendingUpdate | |
| enum class | Kind : uint8_t { kExpireSnapshots , kSetSnapshot , kUpdateLocation , kUpdatePartitionSpec , kUpdatePartitionStatistics , kUpdateProperties , kUpdateSchema , kUpdateSnapshot , kUpdateSnapshotReference , kUpdateSortOrder , kUpdateStatistics } |
Static Protected Member Functions inherited from iceberg::MergingSnapshotUpdate | |
| static Status | ValidateAddedDataFiles (const TableMetadata &metadata, std::optional< int64_t > starting_snapshot_id, std::shared_ptr< Expression > filter, const std::shared_ptr< Snapshot > &parent, std::shared_ptr< FileIO > io, bool case_sensitive=true) |
| Return an error if any snapshot after starting_snapshot_id, or from the beginning if unset, added a data file matching the given filter expression. | |
| static Status | ValidateAddedDataFiles (const TableMetadata &metadata, std::optional< int64_t > starting_snapshot_id, const PartitionSet &partition_set, const std::shared_ptr< Snapshot > &parent, std::shared_ptr< FileIO > io) |
| Return an error if any snapshot after starting_snapshot_id, or from the beginning if unset, added a data file in any partition of the given partition set. | |
| static Status | ValidateDataFilesExist (const TableMetadata &metadata, std::optional< int64_t > starting_snapshot_id, const std::unordered_set< std::string > &file_paths, bool skip_deletes, std::shared_ptr< Expression > filter, const std::shared_ptr< Snapshot > &parent, std::shared_ptr< FileIO > io, bool case_sensitive=true) |
| Return an error if any snapshot after starting_snapshot_id, or from the beginning if unset, removed a file whose path is in file_paths (and skip_deletes is false). | |
| static Status | ValidateNoNewDeletesForDataFiles (const TableMetadata &metadata, std::optional< int64_t > starting_snapshot_id, const DataFileSet &replaced_files, const std::shared_ptr< Snapshot > &parent, std::shared_ptr< FileIO > io, bool ignore_equality_deletes=false) |
| Return an error if any snapshot after starting_snapshot_id, or from the beginning if unset, added a delete file that covers a file in replaced_files. | |
| static Status | ValidateNoNewDeletesForDataFiles (const TableMetadata &metadata, std::optional< int64_t > starting_snapshot_id, std::shared_ptr< Expression > data_filter, const DataFileSet &replaced_files, const std::shared_ptr< Snapshot > &parent, std::shared_ptr< FileIO > io, bool case_sensitive=true) |
| Return an error if any snapshot after starting_snapshot_id, or from the beginning if unset, added a delete file matching the data filter that covers a file in replaced_files. | |
| static Status | ValidateNoNewDeleteFiles (const TableMetadata &metadata, std::optional< int64_t > starting_snapshot_id, std::shared_ptr< Expression > data_filter, const std::shared_ptr< Snapshot > &parent, std::shared_ptr< FileIO > io, bool case_sensitive=true) |
| Return an error if any snapshot after starting_snapshot_id, or from the beginning if unset, added a delete file matching the given row filter. | |
| static Status | ValidateNoNewDeleteFiles (const TableMetadata &metadata, std::optional< int64_t > starting_snapshot_id, const PartitionSet &partition_set, const std::shared_ptr< Snapshot > &parent, std::shared_ptr< FileIO > io) |
| Return an error if any snapshot after starting_snapshot_id, or from the beginning if unset, added a delete file matching any partition in the given partition set. | |
| static Status | ValidateDeletedDataFiles (const TableMetadata &metadata, std::optional< int64_t > starting_snapshot_id, std::shared_ptr< Expression > data_filter, const std::shared_ptr< Snapshot > &parent, std::shared_ptr< FileIO > io, bool case_sensitive=true) |
| Return an error if any snapshot after starting_snapshot_id, or from the beginning if unset, deleted a data file matching the given row filter. | |
| static Status | ValidateDeletedDataFiles (const TableMetadata &metadata, std::optional< int64_t > starting_snapshot_id, const PartitionSet &partition_set, const std::shared_ptr< Snapshot > &parent, std::shared_ptr< FileIO > io) |
| Return an error if any snapshot after starting_snapshot_id, or from the beginning if unset, deleted a data file in any partition of the given partition set. | |
| static Result< std::unique_ptr< DeleteFileIndex > > | AddedDeleteFiles (const TableMetadata &metadata, std::optional< int64_t > starting_snapshot_id, std::shared_ptr< Expression > data_filter, std::shared_ptr< PartitionSet > partition_set, const std::shared_ptr< Snapshot > &parent, std::shared_ptr< FileIO > io, bool case_sensitive=true) |
| Build a DeleteFileIndex of delete files added since starting_snapshot_id. | |
| static Status | ValidateAddedDVs (const TableMetadata &metadata, std::optional< int64_t > starting_snapshot_id, std::shared_ptr< Expression > conflict_filter, const std::unordered_set< std::string > &referenced_data_files, const std::shared_ptr< Snapshot > &parent, std::shared_ptr< FileIO > io, bool case_sensitive=true) |
| Return an error if any snapshot after starting_snapshot_id, or from the beginning if unset, added a deletion vector that conflicts with DVs being written. | |
Protected Attributes inherited from iceberg::SnapshotUpdate | |
| SnapshotSummaryBuilder | summary_ |
Protected Attributes inherited from iceberg::PendingUpdate | |
| std::shared_ptr< TransactionContext > | ctx_ |
Protected Attributes inherited from iceberg::ErrorCollector | |
| std::vector< Error > | errors_ |
API for overwriting files in a table by partition.
The default validation mode is idempotent, meaning the overwrite is correct and should be committed regardless of other concurrent changes to the table. Alternatively, call ValidateNoConflictingDeletes() and ValidateNoConflictingData() to ensure that no conflicting delete files or data files, respectively, have been written since the snapshot passed to ValidateFromSnapshot().
This API accumulates file additions and produces a new Snapshot by replacing all files in partitions containing new data with the new additions. This operation implements dynamic partition replacement.
When committing, these changes are applied to the latest table snapshot. Commit conflicts are resolved by applying the changes to the new latest snapshot and reattempting the commit.
| ReplacePartitions & iceberg::ReplacePartitions::AddFile | ( | const std::shared_ptr< DataFile > & | file | ) |
Add a data file to the table.
| file | A data file with partition_spec_id set |
|
overrideprotectedvirtual |
Apply changes, translating the fail-any-delete error.
Rewrites the delete error raised when ValidateAppendOnly() is set into the partition-conflict message.
Reimplemented from iceberg::MergingSnapshotUpdate.
|
static |
Create a new ReplacePartitions instance.
| table_name | The name of the table |
| ctx | The transaction context |
|
overridevirtual |
A string that describes the action that produced the new snapshot.
Implements iceberg::SnapshotUpdate.
|
overrideprotectedvirtual |
Validate the current metadata.
Child operations can override this to add custom validation.
| current_metadata | Current table metadata to validate |
| snapshot | Ending snapshot on the lineage which is being validated |
Reimplemented from iceberg::SnapshotUpdate.
| ReplacePartitions & iceberg::ReplacePartitions::ValidateAppendOnly | ( | ) |
Validate that no partitions will be replaced and the operation is append-only.
| ReplacePartitions & iceberg::ReplacePartitions::ValidateFromSnapshot | ( | int64_t | snapshot_id | ) |
Set the snapshot ID used in validations for this operation.
All validations check changes after this snapshot ID. If this is not called, validation occurs from the beginning of the table's history.
This method should be called before this operation is committed. If a concurrent operation committed a data or delete file or removed a data file after the given snapshot ID that might contain rows matching a partition marked for deletion, validation detects the conflict and fails.
| snapshot_id | A snapshot ID set to when this operation started to read the table |
| ReplacePartitions & iceberg::ReplacePartitions::ValidateNoConflictingData | ( | ) |
Enable validation that data added concurrently does not conflict with this operation.
Validating concurrent data files is required during non-idempotent replace partition operations. This checks whether a concurrent operation inserted data in any partition being overwritten and aborts the replace to avoid removing rows added concurrently.
| ReplacePartitions & iceberg::ReplacePartitions::ValidateNoConflictingDeletes | ( | ) |
Enable validation that deletes that happened concurrently do not conflict with this operation.
Validating concurrent deletes is required during non-idempotent replace partition operations. This checks whether a concurrent operation deleted data in any partition being overwritten and aborts the replace to avoid undeleting rows removed concurrently.