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

liaoxin01 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/master by this push:
     new 79fcbea9b41 [fix](recycler) Revert shuffle recycle rowsets for 
recycler (#67946)
79fcbea9b41 is described below

commit 79fcbea9b41396380550257b5f700383b5036b9d
Author: Yixuan Wang <[email protected]>
AuthorDate: Tue Sep 15 11:06:30 2026 +0800

    [fix](recycler) Revert shuffle recycle rowsets for recycler (#67946)
    
    ### What problem does this PR solve?
    
    Issue Number: close #xxx
    
    Related PR: https://github.com/apache/doris/pull/63295
    
    Problem Summary:
    
    Only need to update the master version; version 4.1 has been changed.
---
 cloud/src/common/config.h       |   4 -
 cloud/src/recycler/recycler.cpp | 179 ++++-------------
 cloud/test/recycler_test.cpp    | 417 +---------------------------------------
 3 files changed, 38 insertions(+), 562 deletions(-)

diff --git a/cloud/src/common/config.h b/cloud/src/common/config.h
index faae59ad7e3..4e4e6340fd4 100644
--- a/cloud/src/common/config.h
+++ b/cloud/src/common/config.h
@@ -114,10 +114,6 @@ CONF_mInt32(instance_recycler_worker_pool_size, "32");
 // Max number of delete tasks per batch when recycling objects.
 // Each task deletes up to 1000 files. Controls memory usage during 
large-scale deletion.
 CONF_Int32(recycler_max_tasks_per_batch, "1000");
-// Max expired recycle_rowset entries to process for one tablet in one 
recycle_rowsets scan.
-// Remaining entries are left for later scans so deletion can spread across 
tablet prefixes.
-CONF_mInt32(recycle_rowsets_per_tablet_batch_size, "1000");
-CONF_mInt32(recycle_rowsets_delete_batch_size, "300000");
 // The worker pool size for http api `statistics_recycle` worker pool
 CONF_mInt32(instance_recycler_statistics_recycle_worker_pool_size, "5");
 CONF_Bool(enable_checker, "false");
diff --git a/cloud/src/recycler/recycler.cpp b/cloud/src/recycler/recycler.cpp
index 680dfd09a76..19c438c73b9 100644
--- a/cloud/src/recycler/recycler.cpp
+++ b/cloud/src/recycler/recycler.cpp
@@ -5787,61 +5787,15 @@ int InstanceRecycler::recycle_rowsets() {
                 .tag("expired_rowset_meta_size", expired_rowset_size);
     };
 
-    struct RecycleRowsetEntry {
-        std::string key;
-        doris::RowsetMetaCloudPB meta;
-    };
-    struct RecycleRowsetDeleteJob {
-        std::vector<std::string> keys;
-        std::map<std::string, doris::RowsetMetaCloudPB> rowsets;
-    };
-    // Store the scanned recycle key with rowset meta. The scanned key is the 
actual KV key to delete.
-    std::vector<RecycleRowsetEntry> rowsets;
-    int64_t current_tablet_id = -1;
-    int64_t recycled_rowset_count_for_current_tablet = 0;
-    bool current_tablet_skip_logged = false;
-    std::string next_scan_begin;
-    const int64_t rowset_batch_size_per_tablet =
-            std::max(1, config::recycle_rowsets_per_tablet_batch_size);
-    const int64_t delete_rowset_batch_size =
-            std::min(500000, config::recycle_rowsets_delete_batch_size);
-    auto try_reserve_tablet_recycle_slot = [&](int64_t tablet_id) -> bool {
-        if (current_tablet_id != tablet_id) {
-            current_tablet_id = tablet_id;
-            recycled_rowset_count_for_current_tablet = 0;
-            current_tablet_skip_logged = false;
-        }
-        if (recycled_rowset_count_for_current_tablet >= 
rowset_batch_size_per_tablet) {
-            if (!current_tablet_skip_logged) {
-                LOG_INFO(
-                        "skip recycle rowsets for tablet because per-tablet 
batch limit is reached")
-                        .tag("instance_id", instance_id_)
-                        .tag("tablet_id", tablet_id)
-                        .tag("limit", rowset_batch_size_per_tablet);
-                current_tablet_skip_logged = true;
-            }
-            const int64_t next_tablet_id = tablet_id == INT64_MAX ? INT64_MAX 
: tablet_id + 1;
-            recycle_rowset_key({instance_id_, next_tablet_id, ""}, 
&next_scan_begin);
-            return false;
-        }
-        ++recycled_rowset_count_for_current_tablet;
-        return true;
-    };
-    auto next_scan_begin_getter = [&](std::string* begin) -> bool {
-        if (next_scan_begin.empty()) {
-            return false;
-        }
-        *begin = std::move(next_scan_begin);
-        next_scan_begin.clear();
-        return true;
-    };
-
+    std::vector<std::string> rowset_keys;
     std::vector<std::string> rowset_keys_to_mark_recycled;
     std::vector<std::string> rowset_keys_to_abort_job;
+    // rowset_id -> rowset_meta
+    // store rowset id and meta for statistics rs size when delete
+    std::map<std::string, doris::RowsetMetaCloudPB> rowsets;
 
     std::mutex async_recycled_rowset_keys_mutex;
     std::vector<std::string> async_recycled_rowset_keys;
-    std::vector<std::string> rowset_keys_without_data;
     auto worker_pool = std::make_unique<SimpleThreadPool>(
             config::instance_recycler_worker_pool_size, "recycle_rowsets");
     worker_pool->start();
@@ -5887,7 +5841,7 @@ int InstanceRecycler::recycle_rowsets() {
         if (delete_versioned_delete_bitmap_kvs(partition_id, tablet_id, 
rowset_id) != 0) {
             return -1;
         }
-        rowset_keys_without_data.push_back(std::move(key));
+        rowset_keys.push_back(std::move(key));
         return 0;
     };
 
@@ -5915,12 +5869,6 @@ int InstanceRecycler::recycle_rowsets() {
         ++num_expired;
         expired_rowset_size += v.size();
 
-        int64_t tablet_id =
-                rowset.has_type() ? rowset.rowset_meta().tablet_id() : 
rowset.tablet_id();
-        if (!try_reserve_tablet_recycle_slot(tablet_id)) {
-            return 0;
-        }
-
         if (!rowset.has_type()) {                         // old version 
`RecycleRowsetPB`
             if (!rowset.has_resource_id()) [[unlikely]] { // impossible
                 // in old version, keep this key-value pair and it needs to be 
checked manually
@@ -5931,8 +5879,8 @@ int InstanceRecycler::recycle_rowsets() {
                 // old version `RecycleRowsetPB` may has empty resource_id, 
just remove the kv.
                 LOG(INFO) << "delete the recycle rowset kv that has empty 
resource_id, key="
                           << hex(k) << " value=" << proto_to_json(rowset);
-                rowset_keys_without_data.emplace_back(k);
-                return 0;
+                rowset_keys.emplace_back(k);
+                return -1;
             }
             // decode rowset_id
             auto k1 = k;
@@ -6010,23 +5958,38 @@ int InstanceRecycler::recycle_rowsets() {
             }
         } else {
             num_compacted += rowset.type() == RecycleRowsetPB::COMPACT;
-            if (rowset_meta->num_segments() > 0) { // Skip empty rowset
-                rowsets.emplace_back(std::string(k), std::move(*rowset_meta));
-            } else {
+            rowset_keys.emplace_back(k);
+            rowsets.emplace(rowset_meta->rowset_id_v2(), 
std::move(*rowset_meta));
+            if (rowset_meta->num_segments() <= 0) { // Skip empty rowset
                 ++num_empty_rowset;
-                rowset_keys_without_data.emplace_back(k);
             }
         }
         return 0;
     };
 
-    auto submit_delete_rowset_data_job = [&](std::vector<std::string> 
rowset_keys,
-                                             std::map<std::string, 
RowsetMetaCloudPB> rowsets) {
-        worker_pool->submit([&, rowset_keys_to_delete = std::move(rowset_keys),
-                             rowsets_to_delete = std::move(rowsets)]() {
+    auto loop_done = [&]() -> int {
+        std::vector<std::string> rowset_keys_to_delete;
+        // rowset_id -> rowset_meta
+        // store rowset id and meta for statistics rs size when delete
+        std::map<std::string, doris::RowsetMetaCloudPB> rowsets_to_delete;
+        std::vector<std::string> mark_keys_to_process;
+        std::vector<std::string> abort_job_keys_to_process;
+        rowset_keys_to_delete.swap(rowset_keys);
+        rowsets_to_delete.swap(rowsets);
+        mark_keys_to_process.swap(rowset_keys_to_mark_recycled);
+        abort_job_keys_to_process.swap(rowset_keys_to_abort_job);
+        if (!mark_keys_to_process.empty()) {
+            submit_batch_mark_rowsets_as_recycled_job<RecycleRowsetPB>(
+                    *worker_pool, std::move(mark_keys_to_process));
+        }
+        if (!abort_job_keys_to_process.empty()) {
+            submit_recycle_prepare_rowsets_job(*worker_pool, 
std::move(abort_job_keys_to_process),
+                                               &num_recycled);
+        }
+        worker_pool->submit([&, rowset_keys_to_delete = 
std::move(rowset_keys_to_delete),
+                             rowsets_to_delete = 
std::move(rowsets_to_delete)]() mutable {
             std::vector<std::vector<std::string>> 
versioned_delete_bitmap_key_groups;
-            if (!rowsets_to_delete.empty() &&
-                delete_rowset_data(rowsets_to_delete, 
RowsetRecyclingState::FORMAL_ROWSET,
+            if (delete_rowset_data(rowsets_to_delete, 
RowsetRecyclingState::FORMAL_ROWSET,
                                    metrics_context, 
&versioned_delete_bitmap_key_groups) != 0) {
                 LOG(WARNING) << "failed to delete rowset data, instance_id=" 
<< instance_id_;
                 return;
@@ -6042,74 +6005,8 @@ int InstanceRecycler::recycle_rowsets() {
                 LOG(WARNING) << "failed to delete recycle rowset kv, 
instance_id=" << instance_id_;
                 return;
             }
-
             num_recycled.fetch_add(rowset_keys_to_delete.size(), 
std::memory_order_relaxed);
         });
-    };
-
-    bool scan_finished = false;
-    auto loop_done = [&]() -> int {
-        std::vector<std::string> mark_keys_to_process;
-        std::vector<std::string> abort_job_keys_to_process;
-        mark_keys_to_process.swap(rowset_keys_to_mark_recycled);
-        abort_job_keys_to_process.swap(rowset_keys_to_abort_job);
-        if (!mark_keys_to_process.empty()) {
-            submit_batch_mark_rowsets_as_recycled_job<RecycleRowsetPB>(
-                    *worker_pool, std::move(mark_keys_to_process));
-        }
-        if (!abort_job_keys_to_process.empty()) {
-            submit_recycle_prepare_rowsets_job(*worker_pool, 
std::move(abort_job_keys_to_process),
-                                               &num_recycled);
-        }
-        if (!scan_finished && rowsets.size() < delete_rowset_batch_size) {
-            return 0;
-        }
-
-        DORIS_CLOUD_DEFER {
-            // if return -1 in loop done, rowset info in memory is not cleared,
-            // it can lead to memory accumulation
-            rowset_keys_without_data.clear();
-            rowsets.clear();
-        };
-        std::random_device rd;
-        std::mt19937 g(rd());
-        std::ranges::shuffle(rowsets, g);
-
-        std::vector<std::string> rowset_keys_to_delete;
-        rowset_keys_to_delete.reserve(rowset_batch_size_per_tablet);
-        // rowset_id -> rowset_meta
-        // store rowset id and meta for statistics rs size when delete
-        std::map<std::string, doris::RowsetMetaCloudPB> rowsets_to_delete;
-        std::vector<std::vector<std::string>> 
versioned_delete_bitmap_key_groups;
-
-        size_t rowsets_per_batch_size = 0;
-        for (auto& rowset : rowsets) {
-            rowset_keys_to_delete.emplace_back(std::move(rowset.key));
-            rowsets_to_delete.emplace(rowset.meta.rowset_id_v2(), 
std::move(rowset.meta));
-            if (++rowsets_per_batch_size < rowset_batch_size_per_tablet) {
-                continue;
-            }
-
-            submit_delete_rowset_data_job(std::move(rowset_keys_to_delete),
-                                          std::move(rowsets_to_delete));
-            rowsets_per_batch_size = 0;
-            rowset_keys_to_delete.clear();
-            rowsets_to_delete.clear();
-        }
-
-        if (!rowset_keys_to_delete.empty() || !rowsets_to_delete.empty()) {
-            submit_delete_rowset_data_job(std::move(rowset_keys_to_delete),
-                                          std::move(rowsets_to_delete));
-        }
-
-        for (size_t i = 0; i < rowset_keys_without_data.size(); i += 
rowset_batch_size_per_tablet) {
-            auto begin = rowset_keys_without_data.begin() + i;
-            auto end = rowset_keys_without_data.begin() +
-                       std::min(i + rowset_batch_size_per_tablet, 
rowset_keys_without_data.size());
-            std::vector<std::string> 
rowset_keys_to_remove(std::make_move_iterator(begin),
-                                                           
std::make_move_iterator(end));
-            submit_delete_rowset_data_job(std::move(rowset_keys_to_remove), 
{});
-        }
         return 0;
     };
 
@@ -6117,16 +6014,8 @@ int InstanceRecycler::recycle_rowsets() {
         scan_and_statistics_rowsets();
     }
     // recycle_func and loop_done for scan and recycle
-    int ret = scan_and_recycle(recyc_rs_key0, recyc_rs_key1, 
std::move(handle_rowset_kv), loop_done,
-                               std::move(next_scan_begin_getter));
-    scan_finished = true;
-    // if the size of rowsets is always less than delete_rowset_batch_size
-    // it need to submit the task directly
-    // else if the size of rowsets is greater than delete_rowset_batch_size,
-    // but there are residual, whether due to failed or unsuccessful cleanup, 
this behavior is idempotent
-    if (loop_done() != 0) {
-        ret = -1;
-    }
+    int ret = scan_and_recycle(recyc_rs_key0, recyc_rs_key1, 
std::move(handle_rowset_kv),
+                               std::move(loop_done));
 
     worker_pool->stop();
 
diff --git a/cloud/test/recycler_test.cpp b/cloud/test/recycler_test.cpp
index a93359785d3..3946f9bc07c 100644
--- a/cloud/test/recycler_test.cpp
+++ b/cloud/test/recycler_test.cpp
@@ -326,16 +326,6 @@ static int create_recycle_rowset(TxnKv* txn_kv, 
StorageVaultAccessor* accessor,
     return 0;
 }
 
-static size_t count_recycle_rowsets(TxnKv* txn_kv) {
-    std::unique_ptr<Transaction> txn;
-    EXPECT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
-    std::unique_ptr<RangeGetIterator> it;
-    auto begin_key = recycle_key_prefix(instance_id);
-    auto end_key = recycle_key_prefix(instance_id + '\xff');
-    EXPECT_EQ(txn->get(begin_key, end_key, &it), TxnErrorCode::TXN_OK);
-    return it->size();
-}
-
 static size_t count_recycle_rowsets(TxnKv* txn_kv, int64_t tablet_id) {
     std::unique_ptr<Transaction> txn;
     EXPECT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
@@ -1518,16 +1508,12 @@ TEST(RecyclerTest, recycle_rowsets) {
     check_delete_bitmap_keys_size(txn_kv.get(), tablet_id, 1000);
     check_delete_bitmap_file_size(accessor, tablet_id, 1000);
 
-    for (size_t i = 0; i < 10; i++) {
-        ASSERT_EQ(recycler.recycle_rowsets(), 0);
-    }
+    ASSERT_EQ(recycler.recycle_rowsets(), 0);
+    ASSERT_EQ(recycler.recycle_rowsets(), 0);
 
     // check rowset does not exist on obj store
     std::unique_ptr<ListIterator> list_iter;
     ASSERT_EQ(0, accessor->list_directory(tablet_path_prefix(tablet_id), 
&list_iter));
-    for (auto file = list_iter->next(); file.has_value(); file = 
list_iter->next()) {
-        LOG(INFO) << "file: " << file->path;
-    }
     EXPECT_FALSE(list_iter->has_next());
     // check all recycle rowset kv have been deleted
     std::unique_ptr<Transaction> txn;
@@ -1545,18 +1531,6 @@ TEST(RecyclerTest, recycle_rowsets) {
     check_delete_bitmap_file_size(accessor, tablet_id, 0);
 }
 
-TEST(RecyclerTest, next_recycle_rowset_tablet_key_overwrites_existing_buffer) {
-    std::string next_key = recycle_rowset_key({instance_id, 10002, "rowset"});
-    ASSERT_EQ(InstanceRecycler::next_recycle_rowset_tablet_key(instance_id, 
10002, &next_key), 0);
-
-    std::string_view k1 = next_key;
-    k1.remove_prefix(1);
-    std::vector<std::tuple<std::variant<int64_t, std::string>, int, int>> out;
-    ASSERT_EQ(decode_key(&k1, &out), 0);
-    EXPECT_EQ(std::get<int64_t>(std::get<0>(out[3])), 10003);
-    EXPECT_TRUE(std::get<std::string>(std::get<0>(out[4])).empty());
-}
-
 TEST(RecyclerTest, recycle_rowsets_only_marks_prepare_rowsets_as_recycled) {
     config::retention_seconds = 0;
     auto old_enable_mark = config::enable_mark_delete_rowset_before_recycle;
@@ -2700,77 +2674,6 @@ TEST(RecyclerTest, 
recycle_tmp_rowsets_cross_abort_recheck_batch_boundary) {
     EXPECT_FALSE(list_iter->has_next());
 }
 
-TEST(RecyclerTest, 
recycle_rowsets_tablet_batch_limit_recycles_remaining_in_next_round) {
-    config::retention_seconds = 0;
-    auto txn_kv = std::make_shared<MemTxnKv>();
-    ASSERT_EQ(txn_kv->init(), 0);
-
-    InstanceInfoPB instance;
-    instance.set_instance_id(instance_id);
-    auto obj_info = instance.add_obj_info();
-    obj_info->set_id("recycle_rowsets_batch_limit");
-    obj_info->set_ak(config::test_s3_ak);
-    obj_info->set_sk(config::test_s3_sk);
-    obj_info->set_endpoint(config::test_s3_endpoint);
-    obj_info->set_region(config::test_s3_region);
-    obj_info->set_bucket(config::test_s3_bucket);
-    obj_info->set_prefix("recycle_rowsets_batch_limit");
-
-    auto old_worker_pool_size = config::instance_recycler_worker_pool_size;
-    auto old_max_rowsets_per_tablet = 
config::recycle_rowsets_per_tablet_batch_size;
-    auto old_enable_mark = config::enable_mark_delete_rowset_before_recycle;
-    config::instance_recycler_worker_pool_size = 1;
-    config::recycle_rowsets_per_tablet_batch_size = 3;
-    config::enable_mark_delete_rowset_before_recycle = false;
-    DORIS_CLOUD_DEFER {
-        config::instance_recycler_worker_pool_size = old_worker_pool_size;
-        config::recycle_rowsets_per_tablet_batch_size = 
old_max_rowsets_per_tablet;
-        config::enable_mark_delete_rowset_before_recycle = old_enable_mark;
-    };
-
-    InstanceRecycler recycler(txn_kv, instance, thread_group,
-                              std::make_shared<TxnLazyCommitter>(txn_kv));
-    ASSERT_EQ(recycler.init(), 0);
-    auto accessor = recycler.accessor_map_.begin()->second;
-
-    doris::TabletSchemaCloudPB schema;
-    schema.set_schema_version(1);
-    constexpr int64_t index_id = 10001;
-    constexpr int64_t first_tablet_id = 10002;
-    constexpr int64_t second_tablet_id = 10020;
-    std::vector<std::string> first_tablet_rowset_ids;
-    for (int i = 0; i < 5; ++i) {
-        auto rowset =
-                create_rowset("recycle_rowsets_batch_limit", first_tablet_id, 
index_id, 1, schema);
-        first_tablet_rowset_ids.push_back(rowset.rowset_id_v2());
-        ASSERT_EQ(create_recycle_rowset(txn_kv.get(), accessor.get(), rowset,
-                                        RecycleRowsetPB::COMPACT, true),
-                  0);
-    }
-    for (int i = 0; i < 2; ++i) {
-        auto rowset =
-                create_rowset("recycle_rowsets_batch_limit", second_tablet_id, 
index_id, 1, schema);
-        ASSERT_EQ(create_recycle_rowset(txn_kv.get(), accessor.get(), rowset,
-                                        RecycleRowsetPB::COMPACT, true),
-                  0);
-    }
-
-    ASSERT_EQ(recycler.recycle_rowsets(), 0);
-    for (int i = 0; i < 3; ++i) {
-        EXPECT_EQ(accessor->exists(segment_path(first_tablet_id, 
first_tablet_rowset_ids[i], 0)),
-                  1);
-    }
-    for (int i = 3; i < 5; ++i) {
-        EXPECT_EQ(accessor->exists(segment_path(first_tablet_id, 
first_tablet_rowset_ids[i], 0)),
-                  0);
-    }
-
-    ASSERT_EQ(recycler.recycle_rowsets(), 0);
-    for (const auto& rowset_id : first_tablet_rowset_ids) {
-        EXPECT_EQ(accessor->exists(segment_path(first_tablet_id, rowset_id, 
0)), 1);
-    }
-}
-
 TEST(RecyclerTest, recycle_rowsets_with_data_ref_count) {
     config::retention_seconds = 0;
     auto txn_kv = std::make_shared<MemTxnKv>();
@@ -2862,317 +2765,6 @@ TEST(RecyclerTest, recycle_rowsets_with_data_ref_count) 
{
     check_delete_bitmap_file_size(accessor, tablet_id, 3);
 }
 
-TEST(RecyclerTest, recycle_rowsets_limit_per_tablet_batch) {
-    config::retention_seconds = 0;
-    auto origin_worker_pool_size = config::instance_recycler_worker_pool_size;
-    auto origin_batch_size = config::recycle_rowsets_per_tablet_batch_size;
-    config::instance_recycler_worker_pool_size = 4;
-    config::recycle_rowsets_per_tablet_batch_size = 2;
-    DORIS_CLOUD_DEFER {
-        config::instance_recycler_worker_pool_size = origin_worker_pool_size;
-        config::recycle_rowsets_per_tablet_batch_size = origin_batch_size;
-    };
-
-    auto txn_kv = std::make_shared<MemTxnKv>();
-    ASSERT_EQ(txn_kv->init(), 0);
-
-    InstanceInfoPB instance;
-    instance.set_instance_id(instance_id);
-    auto obj_info = instance.add_obj_info();
-    obj_info->set_id("recycle_rowsets_limit_per_tablet_batch");
-    obj_info->set_ak(config::test_s3_ak);
-    obj_info->set_sk(config::test_s3_sk);
-    obj_info->set_endpoint(config::test_s3_endpoint);
-    obj_info->set_region(config::test_s3_region);
-    obj_info->set_bucket(config::test_s3_bucket);
-    obj_info->set_prefix("recycle_rowsets_limit_per_tablet_batch");
-
-    InstanceRecycler recycler(txn_kv, instance, thread_group,
-                              std::make_shared<TxnLazyCommitter>(txn_kv));
-    ASSERT_EQ(recycler.init(), 0);
-    auto accessor = recycler.accessor_map_.begin()->second;
-
-    doris::TabletSchemaCloudPB schema;
-    schema.set_schema_version(1);
-    schema.set_inverted_index_storage_format(InvertedIndexStorageFormatPB::V1);
-
-    constexpr int index_id = 10001;
-    constexpr int64_t tablet_id0 = 100020;
-    constexpr int64_t tablet_id1 = 100021;
-    for (int64_t tablet_id : {tablet_id0, tablet_id1}) {
-        for (int i = 0; i < 5; ++i) {
-            auto rowset = 
create_rowset("recycle_rowsets_limit_per_tablet_batch", tablet_id,
-                                        index_id, 1, schema);
-            ASSERT_EQ(0, create_recycle_rowset(txn_kv.get(), accessor.get(), 
rowset,
-                                               RecycleRowsetPB::COMPACT, 
true));
-        }
-    }
-
-    auto count_recycle_rowsets = [&](int64_t tablet_id) {
-        std::unique_ptr<Transaction> txn;
-        EXPECT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
-        std::unique_ptr<RangeGetIterator> it;
-        auto begin_key = recycle_rowset_key({instance_id, tablet_id, ""});
-        auto end_key = recycle_rowset_key({instance_id, tablet_id, "\xff"});
-        EXPECT_EQ(txn->get(begin_key, end_key, &it), TxnErrorCode::TXN_OK);
-        return it->size();
-    };
-
-    ASSERT_EQ(recycler.recycle_rowsets(), 0);
-
-    ASSERT_EQ(recycler.recycle_rowsets(), 0);
-    EXPECT_EQ(count_recycle_rowsets(tablet_id0), 1);
-    EXPECT_EQ(count_recycle_rowsets(tablet_id1), 1);
-
-    ASSERT_EQ(recycler.recycle_rowsets(), 0);
-    EXPECT_EQ(count_recycle_rowsets(tablet_id0), 0);
-    EXPECT_EQ(count_recycle_rowsets(tablet_id1), 0);
-}
-
-TEST(RecyclerTest, recycle_rowsets_delete_remaining_rowsets_by_tablet) {
-    config::retention_seconds = 0;
-    auto origin_worker_pool_size = config::instance_recycler_worker_pool_size;
-    auto origin_per_tablet_batch_size = 
config::recycle_rowsets_per_tablet_batch_size;
-    auto origin_delete_batch_size = config::recycle_rowsets_delete_batch_size;
-    config::instance_recycler_worker_pool_size = 4;
-    config::recycle_rowsets_per_tablet_batch_size = 10;
-    config::recycle_rowsets_delete_batch_size = 10;
-    DORIS_CLOUD_DEFER {
-        config::instance_recycler_worker_pool_size = origin_worker_pool_size;
-        config::recycle_rowsets_per_tablet_batch_size = 
origin_per_tablet_batch_size;
-        config::recycle_rowsets_delete_batch_size = origin_delete_batch_size;
-    };
-
-    auto txn_kv = std::make_shared<MemTxnKv>();
-    ASSERT_EQ(txn_kv->init(), 0);
-
-    InstanceInfoPB instance;
-    instance.set_instance_id(instance_id);
-    auto obj_info = instance.add_obj_info();
-    obj_info->set_id("recycle_rowsets_below_batch_threshold");
-    obj_info->set_ak(config::test_s3_ak);
-    obj_info->set_sk(config::test_s3_sk);
-    obj_info->set_endpoint(config::test_s3_endpoint);
-    obj_info->set_region(config::test_s3_region);
-    obj_info->set_bucket(config::test_s3_bucket);
-    obj_info->set_prefix("recycle_rowsets_below_batch_threshold");
-
-    InstanceRecycler recycler(txn_kv, instance, thread_group,
-                              std::make_shared<TxnLazyCommitter>(txn_kv));
-    ASSERT_EQ(recycler.init(), 0);
-    auto accessor = recycler.accessor_map_.begin()->second;
-
-    doris::TabletSchemaCloudPB schema;
-    schema.set_schema_version(1);
-    schema.set_inverted_index_storage_format(InvertedIndexStorageFormatPB::V1);
-
-    constexpr int64_t index_id = 10001;
-    constexpr int64_t tablet_id0 = 100030;
-    constexpr int64_t tablet_id1 = 100031;
-    for (int64_t tablet_id : {tablet_id0, tablet_id1}) {
-        for (int i = 0; i < 2; ++i) {
-            auto rowset = 
create_rowset("recycle_rowsets_below_batch_threshold", tablet_id,
-                                        index_id, 1, schema);
-            ASSERT_EQ(0, create_recycle_rowset(txn_kv.get(), accessor.get(), 
rowset,
-                                               RecycleRowsetPB::COMPACT, 
true));
-        }
-        auto prefix_rowset =
-                create_rowset("recycle_rowsets_below_batch_threshold", 
tablet_id, index_id, 1,
-                              schema, RowsetStatePB::BEGIN_PARTIAL_UPDATE);
-        ASSERT_EQ(0, create_recycle_rowset(txn_kv.get(), accessor.get(), 
prefix_rowset,
-                                           RecycleRowsetPB::COMPACT, true));
-    }
-
-    for (size_t i = 0; i < 10; i++) {
-        ASSERT_EQ(recycler.recycle_rowsets(), 0);
-    }
-    std::unique_ptr<ListIterator> list_iter;
-    ASSERT_EQ(0, accessor->list_directory(tablet_path_prefix(tablet_id0), 
&list_iter));
-    EXPECT_FALSE(list_iter->has_next());
-    ASSERT_EQ(0, accessor->list_directory(tablet_path_prefix(tablet_id1), 
&list_iter));
-    EXPECT_FALSE(list_iter->has_next());
-    EXPECT_EQ(count_recycle_rowsets(txn_kv.get(), tablet_id0), 0);
-    EXPECT_EQ(count_recycle_rowsets(txn_kv.get(), tablet_id1), 0);
-    EXPECT_EQ(count_recycle_rowsets(txn_kv.get()), 0);
-}
-
-TEST(RecyclerTest, recycle_rowsets_delete_full_batches_and_leftover_kvs) {
-    config::retention_seconds = 0;
-    auto origin_worker_pool_size = config::instance_recycler_worker_pool_size;
-    auto origin_per_tablet_batch_size = 
config::recycle_rowsets_per_tablet_batch_size;
-    auto origin_delete_batch_size = config::recycle_rowsets_delete_batch_size;
-    config::instance_recycler_worker_pool_size = 1;
-    config::recycle_rowsets_per_tablet_batch_size = 3;
-    config::recycle_rowsets_delete_batch_size = 7;
-    DORIS_CLOUD_DEFER {
-        config::instance_recycler_worker_pool_size = origin_worker_pool_size;
-        config::recycle_rowsets_per_tablet_batch_size = 
origin_per_tablet_batch_size;
-        config::recycle_rowsets_delete_batch_size = origin_delete_batch_size;
-    };
-
-    auto txn_kv = std::make_shared<MemTxnKv>();
-    ASSERT_EQ(txn_kv->init(), 0);
-
-    InstanceInfoPB instance;
-    instance.set_instance_id(instance_id);
-    auto obj_info = instance.add_obj_info();
-    obj_info->set_id("recycle_rowsets_batched_delete");
-    obj_info->set_ak(config::test_s3_ak);
-    obj_info->set_sk(config::test_s3_sk);
-    obj_info->set_endpoint(config::test_s3_endpoint);
-    obj_info->set_region(config::test_s3_region);
-    obj_info->set_bucket(config::test_s3_bucket);
-    obj_info->set_prefix("recycle_rowsets_batched_delete");
-
-    InstanceRecycler recycler(txn_kv, instance, thread_group,
-                              std::make_shared<TxnLazyCommitter>(txn_kv));
-    ASSERT_EQ(recycler.init(), 0);
-    auto accessor = recycler.accessor_map_.begin()->second;
-
-    doris::TabletSchemaCloudPB schema;
-    schema.set_schema_version(1);
-    schema.set_inverted_index_storage_format(InvertedIndexStorageFormatPB::V1);
-
-    constexpr int64_t index_id = 10001;
-    for (int64_t tablet_id : {100040, 100041, 100042}) {
-        for (int i = 0; i < 2; ++i) {
-            auto rowset =
-                    create_rowset("recycle_rowsets_batched_delete", tablet_id, 
index_id, 1, schema);
-            ASSERT_EQ(0, create_recycle_rowset(txn_kv.get(), accessor.get(), 
rowset,
-                                               RecycleRowsetPB::COMPACT, 
true));
-        }
-        auto empty_rowset =
-                create_rowset("recycle_rowsets_batched_delete", tablet_id, 
index_id, 0, schema);
-        ASSERT_EQ(0, create_recycle_rowset(txn_kv.get(), accessor.get(), 
empty_rowset,
-                                           RecycleRowsetPB::COMPACT, true));
-    }
-
-    for (size_t i = 0; i < 3; i++) {
-        ASSERT_EQ(recycler.recycle_rowsets(), 0);
-    }
-    std::unique_ptr<ListIterator> list_iter;
-    for (int64_t tablet_id : {100040, 100041, 100042}) {
-        ASSERT_EQ(0, accessor->list_directory(tablet_path_prefix(tablet_id), 
&list_iter));
-        EXPECT_FALSE(list_iter->has_next());
-    }
-    EXPECT_EQ(count_recycle_rowsets(txn_kv.get()), 0);
-}
-
-TEST(RecyclerTest, 
recycle_rowsets_delete_prefix_rowset_kvs_without_remaining_rowsets) {
-    config::retention_seconds = 0;
-    auto origin_worker_pool_size = config::instance_recycler_worker_pool_size;
-    auto origin_per_tablet_batch_size = 
config::recycle_rowsets_per_tablet_batch_size;
-    auto origin_delete_batch_size = config::recycle_rowsets_delete_batch_size;
-    config::instance_recycler_worker_pool_size = 1;
-    config::recycle_rowsets_per_tablet_batch_size = 10;
-    config::recycle_rowsets_delete_batch_size = 10;
-    DORIS_CLOUD_DEFER {
-        config::instance_recycler_worker_pool_size = origin_worker_pool_size;
-        config::recycle_rowsets_per_tablet_batch_size = 
origin_per_tablet_batch_size;
-        config::recycle_rowsets_delete_batch_size = origin_delete_batch_size;
-    };
-
-    auto txn_kv = std::make_shared<MemTxnKv>();
-    ASSERT_EQ(txn_kv->init(), 0);
-
-    InstanceInfoPB instance;
-    instance.set_instance_id(instance_id);
-    auto obj_info = instance.add_obj_info();
-    obj_info->set_id("recycle_rowsets_prefix_only");
-    obj_info->set_ak(config::test_s3_ak);
-    obj_info->set_sk(config::test_s3_sk);
-    obj_info->set_endpoint(config::test_s3_endpoint);
-    obj_info->set_region(config::test_s3_region);
-    obj_info->set_bucket(config::test_s3_bucket);
-    obj_info->set_prefix("recycle_rowsets_prefix_only");
-
-    InstanceRecycler recycler(txn_kv, instance, thread_group,
-                              std::make_shared<TxnLazyCommitter>(txn_kv));
-    ASSERT_EQ(recycler.init(), 0);
-    auto accessor = recycler.accessor_map_.begin()->second;
-
-    doris::TabletSchemaCloudPB schema;
-    schema.set_schema_version(1);
-    schema.set_inverted_index_storage_format(InvertedIndexStorageFormatPB::V1);
-
-    constexpr int64_t index_id = 10001;
-    constexpr int64_t tablet_id = 100050;
-    for (int i = 0; i < 3; ++i) {
-        auto rowset = create_rowset("recycle_rowsets_prefix_only", tablet_id, 
index_id, 1, schema,
-                                    RowsetStatePB::BEGIN_PARTIAL_UPDATE);
-        ASSERT_EQ(0, create_recycle_rowset(txn_kv.get(), accessor.get(), 
rowset,
-                                           RecycleRowsetPB::COMPACT, true));
-    }
-
-    for (size_t i = 0; i < 3; i++) {
-        ASSERT_EQ(recycler.recycle_rowsets(), 0);
-    }
-    std::unique_ptr<ListIterator> list_iter;
-    ASSERT_EQ(0, accessor->list_directory(tablet_path_prefix(tablet_id), 
&list_iter));
-    EXPECT_FALSE(list_iter->has_next());
-    EXPECT_EQ(count_recycle_rowsets(txn_kv.get()), 0);
-}
-
-TEST(RecyclerTest, 
recycle_rowsets_delete_old_empty_resource_id_kvs_with_normal_rowsets) {
-    config::retention_seconds = 0;
-    auto origin_worker_pool_size = config::instance_recycler_worker_pool_size;
-    auto origin_per_tablet_batch_size = 
config::recycle_rowsets_per_tablet_batch_size;
-    auto origin_delete_batch_size = config::recycle_rowsets_delete_batch_size;
-    config::instance_recycler_worker_pool_size = 4;
-    config::recycle_rowsets_per_tablet_batch_size = 10;
-    config::recycle_rowsets_delete_batch_size = 10;
-    DORIS_CLOUD_DEFER {
-        config::instance_recycler_worker_pool_size = origin_worker_pool_size;
-        config::recycle_rowsets_per_tablet_batch_size = 
origin_per_tablet_batch_size;
-        config::recycle_rowsets_delete_batch_size = origin_delete_batch_size;
-    };
-
-    auto txn_kv = std::make_shared<MemTxnKv>();
-    ASSERT_EQ(txn_kv->init(), 0);
-
-    InstanceInfoPB instance;
-    instance.set_instance_id(instance_id);
-    auto obj_info = instance.add_obj_info();
-    obj_info->set_id("recycle_rowsets_old_empty_resource_id");
-    obj_info->set_ak(config::test_s3_ak);
-    obj_info->set_sk(config::test_s3_sk);
-    obj_info->set_endpoint(config::test_s3_endpoint);
-    obj_info->set_region(config::test_s3_region);
-    obj_info->set_bucket(config::test_s3_bucket);
-    obj_info->set_prefix("recycle_rowsets_old_empty_resource_id");
-
-    InstanceRecycler recycler(txn_kv, instance, thread_group,
-                              std::make_shared<TxnLazyCommitter>(txn_kv));
-    ASSERT_EQ(recycler.init(), 0);
-    auto accessor = recycler.accessor_map_.begin()->second;
-
-    doris::TabletSchemaCloudPB schema;
-    schema.set_schema_version(1);
-    schema.set_inverted_index_storage_format(InvertedIndexStorageFormatPB::V1);
-
-    constexpr int64_t index_id = 10001;
-    constexpr int64_t tablet_id = 100060;
-    for (int i = 0; i < 3; ++i) {
-        auto rowset = create_rowset("recycle_rowsets_old_empty_resource_id", 
tablet_id, index_id, 1,
-                                    schema);
-        ASSERT_EQ(0, create_recycle_rowset(txn_kv.get(), accessor.get(), 
rowset,
-                                           RecycleRowsetPB::COMPACT, true));
-    }
-    for (int i = 0; i < 2; ++i) {
-        auto old_rowset = create_rowset("", tablet_id, index_id, 0, schema);
-        ASSERT_EQ(0, create_recycle_rowset(txn_kv.get(), accessor.get(), 
old_rowset,
-                                           RecycleRowsetPB::UNKNOWN, false));
-    }
-
-    for (size_t i = 0; i < 3; i++) {
-        ASSERT_EQ(recycler.recycle_rowsets(), 0);
-    }
-    std::unique_ptr<ListIterator> list_iter;
-    ASSERT_EQ(0, accessor->list_directory(tablet_path_prefix(tablet_id), 
&list_iter));
-    EXPECT_FALSE(list_iter->has_next());
-    EXPECT_EQ(count_recycle_rowsets(txn_kv.get()), 0);
-}
-
 TEST(RecyclerTest, bench_recycle_rowsets) {
     config::retention_seconds = 0;
     auto txn_kv = std::make_shared<MemTxnKv>();
@@ -3236,9 +2828,8 @@ TEST(RecyclerTest, bench_recycle_rowsets) {
     check_delete_bitmap_keys_size(txn_kv.get(), tablet_id, 1000);
     check_delete_bitmap_file_size(accessor, tablet_id, 1000);
 
-    for (size_t i = 0; i < 10; i++) {
-        ASSERT_EQ(recycler.recycle_rowsets(), 0);
-    }
+    ASSERT_EQ(recycler.recycle_rowsets(), 0);
+    ASSERT_EQ(recycler.recycle_rowsets(), 0);
     ASSERT_EQ(recycler.check_recycle_tasks(), false);
 
     // check rowset does not exist on obj store


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

Reply via email to