This is an automated email from the ASF dual-hosted git repository.

yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/branch-4.1 by this push:
     new 03b6cc30fb0 branch-4.1: [fix](cloud) Exclude covered rowsets from 
compaction minimum timestamps #67617 (#67751)
03b6cc30fb0 is described below

commit 03b6cc30fb0f95f9faa90039152e1e1599fc5fe2
Author: github-actions[bot] 
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Thu Sep 10 21:29:16 2026 +0800

    branch-4.1: [fix](cloud) Exclude covered rowsets from compaction minimum 
timestamps #67617 (#67751)
    
    Cherry-picked from #67617
    
    Co-authored-by: Luwei <[email protected]>
---
 cloud/src/meta-store/meta_reader.cpp       |  13 ++-
 cloud/test/clone_chain_reader_test.cpp     |  52 +++++++++
 cloud/test/meta_reader_test.cpp            | 130 +++++++++++++++++++++
 cloud/test/recycler_operation_log_test.cpp | 179 +++++++++++++++++++++++++++++
 4 files changed, 368 insertions(+), 6 deletions(-)

diff --git a/cloud/src/meta-store/meta_reader.cpp 
b/cloud/src/meta-store/meta_reader.cpp
index 6a62caa19cf..6105678ad1d 100644
--- a/cloud/src/meta-store/meta_reader.cpp
+++ b/cloud/src/meta-store/meta_reader.cpp
@@ -679,7 +679,7 @@ TxnErrorCode MetaReader::get_rowset_metas(Transaction* txn, 
int64_t tablet_id,
                                           int64_t start_version, int64_t 
end_version,
                                           std::vector<RowsetMetaCloudPB>* 
rowset_metas,
                                           bool snapshot) {
-    std::map<int64_t, RowsetMetaCloudPB> rowset_graph;
+    std::map<int64_t, std::pair<RowsetMetaCloudPB, Versionstamp>> rowset_graph;
 
     {
         std::string start_key =
@@ -698,8 +698,7 @@ TxnErrorCode MetaReader::get_rowset_metas(Transaction* txn, 
int64_t tablet_id,
                 versioned::document_get_range<RowsetMetaCloudPB>(txn, 
start_key, end_key, options);
         for (auto&& kvp = iter->next(); kvp.has_value(); kvp = iter->next()) {
             auto&& [key, version, rowset_meta] = *kvp;
-            rowset_graph.emplace(rowset_meta.end_version(), 
std::move(rowset_meta));
-            min_read_versionstamp_ = std::min(min_read_versionstamp_, version);
+            rowset_graph.try_emplace(rowset_meta.end_version(), 
std::move(rowset_meta), version);
             DCHECK(version < snapshot_version_)
                     << "version: " << version.to_string()
                     << ", snapshot_version: " << snapshot_version_.to_string();
@@ -744,12 +743,11 @@ TxnErrorCode MetaReader::get_rowset_metas(Transaction* 
txn, int64_t tablet_id,
                 continue;
             }
 
-            min_read_versionstamp_ = std::min(min_read_versionstamp_, version);
             last_start_version = start_version;
             // erase the rowsets that are covered by this compact rowset
             rowset_graph.erase(rowset_graph.lower_bound(start_version),
                                rowset_graph.upper_bound(end_version));
-            rowset_graph.emplace(end_version, std::move(rowset_meta));
+            rowset_graph.try_emplace(end_version, std::move(rowset_meta), 
version);
         }
         if (!iter->is_valid()) {
             LOG_ERROR("failed to get compacted rowset metas")
@@ -764,7 +762,10 @@ TxnErrorCode MetaReader::get_rowset_metas(Transaction* 
txn, int64_t tablet_id,
 
     rowset_metas->clear();
     rowset_metas->reserve(rowset_graph.size());
-    for (auto&& [version, rowset_meta] : rowset_graph) {
+    // Only returned rowsets contribute dependencies after both scans succeed.
+    for (auto&& [end_version, entry] : rowset_graph) {
+        auto& [rowset_meta, versionstamp] = entry;
+        min_read_versionstamp_ = std::min(min_read_versionstamp_, 
versionstamp);
         rowset_metas->emplace_back(std::move(rowset_meta));
     }
 
diff --git a/cloud/test/clone_chain_reader_test.cpp 
b/cloud/test/clone_chain_reader_test.cpp
index fdc6be6ef6a..5973e54fb02 100644
--- a/cloud/test/clone_chain_reader_test.cpp
+++ b/cloud/test/clone_chain_reader_test.cpp
@@ -1831,6 +1831,58 @@ TEST_F(CloneChainReaderTest, GetRowsetMeta) {
     }
 }
 
+TEST_F(CloneChainReaderTest, GetRowsetMetasMinReadVersionstamp) {
+    constexpr int64_t tablet_id = 15001;
+    std::unique_ptr<Transaction> txn;
+    ASSERT_EQ(txn_kv_->create_txn(&txn), TxnErrorCode::TXN_OK);
+    auto put_rowset = [&](const std::string& instance_id, const std::string& 
id, int64_t start,
+                          int64_t end, Versionstamp version, bool compact) {
+        doris::RowsetMetaCloudPB rowset;
+        rowset.set_rowset_id(0);
+        rowset.set_rowset_id_v2(id);
+        rowset.set_tablet_id(tablet_id);
+        rowset.set_start_version(start);
+        rowset.set_end_version(end);
+        auto key = compact ? versioned::meta_rowset_compact_key({instance_id, 
tablet_id, end})
+                           : versioned::meta_rowset_load_key({instance_id, 
tablet_id, end});
+        ASSERT_TRUE(versioned::document_put(txn.get(), key, version, 
std::move(rowset)));
+    };
+    put_rowset("A", "L1", 1, 1, Versionstamp(70, 1), false);
+    put_rowset("A", "L2", 2, 2, Versionstamp(80, 1), false);
+    put_rowset("A", "L3", 3, 3, Versionstamp(90, 1), false);
+    put_rowset("A", "compact", 2, 3, Versionstamp(150, 2), true);
+    // This source rowset is newer than A's source snapshot (1000) and is 
invisible to C.
+    put_rowset("A", "too_new", 1, 3, Versionstamp(1000, 1), true);
+    put_rowset("C", "L4", 4, 4, Versionstamp(2170, 1), false);
+    ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK);
+
+    CloneChainReader reader(instance_ids_[2], Versionstamp(3000), 
txn_kv_.get(),
+                            resource_mgr_.get());
+    std::vector<doris::RowsetMetaCloudPB> rowsets;
+    ASSERT_EQ(reader.get_rowset_metas(tablet_id, 2, 4, &rowsets), 
TxnErrorCode::TXN_OK);
+    ASSERT_EQ(rowsets.size(), 2);
+    EXPECT_EQ(rowsets[0].rowset_id_v2(), "compact");
+    EXPECT_EQ(rowsets[0].start_version(), 2);
+    EXPECT_EQ(rowsets[0].end_version(), 3);
+    EXPECT_EQ(rowsets[0].reference_instance_id(), "A");
+    EXPECT_EQ(rowsets[1].rowset_id_v2(), "L4");
+    EXPECT_EQ(rowsets[1].start_version(), 4);
+    EXPECT_EQ(rowsets[1].end_version(), 4);
+    EXPECT_EQ(rowsets[1].reference_instance_id(), "C");
+    EXPECT_EQ(reader.min_read_versionstamp(), Versionstamp(150, 2));
+
+    ASSERT_EQ(reader.get_rowset_metas(tablet_id, 1, 1, &rowsets), 
TxnErrorCode::TXN_OK);
+    ASSERT_EQ(rowsets.size(), 1);
+    EXPECT_EQ(rowsets[0].rowset_id_v2(), "L1");
+    EXPECT_EQ(rowsets[0].reference_instance_id(), "A");
+    EXPECT_EQ(reader.min_read_versionstamp(), Versionstamp(70, 1));
+    ASSERT_EQ(reader.get_rowset_metas(tablet_id, 2, 4, &rowsets), 
TxnErrorCode::TXN_OK);
+    EXPECT_EQ(reader.min_read_versionstamp(), Versionstamp(70, 1));
+    ASSERT_EQ(reader.get_rowset_metas(tablet_id, 5, 6, &rowsets), 
TxnErrorCode::TXN_OK);
+    EXPECT_TRUE(rowsets.empty());
+    EXPECT_EQ(reader.min_read_versionstamp(), Versionstamp(70, 1));
+}
+
 TEST_F(CloneChainReaderTest, GetRowsetMetas) {
     std::string instance_id = instance_ids_[2]; // C
     Versionstamp snapshot_version = snapshot_versions_[2];
diff --git a/cloud/test/meta_reader_test.cpp b/cloud/test/meta_reader_test.cpp
index 4a5ba75aae1..dac4f370b28 100644
--- a/cloud/test/meta_reader_test.cpp
+++ b/cloud/test/meta_reader_test.cpp
@@ -23,6 +23,7 @@
 #include <gtest/gtest.h>
 
 #include <memory>
+#include <tuple>
 
 #include "common/config.h"
 #include "common/logging.h"
@@ -1616,6 +1617,135 @@ TEST(MetaReaderTest, GetRowsetMetas) {
     }
 }
 
+TEST(MetaReaderTest, GetRowsetMetasMinReadVersionstamp) {
+    auto txn_kv = std::make_shared<MemTxnKv>();
+    ASSERT_EQ(txn_kv->init(), 0);
+    const std::string instance_id = "rowset_min_version";
+    constexpr int64_t tablet_id = 4001;
+
+    std::unique_ptr<Transaction> txn;
+    ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
+    auto put_rowset = [&](const std::string& id, int64_t start, int64_t end, 
Versionstamp version,
+                          bool compact) {
+        doris::RowsetMetaCloudPB rowset;
+        rowset.set_rowset_id(0);
+        rowset.set_rowset_id_v2(id);
+        rowset.set_tablet_id(tablet_id);
+        rowset.set_start_version(start);
+        rowset.set_end_version(end);
+        auto key = compact ? versioned::meta_rowset_compact_key({instance_id, 
tablet_id, end})
+                           : versioned::meta_rowset_load_key({instance_id, 
tablet_id, end});
+        ASSERT_TRUE(versioned::document_put(txn.get(), key, version, 
std::move(rowset)));
+    };
+    put_rowset("L1", 1, 1, Versionstamp(70, 1), false);
+    put_rowset("L2", 2, 2, Versionstamp(80, 1), false);
+    put_rowset("L3", 3, 3, Versionstamp(90, 1), false);
+    put_rowset("A", 2, 3, Versionstamp(150, 2), true);
+    put_rowset("L4", 4, 4, Versionstamp(170, 1), false);
+    // Compact keys are scanned by descending end version: A also hides this 
older compact.
+    put_rowset("old_compact", 2, 2, Versionstamp(110, 1), true);
+    put_rowset("old_A", 2, 3, Versionstamp(120, 1), true);
+    TabletStatsPB stats;
+    stats.set_num_rowsets(1);
+    versioned_put(txn.get(), versioned::tablet_load_stats_key({instance_id, 
tablet_id}),
+                  Versionstamp(70, 2), stats.SerializeAsString());
+    versioned_put(txn.get(), versioned::tablet_compact_stats_key({instance_id, 
tablet_id}),
+                  Versionstamp(75, 1), stats.SerializeAsString());
+    ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK);
+
+    using Rowset = std::tuple<std::string, int64_t, int64_t>;
+    auto check_rowsets = [&](MetaReader& reader, int64_t start, int64_t end,
+                             const std::vector<Rowset>& expected, Versionstamp 
expected_min) {
+        SCOPED_TRACE(fmt::format("rowset range [{}, {}]", start, end));
+        std::vector<doris::RowsetMetaCloudPB> rowsets;
+        ASSERT_EQ(reader.get_rowset_metas(tablet_id, start, end, &rowsets), 
TxnErrorCode::TXN_OK);
+        std::vector<Rowset> actual;
+        for (const auto& rowset : rowsets) {
+            EXPECT_EQ(rowset.tablet_id(), tablet_id);
+            actual.emplace_back(rowset.rowset_id_v2(), rowset.start_version(),
+                                rowset.end_version());
+        }
+        EXPECT_EQ(actual, expected);
+        EXPECT_EQ(reader.min_read_versionstamp(), expected_min)
+                << "actual metadata version: " << reader.min_read_version()
+                << ", expected: " << expected_min.version();
+    };
+    {
+        MetaReader reader(instance_id, txn_kv.get());
+        // Covered L2@80/L3@90 must not lower the minimum of A@150/L4@170.
+        check_rowsets(reader, 2, 4, {{"A", 2, 3}, {"L4", 4, 4}}, 
Versionstamp(150, 2));
+    }
+    {
+        MetaReader reader(instance_id, txn_kv.get());
+        check_rowsets(reader, 2, 3, {{"A", 2, 3}}, Versionstamp(150, 2));
+    }
+    {
+        MetaReader reader(instance_id, txn_kv.get());
+        // A selected old load still contributes its full metadata 
versionstamp.
+        check_rowsets(reader, 1, 4, {{"L1", 1, 1}, {"A", 2, 3}, {"L4", 4, 4}}, 
Versionstamp(70, 1));
+    }
+    {
+        MetaReader reader(instance_id, txn_kv.get(), Versionstamp(100, 0));
+        check_rowsets(reader, 2, 3, {{"L2", 2, 2}, {"L3", 3, 3}}, 
Versionstamp(80, 1));
+    }
+    for (bool load_stats : {true, false}) {
+        SCOPED_TRACE(load_stats ? "prior load stats" : "prior compact stats");
+        MetaReader reader(instance_id, txn_kv.get());
+        auto err = load_stats ? reader.get_tablet_load_stats(tablet_id, 
&stats, nullptr)
+                              : reader.get_tablet_compact_stats(tablet_id, 
&stats, nullptr);
+        ASSERT_EQ(err, TxnErrorCode::TXN_OK);
+        auto previous_min = load_stats ? Versionstamp(70, 2) : 
Versionstamp(75, 1);
+        EXPECT_EQ(reader.min_read_versionstamp(), previous_min);
+        check_rowsets(reader, 2, 4, {{"A", 2, 3}, {"L4", 4, 4}}, previous_min);
+        check_rowsets(reader, 5, 6, {}, previous_min);
+    }
+    {
+        MetaReader reader(instance_id, txn_kv.get());
+        check_rowsets(reader, 5, 6, {}, Versionstamp::max());
+        check_rowsets(reader, 2, 3, {{"A", 2, 3}}, Versionstamp(150, 2));
+        check_rowsets(reader, 4, 4, {{"L4", 4, 4}}, Versionstamp(150, 2));
+        check_rowsets(reader, 1, 1, {{"L1", 1, 1}}, Versionstamp(70, 1));
+        check_rowsets(reader, 2, 4, {{"A", 2, 3}, {"L4", 4, 4}}, 
Versionstamp(70, 1));
+        check_rowsets(reader, 5, 6, {}, Versionstamp(70, 1));
+    }
+    {
+        // Equal transaction versions must still compare the two-byte order 
component.
+        ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
+        put_rowset("L5", 5, 5, Versionstamp(150, 1), false);
+        ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK);
+        MetaReader reader(instance_id, txn_kv.get());
+        check_rowsets(reader, 2, 5, {{"A", 2, 3}, {"L4", 4, 4}, {"L5", 5, 5}},
+                      Versionstamp(150, 1));
+    }
+}
+
+TEST(MetaReaderTest, GetRowsetMetasScanFailure) {
+    auto txn_kv = std::make_shared<MemTxnKv>();
+    ASSERT_EQ(txn_kv->init(), 0);
+    const std::string instance_id = "rowset_scan_failure";
+    constexpr int64_t tablet_id = 4001;
+    std::unique_ptr<Transaction> txn;
+    ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
+    doris::RowsetMetaCloudPB load;
+    load.set_rowset_id(0);
+    load.set_start_version(2);
+    load.set_end_version(2);
+    ASSERT_TRUE(versioned::document_put(
+            txn.get(), versioned::meta_rowset_load_key({instance_id, 
tablet_id, 2}),
+            Versionstamp(80, 1), std::move(load)));
+    // A malformed compact document fails after the load scan has succeeded.
+    versioned_put(txn.get(), versioned::meta_rowset_compact_key({instance_id, 
tablet_id, 3}),
+                  Versionstamp(150, 1), "invalid protobuf");
+    ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK);
+    MetaReader reader(instance_id, txn_kv.get());
+    std::vector<doris::RowsetMetaCloudPB> rowsets(1);
+    rowsets[0].set_rowset_id_v2("unchanged");
+    ASSERT_EQ(reader.get_rowset_metas(tablet_id, 2, 3, &rowsets), 
TxnErrorCode::TXN_INVALID_DATA);
+    ASSERT_EQ(rowsets.size(), 1);
+    EXPECT_EQ(rowsets[0].rowset_id_v2(), "unchanged");
+    EXPECT_EQ(reader.min_read_versionstamp(), Versionstamp::max());
+}
+
 TEST(MetaReaderTest, GetPartitionPendingTxnId) {
     auto txn_kv = std::make_shared<MemTxnKv>();
     ASSERT_EQ(txn_kv->init(), 0);
diff --git a/cloud/test/recycler_operation_log_test.cpp 
b/cloud/test/recycler_operation_log_test.cpp
index 1a840db74df..f4858f9ec3f 100644
--- a/cloud/test/recycler_operation_log_test.cpp
+++ b/cloud/test/recycler_operation_log_test.cpp
@@ -1070,6 +1070,185 @@ doris::RowsetMetaCloudPB create_rowset(int64_t txn_id, 
int64_t tablet_id, int pa
     return rowset;
 }
 
+TEST(RecycleOperationLogTest, CoveredLoadDoesNotProtectLaterCompaction) {
+    auto meta_service = get_meta_service(false);
+    const std::string test_instance_id = "covered_load_compaction";
+    auto txn_kv = std::dynamic_pointer_cast<MemTxnKv>(meta_service->txn_kv());
+    ASSERT_NE(txn_kv, nullptr);
+    auto* sp = SyncPoint::get_instance();
+    DORIS_CLOUD_DEFER {
+        sp->clear_all_call_backs();
+        sp->disable_processing();
+    };
+    sp->set_call_back("get_instance_id", [&](auto&& args) {
+        auto* ret = try_any_cast_ret<std::string>(args);
+        ret->first = test_instance_id;
+        ret->second = true;
+    });
+    sp->enable_processing();
+
+    constexpr int64_t table_id = 20001, index_id = 20002, partition_id = 
20003, tablet_id = 20004;
+    InstanceInfoPB instance;
+    instance.set_instance_id(test_instance_id);
+    instance.set_multi_version_status(MULTI_VERSION_ENABLED);
+    std::unique_ptr<Transaction> txn;
+    ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
+    txn->put(instance_key(test_instance_id), instance.SerializeAsString());
+    ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK);
+    meta_service->resource_mgr()->refresh_instance(test_instance_id);
+    
ASSERT_TRUE(meta_service->resource_mgr()->is_version_read_enabled(test_instance_id));
+    create_tablet(meta_service.get(), table_id, index_id, partition_id, 
tablet_id);
+
+    auto l2 = create_rowset(102, tablet_id, partition_id, 2, 10);
+    auto l3 = create_rowset(103, tablet_id, partition_id, 3, 10);
+    auto l4 = create_rowset(104, tablet_id, partition_id, 4, 10);
+    ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
+    ASSERT_TRUE(versioned::document_put(
+            txn.get(), versioned::meta_rowset_load_key({test_instance_id, 
tablet_id, 2}),
+            Versionstamp(80), doris::RowsetMetaCloudPB(l2)));
+    ASSERT_TRUE(versioned::document_put(
+            txn.get(), versioned::meta_rowset_load_key({test_instance_id, 
tablet_id, 3}),
+            Versionstamp(90), doris::RowsetMetaCloudPB(l3)));
+    SnapshotPB snapshot;
+    snapshot.set_status(SNAPSHOT_NORMAL);
+    versioned_put(txn.get(), versioned::snapshot_full_key(test_instance_id), 
Versionstamp(100),
+                  snapshot.SerializeAsString());
+    // Keep both stats dependencies at/after 150 so only the selected inputs 
decide protection.
+    TabletStatsPB compact_stats;
+    compact_stats.set_num_rows(20);
+    compact_stats.set_num_rowsets(2);
+    compact_stats.set_num_segments(2);
+    compact_stats.set_data_size(2200);
+    compact_stats.set_index_size(200);
+    compact_stats.set_segment_size(2000);
+    compact_stats.set_cumulative_point(2);
+    versioned_put(txn.get(), 
versioned::tablet_compact_stats_key({test_instance_id, tablet_id}),
+                  Versionstamp(150), compact_stats.SerializeAsString());
+    versioned_put(txn.get(), 
versioned::tablet_load_stats_key({test_instance_id, tablet_id}),
+                  Versionstamp(150), TabletStatsPB().SerializeAsString());
+    ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK);
+
+    std::vector<std::string> output_rowset_ids;
+    auto compact = [&](int64_t end, int64_t commit_version, int base_count, 
OperationLogPB* log,
+                       Versionstamp* log_version) {
+        SCOPED_TRACE(commit_version);
+        const int64_t txn_id = 30000 + end;
+        const int64_t rows = (end - 1) * 10;
+        auto output = create_rowset(txn_id, tablet_id, partition_id, 2, rows);
+        output.set_end_version(end);
+        output_rowset_ids.push_back(output.rowset_id_v2());
+        ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
+        txn->put(meta_rowset_tmp_key({test_instance_id, txn_id, tablet_id}),
+                 output.SerializeAsString());
+        ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK);
+
+        StartTabletJobRequest start;
+        auto* idx = start.mutable_job()->mutable_idx();
+        idx->set_table_id(table_id);
+        idx->set_index_id(index_id);
+        idx->set_partition_id(partition_id);
+        idx->set_tablet_id(tablet_id);
+        auto* job = start.mutable_job()->add_compaction();
+        job->set_id(fmt::format("base_{}", base_count));
+        job->set_initiator("test_be");
+        job->set_type(TabletCompactionJobPB::BASE);
+        job->set_base_compaction_cnt(base_count);
+        job->set_cumulative_compaction_cnt(0);
+        job->add_input_versions(2);
+        job->add_input_versions(end);
+        job->set_expiration(time(nullptr) + 3600);
+        job->set_lease(time(nullptr) + 3600);
+        brpc::Controller cntl;
+        StartTabletJobResponse start_response;
+        meta_service->start_tablet_job(&cntl, &start, &start_response, 
nullptr);
+        ASSERT_EQ(start_response.status().code(), MetaServiceCode::OK)
+                << start_response.status().msg();
+
+        FinishTabletJobRequest finish;
+        finish.set_action(FinishTabletJobRequest::COMMIT);
+        finish.mutable_job()->CopyFrom(start.job());
+        job = finish.mutable_job()->mutable_compaction(0);
+        job->add_txn_id(txn_id);
+        job->add_output_versions(end);
+        job->add_output_rowset_ids(output.rowset_id_v2());
+        job->set_output_cumulative_point(2);
+        job->set_num_input_rows(rows);
+        job->set_num_output_rows(rows);
+        job->set_num_input_rowsets(2);
+        job->set_num_output_rowsets(1);
+        job->set_num_input_segments(2);
+        job->set_num_output_segments(1);
+        job->set_size_input_rowsets(rows * 110);
+        job->set_size_output_rowsets(rows * 110);
+        job->set_index_size_input_rowsets(rows * 10);
+        job->set_index_size_output_rowsets(rows * 10);
+        job->set_segment_size_input_rowsets(rows * 100);
+        job->set_segment_size_output_rowsets(rows * 100);
+        txn_kv->update_commit_version(commit_version - 1);
+        FinishTabletJobResponse finish_response;
+        meta_service->finish_tablet_job(&cntl, &finish, &finish_response, 
nullptr);
+        ASSERT_EQ(finish_response.status().code(), MetaServiceCode::OK)
+                << finish_response.status().msg();
+        ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
+        ASSERT_EQ(read_operation_log(txn.get(), 
versioned::log_key(test_instance_id), log_version,
+                                     log),
+                  TxnErrorCode::TXN_OK);
+        ASSERT_EQ(log_version->version(), commit_version);
+        ASSERT_TRUE(log->has_compaction());
+        ASSERT_EQ(log->compaction().recycle_rowsets_size(), 2);
+    };
+
+    OperationLogPB first_log, second_log;
+    Versionstamp first_version, second_version;
+    ASSERT_NO_FATAL_FAILURE(compact(3, 150, 0, &first_log, &first_version));
+    EXPECT_EQ(first_log.min_timestamp(), 80);
+    
EXPECT_EQ(first_log.compaction().recycle_rowsets(0).rowset_meta().rowset_id_v2(),
+              l2.rowset_id_v2());
+    
EXPECT_EQ(first_log.compaction().recycle_rowsets(1).rowset_meta().rowset_id_v2(),
+              l3.rowset_id_v2());
+
+    ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
+    ASSERT_TRUE(versioned::document_put(
+            txn.get(), versioned::meta_rowset_load_key({test_instance_id, 
tablet_id, 4}),
+            Versionstamp(170), doris::RowsetMetaCloudPB(l4)));
+    TabletStatsPB load_stats;
+    load_stats.set_num_rows(10);
+    load_stats.set_num_rowsets(1);
+    load_stats.set_num_segments(1);
+    load_stats.set_data_size(1100);
+    load_stats.set_index_size(100);
+    load_stats.set_segment_size(1000);
+    versioned_put(txn.get(), 
versioned::tablet_load_stats_key({test_instance_id, tablet_id}),
+                  Versionstamp(170), load_stats.SerializeAsString());
+    ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK);
+    ASSERT_NO_FATAL_FAILURE(compact(4, 200, 1, &second_log, &second_version));
+    EXPECT_EQ(second_log.min_timestamp(), 150);
+    const auto& inputs = second_log.compaction().recycle_rowsets();
+    EXPECT_EQ(inputs[0].rowset_meta().start_version(), 2);
+    EXPECT_EQ(inputs[0].rowset_meta().end_version(), 3);
+    EXPECT_EQ(inputs[0].rowset_meta().rowset_id_v2(), output_rowset_ids[0]);
+    EXPECT_EQ(inputs[1].rowset_meta().rowset_id_v2(), l4.rowset_id_v2());
+
+    // Exercise the real Recycler snapshot-reference decision with the actual 
MS logs.
+    OperationLogRecycleChecker checker(test_instance_id, txn_kv.get(), 
instance);
+    ASSERT_EQ(checker.init(), 0);
+    OperationLogReferenceInfo first_reference, second_reference;
+    EXPECT_FALSE(checker.can_recycle(first_version, first_log.min_timestamp(), 
&first_reference));
+    EXPECT_TRUE(first_reference.referenced_by_snapshot);
+    EXPECT_EQ(first_reference.referenced_snapshot_timestamp, 
Versionstamp(100));
+    EXPECT_TRUE(checker.can_recycle(second_version, 
second_log.min_timestamp(), &second_reference));
+    EXPECT_FALSE(second_reference.referenced_by_snapshot);
+
+    MetaReader snapshot_reader(test_instance_id, txn_kv.get(), 
Versionstamp(100));
+    std::vector<doris::RowsetMetaCloudPB> snapshot_rowsets;
+    ASSERT_EQ(snapshot_reader.get_rowset_metas(tablet_id, 2, 3, 
&snapshot_rowsets),
+              TxnErrorCode::TXN_OK);
+    ASSERT_EQ(snapshot_rowsets.size(), 2);
+    EXPECT_EQ(snapshot_rowsets[0].rowset_id_v2(), l2.rowset_id_v2());
+    EXPECT_EQ(snapshot_rowsets[1].rowset_id_v2(), l3.rowset_id_v2());
+    EXPECT_EQ(snapshot_reader.min_read_versionstamp(), Versionstamp(80));
+}
+
 TEST(RecycleOperationLogTest, RecycleCompactionLog) {
     // Ensure strip behavior is enabled for this test
     auto old_flag = config::enable_recycle_rowset_strip_key_bounds;


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to