Commit#

Interface#

class FileStoreCommit#

Interface for commit operations in a file store.

The FileStoreCommit class provides interfaces for committing changes, expiring old snapshots, dropping partitions, and retrieving commit metrics.

Note

Direct file-system commits support append-only and primary-key tables on non-object-store paths. CommitContextBuilder::WithCatalog() also supports object stores when the catalog manages versions. UseRESTCatalogCommit() only prepares a request for the caller to send.

Public Functions

virtual ~FileStoreCommit() = default#
virtual Status Commit(const std::vector<std::shared_ptr<CommitMessage>> &commit_messages, int64_t commit_identifier = BATCH_WRITE_COMMIT_IDENTIFIER, std::optional<int64_t> watermark = std::nullopt) = 0#

Commit changes to the file store.

Parameters:
  • commit_messages – A vector of commit messages to be committed.

  • commit_identifier – An optional identifier for the commit operation. Default is BATCH_WRITE_COMMIT_IDENTIFIER.

  • watermark – An optional event-time watermark used to indicate the progress of data processing. Default is std::nullopt.

Returns:

Status indicating the success or failure of the commit operation.

virtual Result<int64_t> CommitWithProgress(const std::vector<RealtimeCommitProgress> &realtime_commits, int64_t commit_identifier, std::optional<int64_t> watermark) = 0#

Commit sealed real-time segments and persist their partition-bucket offset progress.

Entries for each partition-bucket must be ordered and non-overlapping after the offset recorded by the latest committed snapshot. Input entries may be unordered; this method orders them by partition, bucket, and offset before validation. Offset gaps are allowed. The resulting snapshot atomically publishes the data files and the updated offset map.

For each partition-bucket, the upstream coordinator must include the complete prefix of prepared-but-uncommitted entries through the requested progress. Because offsets may be sparse, this method cannot distinguish a valid offset gap from an omitted prepared entry. Omitting an earlier entry may advance committed progress past unpublished files and allow their real-time data to be reclaimed.

Snapshot conflicts are retried internally using the configured commit retry limit, timeout, and backoff. Each attempt reloads the latest snapshot and rebases both file changes and offset progress. A retry succeeds idempotently when both the identifier and all requested ranges are already committed. Inconsistent identifiers, overlapping offset progress, and file or index conflicts fail without further retry.

An error is terminal for the writer state which produced realtime_commits. The caller must discard its RealtimeContext and FileStoreWrite, load the current latest snapshot’s durable offsets, recreate both objects, and replay input from those exclusive offsets. Conflicts from CommitContextBuilder::WithCatalog() are retried internally. With UseRESTCatalogCommit(), the caller sends the request and handles server conflicts. Concurrent rollback or partition deletion must be fenced by the upstream coordinator. Configure catalog writers with WriteContextBuilder::WithCatalog() too.

Parameters:
  • realtime_commits – Commit messages and left-closed, right-open offset ranges to commit.

  • commit_identifier – Identifier of the streaming commit operation.

  • watermark – Optional event-time watermark.

Returns:

The id of the latest snapshot containing the committed progress. On retry, this may be a snapshot produced by a later commit and is suitable for refreshing a real-time context.

virtual Result<int32_t> FilterAndCommit(const std::map<int64_t, std::vector<std::shared_ptr<CommitMessage>>> &commit_identifier_and_messages, std::optional<int64_t> watermark = std::nullopt) = 0#

Filter out all std::vector<CommitMessage> which have been committed and commit the remaining ones.

Compared to commit, this method will first check if a commit_identifier has been committed, so this method might be slower. A common usage of this method is to retry the commit process after a failure.

Parameters:
  • commit_identifier_and_messages – A map containing all CommitMessages in question. The key is the commit_identifier.

  • watermark – An optional event-time watermark used to indicate the progress of data processing. Default is std::nullopt.

Returns:

Number of std::vector<CommitMessage> committed.

virtual Status Overwrite(const std::map<std::string, std::string> &partition, const std::vector<std::shared_ptr<CommitMessage>> &commit_messages, int64_t commit_identifier, std::optional<int64_t> watermark = std::nullopt) = 0#

