Skip to content

Commit d46fd39

Browse files
committed
fix(scan): ignore cached manifest entries of a rewritten snapshot
The snapshot live manifest entry cache took a matching snapshot id as a hit. A rollback deletes the newest snapshots and the commits that follow reuse their ids, so a process that kept the cache planned the rewritten snapshot from the entries of the deleted one: reads of the deleted files failed and rows committed after the rollback were missing until the ids moved past the cached ones. Record the delta manifest list the entries were built from with every cached snapshot and require it to match on a hit. Its name is unique per commit, so a rewritten snapshot misses and is rebuilt in place. The serialized form changes with it; the magic moves from SMEC to SMED so bytes of the old layout are rebuilt rather than misread. (cherry picked from commit 3b76467)
1 parent ded7cfb commit d46fd39

5 files changed

Lines changed: 207 additions & 20 deletions

File tree

‎src/paimon/core/manifest/snapshot_live_manifest_entries.cpp‎

Lines changed: 25 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@
2121

2222
#include <algorithm>
2323
#include <cstring>
24+
#include <string>
2425
#include <utility>
2526

2627
#include "paimon/common/io/memory_segment_output_stream.h"
@@ -34,7 +35,9 @@
3435
namespace paimon {
3536
namespace {
3637

37-
constexpr int32_t kMagic = 0x534d4543; // SMEC
38+
// Bumped from SMEC when the delta manifest list was added to every cached snapshot; older bytes
39+
// fail the magic check and are rebuilt.
40+
constexpr int32_t kMagic = 0x534d4544; // SMED
3841

3942
size_t NormalizeMaxSnapshots(int32_t max_snapshots) {
4043
return static_cast<size_t>(std::max(0, max_snapshots));
@@ -68,15 +71,27 @@ std::optional<SnapshotLiveManifestEntries::Entry> SnapshotLiveManifestEntries::L
6871
return std::optional<Entry>();
6972
}
7073
--iter;
71-
return Entry{iter->first, iter->second};
74+
return Entry{iter->first, iter->second.delta_manifest_list, iter->second.entries};
7275
}
7376

74-
void SnapshotLiveManifestEntries::Put(int64_t snapshot_id, std::vector<ManifestEntry>&& entries) {
77+
std::optional<SnapshotLiveManifestEntries::Entry> SnapshotLiveManifestEntries::Find(
78+
int64_t snapshot_id, const std::string& delta_manifest_list) const {
79+
auto iter = entries_by_snapshot_.find(snapshot_id);
80+
if (iter == entries_by_snapshot_.end() ||
81+
iter->second.delta_manifest_list != delta_manifest_list) {
82+
return std::optional<Entry>();
83+
}
84+
return Entry{iter->first, iter->second.delta_manifest_list, iter->second.entries};
85+
}
86+
87+
void SnapshotLiveManifestEntries::Put(int64_t snapshot_id, const std::string& delta_manifest_list,
88+
std::vector<ManifestEntry>&& entries) {
7589
if (NormalizeMaxSnapshots(max_snapshots_) == 0) {
7690
return;
7791
}
7892
entries_by_snapshot_[snapshot_id] =
79-
std::make_shared<const std::vector<ManifestEntry>>(std::move(entries));
93+
Value{delta_manifest_list,
94+
std::make_shared<const std::vector<ManifestEntry>>(std::move(entries))};
8095
EvictIfNeeded();
8196
}
8297

@@ -91,9 +106,10 @@ Result<std::shared_ptr<Bytes>> SnapshotLiveManifestEntries::Serialize(
91106
out.WriteValue<int32_t>(static_cast<int32_t>(entries_by_snapshot_.size()));
92107

93108
ManifestEntrySerializer serializer(pool);
94-
for (const auto& [snapshot_id, entries] : entries_by_snapshot_) {
109+
for (const auto& [snapshot_id, value] : entries_by_snapshot_) {
95110
out.WriteValue<int64_t>(snapshot_id);
96-
PAIMON_RETURN_NOT_OK(serializer.SerializeList(*entries, &out));
111+
out.WriteString(value.delta_manifest_list);
112+
PAIMON_RETURN_NOT_OK(serializer.SerializeList(*value.entries, &out));
97113
}
98114
return ToBytes(out, pool);
99115
}
@@ -121,9 +137,11 @@ Result<SnapshotLiveManifestEntries> SnapshotLiveManifestEntries::Deserialize(
121137
ManifestEntrySerializer serializer(pool);
122138
for (int32_t i = 0; i < snapshot_count; i++) {
123139
PAIMON_ASSIGN_OR_RAISE(int64_t snapshot_id, in.ReadValue<int64_t>());
140+
PAIMON_ASSIGN_OR_RAISE(std::string delta_manifest_list, in.ReadString());
124141
PAIMON_ASSIGN_OR_RAISE(std::vector<ManifestEntry> entries, serializer.DeserializeList(&in));
125142
snapshot_live_manifest_entries.entries_by_snapshot_[snapshot_id] =
126-
std::make_shared<const std::vector<ManifestEntry>>(std::move(entries));
143+
Value{std::move(delta_manifest_list),
144+
std::make_shared<const std::vector<ManifestEntry>>(std::move(entries))};
127145
}
128146
snapshot_live_manifest_entries.EvictIfNeeded();
129147
return snapshot_live_manifest_entries;

‎src/paimon/core/manifest/snapshot_live_manifest_entries.h‎

Lines changed: 18 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@
2424
#include <map>
2525
#include <memory>
2626
#include <optional>
27+
#include <string>
2728
#include <vector>
2829

2930
#include "paimon/core/manifest/manifest_entry.h"
@@ -38,17 +39,27 @@ class MemorySegment;
3839
///
3940
/// This value object owns merged live manifest entries by snapshot id. It does not own or access a
4041
/// cache; callers are responsible for storing the serialized bytes in the cache layer.
42+
///
43+
/// A snapshot id alone does not identify a snapshot: after a rollback, or after a table is dropped
44+
/// and recreated at the same path, later commits reuse the ids of the deleted snapshots. Every
45+
/// cached snapshot therefore also records the delta manifest list it was built from, whose name is
46+
/// unique per commit, and `Find()` only reports a hit when both match.
4147
class SnapshotLiveManifestEntries {
4248
public:
4349
struct Entry {
4450
int64_t snapshot_id;
51+
std::string delta_manifest_list;
4552
std::shared_ptr<const std::vector<ManifestEntry>> entries;
4653
};
4754

4855
explicit SnapshotLiveManifestEntries(int32_t max_snapshots);
4956

5057
std::optional<Entry> LatestBeforeOrEqual(int64_t snapshot_id) const;
51-
void Put(int64_t snapshot_id, std::vector<ManifestEntry>&& entries);
58+
/// Returns the entries cached for exactly this snapshot: the same id, built from the same delta
59+
/// manifest list. A snapshot with the same id but another delta manifest list is a miss.
60+
std::optional<Entry> Find(int64_t snapshot_id, const std::string& delta_manifest_list) const;
61+
void Put(int64_t snapshot_id, const std::string& delta_manifest_list,
62+
std::vector<ManifestEntry>&& entries);
5263
size_t Size() const;
5364

5465
Result<std::shared_ptr<Bytes>> Serialize(const std::shared_ptr<MemoryPool>& pool) const;
@@ -59,7 +70,12 @@ class SnapshotLiveManifestEntries {
5970
private:
6071
void EvictIfNeeded();
6172

62-
std::map<int64_t, std::shared_ptr<const std::vector<ManifestEntry>>> entries_by_snapshot_;
73+
struct Value {
74+
std::string delta_manifest_list;
75+
std::shared_ptr<const std::vector<ManifestEntry>> entries;
76+
};
77+
78+
std::map<int64_t, Value> entries_by_snapshot_;
6379
int32_t max_snapshots_;
6480
};
6581

‎src/paimon/core/operation/file_store_scan.cpp‎

Lines changed: 7 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -325,9 +325,10 @@ Status FileStoreScan::ReadManifestEntries(const std::vector<ManifestFileMeta>& m
325325
}
326326

327327
// Cache merged live manifest entries for one bucket before applying scan filters. Each cache value
328-
// keeps a bounded number of snapshot results for the same table/branch/bucket. Exact snapshot hits
329-
// can be returned directly; cache misses rebuild the target snapshot bucket from the target
330-
// snapshot's data manifests.
328+
// keeps a bounded number of snapshot results for the same table/branch/bucket. A hit needs the same
329+
// snapshot id built from the same delta manifest list, because a rollback (or a table recreated at
330+
// the same path) reuses snapshot ids for different content; cache misses rebuild the target
331+
// snapshot bucket from the target snapshot's data manifests.
331332
Status FileStoreScan::ReadManifestEntriesWithCache(
332333
const Snapshot& snapshot, const std::vector<ManifestFileMeta>& all_manifest_metas,
333334
int32_t bucket, std::vector<ManifestEntry>* manifest_entries, bool* cache_hit) const {
@@ -339,8 +340,8 @@ Status FileStoreScan::ReadManifestEntriesWithCache(
339340
metrics_->ObserveHistogram(ScanMetrics::SNAPSHOT_CACHE_LOAD_DURATION,
340341
static_cast<double>(cache_load_duration_ms));
341342
std::optional<SnapshotLiveManifestEntries::Entry> cached =
342-
cached_entries.LatestBeforeOrEqual(snapshot.Id());
343-
if (cached && cached->snapshot_id == snapshot.Id()) {
343+
cached_entries.Find(snapshot.Id(), snapshot.DeltaManifestList());
344+
if (cached) {
344345
*cache_hit = true;
345346
*manifest_entries = *cached->entries;
346347
return Status::OK();
@@ -358,7 +359,7 @@ Status FileStoreScan::ReadManifestEntriesWithCache(
358359
PAIMON_RETURN_NOT_OK(
359360
ReadAndMergeBucketFileEntries(bucket_manifest_metas, bucket, manifest_entries));
360361
std::vector<ManifestEntry> cache_entries = *manifest_entries;
361-
cached_entries.Put(snapshot.Id(), std::move(cache_entries));
362+
cached_entries.Put(snapshot.Id(), snapshot.DeltaManifestList(), std::move(cache_entries));
362363
Duration cache_store_duration;
363364
PAIMON_RETURN_NOT_OK(StoreSnapshotLiveManifestEntries(bucket, cached_entries));
364365
const uint64_t cache_store_duration_ms = cache_store_duration.Get();

‎src/paimon/core/operation/file_store_scan_test.cpp‎

Lines changed: 24 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -203,17 +203,33 @@ TEST_F(FileStoreScanTest, TestSnapshotLiveManifestEntries) {
203203
snapshot1.emplace_back(FileKind::Add(), BinaryRow::EmptyRow(), /*bucket=*/0,
204204
/*total_buckets=*/1, file1);
205205
SnapshotLiveManifestEntries entries(/*max_snapshots=*/2);
206-
entries.Put(/*snapshot_id=*/1, std::move(snapshot1));
206+
entries.Put(/*snapshot_id=*/1, "manifest-list-1", std::move(snapshot1));
207207
ASSERT_EQ(entries.Size(), 1);
208208
auto hit = entries.LatestBeforeOrEqual(/*snapshot_id=*/1);
209209
ASSERT_TRUE(hit);
210210
ASSERT_EQ(hit->snapshot_id, 1);
211+
ASSERT_EQ(hit->delta_manifest_list, "manifest-list-1");
211212
ASSERT_EQ(hit->entries->size(), 1);
212213
ASSERT_EQ((*hit->entries)[0].FileName(), "file-1");
213214
auto latest_before_2 = entries.LatestBeforeOrEqual(/*snapshot_id=*/2);
214215
ASSERT_TRUE(latest_before_2);
215216
ASSERT_EQ(latest_before_2->snapshot_id, 1);
216217

218+
// Find() needs the same id built from the same delta manifest list.
219+
auto found = entries.Find(/*snapshot_id=*/1, "manifest-list-1");
220+
ASSERT_TRUE(found);
221+
ASSERT_EQ((*found->entries)[0].FileName(), "file-1");
222+
ASSERT_FALSE(entries.Find(/*snapshot_id=*/1, "manifest-list-1-rewritten"));
223+
ASSERT_FALSE(entries.Find(/*snapshot_id=*/2, "manifest-list-1"));
224+
225+
// A rewritten snapshot with the same id replaces the stale entries in place.
226+
entries.Put(/*snapshot_id=*/1, "manifest-list-1-rewritten", {});
227+
ASSERT_EQ(entries.Size(), 1);
228+
ASSERT_FALSE(entries.Find(/*snapshot_id=*/1, "manifest-list-1"));
229+
auto rewritten = entries.Find(/*snapshot_id=*/1, "manifest-list-1-rewritten");
230+
ASSERT_TRUE(rewritten);
231+
ASSERT_TRUE(rewritten->entries->empty());
232+
217233
std::vector<ManifestEntry> snapshot3;
218234
ASSERT_OK_AND_ASSIGN(
219235
auto file3,
@@ -225,13 +241,13 @@ TEST_F(FileStoreScanTest, TestSnapshotLiveManifestEntries) {
225241
/*write_cols=*/std::nullopt));
226242
snapshot3.emplace_back(FileKind::Add(), BinaryRow::EmptyRow(), /*bucket=*/0,
227243
/*total_buckets=*/1, file3);
228-
entries.Put(/*snapshot_id=*/3, std::move(snapshot3));
244+
entries.Put(/*snapshot_id=*/3, "manifest-list-3", std::move(snapshot3));
229245

230246
auto latest_before_4 = entries.LatestBeforeOrEqual(/*snapshot_id=*/4);
231247
ASSERT_TRUE(latest_before_4);
232248
ASSERT_EQ(latest_before_4->snapshot_id, 3);
233249

234-
entries.Put(/*snapshot_id=*/5, {});
250+
entries.Put(/*snapshot_id=*/5, "manifest-list-5", {});
235251
ASSERT_EQ(entries.Size(), 2);
236252
ASSERT_FALSE(entries.LatestBeforeOrEqual(/*snapshot_id=*/1));
237253
ASSERT_TRUE(entries.LatestBeforeOrEqual(/*snapshot_id=*/3));
@@ -251,8 +267,8 @@ TEST_F(FileStoreScanTest, TestSnapshotLiveManifestEntriesSerialization) {
251267
manifest_entries.emplace_back(FileKind::Add(), BinaryRow::EmptyRow(), /*bucket=*/0,
252268
/*total_buckets=*/1, file1);
253269
SnapshotLiveManifestEntries entries(/*max_snapshots=*/2);
254-
entries.Put(/*snapshot_id=*/1, std::move(manifest_entries));
255-
entries.Put(/*snapshot_id=*/3, {});
270+
entries.Put(/*snapshot_id=*/1, "manifest-list-1", std::move(manifest_entries));
271+
entries.Put(/*snapshot_id=*/3, "manifest-list-3", {});
256272

257273
ASSERT_OK_AND_ASSIGN(auto bytes, entries.Serialize(GetDefaultPool()));
258274
ASSERT_OK_AND_ASSIGN(auto deserialized,
@@ -262,9 +278,12 @@ TEST_F(FileStoreScanTest, TestSnapshotLiveManifestEntriesSerialization) {
262278
auto hit = deserialized.LatestBeforeOrEqual(/*snapshot_id=*/2);
263279
ASSERT_TRUE(hit);
264280
ASSERT_EQ(hit->snapshot_id, 1);
281+
ASSERT_EQ(hit->delta_manifest_list, "manifest-list-1");
265282
ASSERT_EQ(hit->entries->size(), 1);
266283
ASSERT_EQ((*hit->entries)[0].FileName(), "file-1");
267284
ASSERT_EQ(deserialized.LatestBeforeOrEqual(/*snapshot_id=*/4)->snapshot_id, 3);
285+
ASSERT_TRUE(deserialized.Find(/*snapshot_id=*/3, "manifest-list-3"));
286+
ASSERT_FALSE(deserialized.Find(/*snapshot_id=*/3, "manifest-list-1"));
268287

269288
ASSERT_OK_AND_ASSIGN(auto evicted_deserialized, SnapshotLiveManifestEntries::Deserialize(
270289
MemorySegment::Wrap(bytes),

‎src/paimon/core/table/source/table_scan_test.cpp‎

Lines changed: 133 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,16 +19,33 @@
1919

2020
#include "paimon/table/source/table_scan.h"
2121

22+
#include <algorithm>
23+
#include <cstdint>
24+
#include <map>
2225
#include <memory>
2326
#include <string>
2427
#include <utility>
2528
#include <vector>
2629

30+
#include "arrow/api.h"
31+
#include "fmt/format.h"
2732
#include "gtest/gtest.h"
33+
#include "paimon/common/io/cache/lru_cache.h"
34+
#include "paimon/common/utils/checked_cast.h"
35+
#include "paimon/common/utils/path_util.h"
36+
#include "paimon/core/snapshot.h"
37+
#include "paimon/core/utils/snapshot_manager.h"
2838
#include "paimon/defs.h"
39+
#include "paimon/fs/file_system.h"
2940
#include "paimon/metrics.h"
41+
#include "paimon/read_context.h"
42+
#include "paimon/record_batch.h"
3043
#include "paimon/scan_context.h"
3144
#include "paimon/status.h"
45+
#include "paimon/table/source/scan_metrics.h"
46+
#include "paimon/table/source/table_read.h"
47+
#include "paimon/testing/utils/read_result_collector.h"
48+
#include "paimon/testing/utils/test_helper.h"
3249
#include "paimon/testing/utils/testharness.h"
3350

3451
namespace paimon::test {
@@ -122,4 +139,120 @@ TEST(TableScanTest, TestReadOptimizedPrimaryKeyStreamingScanUnsupported) {
122139
"key table");
123140
}
124141

142+
namespace {
143+
144+
struct CachedScanResult {
145+
std::vector<std::string> rows;
146+
uint64_t cache_hit = 0;
147+
};
148+
149+
// Plans bucket 0 through the snapshot live manifest entry cache and reads every row back as
150+
// "k=v", sorted.
151+
Result<CachedScanResult> ScanBucketThroughCache(const std::string& table_path,
152+
const std::map<std::string, std::string>& options,
153+
const std::shared_ptr<Cache>& cache) {
154+
ScanContextBuilder scan_builder(table_path);
155+
scan_builder.SetOptions(options).WithCache(cache).SetBucketFilter(0);
156+
PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<ScanContext> scan_context, scan_builder.Finish());
157+
PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<TableScan> table_scan,
158+
TableScan::Create(std::move(scan_context)));
159+
PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<Plan> plan, table_scan->CreatePlan());
160+
CachedScanResult result;
161+
std::shared_ptr<Metrics> metrics = table_scan->GetMetrics();
162+
PAIMON_ASSIGN_OR_RAISE(uint64_t cache_enabled,
163+
metrics->GetCounter(ScanMetrics::LAST_SNAPSHOT_CACHE_ENABLED));
164+
if (cache_enabled != 1) {
165+
return Status::Invalid("the snapshot live manifest entry cache is not in use");
166+
}
167+
PAIMON_ASSIGN_OR_RAISE(result.cache_hit,
168+
metrics->GetCounter(ScanMetrics::LAST_SNAPSHOT_CACHE_HIT));
169+
170+
ReadContextBuilder read_builder(table_path);
171+
read_builder.SetOptions(options).WithCache(cache);
172+
PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<ReadContext> read_context, read_builder.Finish());
173+
PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<TableRead> table_read,
174+
TableRead::Create(std::move(read_context)));
175+
PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<BatchReader> reader,
176+
table_read->CreateReader(plan->Splits()));
177+
PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<arrow::ChunkedArray> chunks,
178+
ReadResultCollector::CollectResult(std::move(reader)));
179+
for (const std::shared_ptr<arrow::Array>& chunk : chunks->chunks()) {
180+
const auto& rows = checked_cast<const arrow::StructArray&>(*chunk);
181+
std::shared_ptr<arrow::Array> key_column = rows.GetFieldByName("k");
182+
std::shared_ptr<arrow::Array> value_column = rows.GetFieldByName("v");
183+
if (key_column == nullptr || key_column->type_id() != arrow::Type::STRING ||
184+
value_column == nullptr || value_column->type_id() != arrow::Type::STRING) {
185+
return Status::Invalid("read result misses string columns k and v");
186+
}
187+
auto keys = checked_pointer_cast<arrow::StringArray>(key_column);
188+
auto values = checked_pointer_cast<arrow::StringArray>(value_column);
189+
for (int64_t i = 0; i < rows.length(); ++i) {
190+
result.rows.push_back(fmt::format("{}={}", keys->GetString(i), values->GetString(i)));
191+
}
192+
}
193+
std::sort(result.rows.begin(), result.rows.end());
194+
return result;
195+
}
196+
197+
} // namespace
198+
199+
// A rollback deletes the newest snapshots, and the commits that follow reuse their ids for
200+
// different content. Entries cached for the deleted snapshot must not serve the rewritten one.
201+
TEST(TableScanTest, TestManifestEntryCacheIgnoresRewrittenSnapshot) {
202+
std::unique_ptr<UniqueTestDirectory> dir = UniqueTestDirectory::Create("local");
203+
std::map<std::string, std::string> options = {
204+
{Options::BUCKET, "1"},
205+
{Options::SCAN_MANIFEST_ENTRY_CACHE_MAX_SNAPSHOTS, "2"},
206+
{Options::SCAN_MANIFEST_ENTRY_LAZY_DECODE_ENABLED, "true"}};
207+
arrow::FieldVector fields = {arrow::field("k", arrow::utf8(), /*nullable=*/false),
208+
arrow::field("v", arrow::utf8())};
209+
ASSERT_OK_AND_ASSIGN(std::unique_ptr<TestHelper> helper,
210+
TestHelper::Create(dir->Str(), arrow::schema(fields),
211+
/*partition_keys=*/{}, /*primary_keys=*/{"k"}, options,
212+
/*is_streaming_mode=*/true));
213+
std::string table_path = PathUtil::JoinPath(dir->Str(), "foo.db/bar");
214+
auto row_type = arrow::struct_(fields);
215+
auto write = [&](TestHelper* writer, const std::string& rows, int64_t commit_identifier) {
216+
ASSERT_OK_AND_ASSIGN(std::unique_ptr<RecordBatch> batch,
217+
TestHelper::MakeRecordBatch(row_type, rows, /*partition_map=*/{},
218+
/*bucket=*/0, /*row_kinds=*/{}));
219+
ASSERT_OK(writer->WriteAndCommit(std::move(batch), commit_identifier, std::nullopt));
220+
};
221+
write(helper.get(), R"([["k1", "v1"], ["k2", "v2"]])", 1);
222+
write(helper.get(), R"([["k3", "v3"]])", 2);
223+
write(helper.get(), R"([["k1", "v1b"]])", 3);
224+
225+
auto cache = std::make_shared<LruCache>(/*max_weight=*/64 * 1024 * 1024);
226+
ASSERT_OK_AND_ASSIGN(CachedScanResult before,
227+
ScanBucketThroughCache(table_path, options, cache));
228+
ASSERT_EQ(before.rows, (std::vector<std::string>{"k1=v1b", "k2=v2", "k3=v3"}));
229+
ASSERT_EQ(before.cache_hit, 0);
230+
231+
// Roll back to snapshot 1 the way `rollback_to` does: drop the newer snapshot files and move
232+
// the LATEST hint. The next commits create a new snapshot 2 and 3.
233+
std::shared_ptr<FileSystem> fs = dir->GetFileSystem();
234+
SnapshotManager snapshot_manager(fs, table_path);
235+
ASSERT_OK_AND_ASSIGN(Snapshot old_snapshot_3, snapshot_manager.LoadSnapshot(3));
236+
ASSERT_OK(fs->Delete(snapshot_manager.SnapshotPath(3), /*recursive=*/false));
237+
ASSERT_OK(fs->Delete(snapshot_manager.SnapshotPath(2), /*recursive=*/false));
238+
ASSERT_OK(snapshot_manager.CommitLatestHint(1));
239+
ASSERT_OK_AND_ASSIGN(std::unique_ptr<TestHelper> writer_after_rollback,
240+
TestHelper::Create(table_path, options, /*is_streaming_mode=*/true));
241+
write(writer_after_rollback.get(), R"([["k4", "v4"]])", 4);
242+
write(writer_after_rollback.get(), R"([["k1", "v1c"]])", 5);
243+
ASSERT_OK_AND_ASSIGN(Snapshot new_snapshot_3, snapshot_manager.LoadSnapshot(3));
244+
ASSERT_NE(new_snapshot_3.DeltaManifestList(), old_snapshot_3.DeltaManifestList());
245+
246+
ASSERT_OK_AND_ASSIGN(CachedScanResult after,
247+
ScanBucketThroughCache(table_path, options, cache));
248+
ASSERT_EQ(after.rows, (std::vector<std::string>{"k1=v1c", "k2=v2", "k4=v4"}));
249+
ASSERT_EQ(after.cache_hit, 0);
250+
251+
// The rewritten snapshot 3 is cached in its own right afterwards.
252+
ASSERT_OK_AND_ASSIGN(CachedScanResult again,
253+
ScanBucketThroughCache(table_path, options, cache));
254+
ASSERT_EQ(again.rows, after.rows);
255+
ASSERT_EQ(again.cache_hit, 1);
256+
}
257+
125258
} // namespace paimon::test

0 commit comments

Comments
 (0)