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]