Overwrite from manifest committable and partition.

Note

A full-table overwrite clears all committed real-time progress. A partition overwrite removes progress only for matching partitions. In either case, active real-time writers and their RealtimeContext instances must be recreated before further real-time operations.

Parameters:
  • partition – A single partition maps each partition key to a partition value. Depending on the user-defined statement, the partition might not include all partition keys. Also note that this partition does not necessarily equal to the partitions of the newly added key-values. This is just the partition to be cleaned up.

  • commit_messages – Description of the commit messages.

  • commit_identifier – Unique identifier.

  • watermark – An optional event-time watermark used to indicate the progress of data processing. Default is std::nullopt.

Returns:

Result of the operation.

virtual Result<int32_t> FilterAndOverwrite(const std::map<std::string, std::string> &partition, const std::vector<std::shared_ptr<CommitMessage>> &commit_messages, int64_t commit_identifier, std::optional<int64_t> watermark = std::nullopt) = 0#

This is a temporary interface for internal use.

It will be removed in a future version. Please do not rely on it for long-term use.

Note

A full-table overwrite clears all committed real-time progress. A partition overwrite removes progress only for matching partitions. In either case, active real-time writers and their RealtimeContext instances must be recreated before further real-time operations.

Parameters:
  • partition – Description of the partition.

  • commit_messages – Description of the commit messages.

  • commit_identifier – Unique identifier.

  • watermark – An optional event-time watermark used to indicate the progress of data processing. Default is std::nullopt.

Returns:

Result of the operation.

virtual Result<std::string> GetLastCommitTableRequest() = 0#

Returns the request from the latest commit attempt, including failed attempts.

Each attempt clears the request of the one before it, so an attempt that failed before building one, as one naming a branch a catalog cannot address does, leaves this returning an error. Catalog commits send requests automatically; UseRESTCatalogCommit() only prepares them.

Note

Temporary interface for internal use, will be removed in the future.

Returns:

JSON with tableId, baseSnapshotUuid, snapshot and statistics. Unset tableId and an absent or legacy base snapshot UUID are serialized as null.

virtual Result<int32_t> Expire() = 0#

Expire old snapshot in the file store.

Protects files referenced by retained snapshots, including files restored by rollback. Catalog commits require the current snapshot and retained history to be published to the file system before deletion. An unpublished or mismatched current snapshot skips expiration; catalog or retained-metadata read errors propagate without deleting files. Retained snapshots with index manifests return NotImplemented before deletion.

Note

Coordinate rollback and expiration so they do not run concurrently.

Note

Returns NotImplemented before any snapshot is read or any file is deleted when this commit is on a branch other than main, or the table path holds such a branch under branch/branch-<name>: branches share the table’s data files, while expiration only reads the retained snapshots of its own branch. A branch held only by a catalog, or created while expiration runs, is not found; do not expire a table with such a branch, and serialize branch creation and expiration.

Returns:

Result<int32_t> indicating the number of expired items or an error status.

virtual Status DropPartition(const std::vector<std::map<std::string, std::string>> &partitions, int64_t commit_identifier) = 0#

Drop specified partitions from the file store.

Note

A partition drop removes committed real-time progress only for matching partitions. Active real-time writers and their RealtimeContext instances must be recreated before further real-time operations.

Parameters:
  • partitions – A vector of partitions to be dropped.

  • commit_identifier – An identifier for the commit operation.

Returns:

Status indicating the success or failure of the drop partition operation.

virtual Status TruncateTable(int64_t commit_identifier) = 0#

Truncate the whole table by overwriting all partitions with empty data.

The generated snapshot has commit kind OVERWRITE.

Note

Truncation clears all committed real-time progress. Active real-time writers and their RealtimeContext instances must be recreated before further real-time operations.

Parameters:

commit_identifier – An identifier for the commit operation.

Returns:

Status indicating the success or failure of the truncate operation.

virtual Status Abort(const std::vector<std::shared_ptr<CommitMessage>> &commit_messages) = 0#

