This is an automated email from the ASF dual-hosted git repository.
SteNicholas pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-cpp.git
The following commit(s) were added to refs/heads/main by this push:
new 70165ff1 feat(realtime): add commit retries, metrics, and PK predicate
pruning (#324)
70165ff1 is described below
commit 70165ff127b9129234c81deab745332d0e12cdfa
Author: lxy <[email protected]>
AuthorDate: Fri Sep 11 11:43:27 2026 +0800
feat(realtime): add commit retries, metrics, and PK predicate pruning (#324)
---
include/paimon/file_store_commit.h | 16 +-
include/paimon/realtime/realtime_context.h | 18 +
include/paimon/realtime/realtime_store.h | 17 +-
.../core/operation/abstract_file_store_write.cpp | 9 +-
.../append_only_file_store_write_test.cpp | 18 +
.../core/operation/file_store_commit_impl.cpp | 133 +++--
src/paimon/core/operation/file_store_commit_impl.h | 18 +-
.../core/operation/file_store_commit_impl_test.cpp | 90 +++
.../core/operation/key_value_file_store_write.cpp | 3 +-
.../operation/key_value_file_store_write_test.cpp | 14 +
.../core/operation/merge_file_split_read.cpp | 2 +-
src/paimon/core/realtime/arrow_realtime_store.cpp | 103 ++--
src/paimon/core/realtime/arrow_realtime_store.h | 15 +
.../core/realtime/arrow_realtime_store_test.cpp | 43 ++
.../realtime/primary_key_realtime_store_test.cpp | 33 +-
src/paimon/core/realtime/realtime_context_impl.cpp | 38 ++
src/paimon/core/realtime/realtime_context_impl.h | 4 +-
src/paimon/core/realtime/realtime_context_test.cpp | 40 +-
.../core/table/source/key_value_table_read.cpp | 18 +-
test/inte/realtime_write_inte_test.cpp | 607 ++++++++++++++++++++-
20 files changed, 1122 insertions(+), 117 deletions(-)
diff --git a/include/paimon/file_store_commit.h
b/include/paimon/file_store_commit.h
index 18cdaf20..553112b3 100644
--- a/include/paimon/file_store_commit.h
+++ b/include/paimon/file_store_commit.h
@@ -85,10 +85,18 @@ class PAIMON_EXPORT FileStoreCommit {
/// Omitting an earlier entry may advance committed progress past
unpublished files and allow
/// their real-time data to be reclaimed.
///
- /// If this method returns an error, the caller may retry with the same
arguments. Each call
- /// reloads the latest committed state. As in `FilterAndCommit`, a retry's
identifier is
- /// considered committed when it is not newer than the latest identifier
for `commit_user`.
- /// The requested offset ranges must also be covered by the latest
committed progress.
+ /// 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.
+ /// External conflicts returned after submitting a REST catalog request
are not retried by this
+ /// method. Concurrent rollback or partition deletion from another process
must be fenced by
+ /// the upstream coordinator.
///
/// @param realtime_commits Commit messages and left-closed, right-open
offset ranges to
/// commit.
diff --git a/include/paimon/realtime/realtime_context.h
b/include/paimon/realtime/realtime_context.h
index 200e4ba4..ca410bb6 100644
--- a/include/paimon/realtime/realtime_context.h
+++ b/include/paimon/realtime/realtime_context.h
@@ -25,6 +25,7 @@
#include <string>
#include <utility>
+#include "paimon/metrics.h"
#include "paimon/result.h"
#include "paimon/visibility.h"
@@ -65,6 +66,20 @@ struct PAIMON_EXPORT RealtimePartitionBucket {
/// Exclusive committed end offset for each partition-bucket.
using RealtimeOffsetMap = std::map<RealtimePartitionBucket, int64_t>;
+/// Gauge names exposed by `RealtimeContext::GetMetrics` and real-time
file-store writers.
+class PAIMON_EXPORT RealtimeMetrics {
+ public:
+ RealtimeMetrics() = delete;
+ ~RealtimeMetrics() = delete;
+
+ static constexpr char kBuildingMemoryBytes[] =
"realtimeBuildingMemoryBytes";
+ static constexpr char kSealedMemoryBytes[] = "realtimeSealedMemoryBytes";
+ static constexpr char kTotalMemoryBytes[] = "realtimeTotalMemoryBytes";
+ static constexpr char kBuildingRowCount[] = "realtimeBuildingRowCount";
+ static constexpr char kSealedRowCount[] = "realtimeSealedRowCount";
+ static constexpr char kTotalRowCount[] = "realtimeTotalRowCount";
+};
+
/// Framework-managed context that owns the `RealtimeStore` instances used by
real-time operations.
///
/// Applications share one context between `WriteContext`, `ScanContext`, and
`ReadContext`. The
@@ -89,6 +104,9 @@ class PAIMON_EXPORT RealtimeContext {
static Result<std::shared_ptr<RealtimeContext>> Create(
const std::shared_ptr<RealtimeStoreFactory>& factory);
+ /// Returns current memory and physical row-count gauges aggregated over
all stores.
+ virtual std::shared_ptr<Metrics> GetMetrics() const = 0;
+
virtual ~RealtimeContext();
protected:
diff --git a/include/paimon/realtime/realtime_store.h
b/include/paimon/realtime/realtime_store.h
index 115eddce..6cc1b085 100644
--- a/include/paimon/realtime/realtime_store.h
+++ b/include/paimon/realtime/realtime_store.h
@@ -60,7 +60,7 @@ struct PAIMON_EXPORT RealtimeStoreCreateRequest {
std::shared_ptr<MemoryPool> memory_pool;
/// Table mode implemented by the store.
RealtimeStoreMode mode = RealtimeStoreMode::APPEND_ONLY;
- /// Statistics collected by append-only stores.
+ /// Statistics collected by the store for query pruning.
StatisticsMode statistics_mode = StatisticsMode::NONE;
};
@@ -80,6 +80,18 @@ struct PAIMON_EXPORT RealtimeWriteBatch {
OffsetRange offset_range;
};
+/// Current memory and row counts tracked by a `RealtimeStore`.
+///
+/// Row counts are physical stored rows. For primary-key stores they include
old versions and
+/// delete records, rather than the rows visible after merge-on-read. A
segment removed by
+/// `AdvanceCommittedOffset` is no longer included, even if an older read view
still pins it.
+struct PAIMON_EXPORT RealtimeStoreDataUsage {
+ uint64_t building_memory_bytes = 0;
+ uint64_t sealed_memory_bytes = 0;
+ uint64_t building_row_count = 0;
+ uint64_t sealed_row_count = 0;
+};
+
/// Opaque handle to an immutable segment returned by
`RealtimeStore::SealForCommit`.
///
/// A plugin may store the segment in memory or in spill files. Callers use
this handle only to
@@ -174,6 +186,9 @@ class PAIMON_EXPORT RealtimeStore {
/// or retain them. Existing read views continue to keep referenced
resources alive.
virtual Status AdvanceCommittedOffset(int64_t committed_end_offset) = 0;
+ /// Returns one consistent snapshot of current building and sealed data
usage.
+ virtual RealtimeStoreDataUsage GetDataUsage() const = 0;
+
/// Returns the number of bytes currently retained by building and sealed
segments.
virtual uint64_t GetMemoryUsage() const = 0;
};
diff --git a/src/paimon/core/operation/abstract_file_store_write.cpp
b/src/paimon/core/operation/abstract_file_store_write.cpp
index c4da29bf..8bc9d11a 100644
--- a/src/paimon/core/operation/abstract_file_store_write.cpp
+++ b/src/paimon/core/operation/abstract_file_store_write.cpp
@@ -427,7 +427,14 @@ Status AbstractFileStoreWrite::Close() {
}
std::shared_ptr<Metrics> AbstractFileStoreWrite::GetMetrics() const {
- return metrics_;
+ std::shared_ptr<RealtimeContext> context = GetRealtimeContext();
+ if (!context) {
+ return metrics_;
+ }
+ auto result = std::make_shared<MetricsImpl>();
+ result->Merge(metrics_);
+ result->Merge(context->GetMetrics());
+ return result;
}
Status AbstractFileStoreWrite::CheckRealtimeWriteUsable() const {
diff --git a/src/paimon/core/operation/append_only_file_store_write_test.cpp
b/src/paimon/core/operation/append_only_file_store_write_test.cpp
index 612f931f..bfeaefce 100644
--- a/src/paimon/core/operation/append_only_file_store_write_test.cpp
+++ b/src/paimon/core/operation/append_only_file_store_write_test.cpp
@@ -311,11 +311,29 @@ TEST_F(AppendOnlyFileStoreWriteTest,
TestRealtimeWriteTracksExternalOffsetRange)
[10, 1, "a"],
[20, 2, "b"]
])")));
+ std::shared_ptr<Metrics> building_metrics = file_store_write->GetMetrics();
+ ASSERT_OK_AND_ASSIGN(double building_rows,
+
building_metrics->GetGauge(RealtimeMetrics::kBuildingRowCount));
+ ASSERT_OK_AND_ASSIGN(double sealed_rows,
+
building_metrics->GetGauge(RealtimeMetrics::kSealedRowCount));
+ ASSERT_OK_AND_ASSIGN(double total_rows,
+
building_metrics->GetGauge(RealtimeMetrics::kTotalRowCount));
+ ASSERT_EQ(2, building_rows);
+ ASSERT_EQ(0, sealed_rows);
+ ASSERT_EQ(2, total_rows);
ASSERT_NOK_WITH_MSG(
file_store_write->PrepareCommit(/*wait_compaction=*/false,
/*commit_identifier=*/0),
"real-time writer must use PrepareCommitWithProgress");
ASSERT_OK_AND_ASSIGN(auto first_prepared,
file_store_write->PrepareCommitWithProgress(/*commit_identifier=*/0));
+ std::shared_ptr<Metrics> sealed_metrics = file_store_write->GetMetrics();
+ ASSERT_OK_AND_ASSIGN(building_rows,
+
sealed_metrics->GetGauge(RealtimeMetrics::kBuildingRowCount));
+ ASSERT_OK_AND_ASSIGN(sealed_rows,
sealed_metrics->GetGauge(RealtimeMetrics::kSealedRowCount));
+ ASSERT_OK_AND_ASSIGN(total_rows,
sealed_metrics->GetGauge(RealtimeMetrics::kTotalRowCount));
+ ASSERT_EQ(0, building_rows);
+ ASSERT_EQ(2, sealed_rows);
+ ASSERT_EQ(2, total_rows);
ASSERT_EQ(1, first_prepared.size());
ASSERT_TRUE(first_prepared[0].partition_bucket.partition.empty());
ASSERT_EQ(0, first_prepared[0].partition_bucket.bucket);
diff --git a/src/paimon/core/operation/file_store_commit_impl.cpp
b/src/paimon/core/operation/file_store_commit_impl.cpp
index 736825f7..adac3e15 100644
--- a/src/paimon/core/operation/file_store_commit_impl.cpp
+++ b/src/paimon/core/operation/file_store_commit_impl.cpp
@@ -523,6 +523,56 @@ Result<std::vector<std::shared_ptr<ManifestCommittable>>>
FileStoreCommitImpl::F
}
}
+Result<std::optional<Snapshot>>
FileStoreCommitImpl::LatestSnapshotOfCommitUserAtOrBefore(
+ const std::optional<Snapshot>& latest_snapshot) const {
+ if (!latest_snapshot) {
+ return std::optional<Snapshot>();
+ }
+ PAIMON_ASSIGN_OR_RAISE(std::optional<int64_t> earliest_snapshot_id,
+ snapshot_manager_->EarliestSnapshotId());
+ if (!earliest_snapshot_id || earliest_snapshot_id.value() >
latest_snapshot->Id()) {
+ return std::optional<Snapshot>();
+ }
+ for (int64_t snapshot_id = latest_snapshot->Id(); snapshot_id >=
earliest_snapshot_id.value();
+ --snapshot_id) {
+ PAIMON_ASSIGN_OR_RAISE(Snapshot snapshot,
snapshot_manager_->LoadSnapshot(snapshot_id));
+ if (snapshot.CommitUser() == commit_user_) {
+ return std::optional<Snapshot>(std::move(snapshot));
+ }
+ }
+ return std::optional<Snapshot>();
+}
+
+Result<std::optional<int64_t>> FileStoreCommitImpl::ResolveRealtimeCommit(
+ const std::optional<Snapshot>& latest_snapshot, int64_t identifier,
+ const std::map<RealtimePartitionBucket, OffsetRange>& realtime_ranges)
const {
+ PAIMON_ASSIGN_OR_RAISE(std::optional<Snapshot> latest_snapshot_of_user,
+
LatestSnapshotOfCommitUserAtOrBefore(latest_snapshot));
+ const bool identifier_committed =
+ latest_snapshot_of_user && identifier <=
latest_snapshot_of_user->CommitIdentifier();
+ const std::optional<Snapshot>& offsets_snapshot =
+ identifier_committed ? latest_snapshot_of_user : latest_snapshot;
+ PAIMON_ASSIGN_OR_RAISE(RealtimeOffsetMap committed_offsets,
+
RealtimeCommitProperties::ReadOffsets(offsets_snapshot, fs_));
+ PAIMON_ASSIGN_OR_RAISE(bool ranges_committed,
RealtimeCommitProperties::AreRangesCommitted(
+ committed_offsets,
realtime_ranges));
+ if (!identifier_committed) {
+ if (ranges_committed) {
+ return Status::Invalid(
+ "real-time offset ranges were committed by another commit user
or identifier");
+ }
+ return std::optional<int64_t>();
+ }
+ if (!ranges_committed) {
+ return Status::Invalid(
+ "real-time commit identifier was committed without the requested
offset ranges");
+ }
+ if (!latest_snapshot) {
+ return Status::Invalid("real-time commit ranges are covered without a
snapshot");
+ }
+ return std::optional<int64_t>(latest_snapshot->Id());
+}
+
Status FileStoreCommitImpl::Overwrite(
const std::map<std::string, std::string>& partition,
const std::vector<std::shared_ptr<CommitMessage>>& commit_messages,
int64_t identifier,
@@ -696,7 +746,7 @@ Status FileStoreCommitImpl::ExecuteOverwrite(
}
if (with_compact) {
- PAIMON_ASSIGN_OR_RAISE(int32_t cnt,
+ PAIMON_ASSIGN_OR_RAISE(TryCommitResult result,
TryCommit(changes->compact_table_files,
/*changelog_files=*/{},
changes->compact_index_files,
identifier, watermark,
committable->Properties(),
/*realtime_ranges=*/{},
@@ -705,7 +755,7 @@ Status FileStoreCommitImpl::ExecuteOverwrite(
/*removed_realtime_partitions=*/{},
/*detect_conflicts=*/true,
/*retry_on_conflict=*/true));
- *attempt += cnt;
+ *attempt += result.attempts;
*generated_snapshot += 1;
}
@@ -807,11 +857,13 @@ Result<int32_t> FileStoreCommitImpl::TryOverwrite(
// ExecuteOverwrite has already resolved dynamic overwrite to the concrete
affected
// partitions. Only an empty final partition list denotes a full-table
replacement.
const bool reset_all_realtime_progress = partitions.empty();
- return TryCommit(changes_provider, commit_identifier, watermark,
properties,
- /*realtime_ranges=*/{}, Snapshot::CommitKind::Overwrite(),
- reset_all_realtime_progress, partitions,
- /*detect_conflicts=*/true,
- /*retry_on_conflict=*/true);
+ PAIMON_ASSIGN_OR_RAISE(TryCommitResult result,
+ TryCommit(changes_provider, commit_identifier,
watermark, properties,
+ /*realtime_ranges=*/{},
Snapshot::CommitKind::Overwrite(),
+ reset_all_realtime_progress, partitions,
+ /*detect_conflicts=*/true,
+ /*retry_on_conflict=*/true));
+ return result.attempts;
}
Status FileStoreCommitImpl::Commit(
@@ -854,19 +906,22 @@ Status FileStoreCommitImpl::Commit(
}
PAIMON_ASSIGN_OR_RAISE(
- int32_t cnt,
+ TryCommitResult result,
TryCommit(changes.append_table_files, changes.append_changelog,
changes.append_index_files, committable->Identifier(),
committable->Watermark(), committable->Properties(),
realtime_ranges,
commit_kind,
/*reset_all_realtime_progress=*/false,
/*removed_realtime_partitions=*/{}, check_append_files,
retry_on_conflict));
- attempt += cnt;
+ attempt += result.attempts;
+ if (result.already_committed) {
+ return Status::OK();
+ }
generated_snapshot += 1;
}
if (changes.HasCompactChanges()) {
- PAIMON_ASSIGN_OR_RAISE(int32_t cnt,
+ PAIMON_ASSIGN_OR_RAISE(TryCommitResult result,
TryCommit(changes.compact_table_files,
changes.compact_changelog,
changes.compact_index_files,
committable->Identifier(),
committable->Watermark(),
committable->Properties(),
@@ -874,7 +929,7 @@ Status FileStoreCommitImpl::Commit(
/*reset_all_realtime_progress=*/false,
/*removed_realtime_partitions=*/{},
/*detect_conflicts=*/true,
retry_on_conflict));
- attempt += cnt;
+ attempt += result.attempts;
generated_snapshot += 1;
}
return Status::OK();
@@ -936,40 +991,25 @@ Result<int64_t> FileStoreCommitImpl::CommitWithProgress(
std::shared_ptr<ManifestCommittable> committable =
CreateManifestCommittable(identifier, commit_messages, watermark,
/*properties=*/{});
- PAIMON_ASSIGN_OR_RAISE(std::vector<std::shared_ptr<ManifestCommittable>>
pending_committables,
- FilterCommitted({committable}));
- const bool identifier_committed = pending_committables.empty();
-
PAIMON_ASSIGN_OR_RAISE(std::optional<Snapshot> latest_snapshot,
snapshot_manager_->LatestSnapshot());
- PAIMON_ASSIGN_OR_RAISE(RealtimeOffsetMap committed_offsets,
-
RealtimeCommitProperties::ReadOffsets(latest_snapshot, fs_));
- PAIMON_ASSIGN_OR_RAISE(bool ranges_committed,
RealtimeCommitProperties::AreRangesCommitted(
- committed_offsets,
realtime_ranges));
- if (ranges_committed != identifier_committed) {
- return Status::Invalid(
- ranges_committed
- ? "real-time offset ranges were committed by another commit
user or identifier"
- : "real-time commit identifier was committed without the
requested offset ranges");
- }
- if (ranges_committed) {
- if (!latest_snapshot) {
- return Status::Invalid("real-time commit ranges are covered
without a snapshot");
- }
- return latest_snapshot->Id();
+ PAIMON_ASSIGN_OR_RAISE(std::optional<int64_t> committed_snapshot_id,
+ ResolveRealtimeCommit(latest_snapshot, identifier,
realtime_ranges));
+ if (committed_snapshot_id) {
+ return committed_snapshot_id.value();
}
- PAIMON_RETURN_NOT_OK(CheckFilesExistence(pending_committables));
- const int64_t previous_snapshot_id = last_committed_snapshot_id_;
+ PAIMON_RETURN_NOT_OK(CheckFilesExistence({committable}));
+ last_committed_snapshot_id_ = -1;
PAIMON_RETURN_NOT_OK(Commit(committable, /*check_append_files=*/false,
- /*retry_on_conflict=*/false, realtime_ranges));
- if (last_committed_snapshot_id_ <= previous_snapshot_id) {
+ /*retry_on_conflict=*/true, realtime_ranges));
+ if (last_committed_snapshot_id_ < Snapshot::FIRST_SNAPSHOT_ID) {
return Status::Invalid("real-time commit did not produce a snapshot");
}
return last_committed_snapshot_id_;
}
-Result<int32_t> FileStoreCommitImpl::TryCommit(
+Result<FileStoreCommitImpl::TryCommitResult> FileStoreCommitImpl::TryCommit(
const std::vector<ManifestEntry>& delta_files,
const std::vector<ManifestEntry>& changelog_files,
const std::vector<IndexManifestEntry>& index_entries, int64_t identifier,
@@ -985,7 +1025,7 @@ Result<int32_t> FileStoreCommitImpl::TryCommit(
detect_conflicts, retry_on_conflict);
}
-Result<int32_t> FileStoreCommitImpl::TryCommit(
+Result<FileStoreCommitImpl::TryCommitResult> FileStoreCommitImpl::TryCommit(
const std::shared_ptr<CommitChangesProvider>& changes_provider, int64_t
identifier,
std::optional<int64_t> watermark, const std::map<std::string,
std::string>& properties,
const std::map<RealtimePartitionBucket, OffsetRange>& realtime_ranges,
@@ -997,6 +1037,15 @@ Result<int32_t> FileStoreCommitImpl::TryCommit(
while (true) {
PAIMON_ASSIGN_OR_RAISE(std::optional<Snapshot> latest_snapshot,
snapshot_manager_->LatestSnapshot());
+ if (!realtime_ranges.empty()) {
+ PAIMON_ASSIGN_OR_RAISE(
+ std::optional<int64_t> committed_snapshot_id,
+ ResolveRealtimeCommit(latest_snapshot, identifier,
realtime_ranges));
+ if (committed_snapshot_id) {
+ last_committed_snapshot_id_ = committed_snapshot_id.value();
+ return TryCommitResult{retry_count,
/*already_committed=*/true};
+ }
+ }
PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<CommitChanges> commit_changes,
changes_provider->Provide(latest_snapshot));
using SnapshotProperties = std::map<std::string, std::string>;
@@ -1014,8 +1063,16 @@ Result<int32_t> FileStoreCommitImpl::TryCommit(
if (commit_success) {
break;
}
+ if (!realtime_ranges.empty()) {
+ auto offsets_iter =
snapshot_properties.find(RealtimeCommitProperties::kOffsetsKey);
+ if (offsets_iter != snapshot_properties.end()) {
+ // AtomicStore reported a known conflict, so this attempt's
offsets file is not
+ // referenced by a snapshot and can be removed before rebasing
the next attempt.
+ [[maybe_unused]] Status cleanup_status =
+ fs_->Delete(offsets_iter->second, /*recursive=*/false);
+ }
+ }
if (!retry_on_conflict) {
- // TODO(xinyu.lxy): Support failure recovery and idempotent retry
for real-time commits.
return Status::Invalid("real-time commit failed due to snapshot
conflict");
}
int64_t current_millis = DateTimeUtils::GetCurrentUTCTimeUs() / 1000;
@@ -1029,7 +1086,7 @@ Result<int32_t> FileStoreCommitImpl::TryCommit(
retry_waiter_.RetryWait(retry_count);
retry_count++;
}
- return retry_count + 1;
+ return TryCommitResult{retry_count + 1, /*already_committed=*/false};
}
Status FileStoreCommitImpl::CheckSameBucketFromSnapshot(
diff --git a/src/paimon/core/operation/file_store_commit_impl.h
b/src/paimon/core/operation/file_store_commit_impl.h
index ca4de22f..31006ebb 100644
--- a/src/paimon/core/operation/file_store_commit_impl.h
+++ b/src/paimon/core/operation/file_store_commit_impl.h
@@ -147,6 +147,11 @@ class FileStoreCommitImpl : public FileStoreCommit {
Status Init(std::unique_ptr<CommitContext> ctx);
private:
+ struct TryCommitResult {
+ int32_t attempts;
+ bool already_committed;
+ };
+
Status Commit(const std::shared_ptr<ManifestCommittable>&
manifest_committable,
bool check_append_files, bool retry_on_conflict,
const std::map<RealtimePartitionBucket, OffsetRange>&
realtime_ranges);
@@ -175,6 +180,15 @@ class FileStoreCommitImpl : public FileStoreCommit {
Result<std::vector<std::shared_ptr<ManifestCommittable>>> FilterCommitted(
const std::vector<std::shared_ptr<ManifestCommittable>>& committables);
+ Result<std::optional<Snapshot>> LatestSnapshotOfCommitUserAtOrBefore(
+ const std::optional<Snapshot>& latest_snapshot) const;
+
+ /// Returns the containing snapshot id when the commit is already
complete, or nullopt when
+ /// neither its identifier nor offset ranges have been committed.
+ Result<std::optional<int64_t>> ResolveRealtimeCommit(
+ const std::optional<Snapshot>& latest_snapshot, int64_t identifier,
+ const std::map<RealtimePartitionBucket, OffsetRange>& realtime_ranges)
const;
+
std::shared_ptr<ManifestCommittable> CreateManifestCommittable(
int64_t identifier, const std::vector<std::shared_ptr<CommitMessage>>&
commit_messages,
std::optional<int64_t> watermark, const std::map<std::string,
std::string>& properties);
@@ -185,7 +199,7 @@ class FileStoreCommitImpl : public FileStoreCommit {
void ReportCommit(const ManifestEntryChanges& changes, int64_t
commit_duration,
int32_t generated_snapshot, int32_t attempt);
- Result<int32_t> TryCommit(
+ Result<TryCommitResult> TryCommit(
const std::vector<ManifestEntry>& delta_files,
const std::vector<ManifestEntry>& changelog_files,
const std::vector<IndexManifestEntry>& index_entries, int64_t
identifier,
@@ -195,7 +209,7 @@ class FileStoreCommitImpl : public FileStoreCommit {
const std::vector<std::map<std::string, std::string>>&
removed_realtime_partitions,
bool detect_conflicts, bool retry_on_conflict);
- Result<int32_t> TryCommit(
+ Result<TryCommitResult> TryCommit(
const std::shared_ptr<CommitChangesProvider>& changes_provider,
int64_t identifier,
std::optional<int64_t> watermark, const std::map<std::string,
std::string>& properties,
const std::map<RealtimePartitionBucket, OffsetRange>& realtime_ranges,
diff --git a/src/paimon/core/operation/file_store_commit_impl_test.cpp
b/src/paimon/core/operation/file_store_commit_impl_test.cpp
index d55f21cb..7c756c3d 100644
--- a/src/paimon/core/operation/file_store_commit_impl_test.cpp
+++ b/src/paimon/core/operation/file_store_commit_impl_test.cpp
@@ -62,6 +62,7 @@
#include "paimon/core/manifest/manifest_file.h"
#include "paimon/core/manifest/manifest_file_meta.h"
#include "paimon/core/manifest/manifest_list.h"
+#include "paimon/core/operation/commit/realtime_commit_properties.h"
#include "paimon/core/operation/metrics/commit_metrics.h"
#include "paimon/core/partition/partition_statistics.h"
#include "paimon/core/schema/table_schema.h"
@@ -600,6 +601,95 @@ TEST_F(FileStoreCommitImplTest,
TestCommitWithConflictSnapshotAndRetryOnce) {
ASSERT_TRUE(exist);
}
+TEST_F(FileStoreCommitImplTest,
TestRealtimeCommitRebasesAfterSnapshotConflict) {
+ const std::vector<std::string> data_files = {
+ "/f1=10/bucket-0/data-51a45441-6037-4af3-b67b-5cefd75dc6f2-0.orc",
+ "/f1=10/bucket-1/data-6828284c-e707-49b5-af6b-69be79af120c-0.orc",
+ "/f1=20/bucket-0/data-8dc7f04c-3c98-48b2-9d56-834d746c4a40-0.orc",
+ "/f1=10/bucket-1/data-fd1d2255-43f2-4534-b4cc-08b29e662940-0.orc",
+ "/f1=20/bucket-0/data-7b3f4cc7-116b-4d2f-9c62-5dadc1f11bcb-0.orc"};
+ ASSERT_OK(PrepareFakeFiles(data_files));
+
+ std::vector<std::shared_ptr<CommitMessage>> realtime_messages =
+ GetCommitMessages(paimon::test::GetDataDir() +
+
"/orc/append_09.db/append_09/commit_messages/commit_messages-01",
+ /*version=*/3);
+ std::vector<std::shared_ptr<CommitMessage>> concurrent_messages =
+ GetCommitMessages(paimon::test::GetDataDir() +
+
"/orc/append_09.db/append_09/commit_messages/commit_messages-02",
+ /*version=*/3);
+ ASSERT_EQ(3, realtime_messages.size());
+ ASSERT_EQ(2, concurrent_messages.size());
+
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<FileSystem> fs,
+ FileSystemFactory::Get("gmock_fs", table_path_, {}));
+ CommitContextBuilder realtime_builder(table_path_, "realtime_commit_user");
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<CommitContext> realtime_context,
+ realtime_builder.AddOption(Options::REALTIME_ENABLED,
"true")
+ .AddOption(Options::COMMIT_MAX_RETRIES, "1")
+ .AddOption(Options::COMMIT_MIN_RETRY_WAIT, "1ms")
+ .AddOption(Options::COMMIT_MAX_RETRY_WAIT, "1ms")
+ .WithFileSystem(fs)
+ .Finish());
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreCommit> realtime_commit,
+ FileStoreCommit::Create(std::move(realtime_context)));
+ auto realtime_commit_impl = std::dynamic_pointer_cast<FileStoreCommitImpl>(
+ std::shared_ptr<FileStoreCommit>(std::move(realtime_commit)));
+
+ std::vector<RealtimeCommitProgress> realtime_progress;
+ for (const std::shared_ptr<CommitMessage>& message : realtime_messages) {
+ std::shared_ptr<CommitMessageImpl> message_impl =
+ std::dynamic_pointer_cast<CommitMessageImpl>(message);
+ ASSERT_NE(nullptr, message_impl);
+ std::map<std::string, std::string> partition;
+ ASSERT_OK_AND_ASSIGN(partition,
+
realtime_commit_impl->PartitionToMap(message_impl->Partition()));
+ realtime_progress.push_back(RealtimeCommitProgress{
+ message, RealtimePartitionBucket(std::move(partition),
message_impl->Bucket()),
+ OffsetRange(/*begin=*/0, /*end=*/1)});
+ }
+
+ CommitContextBuilder concurrent_builder(table_path_,
"concurrent_commit_user");
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<CommitContext> concurrent_context,
+
concurrent_builder.WithFileSystem(file_system_).Finish());
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreCommit> concurrent_commit,
+
FileStoreCommit::Create(std::move(concurrent_context)));
+
+ auto* mock_fs = dynamic_cast<GmockFileSystem*>(fs.get());
+ ASSERT_NE(nullptr, mock_fs);
+ EXPECT_CALL(*mock_fs, Exists(testing::_))
+ .Times(testing::AnyNumber())
+ .WillRepeatedly(testing::Invoke(
+ [mock_fs](const std::string& path) { return
mock_fs->LocalFileSystem::Exists(path); }));
+ const std::string first_snapshot_path = PathUtil::JoinPath(table_path_,
"snapshot/snapshot-1");
+ EXPECT_CALL(*mock_fs, Exists(testing::StrEq(first_snapshot_path)))
+ .WillOnce(testing::Invoke([&](const std::string&) -> Result<bool> {
+ PAIMON_RETURN_NOT_OK(concurrent_commit->Commit(
+ concurrent_messages, /*commit_identifier=*/10,
/*watermark=*/std::nullopt));
+ return true;
+ }))
+ .RetiresOnSaturation();
+
+ ASSERT_OK_AND_ASSIGN(int64_t snapshot_id,
realtime_commit_impl->CommitWithProgress(
+ realtime_progress,
/*commit_identifier=*/1,
+ /*watermark=*/std::nullopt));
+ ASSERT_EQ(2, snapshot_id);
+ ASSERT_OK_AND_ASSIGN(Snapshot snapshot,
+
realtime_commit_impl->snapshot_manager_->LoadSnapshot(snapshot_id));
+ ASSERT_OK_AND_ASSIGN(std::vector<ManifestEntry> entries,
+
realtime_commit_impl->ReadAddManifestEntries(snapshot));
+ ASSERT_EQ(5, entries.size());
+ ASSERT_OK_AND_ASSIGN(RealtimeOffsetMap offsets,
RealtimeCommitProperties::ReadOffsets(
+
std::optional<Snapshot>(snapshot), fs));
+ ASSERT_EQ(3, offsets.size());
+ for (const RealtimeCommitProgress& progress : realtime_progress) {
+ ASSERT_EQ(1, offsets.at(progress.partition_bucket));
+ }
+ ASSERT_OK_AND_ASSIGN(uint64_t attempts,
realtime_commit_impl->GetCommitMetrics()->GetCounter(
+
CommitMetrics::LAST_COMMIT_ATTEMPTS));
+ ASSERT_EQ(2, attempts);
+}
+
TEST_F(FileStoreCommitImplTest,
TestCommitWithAtomicWriteSnapshotTimeoutAndActuallySucceed) {
std::string test_data_path = paimon::test::GetDataDir() +
"/parquet/append_09.db/append_09/";
auto dir = UniqueTestDirectory::Create();
diff --git a/src/paimon/core/operation/key_value_file_store_write.cpp
b/src/paimon/core/operation/key_value_file_store_write.cpp
index cfa62873..d6b6a7a4 100644
--- a/src/paimon/core/operation/key_value_file_store_write.cpp
+++ b/src/paimon/core/operation/key_value_file_store_write.cpp
@@ -143,7 +143,8 @@ Result<std::shared_ptr<BatchWriter>>
KeyValueFileStoreWrite::CreateWriter(
RealtimeStoreState store_state,
realtime_context_impl->GetOrCreateRealtimeStore(
RealtimeStoreCreateRequest{std::move(c_write_schema),
options_.ToMap(), pool_,
- RealtimeStoreMode::PRIMARY_KEY},
+ RealtimeStoreMode::PRIMARY_KEY,
+
options_.GetRealtimeStoreStatisticsMode()},
RealtimePartitionBucket(partition_map, bucket)));
realtime_store_state = std::move(store_state);
compact_manager = std::make_shared<NoopCompactManager>();
diff --git a/src/paimon/core/operation/key_value_file_store_write_test.cpp
b/src/paimon/core/operation/key_value_file_store_write_test.cpp
index b71708db..988fb3f5 100644
--- a/src/paimon/core/operation/key_value_file_store_write_test.cpp
+++ b/src/paimon/core/operation/key_value_file_store_write_test.cpp
@@ -468,12 +468,26 @@ TEST_F(KeyValueFileStoreWriteTest, TestRealtimeWrite) {
ASSERT_NOK_WITH_MSG(writer->Write(MakeBatch(realtime_schema, R"([[30, 3,
"backwards"]])")),
"offset moved backwards or was duplicated");
ASSERT_NOK(writer->Write(MakeBatch(realtime_schema, R"([[null, 3,
"null-offset"]])")));
+ std::shared_ptr<Metrics> building_metrics = writer->GetMetrics();
+ ASSERT_OK_AND_ASSIGN(double building_rows,
+
building_metrics->GetGauge(RealtimeMetrics::kBuildingRowCount));
+ ASSERT_OK_AND_ASSIGN(double total_rows,
+
building_metrics->GetGauge(RealtimeMetrics::kTotalRowCount));
+ // The real-time gauges count physical versions, including the delete
record.
+ ASSERT_EQ(3, building_rows);
+ ASSERT_EQ(3, total_rows);
ASSERT_OK_AND_ASSIGN(auto store_rows,
ReadRealtimePrimaryKeyStoreRows(realtime_context));
ASSERT_EQ(
(decltype(store_rows){{0, 1, "old", 0, 10}, {2, 1, "new", 2, 30}, {3,
2, "two", 1, 20}}),
store_rows);
ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> progresses,
writer->PrepareCommitWithProgress(0));
+ std::shared_ptr<Metrics> sealed_metrics = writer->GetMetrics();
+ ASSERT_OK_AND_ASSIGN(double sealed_rows,
+
sealed_metrics->GetGauge(RealtimeMetrics::kSealedRowCount));
+ ASSERT_OK_AND_ASSIGN(total_rows,
sealed_metrics->GetGauge(RealtimeMetrics::kTotalRowCount));
+ ASSERT_EQ(3, sealed_rows);
+ ASSERT_EQ(3, total_rows);
ASSERT_EQ(1, progresses.size());
ASSERT_EQ(OffsetRange(10, 31), progresses[0].offset_range);
std::shared_ptr<CommitMessageImpl> commit_message =
diff --git a/src/paimon/core/operation/merge_file_split_read.cpp
b/src/paimon/core/operation/merge_file_split_read.cpp
index 746d4d3a..7c6b1724 100644
--- a/src/paimon/core/operation/merge_file_split_read.cpp
+++ b/src/paimon/core/operation/merge_file_split_read.cpp
@@ -262,7 +262,7 @@ class MergeFileSplitRead::RealtimeReaderBuilder {
partition,
owner_->context_->GetPredicate(),
data_file_path_factory,
readers));
return CollectRawDiskReaders(level0_files, level0_deletion_files,
partition,
- /*predicate=*/nullptr,
data_file_path_factory, readers);
+ owner_->predicate_for_keys_,
data_file_path_factory, readers);
}
Status CollectRawDiskReaders(const
std::vector<std::shared_ptr<DataFileMeta>>& data_files,
diff --git a/src/paimon/core/realtime/arrow_realtime_store.cpp
b/src/paimon/core/realtime/arrow_realtime_store.cpp
index 3a819024..4950ecfb 100644
--- a/src/paimon/core/realtime/arrow_realtime_store.cpp
+++ b/src/paimon/core/realtime/arrow_realtime_store.cpp
@@ -235,7 +235,9 @@ class ArrowRealtimeStore::AppendQueryBatchReader : public
BatchReader {
continue;
}
const StoredBatch& stored = batches[next_batch_++];
- PAIMON_ASSIGN_OR_RAISE(bool may_match, MayMatch(stored));
+ PAIMON_ASSIGN_OR_RAISE(bool may_match,
ArrowRealtimeStore::MayMatchStatistics(
+ stored, read_schema_,
predicate_filter_,
+ statistics_mapping_,
memory_pool_));
if (!may_match) {
continue;
}
@@ -260,26 +262,6 @@ class ArrowRealtimeStore::AppendQueryBatchReader : public
BatchReader {
view_ = nullptr;
}
- private:
- Result<bool> MayMatch(const StoredBatch& stored) const {
- if (!predicate_filter_ || !stored.statistics) {
- return true;
- }
- const BatchStatistics& statistics = stored.statistics.value();
- std::shared_ptr<InternalRow> min_row = std::make_shared<ColumnarRow>(
- statistics.min_values, statistics.min_values->fields(),
memory_pool_, /*row_id=*/0);
- std::shared_ptr<InternalRow> max_row = std::make_shared<ColumnarRow>(
- statistics.max_values, statistics.max_values->fields(),
memory_pool_, /*row_id=*/0);
- ProjectedRow projected_min(min_row, statistics_mapping_);
- ProjectedRow projected_max(max_row, statistics_mapping_);
- std::shared_ptr<InternalArray> null_counts =
- std::make_shared<ColumnarArray>(statistics.null_counts.get(),
memory_pool_,
- /*offset=*/0,
statistics.null_counts->length());
- ProjectedArray projected_null_counts(null_counts, statistics_mapping_);
- return predicate_filter_->Test(read_schema_, stored.data->length(),
projected_min,
- projected_max, projected_null_counts);
- }
-
private:
std::shared_ptr<ReadView> view_;
std::shared_ptr<arrow::Schema> read_schema_;
@@ -304,7 +286,7 @@ ArrowRealtimeStore::ArrowRealtimeStore(const
std::shared_ptr<arrow::Schema>& wri
Result<std::optional<ArrowRealtimeStore::BatchStatistics>>
ArrowRealtimeStore::CollectStatistics(
const std::shared_ptr<arrow::StructArray>& data) const {
- if (mode_ == RealtimeStoreMode::PRIMARY_KEY || statistics_mode_ ==
StatisticsMode::NONE) {
+ if (statistics_mode_ == StatisticsMode::NONE) {
return std::optional<BatchStatistics>();
}
@@ -358,8 +340,31 @@ Result<std::optional<ArrowRealtimeStore::BatchStatistics>>
ArrowRealtimeStore::C
PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(
std::shared_ptr<arrow::StructArray> max_values_struct,
arrow::StructArray::Make(max_values, write_schema_->fields()));
- return std::optional<BatchStatistics>(BatchStatistics{
- std::move(min_values_struct), std::move(max_values_struct),
std::move(null_counts_array)});
+ return std::optional<BatchStatistics>(BatchStatistics{arrow_pool_,
std::move(min_values_struct),
+
std::move(max_values_struct),
+
std::move(null_counts_array)});
+}
+
+Result<bool> ArrowRealtimeStore::MayMatchStatistics(
+ const StoredBatch& stored, const std::shared_ptr<arrow::Schema>&
read_schema,
+ const std::shared_ptr<PredicateFilter>& predicate_filter,
+ const std::vector<int32_t>& statistics_mapping,
+ const std::shared_ptr<MemoryPool>& memory_pool) {
+ if (!predicate_filter || !stored.statistics) {
+ return true;
+ }
+ const BatchStatistics& statistics = stored.statistics.value();
+ std::shared_ptr<InternalRow> min_row = std::make_shared<ColumnarRow>(
+ statistics.min_values, statistics.min_values->fields(), memory_pool,
/*row_id=*/0);
+ std::shared_ptr<InternalRow> max_row = std::make_shared<ColumnarRow>(
+ statistics.max_values, statistics.max_values->fields(), memory_pool,
/*row_id=*/0);
+ ProjectedRow projected_min(min_row, statistics_mapping);
+ ProjectedRow projected_max(max_row, statistics_mapping);
+ std::shared_ptr<InternalArray> null_counts =
std::make_shared<ColumnarArray>(
+ statistics.null_counts.get(), memory_pool, /*offset=*/0,
statistics.null_counts->length());
+ ProjectedArray projected_null_counts(null_counts, statistics_mapping);
+ return predicate_filter->Test(read_schema, stored.data->length(),
projected_min, projected_max,
+ projected_null_counts);
}
Status ArrowRealtimeStore::Write(RealtimeWriteBatch&& write_batch) {
@@ -392,6 +397,7 @@ Status ArrowRealtimeStore::Write(RealtimeWriteBatch&&
write_batch) {
ArrowUtils::GetArrayMemoryUsage(statistics->null_counts->data());
}
building_memory_usage_ += memory_usage;
+ building_row_count_ += static_cast<uint64_t>(struct_array->length());
building_batches_.push_back(StoredBatch{std::move(struct_array),
write_batch.offset_range,
std::move(statistics),
memory_usage});
if (!building_range_) {
@@ -409,9 +415,12 @@
Result<std::optional<std::shared_ptr<RealtimeSegmentHandle>>> ArrowRealtimeStore
}
auto segment = std::make_shared<Segment>(building_range_.value(),
std::move(building_batches_));
sealed_segments_.push_back(segment);
+ sealed_memory_usage_ += building_memory_usage_;
+ sealed_row_count_ += building_row_count_;
building_batches_.clear();
building_range_.reset();
building_memory_usage_ = 0;
+ building_row_count_ = 0;
return
std::optional<std::shared_ptr<RealtimeSegmentHandle>>(std::move(segment));
}
@@ -459,9 +468,24 @@ Result<std::vector<std::unique_ptr<BatchReader>>>
ArrowRealtimeStore::CreateQuer
if (!arrow_view->GetOffsetRange()) {
return readers;
}
+ std::shared_ptr<PredicateFilter> predicate_filter;
+ if (context.predicate) {
+ predicate_filter =
std::dynamic_pointer_cast<PredicateFilter>(context.predicate);
+ }
+ std::vector<int32_t> statistics_mapping;
+ statistics_mapping.reserve(read_schema->num_fields());
+ for (const std::shared_ptr<arrow::Field>& field : read_schema->fields()) {
+
statistics_mapping.push_back(write_schema_->GetFieldIndex(field->name()));
+ }
if (mode_ == RealtimeStoreMode::PRIMARY_KEY) {
for (const std::shared_ptr<Segment>& segment :
arrow_view->GetSegments()) {
for (const StoredBatch& batch : segment->GetBatches()) {
+ PAIMON_ASSIGN_OR_RAISE(bool may_match,
+ MayMatchStatistics(batch, read_schema,
predicate_filter,
+ statistics_mapping,
memory_pool_));
+ if (!may_match) {
+ continue;
+ }
PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<arrow::StructArray>
projected,
ProjectBatch(batch.data, read_schema,
arrow_pool_.get()));
readers.push_back(std::make_unique<StoredBatchReader>(projected, arrow_pool_));
@@ -470,15 +494,6 @@ Result<std::vector<std::unique_ptr<BatchReader>>>
ArrowRealtimeStore::CreateQuer
return readers;
}
- std::shared_ptr<PredicateFilter> predicate_filter;
- if (context.predicate) {
- predicate_filter =
std::dynamic_pointer_cast<PredicateFilter>(context.predicate);
- }
- std::vector<int32_t> statistics_mapping;
- statistics_mapping.reserve(read_schema->num_fields());
- for (const std::shared_ptr<arrow::Field>& field : read_schema->fields()) {
-
statistics_mapping.push_back(write_schema_->GetFieldIndex(field->name()));
- }
std::unique_ptr<BatchReader> reader =
std::make_unique<AppendQueryBatchReader>(
arrow_view, read_schema, predicate_filter,
std::move(statistics_mapping), arrow_pool_,
memory_pool_);
@@ -491,22 +506,34 @@ Status ArrowRealtimeStore::AdvanceCommittedOffset(int64_t
committed_end_offset)
// TODO(xinyu.lxy): Consider deferring segment destruction to a
reclamation queue. Existing
// read views may pin reclaimed batches, so the last query releasing a
view can otherwise pay
// the full buffer destruction cost and observe higher tail latency.
+ uint64_t reclaimed_memory_usage = 0;
+ uint64_t reclaimed_row_count = 0;
+ for (const std::shared_ptr<Segment>& segment : sealed_segments_) {
+ if (segment->GetOffsetRange().end <= committed_end_offset) {
+ reclaimed_memory_usage += segment->GetMemoryUsage();
+ reclaimed_row_count +=
static_cast<uint64_t>(segment->GetRowCount());
+ }
+ }
sealed_segments_.erase(
std::remove_if(sealed_segments_.begin(), sealed_segments_.end(),
[committed_end_offset](const std::shared_ptr<Segment>&
segment) {
return segment->GetOffsetRange().end <=
committed_end_offset;
}),
sealed_segments_.end());
+ sealed_memory_usage_ -= reclaimed_memory_usage;
+ sealed_row_count_ -= reclaimed_row_count;
return Status::OK();
}
+RealtimeStoreDataUsage ArrowRealtimeStore::GetDataUsage() const {
+ std::lock_guard<std::mutex> lock(mutex_);
+ return RealtimeStoreDataUsage{building_memory_usage_,
sealed_memory_usage_, building_row_count_,
+ sealed_row_count_};
+}
+
uint64_t ArrowRealtimeStore::GetMemoryUsage() const {
std::lock_guard<std::mutex> lock(mutex_);
- uint64_t result = building_memory_usage_;
- for (const std::shared_ptr<Segment>& segment : sealed_segments_) {
- result += segment->GetMemoryUsage();
- }
- return result;
+ return building_memory_usage_ + sealed_memory_usage_;
}
} // namespace paimon
diff --git a/src/paimon/core/realtime/arrow_realtime_store.h
b/src/paimon/core/realtime/arrow_realtime_store.h
index 4032568a..630744f5 100644
--- a/src/paimon/core/realtime/arrow_realtime_store.h
+++ b/src/paimon/core/realtime/arrow_realtime_store.h
@@ -36,6 +36,7 @@ class StructArray;
namespace paimon {
class MemoryPool;
+class PredicateFilter;
/// Internal Arrow-backed implementation of the default `RealtimeStore`.
class ArrowRealtimeStore final : public RealtimeStore {
@@ -60,10 +61,15 @@ class ArrowRealtimeStore final : public RealtimeStore {
Status AdvanceCommittedOffset(int64_t committed_end_offset) override;
+ RealtimeStoreDataUsage GetDataUsage() const override;
+
uint64_t GetMemoryUsage() const override;
private:
struct BatchStatistics {
+ // Arrow buffers retain a raw MemoryPool pointer. Keep their allocator
alive until all
+ // statistics arrays have been released.
+ std::shared_ptr<arrow::MemoryPool> arrow_pool;
std::shared_ptr<arrow::StructArray> min_values;
std::shared_ptr<arrow::StructArray> max_values;
std::shared_ptr<arrow::Array> null_counts;
@@ -85,6 +91,12 @@ class ArrowRealtimeStore final : public RealtimeStore {
Result<std::optional<BatchStatistics>> CollectStatistics(
const std::shared_ptr<arrow::StructArray>& data) const;
+ static Result<bool> MayMatchStatistics(const StoredBatch& stored,
+ const
std::shared_ptr<arrow::Schema>& read_schema,
+ const
std::shared_ptr<PredicateFilter>& predicate_filter,
+ const std::vector<int32_t>&
statistics_mapping,
+ const std::shared_ptr<MemoryPool>&
memory_pool);
+
std::shared_ptr<arrow::Schema> write_schema_;
std::shared_ptr<MemoryPool> memory_pool_;
std::shared_ptr<arrow::MemoryPool> arrow_pool_;
@@ -95,6 +107,9 @@ class ArrowRealtimeStore final : public RealtimeStore {
std::vector<std::shared_ptr<Segment>> sealed_segments_;
std::optional<OffsetRange> building_range_;
uint64_t building_memory_usage_ = 0;
+ uint64_t sealed_memory_usage_ = 0;
+ uint64_t building_row_count_ = 0;
+ uint64_t sealed_row_count_ = 0;
};
} // namespace paimon
diff --git a/src/paimon/core/realtime/arrow_realtime_store_test.cpp
b/src/paimon/core/realtime/arrow_realtime_store_test.cpp
index 70a1a20a..430f8b27 100644
--- a/src/paimon/core/realtime/arrow_realtime_store_test.cpp
+++ b/src/paimon/core/realtime/arrow_realtime_store_test.cpp
@@ -157,6 +157,49 @@ TEST_F(ArrowRealtimeStoreTest, TestWriteValidationAndSeal)
{
ASSERT_GT(store_->GetMemoryUsage(), 0);
}
+TEST_F(ArrowRealtimeStoreTest,
TestDataUsageTracksBuildingSealedAndReclaimedData) {
+ RealtimeStoreDataUsage usage = store_->GetDataUsage();
+ ASSERT_EQ(0, usage.building_memory_bytes);
+ ASSERT_EQ(0, usage.sealed_memory_bytes);
+ ASSERT_EQ(0, usage.building_row_count);
+ ASSERT_EQ(0, usage.sealed_row_count);
+
+ ASSERT_OK(store_->Write(
+ RealtimeWriteBatch{MakeBatch(R"([[0, 0, "a"], [1, 1, "b"]])"),
OffsetRange(0, 2)}));
+ usage = store_->GetDataUsage();
+ ASSERT_GT(usage.building_memory_bytes, 0);
+ ASSERT_EQ(0, usage.sealed_memory_bytes);
+ ASSERT_EQ(2, usage.building_row_count);
+ ASSERT_EQ(0, usage.sealed_row_count);
+ ASSERT_EQ(usage.building_memory_bytes, store_->GetMemoryUsage());
+ const uint64_t first_segment_memory = usage.building_memory_bytes;
+
+ ASSERT_OK_AND_ASSIGN(std::optional<std::shared_ptr<RealtimeSegmentHandle>>
segment,
+ store_->SealForCommit());
+ ASSERT_TRUE(segment.has_value());
+ usage = store_->GetDataUsage();
+ ASSERT_EQ(0, usage.building_memory_bytes);
+ ASSERT_EQ(first_segment_memory, usage.sealed_memory_bytes);
+ ASSERT_EQ(0, usage.building_row_count);
+ ASSERT_EQ(2, usage.sealed_row_count);
+
+ ASSERT_OK(store_->Write(RealtimeWriteBatch{MakeBatch(R"([[2, 2, "c"]])"),
OffsetRange(2, 3)}));
+ usage = store_->GetDataUsage();
+ ASSERT_GT(usage.building_memory_bytes, 0);
+ ASSERT_EQ(first_segment_memory, usage.sealed_memory_bytes);
+ ASSERT_EQ(1, usage.building_row_count);
+ ASSERT_EQ(2, usage.sealed_row_count);
+ ASSERT_EQ(usage.building_memory_bytes + usage.sealed_memory_bytes,
store_->GetMemoryUsage());
+
+ ASSERT_OK(store_->AdvanceCommittedOffset(/*committed_end_offset=*/2));
+ usage = store_->GetDataUsage();
+ ASSERT_GT(usage.building_memory_bytes, 0);
+ ASSERT_EQ(0, usage.sealed_memory_bytes);
+ ASSERT_EQ(1, usage.building_row_count);
+ ASSERT_EQ(0, usage.sealed_row_count);
+ ASSERT_EQ(usage.building_memory_bytes, store_->GetMemoryUsage());
+}
+
TEST_F(ArrowRealtimeStoreTest, TestQueryReaderClipsCommittedOffsetWithBitmap) {
ASSERT_OK(store_->Write(RealtimeWriteBatch{
MakeBatch(R"([[10, 10, "a"], [11, 11, "b"], [12, 12, "c"]])"),
OffsetRange(10, 13)}));
diff --git a/src/paimon/core/realtime/primary_key_realtime_store_test.cpp
b/src/paimon/core/realtime/primary_key_realtime_store_test.cpp
index d62f012d..e4e85ede 100644
--- a/src/paimon/core/realtime/primary_key_realtime_store_test.cpp
+++ b/src/paimon/core/realtime/primary_key_realtime_store_test.cpp
@@ -34,6 +34,8 @@
#include "paimon/core/realtime/realtime_schema_layout.h"
#include "paimon/macros.h"
#include "paimon/memory/memory_pool.h"
+#include "paimon/predicate/literal.h"
+#include "paimon/predicate/predicate_builder.h"
#include "paimon/realtime/arrow_realtime_store_factory.h"
#include "paimon/testing/utils/read_result_collector.h"
#include "paimon/testing/utils/testharness.h"
@@ -70,10 +72,11 @@ std::shared_ptr<arrow::Schema> NestedStoreWriteSchema() {
1)});
}
-std::shared_ptr<ArrowRealtimeStore> CreateStore(const
std::shared_ptr<arrow::Schema>& schema,
- const
std::shared_ptr<MemoryPool>& pool) {
+std::shared_ptr<ArrowRealtimeStore> CreateStore(
+ const std::shared_ptr<arrow::Schema>& schema, const
std::shared_ptr<MemoryPool>& pool,
+ StatisticsMode statistics_mode = StatisticsMode::NONE) {
return std::make_shared<ArrowRealtimeStore>(schema,
RealtimeStoreMode::PRIMARY_KEY,
- StatisticsMode::NONE, pool,
GetArrowPool(pool));
+ statistics_mode, pool,
GetArrowPool(pool));
}
std::unique_ptr<RecordBatch> MakeBatch(const std::string& json) {
@@ -307,6 +310,30 @@ TEST(PrimaryKeyRealtimeStoreTest,
TestQueryReaderPerStoredBatch) {
ASSERT_TRUE(expected->Equals(actual));
}
+TEST(PrimaryKeyRealtimeStoreTest,
TestFullStatisticsPrunesNonMatchingStoredBatch) {
+ std::shared_ptr<ArrowRealtimeStore> store =
+ CreateStore(StoreWriteSchema(), GetDefaultPool(),
StatisticsMode::FULL);
+ ASSERT_OK(store->Write(RealtimeWriteBatch{
+ MakeBatch(R"([[1, 0, 0, 1, "one"], [2, 0, 1, 2, "two"]])"),
OffsetRange(0, 2)}));
+ ASSERT_OK(store->Write(RealtimeWriteBatch{
+ MakeBatch(R"([[3, 0, 2, 10, "ten"], [4, 0, 3, 11, "eleven"]])"),
OffsetRange(2, 4)}));
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeReadView> view,
store->AcquireReadView());
+
+ auto c_schema = std::make_unique<ArrowSchema>();
+ ASSERT_TRUE(arrow::ExportSchema(*StoreWriteSchema(), c_schema.get()).ok());
+ std::shared_ptr<Predicate> predicate = PredicateBuilder::Equal(
+ /*field_index=*/3, /*field_name=*/"id", FieldType::BIGINT,
Literal(int64_t{10}));
+ RealtimeQueryContext context{c_schema.get(), predicate};
+ ASSERT_OK_AND_ASSIGN(std::vector<std::unique_ptr<BatchReader>> readers,
+ store->CreateQueryReaders(view, context));
+ ASSERT_EQ(1, readers.size());
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<arrow::Array> actual,
ReadArray(std::move(readers)));
+ ASSERT_OK_AND_ASSIGN(
+ std::shared_ptr<arrow::Array> expected,
+ ArrayFromJson(actual->type(), R"([[3, 0, 2, 10, "ten"], [4, 0, 3, 11,
"eleven"]])"));
+ ASSERT_TRUE(expected->Equals(actual));
+}
+
TEST(PrimaryKeyRealtimeStoreTest, TestQueryBatchOutlivesStoreAndReader) {
const std::shared_ptr<arrow::Schema> stored_schema = StoreWriteSchema();
std::shared_ptr<MemoryPool> pool = GetMemoryPool();
diff --git a/src/paimon/core/realtime/realtime_context_impl.cpp
b/src/paimon/core/realtime/realtime_context_impl.cpp
index 0f1a4cf8..72a5dd13 100644
--- a/src/paimon/core/realtime/realtime_context_impl.cpp
+++ b/src/paimon/core/realtime/realtime_context_impl.cpp
@@ -39,6 +39,7 @@
#include "arrow/c/helpers.h"
#include "fmt/format.h"
#include "paimon/arrow/abi.h"
+#include "paimon/common/metrics/metrics_impl.h"
#include "paimon/common/utils/arrow/status_utils.h"
#include "paimon/common/utils/scope_guard.h"
#include "paimon/common/utils/uuid.h"
@@ -110,6 +111,43 @@ Status RealtimeContextImpl::CheckUsable() const {
return Status::OK();
}
+std::shared_ptr<Metrics> RealtimeContextImpl::GetMetrics() const {
+ std::vector<std::shared_ptr<RealtimeStore>> stores;
+ {
+ std::lock_guard<std::mutex> lock(mutex_);
+ stores.reserve(stores_.size());
+ for (const auto& [_, entry] : stores_) {
+ stores.push_back(entry.store);
+ }
+ }
+
+ RealtimeStoreDataUsage total_usage;
+ for (const std::shared_ptr<RealtimeStore>& store : stores) {
+ const RealtimeStoreDataUsage usage = store->GetDataUsage();
+ total_usage.building_memory_bytes += usage.building_memory_bytes;
+ total_usage.sealed_memory_bytes += usage.sealed_memory_bytes;
+ total_usage.building_row_count += usage.building_row_count;
+ total_usage.sealed_row_count += usage.sealed_row_count;
+ }
+
+ auto metrics = std::make_shared<MetricsImpl>();
+ metrics->SetGauge(RealtimeMetrics::kBuildingMemoryBytes,
+ static_cast<double>(total_usage.building_memory_bytes));
+ metrics->SetGauge(RealtimeMetrics::kSealedMemoryBytes,
+ static_cast<double>(total_usage.sealed_memory_bytes));
+ metrics->SetGauge(
+ RealtimeMetrics::kTotalMemoryBytes,
+ static_cast<double>(total_usage.building_memory_bytes +
total_usage.sealed_memory_bytes));
+ metrics->SetGauge(RealtimeMetrics::kBuildingRowCount,
+ static_cast<double>(total_usage.building_row_count));
+ metrics->SetGauge(RealtimeMetrics::kSealedRowCount,
+ static_cast<double>(total_usage.sealed_row_count));
+ metrics->SetGauge(
+ RealtimeMetrics::kTotalRowCount,
+ static_cast<double>(total_usage.building_row_count +
total_usage.sealed_row_count));
+ return metrics;
+}
+
Result<RealtimeStoreState> RealtimeContextImpl::GetOrCreateRealtimeStore(
RealtimeStoreCreateRequest&& request, const RealtimePartitionBucket&
partition_bucket) {
PAIMON_RETURN_NOT_OK(CheckUsable());
diff --git a/src/paimon/core/realtime/realtime_context_impl.h
b/src/paimon/core/realtime/realtime_context_impl.h
index c8f8bbaa..6c7c8d0b 100644
--- a/src/paimon/core/realtime/realtime_context_impl.h
+++ b/src/paimon/core/realtime/realtime_context_impl.h
@@ -80,6 +80,8 @@ class PAIMON_EXPORT RealtimeContextImpl final : public
RealtimeContext {
Status CheckUsable() const;
+ std::shared_ptr<Metrics> GetMetrics() const override;
+
Result<RealtimeStoreState> GetOrCreateRealtimeStore(
RealtimeStoreCreateRequest&& request, const RealtimePartitionBucket&
partition_bucket);
@@ -122,7 +124,7 @@ class PAIMON_EXPORT RealtimeContextImpl final : public
RealtimeContext {
void CleanupReadViews();
std::shared_ptr<RealtimeStoreFactory> factory_;
- std::mutex mutex_;
+ mutable std::mutex mutex_;
std::mutex progress_mutex_;
std::map<RealtimePartitionBucket, StoreEntry> stores_;
// Full-table progress used as the initial offset when a store is created
lazily.
diff --git a/src/paimon/core/realtime/realtime_context_test.cpp
b/src/paimon/core/realtime/realtime_context_test.cpp
index 405dc100..af2fe717 100644
--- a/src/paimon/core/realtime/realtime_context_test.cpp
+++ b/src/paimon/core/realtime/realtime_context_test.cpp
@@ -74,14 +74,18 @@ class TestingRealtimeStore : public RealtimeStore {
committed_offsets.push_back(committed_offset);
return Status::OK();
}
+ RealtimeStoreDataUsage GetDataUsage() const override {
+ return data_usage;
+ }
uint64_t GetMemoryUsage() const override {
- return 0;
+ return data_usage.building_memory_bytes +
data_usage.sealed_memory_bytes;
}
int32_t acquire_count = 0;
int32_t advance_count = 0;
bool fail_next_advance = false;
bool return_null_read_view = false;
+ RealtimeStoreDataUsage data_usage;
std::vector<int64_t> committed_offsets;
};
@@ -171,6 +175,40 @@ TEST(RealtimeContextTest,
TestReusesStoreAndCapturesRegisteredViews) {
ASSERT_EQ(1, factory->stores[2]->acquire_count);
}
+TEST(RealtimeContextTest, TestMetricsAggregateAllStores) {
+ auto factory = std::make_shared<TestingRealtimeStoreFactory>();
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContextImpl> context,
CreateContext(factory));
+ ASSERT_OK(GetOrCreateAppendStore(context, {{"dt", "2026-08-02"}}, 0,
MakeWriteSchema(), {},
+ GetDefaultPool()));
+ ASSERT_OK(GetOrCreateAppendStore(context, {{"dt", "2026-08-03"}}, 1,
MakeWriteSchema(), {},
+ GetDefaultPool()));
+ ASSERT_EQ(2, factory->stores.size());
+ factory->stores[0]->data_usage =
+ RealtimeStoreDataUsage{/*building_memory_bytes=*/10,
/*sealed_memory_bytes=*/20,
+ /*building_row_count=*/1,
/*sealed_row_count=*/2};
+ factory->stores[1]->data_usage =
+ RealtimeStoreDataUsage{/*building_memory_bytes=*/30,
/*sealed_memory_bytes=*/40,
+ /*building_row_count=*/3,
/*sealed_row_count=*/4};
+
+ std::shared_ptr<Metrics> metrics = context->GetMetrics();
+ ASSERT_OK_AND_ASSIGN(double building_memory,
+
metrics->GetGauge(RealtimeMetrics::kBuildingMemoryBytes));
+ ASSERT_OK_AND_ASSIGN(double sealed_memory,
+
metrics->GetGauge(RealtimeMetrics::kSealedMemoryBytes));
+ ASSERT_OK_AND_ASSIGN(double total_memory,
+
metrics->GetGauge(RealtimeMetrics::kTotalMemoryBytes));
+ ASSERT_OK_AND_ASSIGN(double building_rows,
+
metrics->GetGauge(RealtimeMetrics::kBuildingRowCount));
+ ASSERT_OK_AND_ASSIGN(double sealed_rows,
metrics->GetGauge(RealtimeMetrics::kSealedRowCount));
+ ASSERT_OK_AND_ASSIGN(double total_rows,
metrics->GetGauge(RealtimeMetrics::kTotalRowCount));
+ ASSERT_EQ(40, building_memory);
+ ASSERT_EQ(60, sealed_memory);
+ ASSERT_EQ(100, total_memory);
+ ASSERT_EQ(4, building_rows);
+ ASSERT_EQ(6, sealed_rows);
+ ASSERT_EQ(10, total_rows);
+}
+
TEST(RealtimeContextTest, TestRejectsMismatchedModeOnStoreReuse) {
auto factory = std::make_shared<TestingRealtimeStoreFactory>();
ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContextImpl> context,
CreateContext(factory));
diff --git a/src/paimon/core/table/source/key_value_table_read.cpp
b/src/paimon/core/table/source/key_value_table_read.cpp
index c25011bc..9feae79e 100644
--- a/src/paimon/core/table/source/key_value_table_read.cpp
+++ b/src/paimon/core/table/source/key_value_table_read.cpp
@@ -19,6 +19,8 @@
#include "paimon/core/table/source/key_value_table_read.h"
+#include <map>
+#include <string>
#include <unordered_set>
#include <utility>
#include <vector>
@@ -46,6 +48,7 @@
#include "paimon/core/table/source/realtime_split.h"
#include "paimon/core/utils/nested_projection_utils.h"
#include "paimon/core/utils/primary_key_table_utils.h"
+#include "paimon/predicate/predicate_utils.h"
#include "paimon/status.h"
namespace paimon {
@@ -104,7 +107,20 @@ Result<std::vector<std::unique_ptr<KeyValueRecordReader>>>
CreateMemoryReaders(
auto c_schema = std::make_unique<ArrowSchema>();
PAIMON_RETURN_NOT_OK_FROM_ARROW(arrow::ExportSchema(*store_read_schema,
c_schema.get()));
ScopeGuard schema_guard([schema = c_schema.get()]() {
ArrowSchemaRelease(schema); });
- RealtimeQueryContext query_context{c_schema.get(), /*predicate=*/nullptr};
+ std::map<std::string, int32_t> primary_key_name_to_index;
+ for (const std::shared_ptr<arrow::Field>& key_field :
key_schema->fields()) {
+ int32_t field_index =
store_read_schema->GetFieldIndex(key_field->name());
+ if (field_index < 0) {
+ return Status::Invalid(
+ "primary key field is missing from real-time store read
schema: ",
+ key_field->name());
+ }
+ primary_key_name_to_index.emplace(key_field->name(), field_index);
+ }
+ PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<Predicate> primary_key_predicate,
+
PredicateUtils::CreatePickedFieldFilter(context->GetPredicate(),
+
primary_key_name_to_index));
+ RealtimeQueryContext query_context{c_schema.get(),
std::move(primary_key_predicate)};
PAIMON_ASSIGN_OR_RAISE(std::vector<std::unique_ptr<BatchReader>>
batch_readers,
memory.store->CreateQueryReaders(memory.read_view,
query_context));
PAIMON_ASSIGN_OR_RAISE(
diff --git a/test/inte/realtime_write_inte_test.cpp
b/test/inte/realtime_write_inte_test.cpp
index bd398c1d..ba05af0f 100644
--- a/test/inte/realtime_write_inte_test.cpp
+++ b/test/inte/realtime_write_inte_test.cpp
@@ -51,6 +51,7 @@
#include "paimon/common/utils/scope_guard.h"
#include "paimon/core/core_options.h"
#include "paimon/core/operation/commit/realtime_commit_properties.h"
+#include "paimon/core/operation/metrics/commit_metrics.h"
#include "paimon/core/realtime/realtime_context_impl.h"
#include "paimon/core/realtime/realtime_schema_layout.h"
#include "paimon/core/schema/schema_manager.h"
@@ -65,6 +66,7 @@
#include "paimon/file_store_commit.h"
#include "paimon/file_store_write.h"
#include "paimon/fs/file_system.h"
+#include "paimon/fs/local/local_file_system.h"
#include "paimon/memory/memory_pool.h"
#include "paimon/orphan_files_cleaner.h"
#include "paimon/predicate/function.h"
@@ -138,6 +140,10 @@ class DelegatingRealtimeStore : public RealtimeStore {
return delegate_->AdvanceCommittedOffset(committed_offset);
}
+ RealtimeStoreDataUsage GetDataUsage() const override {
+ return delegate_->GetDataUsage();
+ }
+
uint64_t GetMemoryUsage() const override {
return delegate_->GetMemoryUsage();
}
@@ -209,6 +215,88 @@ class QueryTrackingRealtimeStore final : public
DelegatingRealtimeStore {
std::shared_ptr<std::weak_ptr<RealtimeReadView>> query_view_;
};
+class SnapshotConflictFileSystem final : public LocalFileSystem {
+ public:
+ explicit SnapshotConflictFileSystem(std::string conflict_snapshot_path)
+ : conflict_snapshot_path_(std::move(conflict_snapshot_path)) {}
+
+ void SetPassthroughChecksForCurrentThread(int32_t checks) {
+ std::lock_guard<std::mutex> lock(mutex_);
+ passthrough_checks_[std::this_thread::get_id()] = checks;
+ }
+
+ Result<bool> Exists(const std::string& path) const override {
+ if (path != conflict_snapshot_path_) {
+ return LocalFileSystem::Exists(path);
+ }
+
+ {
+ std::unique_lock<std::mutex> lock(mutex_);
+ if (conflict_snapshot_store_finished_) {
+ lock.unlock();
+ return LocalFileSystem::Exists(path);
+ }
+ if (conflict_snapshot_leader_selected_ &&
+ conflict_snapshot_leader_ == std::this_thread::get_id()) {
+ // Local AtomicStore rechecks the destination through virtual
Exists during rename.
+ lock.unlock();
+ return LocalFileSystem::Exists(path);
+ }
+ auto passthrough_iter =
passthrough_checks_.find(std::this_thread::get_id());
+ if (passthrough_iter != passthrough_checks_.end() &&
passthrough_iter->second > 0) {
+ --passthrough_iter->second;
+ lock.unlock();
+ return LocalFileSystem::Exists(path);
+ }
+ // Make both committers race for the same target snapshot. One
observes it as absent
+ // and commits it; the other waits for that store and then reports
the conflict.
+ ++conflict_snapshot_attempts_;
+ cv_.notify_all();
+ if (!cv_.wait_for(lock, std::chrono::seconds(10),
+ [this]() { return conflict_snapshot_attempts_ >=
2; })) {
+ return Status::IOError("timed out waiting for concurrent
snapshot commit");
+ }
+ if (!conflict_snapshot_leader_selected_) {
+ conflict_snapshot_leader_selected_ = true;
+ conflict_snapshot_leader_ = std::this_thread::get_id();
+ return false;
+ }
+ if (!cv_.wait_for(lock, std::chrono::seconds(10),
+ [this]() { return
conflict_snapshot_store_finished_; })) {
+ return Status::IOError("timed out serializing conflicting
snapshot commits");
+ }
+ }
+ return LocalFileSystem::Exists(path);
+ }
+
+ Status AtomicStore(const std::string& path, const std::string& content)
override {
+ Status status = FileSystem::AtomicStore(path, content);
+ if (path == conflict_snapshot_path_) {
+ {
+ std::lock_guard<std::mutex> lock(mutex_);
+ conflict_snapshot_store_finished_ = true;
+ }
+ cv_.notify_all();
+ }
+ return status;
+ }
+
+ int32_t ConflictSnapshotAttempts() const {
+ std::lock_guard<std::mutex> lock(mutex_);
+ return conflict_snapshot_attempts_;
+ }
+
+ private:
+ std::string conflict_snapshot_path_;
+ mutable std::mutex mutex_;
+ mutable std::condition_variable cv_;
+ mutable std::map<std::thread::id, int32_t> passthrough_checks_;
+ mutable int32_t conflict_snapshot_attempts_ = 0;
+ mutable bool conflict_snapshot_leader_selected_ = false;
+ mutable std::thread::id conflict_snapshot_leader_;
+ mutable bool conflict_snapshot_store_finished_ = false;
+};
+
} // namespace
namespace {
@@ -322,6 +410,52 @@ class RealtimeWriteInteTest : public ::testing::Test {
protected:
using Row = std::tuple<int64_t, std::string, std::string>;
+ struct RealtimeMetricValues {
+ double building_memory = 0;
+ double sealed_memory = 0;
+ double total_memory = 0;
+ double building_rows = 0;
+ double sealed_rows = 0;
+ double total_rows = 0;
+ };
+
+ struct DiskFileLayout {
+ bool has_level0_file = false;
+ bool has_high_level_file = false;
+ bool has_high_level_deletion_vector = false;
+ };
+
+ static Result<DiskFileLayout> InspectDiskFiles(
+ const std::vector<std::shared_ptr<Split>>& splits) {
+ DiskFileLayout layout;
+ for (const std::shared_ptr<Split>& split : splits) {
+ std::shared_ptr<DataSplitImpl> data_split =
+ std::dynamic_pointer_cast<DataSplitImpl>(split);
+ if (!data_split) {
+ return Status::Invalid("expected a data split");
+ }
+ const std::vector<std::shared_ptr<DataFileMeta>>& files =
data_split->DataFiles();
+ const std::vector<std::optional<DeletionFile>>& deletion_files =
+ data_split->DeletionFiles();
+ if (!deletion_files.empty() && deletion_files.size() !=
files.size()) {
+ return Status::Invalid("deletion files must be empty or match
data files");
+ }
+ for (size_t i = 0; i < files.size(); ++i) {
+ if (files[i]->level == 0) {
+ layout.has_level0_file = true;
+ } else if (files[i]->level > 0) {
+ layout.has_high_level_file = true;
+ if (!deletion_files.empty() && deletion_files[i]) {
+ layout.has_high_level_deletion_vector = true;
+ }
+ } else {
+ return Status::Invalid("data file has a negative level");
+ }
+ }
+ }
+ return layout;
+ }
+
void SetUp() override {
pool_ = GetDefaultPool();
dir_ = UniqueTestDirectory::Create("local");
@@ -814,6 +948,27 @@ class RealtimeWriteInteTest : public ::testing::Test {
return memory_usage;
}
+ static Result<RealtimeMetricValues> ReadRealtimeMetrics(
+ const std::shared_ptr<Metrics>& metrics) {
+ if (!metrics) {
+ return Status::Invalid("real-time metrics are null");
+ }
+ RealtimeMetricValues result;
+ PAIMON_ASSIGN_OR_RAISE(result.building_memory,
+
metrics->GetGauge(RealtimeMetrics::kBuildingMemoryBytes));
+ PAIMON_ASSIGN_OR_RAISE(result.sealed_memory,
+
metrics->GetGauge(RealtimeMetrics::kSealedMemoryBytes));
+ PAIMON_ASSIGN_OR_RAISE(result.total_memory,
+
metrics->GetGauge(RealtimeMetrics::kTotalMemoryBytes));
+ PAIMON_ASSIGN_OR_RAISE(result.building_rows,
+
metrics->GetGauge(RealtimeMetrics::kBuildingRowCount));
+ PAIMON_ASSIGN_OR_RAISE(result.sealed_rows,
+
metrics->GetGauge(RealtimeMetrics::kSealedRowCount));
+ PAIMON_ASSIGN_OR_RAISE(result.total_rows,
+
metrics->GetGauge(RealtimeMetrics::kTotalRowCount));
+ return result;
+ }
+
Result<std::vector<int64_t>> ReadPkSequences(
const std::shared_ptr<RealtimeContext>& realtime_context) const {
PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<RealtimeContextImpl>
realtime_context_impl,
@@ -1072,6 +1227,281 @@ TEST_F(RealtimeWriteInteTest, TestAppendCommitAndRead) {
FinalizeCommitAndCheck(writer.get(), /*realtime_commits=*/{},
/*prepare_identifier=*/0, rows);
}
+TEST_F(RealtimeWriteInteTest,
TestRealtimeMetricsTrackBuildingSealedAndCommittedData) {
+ CreateTable(/*partition_keys=*/{});
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> realtime_context,
+ RealtimeContext::Create());
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> writer,
+ CreateRealtimeWriter(realtime_context));
+
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> first_batch,
+ MakeBatch(MakeRows(/*first_id=*/0, /*count=*/3,
/*partition=*/"p0"),
+ /*partitioned=*/false));
+ ASSERT_OK(writer->Write(std::move(first_batch)));
+ ASSERT_OK_AND_ASSIGN(RealtimeMetricValues building_metrics,
+ ReadRealtimeMetrics(writer->GetMetrics()));
+ ASSERT_GT(building_metrics.building_memory, 0);
+ ASSERT_EQ(0, building_metrics.sealed_memory);
+ ASSERT_EQ(building_metrics.building_memory, building_metrics.total_memory);
+ ASSERT_EQ(3, building_metrics.building_rows);
+ ASSERT_EQ(0, building_metrics.sealed_rows);
+ ASSERT_EQ(3, building_metrics.total_rows);
+
+ ASSERT_OK(writer->Seal());
+ ASSERT_OK_AND_ASSIGN(RealtimeMetricValues sealed_metrics,
+ ReadRealtimeMetrics(realtime_context->GetMetrics()));
+ ASSERT_EQ(0, sealed_metrics.building_memory);
+ ASSERT_EQ(building_metrics.total_memory, sealed_metrics.sealed_memory);
+ ASSERT_EQ(sealed_metrics.sealed_memory, sealed_metrics.total_memory);
+ ASSERT_EQ(0, sealed_metrics.building_rows);
+ ASSERT_EQ(3, sealed_metrics.sealed_rows);
+ ASSERT_EQ(3, sealed_metrics.total_rows);
+
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> second_batch,
+ MakeBatch(MakeRows(/*first_id=*/3, /*count=*/2,
/*partition=*/"p0"),
+ /*partitioned=*/false));
+ ASSERT_OK(writer->Write(std::move(second_batch)));
+ ASSERT_OK_AND_ASSIGN(RealtimeMetricValues mixed_metrics,
+ ReadRealtimeMetrics(realtime_context->GetMetrics()));
+ ASSERT_GT(mixed_metrics.building_memory, 0);
+ ASSERT_EQ(sealed_metrics.sealed_memory, mixed_metrics.sealed_memory);
+ ASSERT_EQ(mixed_metrics.building_memory + mixed_metrics.sealed_memory,
+ mixed_metrics.total_memory);
+ ASSERT_EQ(2, mixed_metrics.building_rows);
+ ASSERT_EQ(3, mixed_metrics.sealed_rows);
+ ASSERT_EQ(5, mixed_metrics.total_rows);
+
+ ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> commits,
+
writer->PrepareCommitWithProgress(/*commit_identifier=*/0));
+ ASSERT_EQ(1, commits.size());
+ ASSERT_OK_AND_ASSIGN(RealtimeMetricValues prepared_metrics,
+ ReadRealtimeMetrics(writer->GetMetrics()));
+ ASSERT_EQ(0, prepared_metrics.building_memory);
+ ASSERT_EQ(mixed_metrics.total_memory, prepared_metrics.sealed_memory);
+ ASSERT_EQ(prepared_metrics.sealed_memory, prepared_metrics.total_memory);
+ ASSERT_EQ(0, prepared_metrics.building_rows);
+ ASSERT_EQ(5, prepared_metrics.sealed_rows);
+ ASSERT_EQ(5, prepared_metrics.total_rows);
+
+ ASSERT_OK_AND_ASSIGN(int64_t snapshot_id, Commit(commits,
/*commit_identifier=*/0));
+ ASSERT_OK(writer->RefreshCommittedSnapshot(snapshot_id));
+ ASSERT_OK_AND_ASSIGN(RealtimeMetricValues reclaimed_metrics,
+ ReadRealtimeMetrics(realtime_context->GetMetrics()));
+ ASSERT_EQ(0, reclaimed_metrics.building_memory);
+ ASSERT_EQ(0, reclaimed_metrics.sealed_memory);
+ ASSERT_EQ(0, reclaimed_metrics.total_memory);
+ ASSERT_EQ(0, reclaimed_metrics.building_rows);
+ ASSERT_EQ(0, reclaimed_metrics.sealed_rows);
+ ASSERT_EQ(0, reclaimed_metrics.total_rows);
+ ASSERT_OK(writer->Close());
+}
+
+TEST_F(RealtimeWriteInteTest, TestRealtimeCommitRetriesSnapshotConflict) {
+ options_[Options::BUCKET] = "2";
+ options_[Options::COMMIT_MAX_RETRIES] = "1";
+ options_[Options::COMMIT_MIN_RETRY_WAIT] = "1ms";
+ options_[Options::COMMIT_MAX_RETRY_WAIT] = "1ms";
+ CreateTable(/*partition_keys=*/{});
+
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> bucket0_writer,
CreateRealtimeWriter());
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> bucket0_batch,
+ MakeBatch(MakeRows(/*first_id=*/0, /*count=*/2,
/*partition=*/"p0"),
+ /*partitioned=*/false, /*bucket=*/0));
+ ASSERT_OK(bucket0_writer->Write(std::move(bucket0_batch)));
+ ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> bucket0_commits,
+
bucket0_writer->PrepareCommitWithProgress(/*commit_identifier=*/0));
+ ASSERT_EQ(1, bucket0_commits.size());
+
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> bucket1_writer,
CreateRealtimeWriter());
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> bucket1_batch,
+ MakeBatch(MakeRows(/*first_id=*/10, /*count=*/3,
/*partition=*/"p0"),
+ /*partitioned=*/false, /*bucket=*/1));
+ ASSERT_OK(bucket1_writer->Write(std::move(bucket1_batch)));
+ ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> bucket1_commits,
+
bucket1_writer->PrepareCommitWithProgress(/*commit_identifier=*/0));
+ ASSERT_EQ(1, bucket1_commits.size());
+
+ const std::string first_snapshot_path = PathUtil::JoinPath(table_path_,
"snapshot/snapshot-1");
+ auto conflict_file_system =
std::make_shared<SnapshotConflictFileSystem>(first_snapshot_path);
+ CommitContextBuilder bucket0_builder(table_path_,
"bucket0_realtime_commit_user");
+ ASSERT_OK_AND_ASSIGN(
+ std::unique_ptr<CommitContext> bucket0_context,
+
bucket0_builder.SetOptions(options_).WithFileSystem(conflict_file_system).Finish());
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreCommit> bucket0_commit,
+ FileStoreCommit::Create(std::move(bucket0_context)));
+ CommitContextBuilder bucket1_builder(table_path_,
"bucket1_realtime_commit_user");
+ ASSERT_OK_AND_ASSIGN(
+ std::unique_ptr<CommitContext> bucket1_context,
+
bucket1_builder.SetOptions(options_).WithFileSystem(conflict_file_system).Finish());
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreCommit> bucket1_commit,
+ FileStoreCommit::Create(std::move(bucket1_context)));
+
+ std::optional<Result<int64_t>> bucket0_result;
+ std::optional<Result<int64_t>> bucket1_result;
+ std::thread bucket0_thread([&]() {
+ bucket0_result.emplace(bucket0_commit->CommitWithProgress(
+ bucket0_commits, /*commit_identifier=*/0,
/*watermark=*/std::nullopt));
+ });
+ std::thread bucket1_thread([&]() {
+ bucket1_result.emplace(bucket1_commit->CommitWithProgress(
+ bucket1_commits, /*commit_identifier=*/0,
/*watermark=*/std::nullopt));
+ });
+ bucket0_thread.join();
+ bucket1_thread.join();
+
+ ASSERT_TRUE(bucket0_result.has_value());
+ ASSERT_TRUE(bucket1_result.has_value());
+ ASSERT_OK_AND_ASSIGN(int64_t bucket0_snapshot_id,
std::move(bucket0_result.value()));
+ ASSERT_OK_AND_ASSIGN(int64_t bucket1_snapshot_id,
std::move(bucket1_result.value()));
+ std::vector<int64_t> snapshot_ids = {bucket0_snapshot_id,
bucket1_snapshot_id};
+ std::sort(snapshot_ids.begin(), snapshot_ids.end());
+ ASSERT_EQ((std::vector<int64_t>{1, 2}), snapshot_ids);
+ ASSERT_EQ(2, conflict_file_system->ConflictSnapshotAttempts());
+
+ ASSERT_OK_AND_ASSIGN(uint64_t bucket0_attempts,
bucket0_commit->GetCommitMetrics()->GetCounter(
+
CommitMetrics::LAST_COMMIT_ATTEMPTS));
+ ASSERT_OK_AND_ASSIGN(uint64_t bucket1_attempts,
bucket1_commit->GetCommitMetrics()->GetCounter(
+
CommitMetrics::LAST_COMMIT_ATTEMPTS));
+ std::vector<uint64_t> attempts = {bucket0_attempts, bucket1_attempts};
+ std::sort(attempts.begin(), attempts.end());
+ ASSERT_EQ((std::vector<uint64_t>{1, 2}), attempts);
+
+ ASSERT_OK_AND_ASSIGN(RealtimeOffsetMap committed_offsets,
ReadCommittedOffsets());
+ ASSERT_EQ(2, committed_offsets.size());
+ ASSERT_EQ(2,
committed_offsets.at(RealtimePartitionBucket(/*partition=*/{}, /*bucket=*/0)));
+ ASSERT_EQ(3,
committed_offsets.at(RealtimePartitionBucket(/*partition=*/{}, /*bucket=*/1)));
+ ASSERT_OK_AND_ASSIGN(std::vector<Row> rows, ReadRows());
+ std::sort(rows.begin(), rows.end());
+ std::vector<Row> expected_rows = MakeRows(/*first_id=*/0, /*count=*/2,
/*partition=*/"p0");
+ std::vector<Row> bucket1_rows = MakeRows(/*first_id=*/10, /*count=*/3,
/*partition=*/"p0");
+ expected_rows.insert(expected_rows.end(), bucket1_rows.begin(),
bucket1_rows.end());
+ std::sort(expected_rows.begin(), expected_rows.end());
+ ASSERT_EQ(expected_rows, rows);
+ ASSERT_OK(bucket0_writer->Close());
+ ASSERT_OK(bucket1_writer->Close());
+}
+
+TEST_F(RealtimeWriteInteTest, TestRealtimeCommitRetriesConflictWithCompaction)
{
+ options_[Options::COMMIT_MAX_RETRIES] = "1";
+ options_[Options::COMMIT_MIN_RETRY_WAIT] = "1ms";
+ options_[Options::COMMIT_MAX_RETRY_WAIT] = "1ms";
+ CreatePkTable();
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> realtime_context,
+ RealtimeContext::Create());
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> writer,
+ CreateRealtimeWriter(realtime_context));
+
+ for (int64_t id = 0; id < 2; ++id) {
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> batch,
+ MakeBatch(MakeRows(id, /*count=*/1,
/*partition=*/"p0"),
+ /*partitioned=*/false));
+ ASSERT_OK(writer->Write(std::move(batch)));
+ ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> progress,
+
writer->PrepareCommitWithProgress(/*commit_identifier=*/id));
+ ASSERT_OK_AND_ASSIGN(int64_t snapshot_id, Commit(progress,
/*commit_identifier=*/id));
+ ASSERT_OK(writer->RefreshCommittedSnapshot(snapshot_id));
+ }
+
+ WriteContextBuilder compact_write_builder(table_path_,
"compaction_write_user");
+
compact_write_builder.SetOptions(options_).WithStreamingMode(true).WithTempDirectory(
+ PathUtil::JoinPath(dir_->Str(), "conflict-compact-tmp"));
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<WriteContext> compact_write_context,
+ compact_write_builder.Finish());
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> compact_writer,
+
FileStoreWrite::Create(std::move(compact_write_context)));
+ ASSERT_OK(compact_writer->Compact(/*partition=*/{}, /*bucket=*/0,
+ /*full_compaction=*/true));
+ ASSERT_OK_AND_ASSIGN(
+ std::vector<std::shared_ptr<CommitMessage>> compact_messages,
+ compact_writer->PrepareCommit(/*wait_compaction=*/true,
/*commit_identifier=*/0));
+ ASSERT_EQ(1, compact_messages.size());
+ std::shared_ptr<CommitMessageImpl> compact_message =
+ std::dynamic_pointer_cast<CommitMessageImpl>(compact_messages[0]);
+ ASSERT_NE(nullptr, compact_message);
+ ASSERT_GE(compact_message->GetCompactIncrement().CompactBefore().size(),
2);
+
ASSERT_FALSE(compact_message->GetCompactIncrement().CompactAfter().empty());
+
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> append_batch,
+ MakeBatch(MakeRows(/*first_id=*/2, /*count=*/1,
/*partition=*/"p0"),
+ /*partitioned=*/false));
+ ASSERT_OK(writer->Write(std::move(append_batch)));
+ ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> append_progress,
+
writer->PrepareCommitWithProgress(/*commit_identifier=*/2));
+
+ const std::string conflict_snapshot_path =
+ PathUtil::JoinPath(table_path_, "snapshot/snapshot-3");
+ auto conflict_file_system =
+ std::make_shared<SnapshotConflictFileSystem>(conflict_snapshot_path);
+ CommitContextBuilder append_commit_builder(table_path_, commit_user_);
+ ASSERT_OK_AND_ASSIGN(
+ std::unique_ptr<CommitContext> append_commit_context,
+
append_commit_builder.SetOptions(options_).WithFileSystem(conflict_file_system).Finish());
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreCommit> append_commit,
+
FileStoreCommit::Create(std::move(append_commit_context)));
+ CommitContextBuilder compact_commit_builder(table_path_,
"compaction_commit_user");
+ ASSERT_OK_AND_ASSIGN(
+ std::unique_ptr<CommitContext> compact_commit_context,
+
compact_commit_builder.SetOptions(options_).WithFileSystem(conflict_file_system).Finish());
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreCommit> compact_commit,
+
FileStoreCommit::Create(std::move(compact_commit_context)));
+
+ std::optional<Result<int64_t>> append_result;
+ std::optional<Status> compact_status;
+ std::thread append_thread([&]() {
+ // CommitWithProgress probes the latest snapshot once for idempotency
and once while
+ // starting TryCommit. Synchronize only its destination check.
+
conflict_file_system->SetPassthroughChecksForCurrentThread(/*checks=*/2);
+ append_result.emplace(append_commit->CommitWithProgress(
+ append_progress, /*commit_identifier=*/2,
/*watermark=*/std::nullopt));
+ });
+ std::thread compact_thread([&]() {
+ // A regular commit probes the latest snapshot once before its
destination check.
+
conflict_file_system->SetPassthroughChecksForCurrentThread(/*checks=*/1);
+ compact_status.emplace(compact_commit->Commit(compact_messages,
/*commit_identifier=*/0,
+
/*watermark=*/std::nullopt));
+ });
+ append_thread.join();
+ compact_thread.join();
+
+ ASSERT_TRUE(append_result.has_value());
+ ASSERT_TRUE(compact_status.has_value());
+ ASSERT_OK_AND_ASSIGN(int64_t append_snapshot_id,
std::move(append_result.value()));
+ ASSERT_OK(compact_status.value());
+ ASSERT_TRUE(append_snapshot_id == 3 || append_snapshot_id == 4);
+ ASSERT_EQ(2, conflict_file_system->ConflictSnapshotAttempts());
+
+ ASSERT_OK_AND_ASSIGN(uint64_t append_attempts,
append_commit->GetCommitMetrics()->GetCounter(
+
CommitMetrics::LAST_COMMIT_ATTEMPTS));
+ ASSERT_OK_AND_ASSIGN(uint64_t compact_attempts,
compact_commit->GetCommitMetrics()->GetCounter(
+
CommitMetrics::LAST_COMMIT_ATTEMPTS));
+ std::vector<uint64_t> attempts = {append_attempts, compact_attempts};
+ std::sort(attempts.begin(), attempts.end());
+ ASSERT_EQ((std::vector<uint64_t>{1, 2}), attempts);
+
+ ASSERT_OK_AND_ASSIGN(CoreOptions core_options,
CoreOptions::FromMap(options_));
+ SnapshotManager snapshot_manager(core_options.GetFileSystem(),
table_path_);
+ ASSERT_OK_AND_ASSIGN(Snapshot snapshot3,
snapshot_manager.LoadSnapshot(/*snapshot_id=*/3));
+ ASSERT_OK_AND_ASSIGN(Snapshot snapshot4,
snapshot_manager.LoadSnapshot(/*snapshot_id=*/4));
+ ASSERT_TRUE((snapshot3.GetCommitKind() == Snapshot::CommitKind::Append() &&
+ snapshot4.GetCommitKind() == Snapshot::CommitKind::Compact())
||
+ (snapshot3.GetCommitKind() == Snapshot::CommitKind::Compact()
&&
+ snapshot4.GetCommitKind() == Snapshot::CommitKind::Append()));
+
+ ASSERT_OK_AND_ASSIGN(RealtimeOffsetMap committed_offsets,
ReadCommittedOffsets());
+ ASSERT_EQ(1, committed_offsets.size());
+ ASSERT_EQ(3,
committed_offsets.at(RealtimePartitionBucket(/*partition=*/{}, /*bucket=*/0)));
+ ASSERT_OK_AND_ASSIGN(std::vector<Row> rows, ReadRows());
+ ASSERT_EQ(MakeRows(/*first_id=*/0, /*count=*/3, /*partition=*/"p0"), rows);
+
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<Plan> plan,
+ CreatePlan(/*realtime_context=*/nullptr,
/*predicate=*/nullptr));
+ ASSERT_OK_AND_ASSIGN(DiskFileLayout layout,
InspectDiskFiles(plan->Splits()));
+ ASSERT_TRUE(layout.has_level0_file);
+ ASSERT_TRUE(layout.has_high_level_file);
+ ASSERT_OK(writer->Close());
+ ASSERT_OK(compact_writer->Close());
+}
+
TEST_F(RealtimeWriteInteTest, TestSealOnlySealsRealtimeSegment) {
CreateTable(/*partition_keys=*/{});
ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> realtime_context,
@@ -1213,6 +1643,7 @@ TEST_F(RealtimeWriteInteTest,
TestExplicitSnapshotIdWithSparseTailAndReclaimBoun
}
TEST_F(RealtimeWriteInteTest, TestPkRead) {
+ options_[Options::REALTIME_STORE_STATS_MODE] = "full";
CreatePkTable();
auto saw_query_predicate = std::make_shared<std::atomic<bool>>(false);
auto query_view = std::make_shared<std::weak_ptr<RealtimeReadView>>();
@@ -1261,6 +1692,16 @@ TEST_F(RealtimeWriteInteTest, TestPkRead) {
predicate,
/*enable_predicate_filter=*/true));
ASSERT_EQ(nullptr, filtered_result);
ASSERT_FALSE(saw_query_predicate->load(std::memory_order_acquire));
+
+ std::shared_ptr<Predicate> primary_key_predicate = PredicateBuilder::Equal(
+ /*field_index=*/0, /*field_name=*/"id", FieldType::BIGINT,
Literal(int64_t{3}));
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<Plan> primary_key_plan,
+ CreatePlan(realtime_context, primary_key_predicate));
+ ASSERT_OK_AND_ASSIGN(std::vector<Row> primary_key_rows,
+ ReadRows(primary_key_plan, realtime_context,
primary_key_predicate,
+ /*enable_predicate_filter=*/true));
+ ASSERT_EQ((std::vector<Row>{{3, "three", "p0"}}), primary_key_rows);
+ ASSERT_TRUE(saw_query_predicate->load(std::memory_order_acquire));
ASSERT_OK(PrepareAndClose(writer.get()));
writer.reset();
@@ -2324,6 +2765,7 @@ TEST_F(RealtimeWriteInteTest,
TestRealtimeWriteAcrossAppendCompaction) {
TEST_F(RealtimeWriteInteTest, TestPkDvPredicateAcrossHighLevelLevel0AndMemory)
{
options_[Options::FILE_FORMAT] = "parquet";
options_[Options::DELETION_VECTORS_ENABLED] = "true";
+ options_[Options::REALTIME_STORE_STATS_MODE] = "full";
CreatePkTable();
ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> realtime_context,
RealtimeContext::Create());
@@ -2418,31 +2860,10 @@ TEST_F(RealtimeWriteInteTest,
TestPkDvPredicateAcrossHighLevelLevel0AndMemory) {
ASSERT_NE(nullptr, realtime_split);
ASSERT_GT(realtime_split->MemoryEndOffset(),
realtime_split->CommittedEndOffset());
- bool has_level0_file = false;
- bool has_high_level_file = false;
- bool has_high_level_deletion_vector = false;
- for (const std::shared_ptr<Split>& disk_split :
realtime_split->DiskSplits()) {
- std::shared_ptr<DataSplitImpl> data_split =
- std::dynamic_pointer_cast<DataSplitImpl>(disk_split);
- ASSERT_NE(nullptr, data_split);
- const std::vector<std::shared_ptr<DataFileMeta>>& files =
data_split->DataFiles();
- const std::vector<std::optional<DeletionFile>>& deletion_files =
- data_split->DeletionFiles();
- ASSERT_TRUE(deletion_files.empty() || deletion_files.size() ==
files.size());
- for (size_t i = 0; i < files.size(); ++i) {
- if (files[i]->level == 0) {
- has_level0_file = true;
- } else {
- has_high_level_file = true;
- if (!deletion_files.empty() && deletion_files[i]) {
- has_high_level_deletion_vector = true;
- }
- }
- }
- }
- ASSERT_TRUE(has_level0_file);
- ASSERT_TRUE(has_high_level_file);
- ASSERT_TRUE(has_high_level_deletion_vector);
+ ASSERT_OK_AND_ASSIGN(DiskFileLayout layout,
InspectDiskFiles(realtime_split->DiskSplits()));
+ ASSERT_TRUE(layout.has_level0_file);
+ ASSERT_TRUE(layout.has_high_level_file);
+ ASSERT_TRUE(layout.has_high_level_deletion_vector);
// Predicate pushdown may return false-positive candidates, but a newer
unfiltered L0 or
// memory record must still suppress every matching old version from the
high levels.
@@ -2453,6 +2874,142 @@ TEST_F(RealtimeWriteInteTest,
TestPkDvPredicateAcrossHighLevelLevel0AndMemory) {
{4, "base-four", "p0"},
{7, matching_payload, "p0"}}),
candidates);
+
+ std::shared_ptr<Predicate> matching_key_predicate =
PredicateBuilder::Equal(
+ /*field_index=*/0, /*field_name=*/"id", FieldType::BIGINT,
Literal(int64_t{1}));
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<Plan> matching_key_plan,
+ CreatePlan(realtime_context, matching_key_predicate));
+ ASSERT_OK_AND_ASSIGN(std::vector<Row> matching_key_candidates,
+ ReadRows(matching_key_plan, realtime_context,
matching_key_predicate,
+ /*enable_predicate_filter=*/true));
+ ASSERT_EQ((std::vector<Row>{{1, "level0-current-one", "p0"}}),
matching_key_candidates);
+
+ std::shared_ptr<Predicate> deleted_key_predicate = PredicateBuilder::Equal(
+ /*field_index=*/0, /*field_name=*/"id", FieldType::BIGINT,
Literal(int64_t{2}));
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<Plan> deleted_key_plan,
+ CreatePlan(realtime_context, deleted_key_predicate));
+ ASSERT_OK_AND_ASSIGN(std::vector<Row> deleted_key_candidates,
+ ReadRows(deleted_key_plan, realtime_context,
deleted_key_predicate,
+ /*enable_predicate_filter=*/true));
+ ASSERT_TRUE(deleted_key_candidates.empty());
+ ASSERT_OK(PrepareAndClose(writer.get()));
+}
+
+TEST_F(RealtimeWriteInteTest,
TestPkDvPushesOnlyPrimaryKeyPredicatesIntoLevel0AndMemory) {
+ options_[Options::FILE_FORMAT] = "parquet";
+ options_[Options::DELETION_VECTORS_ENABLED] = "true";
+ options_[Options::REALTIME_STORE_STATS_MODE] = "full";
+ options_[Options::WRITE_BATCH_SIZE] = "1";
+ options_["parquet.page.size"] = "1";
+ options_["parquet.write.enable-page-index"] = "true";
+ options_["parquet.write.max-row-group-length"] = "1";
+ options_["parquet.read.enable-page-index-filter"] = "true";
+ CreatePkTable();
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<RealtimeContext> realtime_context,
+ RealtimeContext::Create());
+ ASSERT_OK_AND_ASSIGN(std::unique_ptr<FileStoreWrite> writer,
+ CreateRealtimeWriter(realtime_context));
+
+ auto write_one = [&](const Row& row, RecordBatch::RowKind row_kind) ->
Status {
+ PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<RecordBatch> batch,
+ MakeBatch({row}, /*partitioned=*/false,
/*bucket=*/0, {row_kind}));
+ return writer->Write(std::move(batch));
+ };
+
+ ASSERT_OK(write_one({1, "high-match", "p0"},
RecordBatch::RowKind::INSERT));
+ ASSERT_OK(write_one({2, "high-other", "p0"},
RecordBatch::RowKind::INSERT));
+ ASSERT_OK(write_one({3, "high-old", "p0"}, RecordBatch::RowKind::INSERT));
+ ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> base_progress,
+
writer->PrepareCommitWithProgress(/*commit_identifier=*/0));
+ ASSERT_OK_AND_ASSIGN(int64_t base_snapshot_id, Commit(base_progress,
/*commit_identifier=*/0));
+ ASSERT_OK(writer->RefreshCommittedSnapshot(base_snapshot_id));
+ ASSERT_OK_AND_ASSIGN(Snapshot full_compact_snapshot,
+ CompactAndCommit(/*partition=*/{}, /*bucket=*/0,
+ /*commit_identifier=*/1));
+ ASSERT_OK(writer->RefreshCommittedSnapshot(full_compact_snapshot.Id()));
+
+ ASSERT_OK(write_one({3, "high-current", "p0"},
RecordBatch::RowKind::UPDATE_AFTER));
+ ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> update_progress,
+
writer->PrepareCommitWithProgress(/*commit_identifier=*/2));
+ ASSERT_OK_AND_ASSIGN(int64_t update_snapshot_id,
+ Commit(update_progress, /*commit_identifier=*/2));
+ ASSERT_OK(writer->RefreshCommittedSnapshot(update_snapshot_id));
+ ASSERT_OK_AND_ASSIGN(Snapshot dv_compact_snapshot,
+ CompactAndCommit(/*partition=*/{}, /*bucket=*/0,
+ /*commit_identifier=*/3,
+ /*full_compaction=*/false));
+ ASSERT_OK(writer->RefreshCommittedSnapshot(dv_compact_snapshot.Id()));
+
+ ASSERT_OK(write_one({4, "level0-match", "p0"},
RecordBatch::RowKind::INSERT));
+ ASSERT_OK(write_one({5, "level0-other", "p0"},
RecordBatch::RowKind::INSERT));
+ ASSERT_OK_AND_ASSIGN(std::vector<RealtimeCommitProgress> level0_progress,
+
writer->PrepareCommitWithProgress(/*commit_identifier=*/4));
+ ASSERT_OK_AND_ASSIGN(int64_t level0_snapshot_id,
+ Commit(level0_progress, /*commit_identifier=*/4));
+ ASSERT_OK(writer->RefreshCommittedSnapshot(level0_snapshot_id));
+
+ ASSERT_OK(write_one({6, "memory-match", "p0"},
RecordBatch::RowKind::INSERT));
+ ASSERT_OK(write_one({7, "memory-other", "p0"},
RecordBatch::RowKind::INSERT));
+
+ ASSERT_OK_AND_ASSIGN(std::shared_ptr<Plan> unfiltered_plan,
+ CreatePlan(realtime_context, /*predicate=*/nullptr));
+ ASSERT_EQ(1, unfiltered_plan->Splits().size());
+ std::shared_ptr<RealtimeSplit> realtime_split =
+ std::dynamic_pointer_cast<RealtimeSplit>(unfiltered_plan->Splits()[0]);
+ ASSERT_NE(nullptr, realtime_split);
+ ASSERT_GT(realtime_split->MemoryEndOffset(),
realtime_split->CommittedEndOffset());
+ ASSERT_OK_AND_ASSIGN(DiskFileLayout layout,
InspectDiskFiles(realtime_split->DiskSplits()));
+ ASSERT_TRUE(layout.has_level0_file);
+ ASSERT_TRUE(layout.has_high_level_file);
+ ASSERT_TRUE(layout.has_high_level_deletion_vector);
+
+ // Keep residual filtering disabled so candidate rows expose the actual
source pushdown:
+ // high-level files receive the full predicate, while L0 and memory
receive only PK predicates.
+ auto read_candidates =
+ [&](const std::shared_ptr<Predicate>& predicate) ->
Result<std::vector<Row>> {
+ PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<Plan> plan,
CreatePlan(realtime_context, predicate));
+ return ReadRows(plan, realtime_context, predicate,
+ /*enable_predicate_filter=*/false);
+ };
+
+ std::shared_ptr<Predicate> level0_key_predicate = PredicateBuilder::Equal(
+ /*field_index=*/0, /*field_name=*/"id", FieldType::BIGINT,
Literal(int64_t{4}));
+ ASSERT_OK_AND_ASSIGN(std::vector<Row> level0_key_candidates,
+ read_candidates(level0_key_predicate));
+ ASSERT_EQ((std::vector<Row>{{4, "level0-match", "p0"}}),
level0_key_candidates);
+
+ std::shared_ptr<Predicate> memory_key_predicate = PredicateBuilder::Equal(
+ /*field_index=*/0, /*field_name=*/"id", FieldType::BIGINT,
Literal(int64_t{6}));
+ ASSERT_OK_AND_ASSIGN(std::vector<Row> memory_key_candidates,
+ read_candidates(memory_key_predicate));
+ ASSERT_EQ((std::vector<Row>{{6, "memory-match", "p0"}}),
memory_key_candidates);
+
+ const std::string high_match = "high-match";
+ std::shared_ptr<Predicate> non_key_predicate = PredicateBuilder::Equal(
+ /*field_index=*/1, /*field_name=*/"payload", FieldType::STRING,
+ Literal(FieldType::STRING, high_match.data(), high_match.size()));
+ ASSERT_OK_AND_ASSIGN(std::vector<Row> non_key_candidates,
read_candidates(non_key_predicate));
+ ASSERT_EQ((std::vector<Row>{{1, "high-match", "p0"},
+ {4, "level0-match", "p0"},
+ {5, "level0-other", "p0"},
+ {6, "memory-match", "p0"},
+ {7, "memory-other", "p0"}}),
+ non_key_candidates);
+
+ ASSERT_OK_AND_ASSIGN(
+ std::shared_ptr<Predicate> mixed_predicate,
+ PredicateBuilder::And(
+ {PredicateBuilder::GreaterOrEqual(/*field_index=*/0,
/*field_name=*/"id",
+ FieldType::BIGINT,
Literal(int64_t{4})),
+ PredicateBuilder::Equal(
+ /*field_index=*/1, /*field_name=*/"payload",
FieldType::STRING,
+ Literal(FieldType::STRING, high_match.data(),
high_match.size()))}));
+ ASSERT_OK_AND_ASSIGN(std::vector<Row> mixed_candidates,
read_candidates(mixed_predicate));
+ ASSERT_EQ((std::vector<Row>{{4, "level0-match", "p0"},
+ {5, "level0-other", "p0"},
+ {6, "memory-match", "p0"},
+ {7, "memory-other", "p0"}}),
+ mixed_candidates);
ASSERT_OK(PrepareAndClose(writer.get()));
}