Abort an unsuccessful commit.

The data and index files described by the given commit messages will be deleted on a best-effort basis (delete failures are ignored).

Parameters:

commit_messages – A vector of commit messages whose files should be cleaned up.

Returns:

Status indicating the success or failure of the abort operation.

virtual Result<bool> RollbackToAsLatest(int64_t target_snapshot_id) = 0#

Roll back to the target snapshot and materialize it as the latest snapshot.

Reads the surviving files of both the current latest snapshot and the target snapshot, then commits an OVERWRITE snapshot whose visible state equals the target.

Note

Rollback restores the real-time progress recorded by the target snapshot. Active real-time writers and their RealtimeContext instances must be recreated before further real-time operations. Coordinate rollback with expiration so the target’s files cannot be deleted while the rollback is being prepared or committed.

Parameters:

target_snapshot_id – The snapshot id to roll back to.

Returns:

Result<bool>; true if the atomic commit succeeded. Returns an error status if there is no latest snapshot or the target snapshot does not exist.

virtual FileStoreCommit &RowIdCheckConflict(std::optional<int64_t> row_id_check_from_snapshot) = 0#

Configure row-id conflict checking from a specific snapshot id.

If set to a snapshot id, commit conflict detection will additionally validate row-id conflicts against snapshots after that id. Passing std::nullopt disables this behavior.

Parameters:

row_id_check_from_snapshot – Snapshot id to start row-id conflict checks from, or std::nullopt to disable.

Returns:

Current commit object for chaining.

virtual std::shared_ptr<Metrics> GetCommitMetrics() const = 0#

Retrieve metrics related to commit operations.

Returns:

A shared pointer to a Metrics object containing commit metrics.

Public Static Functions

static Result<std::unique_ptr<FileStoreCommit>> Create(std::unique_ptr<CommitContext> context)#

Create an instance of FileStoreCommit.

Parameters:

context – A unique pointer to the CommitContext used for commit operations.

Returns:

A Result containing a unique pointer to the FileStoreCommit instance.

class CommitContextBuilder#

CommitContextBuilder used to build a CommitContext, has input validation.

Public Functions

CommitContextBuilder(const std::string &root_path, const std::string &commit_user)#

Constructs a CommitContextBuilder with required parameters.

Parameters:
  • root_path – The root path of the Paimon table.

  • commit_user – The user identifier for the commit operation.

explicit CommitContextBuilder(const std::shared_ptr<FormatTable> &table)#

Constructs a CommitContextBuilder for a format table that is already loaded: the only way to commit to one whose schema lives in a metastore rather than under its location, such as a table a REST catalog serves.

The table carries what such a location does not say, so WithFileSystem() is refused here rather than ignored.

There is no commit user: a format table keeps no snapshot to record one in.

Parameters:

table – The format table to commit to, as Catalog::GetFormatTable() hands it back.

~CommitContextBuilder()#
CommitContextBuilder &SetOptions(const std::map<std::string, std::string> &options)#

Set a configuration options map to set some option entries which are not defined in the table schema or whose values you want to overwrite.

Note

The options map will clear the options added by AddOption() before.

Parameters:

options – The configuration options map.

Returns:

Reference to this builder for method chaining.

CommitContextBuilder &AddOption(const std::string &key, const std::string &value)#

Add a single configuration option which is not defined in the table schema or whose value you want to overwrite.

If you want to add multiple options, call AddOption() multiple times or use SetOptions() instead.

Parameters:
  • key – The option key.

  • value – The option value.

Returns:

Reference to this builder for method chaining.

CommitContextBuilder &IgnoreEmptyCommit(bool ignore_empty_commit)#

Sets whether to ignore empty commits (default is true).

When set to true, commits that don’t contain any actual data changes will be ignored.

Parameters:

ignore_empty_commit – True to ignore empty commits, false otherwise.

Returns:

Reference to this builder for method chaining.

CommitContextBuilder &UseRESTCatalogCommit(bool use_rest_catalog_commit)#

Builds a REST commit request for the caller to send (default is false).

Use WithCatalog() for catalog-managed commits; the two modes are mutually exclusive.

Note

Temporary interface, will be removed in the future.

Parameters:

use_rest_catalog_commit – True to build the request without sending it.

Returns:

Reference to this builder for method chaining.

CommitContextBuilder &WithCatalog(const std::shared_ptr<Catalog> &catalog, const Identifier &identifier)#

Loads the current schema from catalog.

Catalogs with version management also load and commit snapshots, rebasing on conflicts. Other catalogs use file-system snapshot commits. Mutually exclusive with UseRESTCatalogCommit().

Data and manifests remain on the file system. Recovery, rollback and row-id checks require historical snapshots there; row-id checks also require the inspected data files’ schemas. Expire() manages file-system snapshots after confirming that the catalog’s current snapshot and retained history are published there. Configure writers with the same catalog via WriteContextBuilder::WithCatalog(). A branch is addressed by the identifier, as tbl$branch_dev, so that the catalog answers for that branch. Its schema is read from the branch directory rather than from the catalog, just as a read of that branch reads it. Configure the writer with the same identifier.

Parameters:
  • catalog – Non-null catalog, kept alive by this context.

  • identifier – Table to commit to, naming a branch of it as tbl$branch_dev.

Returns:

Reference to this builder for method chaining.

CommitContextBuilder &WithTableId(const std::string &table_id)#

Sets the catalog table UUID sent as tableId to detect table recreation (default is unset).

Capture Table::CatalogUuid() before preparing changes. Table::Uuid() may fall back to the table name. An unset ID is sent as null and validated by the server. Requires WithCatalog() or UseRESTCatalogCommit().

Parameters:

table_id – Catalog UUID of the table to commit to.

Returns:

Reference to this builder for method chaining.

CommitContextBuilder &AppendCommitCheckConflict(bool append_commit_check_conflict)#

Sets whether append commits should perform conflict checking (default is false).

Parameters:

append_commit_check_conflict – True to enable append conflict checks.

Returns:

Reference to this builder for method chaining.

CommitContextBuilder &WithMemoryPool(const std::shared_ptr<MemoryPool> &memory_pool)#

Sets the memory pool to be used for memory allocation during commit operations.

Parameters:

memory_pool – Shared pointer to the memory pool instance.

Returns:

Reference to this builder for method chaining.

CommitContextBuilder &WithExecutor(const std::shared_ptr<Executor> &executor)#

Sets the executor to be used for asynchronous operations during commit.

Parameters:

executor – Shared pointer to the executor instance.

Returns:

Reference to this builder for method chaining.

CommitContextBuilder &WithFileSystem(const std::shared_ptr<FileSystem> &file_system)#

Sets a custom file system instance to be used for all file operations in this commit context.

This bypasses the global file system registry and uses the provided implementation directly.

Note

If not set, use default file system (configured in Options::FILE_SYSTEM)

Parameters:

file_system – The file system to use.

Returns:

Reference to this builder for method chaining.

Result<std::unique_ptr<CommitContext>> Finish()#

Build and return a CommitContext instance with input validation.

Returns:

Result containing the constructed CommitContext or an error status.

class CommitContext#

CommitContext is some configuration for commit operations.

Please do not use this class directly, use CommitContextBuilder to build a CommitContext which has input validation.

Public Functions

CommitContext(const std::string &root_path, const std::string &commit_user, bool ignore_empty_commit, bool use_rest_catalog_commit, bool append_commit_check_conflict, const std::shared_ptr<MemoryPool> &memory_pool, const std::shared_ptr<Executor> &executor, const std::shared_ptr<FileSystem> &specific_file_system, const std::map<std::string, std::string> &options, const std::shared_ptr<FormatTable> &format_table)#

Retained for source compatibility. Prefer CommitContextBuilder for input validation.

CommitContext(const std::string &root_path, const std::string &commit_user, bool ignore_empty_commit, bool use_rest_catalog_commit, const std::shared_ptr<Catalog> &catalog, const std::optional<Identifier> &identifier, const std::optional<std::string> &table_id, bool append_commit_check_conflict, const std::shared_ptr<MemoryPool> &memory_pool, const std::shared_ptr<Executor> &executor, const std::shared_ptr<FileSystem> &specific_file_system, const std::map<std::string, std::string> &options, const std::shared_ptr<FormatTable> &format_table)#
~CommitContext()#
inline const std::string &GetRootPath() const#
inline const std::string &GetCommitUser() const#
inline bool IgnoreEmptyCommit() const#
inline bool UseRESTCatalogCommit() const#
inline const std::shared_ptr<Catalog> &GetCatalog() const#

Returns the configured catalog, or null if unset.

inline const std::optional<Identifier> &GetIdentifier() const#

Returns the table identifier supplied with the catalog.

inline const std::optional<std::string> &GetTableId() const#

Returns the catalog table UUID sent as tableId in commit requests.

inline bool AppendCommitCheckConflict() const#
inline std::shared_ptr<MemoryPool> GetMemoryPool() const#
inline std::shared_ptr<Executor> GetExecutor() const#
inline std::shared_ptr<FileSystem> GetSpecificFileSystem() const#
inline const std::map<std::string, std::string> &GetOptions() const#
inline const std::shared_ptr<FormatTable> &GetFormatTable() const#

The format table this context was built from, or null when it names a table path and the schema under that path says what kind of table it is.

class CommitMessage#

Commit message for partition and bucket.

Supports serialization and deserialization compatible with Java Paimon.

Note

Serialized payloads do not embed their serialization version. Transport CurrentVersion() alongside the payload and pass it explicitly to Deserialize() or DeserializeList().

Public Functions

virtual ~CommitMessage()#

Public Static Functions

static int32_t CurrentVersion()#

Gets the version with which this serializer serializes.

static Result<std::string> Serialize(const std::shared_ptr<CommitMessage> &commit_message, const std::shared_ptr<MemoryPool> &pool)#

Serializes a single commit message to a binary string format.

The serialized format is compatible with Java Paimon. The serialization version is not included in the returned payload.

Parameters:
  • commit_message – The commit message to serialize.

  • pool – Memory pool for memory allocation during serialization.

Returns:

Result containing the serialized string data, or an error if serialization fails.

static Result<std::string> SerializeList(const std::vector<std::shared_ptr<CommitMessage>> &commit_messages, const std::shared_ptr<MemoryPool> &pool)#

Serializes a list of commit messages to a binary string format.

The serialization version is not included in the returned payload.

Parameters:
  • commit_messages – Vector of commit messages to serialize.

  • pool – Memory pool for memory allocation during serialization.

Returns:

Result containing the serialized string data, or an error if serialization fails.

static Result<std::shared_ptr<CommitMessage>> Deserialize(int32_t version, const char *buffer, int32_t length, const std::shared_ptr<MemoryPool> &pool)#

Deserializes a single commit message from binary data.

Parameters:
  • version – The serialization format version used when the data was serialized.

  • buffer – Pointer to the binary data buffer.

  • length – Length of the binary data in bytes.

  • pool – Memory pool for memory allocation during deserialization.

Returns:

Result containing the deserialized CommitMessage, or an error if deserialization fails.

static Result<std::vector<std::shared_ptr<CommitMessage>>> DeserializeList(int32_t version, const char *buffer, int32_t length, const std::shared_ptr<MemoryPool> &pool)#

Deserializes a list of commit messages from binary data.

This is the counterpart to SerializeList() for batch processing.

Parameters:
  • version – The serialization format version used when the data was serialized.

  • buffer – Pointer to the binary data buffer.

  • length – Length of the binary data in bytes.

  • pool – Memory pool for memory allocation during deserialization.

Returns:

Result containing a vector of deserialized CommitMessages, or an error if deserialization fails.

static Result<std::string> ToDebugString(const std::shared_ptr<CommitMessage> &commit_message)#

Converts a commit message to a human-readable debug string.

This is useful for logging, debugging, and troubleshooting purposes.

Parameters:

commit_message – The commit message to convert to string.

Returns:

Result containing the debug string representation, or an error if conversion fails.