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 fb5560e4b05 branch-4.1: [opt](cloud) Add recycler rowset recycling 
benchmarks test #68043 (#68581)
fb5560e4b05 is described below

commit fb5560e4b058b0b1adb7aeeb14a7757ab130beb6
Author: github-actions[bot] 
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Wed Sep 30 22:03:16 2026 +0800

    branch-4.1: [opt](cloud) Add recycler rowset recycling benchmarks test 
#68043 (#68581)
    
    Cherry-picked from #68043
    
    Co-authored-by: Yixuan Wang <[email protected]>
---
 cloud/src/recycler/recycler.cpp        |    2 +-
 cloud/test/CMakeLists.txt              |    1 +
 cloud/test/recycler_benchmark_test.cpp | 1069 ++++++++++++++++++++++++++++++++
 3 files changed, 1071 insertions(+), 1 deletion(-)

diff --git a/cloud/src/recycler/recycler.cpp b/cloud/src/recycler/recycler.cpp
index 72a444f3da4..76a1f048b59 100644
--- a/cloud/src/recycler/recycler.cpp
+++ b/cloud/src/recycler/recycler.cpp
@@ -5651,7 +5651,7 @@ int InstanceRecycler::recycle_rowsets() {
                 LOG(INFO) << "delete the recycle rowset kv that has empty 
resource_id, key="
                           << hex(k) << " value=" << proto_to_json(rowset);
                 rowset_keys.emplace_back(k);
-                return -1;
+                return 0;
             }
             // decode rowset_id
             auto k1 = k;
diff --git a/cloud/test/CMakeLists.txt b/cloud/test/CMakeLists.txt
index 4d768601044..057b0132356 100644
--- a/cloud/test/CMakeLists.txt
+++ b/cloud/test/CMakeLists.txt
@@ -41,6 +41,7 @@ set_target_properties(txn_kv_test PROPERTIES COMPILE_FLAGS 
"-fno-access-control"
 
 add_executable(recycler_test
     recycler_test.cpp
+    recycler_benchmark_test.cpp
     recycler_operation_log_test.cpp
     snapshot_data_size_calculator_test.cpp
     recycle_versioned_keys_test.cpp
diff --git a/cloud/test/recycler_benchmark_test.cpp 
b/cloud/test/recycler_benchmark_test.cpp
new file mode 100644
index 00000000000..9ca31a04f5d
--- /dev/null
+++ b/cloud/test/recycler_benchmark_test.cpp
@@ -0,0 +1,1069 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+#include <butil/guid.h>
+#include <fmt/core.h>
+#include <gen_cpp/olap_file.pb.h>
+#include <gtest/gtest-spi.h>
+#include <gtest/gtest.h>
+
+#include <algorithm>
+#include <atomic>
+#include <chrono>
+#include <cmath>
+#include <cstddef>
+#include <cstdint>
+#include <cstdlib>
+#include <exception>
+#include <iostream>
+#include <limits>
+#include <map>
+#include <memory>
+#include <string>
+#include <utility>
+
+#include "common/bvars.h"
+#include "common/config.h"
+#include "common/defer.h"
+#include "common/simple_thread_pool.h"
+#include "cpp/sync_point.h"
+#include "meta-service/meta_service_schema.h"
+#include "meta-service/txn_lazy_committer.h"
+#include "meta-store/blob_message.h"
+#include "meta-store/keys.h"
+#include "meta-store/mem_txn_kv.h"
+#include "meta-store/txn_kv.h"
+#include "meta-store/txn_kv_error.h"
+#include "recycler/recycler.h"
+#include "recycler/s3_accessor.h"
+#include "recycler/util.h"
+
+namespace doris::cloud {
+namespace {
+
+const std::string kBenchmarkInstanceId = "recycler_benchmark_instance";
+const std::string kBenchmarkResourceId = "recycler_benchmark_resource";
+double recycler_benchmark_duration_tolerance_ratio = 0.8;
+
+std::string get_env(const char* name) {
+    const char* value = std::getenv(name);
+    return value == nullptr ? "" : value;
+}
+
+std::string get_env_with_fallback(const char* primary, const char* fallback) {
+    const char* value = std::getenv(primary);
+    return value == nullptr ? get_env(fallback) : std::string(value);
+}
+
+struct BenchmarkS3Config {
+    bool enabled = false;
+    std::string access_key;
+    std::string secret_key;
+    std::string role_arn;
+    std::string external_id;
+    std::string endpoint;
+    std::string provider;
+    std::string bucket;
+    std::string region;
+    std::string prefix;
+};
+
+BenchmarkS3Config load_benchmark_s3_config() {
+    BenchmarkS3Config config;
+    config.enabled = get_env("ENABLE_S3_CLIENT") == "1";
+    if (!config.enabled) {
+        return config;
+    }
+
+    config.access_key = get_env("S3_AK");
+    config.secret_key = get_env("S3_SK");
+    config.role_arn = get_env("AWS_ROLE_ARN");
+    config.external_id = get_env("AWS_EXTERNAL_ID");
+    config.endpoint = get_env_with_fallback("S3_ENDPOINT", "AWS_ENDPOINT");
+    config.provider = get_env("S3_PROVIDER");
+    config.bucket = get_env_with_fallback("S3_BUCKET", "AWS_BUCKET");
+    config.region = get_env_with_fallback("S3_REGION", "AWS_REGION");
+    config.prefix = get_env_with_fallback("S3_PREFIX", "AWS_PREFIX");
+    return config;
+}
+
+void set_obj_store_provider(const std::string& provider, ObjectStoreInfoPB* 
obj_info) {
+    if (provider == "AZURE") {
+        obj_info->set_provider(ObjectStoreInfoPB_Provider_AZURE);
+    } else if (provider == "GCS") {
+        obj_info->set_provider(ObjectStoreInfoPB_Provider_GCP);
+    } else {
+        obj_info->set_provider(ObjectStoreInfoPB_Provider_S3);
+    }
+}
+
+// Number of recyclable rowsets seeded per branch.
+constexpr int64_t kRowsetsPerBranch = 10000;
+// Commit the seeded recycle rowset KVs in batches to keep each txn small.
+constexpr int64_t kSeedCommitBatch = 2000;
+constexpr int64_t kRowsetsPerPackedFile = 2;
+constexpr int64_t kPackedFileCount = kRowsetsPerBranch / kRowsetsPerPackedFile;
+constexpr int kMaxPackedRecyclePasses = 10;
+static_assert(kRowsetsPerBranch % kRowsetsPerPackedFile == 0);
+static_assert(kSeedCommitBatch % kRowsetsPerPackedFile == 0);
+constexpr int64_t kBenchmarkIndexId = 20000;
+constexpr int32_t kBenchmarkSchemaVersion = 1;
+constexpr int64_t kBenchmarkDbId = 10000;
+
+struct RecycleRowsetConfigGuard {
+    RecycleRowsetConfigGuard()
+            : 
old_enable_mark(config::enable_mark_delete_rowset_before_recycle),
+              
old_enable_abort(config::enable_abort_txn_and_job_for_delete_rowset_before_recycle)
 {
+        config::enable_mark_delete_rowset_before_recycle = true;
+        config::enable_abort_txn_and_job_for_delete_rowset_before_recycle = 
true;
+    }
+
+    ~RecycleRowsetConfigGuard() {
+        config::enable_mark_delete_rowset_before_recycle = old_enable_mark;
+        config::enable_abort_txn_and_job_for_delete_rowset_before_recycle = 
old_enable_abort;
+    }
+
+    bool old_enable_mark;
+    bool old_enable_abort;
+};
+
+doris::TabletSchemaCloudPB make_benchmark_schema() {
+    doris::TabletSchemaCloudPB schema;
+    schema.set_schema_version(kBenchmarkSchemaVersion);
+    schema.set_inverted_index_storage_format(InvertedIndexStorageFormatPB::V1);
+    auto* index = schema.add_index();
+    index->set_index_id(1);
+    index->set_index_type(IndexType::INVERTED);
+    return schema;
+}
+
+int put_benchmark_schema(TxnKv* txn_kv, const std::string& instance_id) {
+    std::unique_ptr<Transaction> txn;
+    if (txn_kv->create_txn(&txn) != TxnErrorCode::TXN_OK) {
+        return -1;
+    }
+    std::string schema_key;
+    meta_schema_key({instance_id, kBenchmarkIndexId, kBenchmarkSchemaVersion}, 
&schema_key);
+    auto schema = make_benchmark_schema();
+    MetaServiceCode code = MetaServiceCode::OK;
+    std::string msg;
+    put_schema_kv(code, msg, txn.get(), schema_key, schema);
+    if (code != MetaServiceCode::OK) {
+        return -1;
+    }
+    return txn->commit() == TxnErrorCode::TXN_OK ? 0 : -1;
+}
+
+int remove_benchmark_schema(TxnKv* txn_kv, const std::string& instance_id) {
+    std::unique_ptr<Transaction> txn;
+    if (txn_kv->create_txn(&txn) != TxnErrorCode::TXN_OK) {
+        return -1;
+    }
+    std::string schema_key;
+    meta_schema_key({instance_id, kBenchmarkIndexId, kBenchmarkSchemaVersion}, 
&schema_key);
+    ValueBuf schema_value;
+    if (blob_get(txn.get(), schema_key, &schema_value) != 
TxnErrorCode::TXN_OK) {
+        return -1;
+    }
+    schema_value.remove(txn.get());
+    return txn->commit() == TxnErrorCode::TXN_OK ? 0 : -1;
+}
+
+// Each value maps to one recyclable branch inside
+// InstanceRecycler::recycle_rowsets()::handle_rowset_kv. Both mark and abort
+// flags are enabled by RecycleRowsetConfigGuard for every scenario.
+enum class RecycleRowsetBranch {
+    // Old-version RecycleRowsetPB (no `type`) whose resource_id is empty: the
+    // recycler removes the KV only, without touching object storage.
+    kLegacyEmptyResource,
+    // Old-version RecycleRowsetPB with a real resource_id: recycled through 
the
+    // single-rowset delete_rowset_data_by_prefix path.
+    kLegacyWithResource,
+    // Already-marked PREPARE rowset without a related txn/job: recycled 
directly.
+    kPrepareDirect,
+    // New PREPARE rowset queued to be marked as recycled first
+    // (enable_mark_delete_rowset_before_recycle = true).
+    kPrepareMark,
+    // New PREPARE rowset carrying a load_id that triggers the abort-txn/job 
path
+    // (enable_abort_txn_and_job_for_delete_rowset_before_recycle = true).
+    kPrepareAbort,
+    // New COMPACT/DROP rowset with segments: recycled via the batched
+    // delete_rowset_data path.
+    kCompactedWithData,
+    // New COMPACT/DROP rowset without segments: treated as an empty rowset and
+    // removed by KV delete only.
+    kCompactedEmpty,
+};
+
+// Build one RowsetMetaCloudPB for the given branch. Kept intentionally minimal
+// and self-contained so this benchmark does not depend on recycler_test.cpp.
+doris::RowsetMetaCloudPB make_rowset_meta(RecycleRowsetBranch branch, int64_t 
tablet_id,
+                                          const std::string& rowset_id) {
+    doris::RowsetMetaCloudPB meta;
+    meta.set_rowset_id(0); // deprecated but required
+    meta.set_rowset_id_v2(rowset_id);
+    meta.set_tablet_id(tablet_id);
+    meta.set_index_id(kBenchmarkIndexId);
+    meta.set_schema_version(kBenchmarkSchemaVersion);
+    meta.mutable_tablet_schema()->CopyFrom(make_benchmark_schema());
+    meta.set_start_version(2);
+    meta.set_end_version(2);
+    meta.set_data_disk_size(1024);
+    meta.set_index_disk_size(512);
+    meta.set_total_disk_size(1536);
+    if (branch == RecycleRowsetBranch::kPrepareDirect) {
+        meta.set_is_recycled(true);
+    }
+    switch (branch) {
+    case RecycleRowsetBranch::kCompactedEmpty:
+        meta.set_num_segments(0);
+        break;
+    default:
+        meta.set_num_segments(1);
+        break;
+    }
+    // Only the branches that reach delete_rowset_data need a resolvable 
resource.
+    if (branch != RecycleRowsetBranch::kLegacyEmptyResource) {
+        meta.set_resource_id(kBenchmarkResourceId);
+    }
+    if (branch == RecycleRowsetBranch::kPrepareAbort) {
+        // A load_id makes make_related_txn_or_job_abort_task emit a TXN abort
+        // task; end_version != 1 is required to enter the abort branch.
+        meta.mutable_load_id()->set_hi(tablet_id);
+        meta.mutable_load_id()->set_lo(1);
+        meta.set_txn_id(tablet_id);
+    }
+    return meta;
+}
+
+// Wrap a RowsetMetaCloudPB into a RecycleRowsetPB shaped for the branch.
+RecycleRowsetPB make_recycle_rowset(RecycleRowsetBranch branch,
+                                    const doris::RowsetMetaCloudPB& meta) {
+    RecycleRowsetPB pb;
+    pb.set_creation_time(1); // long expired once retention is 0 / immediate 
recycle
+    pb.set_expiration(1);
+    switch (branch) {
+    case RecycleRowsetBranch::kLegacyEmptyResource:
+        // Old-version layout: no `type`, resource_id left empty on purpose.
+        pb.set_tablet_id(meta.tablet_id());
+        pb.set_resource_id("");
+        break;
+    case RecycleRowsetBranch::kLegacyWithResource:
+        // Old-version layout: no `type`, resource_id populated.
+        pb.set_tablet_id(meta.tablet_id());
+        pb.set_resource_id(meta.resource_id());
+        break;
+    case RecycleRowsetBranch::kPrepareDirect:
+    case RecycleRowsetBranch::kPrepareMark:
+    case RecycleRowsetBranch::kPrepareAbort:
+        pb.set_type(RecycleRowsetPB::PREPARE);
+        pb.mutable_rowset_meta()->CopyFrom(meta);
+        break;
+    case RecycleRowsetBranch::kCompactedWithData:
+    case RecycleRowsetBranch::kCompactedEmpty:
+        pb.set_type(RecycleRowsetPB::COMPACT);
+        pb.mutable_rowset_meta()->CopyFrom(meta);
+        if (branch == RecycleRowsetBranch::kCompactedWithData) {
+            // Match production rowsets whose schema is stored separately in 
the schema KV.
+            pb.mutable_rowset_meta()->clear_tablet_schema();
+        }
+        break;
+    }
+    return pb;
+}
+
+enum class DeleteBitmapVersion { kNone, kV1, kV2 };
+
+void put_benchmark_delete_bitmap(Transaction* txn, const std::string& 
instance_id,
+                                 const doris::RowsetMetaCloudPB& meta,
+                                 DeleteBitmapVersion version) {
+    switch (version) {
+    case DeleteBitmapVersion::kNone:
+        break;
+    case DeleteBitmapVersion::kV1:
+        txn->put(meta_delete_bitmap_key({instance_id, meta.tablet_id(), 
meta.rowset_id_v2(), 0, 0}),
+                 "delete_bitmap_data");
+        break;
+    case DeleteBitmapVersion::kV2: {
+        // As with segment/index files, metadata is enough to exercise 
deletion of
+        // the standalone .dbm file. Keep bitmap storage independent of data 
packing.
+        DeleteBitmapStoragePB storage;
+        storage.set_store_in_fdb(false);
+        blob_put(txn,
+                 versioned::meta_delete_bitmap_key(
+                         {instance_id, meta.tablet_id(), meta.rowset_id_v2()}),
+                 storage, 0);
+        break;
+    }
+    }
+}
+
+// Seed `count` recycle rowset KVs for one branch. Every rowset gets a distinct
+// tablet_id so the per-tablet recycle batch limit never truncates the 
workload.
+int seed_recycle_rowsets(TxnKv* txn_kv, const std::string& instance_id, 
RecycleRowsetBranch branch,
+                         int64_t count, int64_t tablet_id_base, bool 
write_schema_kv = true,
+                         DeleteBitmapVersion bitmap_version = 
DeleteBitmapVersion::kNone) {
+    if (write_schema_kv && put_benchmark_schema(txn_kv, instance_id) != 0) {
+        return -1;
+    }
+    std::unique_ptr<Transaction> txn;
+    for (int64_t i = 0; i < count; ++i) {
+        if (i % kSeedCommitBatch == 0) {
+            if (txn) {
+                if (txn->commit() != TxnErrorCode::TXN_OK) {
+                    return -1;
+                }
+            }
+            if (txn_kv->create_txn(&txn) != TxnErrorCode::TXN_OK) {
+                return -1;
+            }
+        }
+        int64_t tablet_id = tablet_id_base + i;
+        std::string rowset_id = fmt::format("{:018d}", i);
+        auto meta = make_rowset_meta(branch, tablet_id, rowset_id);
+        auto pb = make_recycle_rowset(branch, meta);
+
+        std::string key;
+        recycle_rowset_key({instance_id, tablet_id, rowset_id}, &key);
+        std::string val;
+        pb.SerializeToString(&val);
+        txn->put(key, val);
+        put_benchmark_delete_bitmap(txn.get(), instance_id, meta, 
bitmap_version);
+
+        if (branch == RecycleRowsetBranch::kPrepareAbort) {
+            const auto txn_id = meta.txn_id();
+            const auto label = fmt::format("recycler_benchmark_txn_{}", 
txn_id);
+            TxnIndexPB index;
+            index.mutable_tablet_index()->set_db_id(kBenchmarkDbId);
+            index.mutable_tablet_index()->set_tablet_id(tablet_id);
+            TxnInfoPB info;
+            info.set_txn_id(txn_id);
+            info.set_db_id(kBenchmarkDbId);
+            info.set_label(label);
+            info.set_status(TxnStatusPB::TXN_STATUS_PREPARED);
+            info.set_prepare_time(1);
+            info.set_timeout_ms(60000);
+            TxnRunningPB running;
+            running.set_timeout_time(info.prepare_time() + info.timeout_ms());
+            TxnLabelPB txn_label;
+            txn_label.add_txn_ids(txn_id);
+            txn->put(txn_index_key({instance_id, txn_id}), 
index.SerializeAsString());
+            txn->put(txn_info_key({instance_id, kBenchmarkDbId, txn_id}), 
info.SerializeAsString());
+            txn->put(txn_running_key({instance_id, kBenchmarkDbId, txn_id}),
+                     running.SerializeAsString());
+            txn->atomic_set_ver_value(txn_label_key({instance_id, 
kBenchmarkDbId, label}),
+                                      txn_label.SerializeAsString());
+        }
+    }
+    if (txn && txn->commit() != TxnErrorCode::TXN_OK) {
+        return -1;
+    }
+    return 0;
+}
+
+std::string benchmark_packed_file_path(int64_t file_id) {
+    return fmt::format("data/packed_file/recycler_benchmark_{}.pack", file_id);
+}
+
+// Seed a complete packed file and its rowsets in the same transaction.
+void put_packed_recycle_rowsets(Transaction* txn, const std::string& 
instance_id,
+                                int64_t tablet_id_base, int64_t file_id,
+                                DeleteBitmapVersion bitmap_version) {
+    const auto path = benchmark_packed_file_path(file_id);
+    PackedFileInfoPB packed_info;
+    packed_info.set_resource_id(kBenchmarkResourceId);
+    packed_info.set_state(PackedFileInfoPB::NORMAL);
+    int64_t offset = 0;
+    for (int64_t j = 0; j < kRowsetsPerPackedFile; ++j) {
+        // Spread references across scan batches so different recycler workers 
can
+        // contend on the same packed file KV, instead of handling them 
serially.
+        const auto rowset_number = file_id + j * kPackedFileCount;
+        const auto tablet_id = tablet_id_base + rowset_number;
+        const auto rowset_id = fmt::format("{:018d}", rowset_number);
+        auto meta = make_rowset_meta(RecycleRowsetBranch::kCompactedWithData, 
tablet_id, rowset_id);
+        const std::pair<std::string, int64_t> files[] = {
+                {segment_path(tablet_id, rowset_id, 0), meta.data_disk_size()},
+                {inverted_index_path_v1(tablet_id, rowset_id, 0, 1, ""), 
meta.index_disk_size()}};
+        for (const auto& [small_path, size] : files) {
+            auto& location = 
(*meta.mutable_packed_slice_locations())[small_path];
+            location.set_packed_file_path(path);
+            location.set_offset(offset);
+            location.set_size(size);
+            auto* slice = packed_info.add_slices();
+            slice->set_path(small_path);
+            slice->set_offset(offset);
+            slice->set_size(size);
+            slice->set_deleted(false);
+            slice->set_tablet_id(tablet_id);
+            slice->set_rowset_id(rowset_id);
+            offset += size;
+        }
+        auto rowset = 
make_recycle_rowset(RecycleRowsetBranch::kCompactedWithData, meta);
+        txn->put(recycle_rowset_key({instance_id, tablet_id, rowset_id}),
+                 rowset.SerializeAsString());
+        put_benchmark_delete_bitmap(txn, instance_id, meta, bitmap_version);
+    }
+    packed_info.set_ref_cnt(packed_info.slices_size());
+    packed_info.set_total_slice_num(packed_info.slices_size());
+    packed_info.set_total_slice_bytes(offset);
+    packed_info.set_remaining_slice_bytes(offset);
+    txn->put(packed_file_key({instance_id, path}), 
packed_info.SerializeAsString());
+}
+
+class RecyclerBenchmarkTest : public ::testing::Test {
+protected:
+    struct RecycleMetrics {
+        int64_t num = 0;
+        int64_t bytes = 0;
+    };
+
+    struct TxnKvCounts {
+        int64_t get = 0;
+        int64_t put = 0;
+        int64_t del = 0;
+
+        int64_t total() const { return get + put + del; }
+
+        TxnKvCounts& operator+=(const TxnKvCounts& rhs) {
+            get += rhs.get;
+            put += rhs.put;
+            del += rhs.del;
+            return *this;
+        }
+    };
+
+    struct BenchmarkResult {
+        RecycleMetrics metrics;
+        TxnKvCounts txn_kv_counts;
+        double elapsed_ms = 0;
+    };
+
+    using RecycleFunction = int (InstanceRecycler::*)();
+
+    static const std::map<std::string, RecycleFunction>& recycle_functions() {
+        static const std::map<std::string, RecycleFunction> functions = {
+                {"recycle_cluster_snapshots", 
&InstanceRecycler::recycle_cluster_snapshots},
+                {"recycle_operation_logs", 
&InstanceRecycler::recycle_operation_logs},
+                {"recycle_indexes", &InstanceRecycler::recycle_indexes},
+                {"recycle_partitions", &InstanceRecycler::recycle_partitions},
+                {"recycle_tmp_rowsets", 
&InstanceRecycler::recycle_tmp_rowsets},
+                {"recycle_rowsets", &InstanceRecycler::recycle_rowsets},
+                {"recycle_packed_files", 
&InstanceRecycler::recycle_packed_files},
+                {"abort_timeout_txn", &InstanceRecycler::abort_timeout_txn},
+                {"recycle_expired_txn_label", 
&InstanceRecycler::recycle_expired_txn_label},
+                {"recycle_copy_jobs", &InstanceRecycler::recycle_copy_jobs},
+                {"recycle_stage", &InstanceRecycler::recycle_stage},
+                {"recycle_expired_stage_objects", 
&InstanceRecycler::recycle_expired_stage_objects},
+                {"recycle_versions", &InstanceRecycler::recycle_versions},
+                {"recycle_restore_jobs", 
&InstanceRecycler::recycle_restore_jobs},
+        };
+        return functions;
+    }
+
+    void SetUp() override {
+        old_force_immediate_recycle_ = config::force_immediate_recycle;
+        old_retention_seconds_ = config::retention_seconds;
+
+        const auto tolerance_env = 
get_env("recycler_benchmark_duration_tolerance_ratio");
+        size_t parsed_size = 0;
+        recycler_benchmark_duration_tolerance_ratio =
+                tolerance_env.empty() ? 0.8 : std::stod(tolerance_env, 
&parsed_size);
+        ASSERT_EQ(parsed_size, tolerance_env.size()) << "invalid recycler 
benchmark tolerance";
+        
ASSERT_TRUE(std::isfinite(recycler_benchmark_duration_tolerance_ratio));
+        ASSERT_GE(recycler_benchmark_duration_tolerance_ratio, 0);
+
+        config::force_immediate_recycle = true;
+        config::retention_seconds = 0;
+
+        instance_.set_instance_id(std::string(kBenchmarkInstanceId));
+        auto* obj_info = instance_.add_obj_info();
+        obj_info->set_id(kBenchmarkResourceId);
+        const auto s3_config = load_benchmark_s3_config();
+        ASSERT_NO_FATAL_FAILURE(configure_obj_info(s3_config, obj_info));
+
+        s3_producer_pool_ = std::make_shared<SimpleThreadPool>(
+                config::recycle_pool_parallelism, 
"recycler_benchmark_s3_producer_pool");
+        recycle_tablet_pool_ = std::make_shared<SimpleThreadPool>(
+                config::recycle_pool_parallelism, 
"recycler_benchmark_recycle_tablet_pool");
+        group_recycle_function_pool_ = std::make_shared<SimpleThreadPool>(
+                config::recycle_pool_parallelism, 
"recycler_benchmark_group_recycle_function_pool");
+        ASSERT_EQ(s3_producer_pool_->start(), 0);
+        ASSERT_EQ(recycle_tablet_pool_->start(), 0);
+        ASSERT_EQ(group_recycle_function_pool_->start(), 0);
+
+        thread_pool_group_ = RecyclerThreadPoolGroup(s3_producer_pool_, 
recycle_tablet_pool_,
+                                                     
group_recycle_function_pool_);
+        ASSERT_NO_FATAL_FAILURE(init_recycler());
+    }
+
+    void init_recycler() {
+        recycler_.reset();
+        txn_lazy_committer_.reset();
+        config::fdb_cluster_file_path = "fdb.cluster";
+        txn_kv_ = std::make_shared<FdbTxnKv>();
+        ASSERT_EQ(txn_kv_->init(), 0);
+        txn_lazy_committer_ = std::make_shared<TxnLazyCommitter>(txn_kv_);
+        recycler_ = std::make_unique<InstanceRecycler>(txn_kv_, instance_, 
thread_pool_group_,
+                                                       txn_lazy_committer_);
+
+        if (load_benchmark_s3_config().enabled) {
+            auto s3_conf = S3Conf::from_obj_store_info(instance_.obj_info(0));
+            ASSERT_TRUE(s3_conf.has_value());
+
+            std::shared_ptr<S3Accessor> accessor;
+            ASSERT_EQ(S3Accessor::create(std::move(*s3_conf), &accessor), 0);
+            s3_accessor_ = std::move(accessor);
+            recycler_->TEST_add_accessor(kBenchmarkResourceId, s3_accessor_);
+        }
+        ASSERT_EQ(recycler_->init(), 0);
+    }
+
+    template <typename Benchmark>
+    void run_benchmark(const std::string& branch, Benchmark benchmark) {
+        ::testing::TestPartResultArray failures;
+        {
+            ::testing::ScopedFakeTestPartResultReporter reporter(
+                    
::testing::ScopedFakeTestPartResultReporter::INTERCEPT_ALL_THREADS, &failures);
+            SCOPED_TRACE(branch);
+            try {
+                // Fatal assertions return from this lambda or the scenario, 
not the whole test.
+                [&] { benchmark(); }();
+            } catch (const std::exception& e) {
+                ADD_FAILURE() << "benchmark threw: " << e.what();
+            } catch (...) {
+                ADD_FAILURE() << "benchmark threw an unknown exception";
+            }
+        }
+        for (int i = 0; i < failures.size(); ++i) {
+            const auto& failure = failures.GetTestPartResult(i);
+            if (failure.failed()) {
+                benchmark_failures_ +=
+                        fmt::format("branch={}\n{}\n", branch, 
::testing::PrintToString(failure));
+            }
+        }
+    }
+
+    static void configure_obj_info(const BenchmarkS3Config& s3_config,
+                                   ObjectStoreInfoPB* obj_info) {
+        if (!s3_config.enabled) {
+            obj_info->set_prefix(kBenchmarkResourceId);
+            return;
+        }
+        ASSERT_FALSE(s3_config.endpoint.empty());
+        ASSERT_FALSE(s3_config.region.empty());
+        ASSERT_FALSE(s3_config.bucket.empty());
+        ASSERT_FALSE(s3_config.prefix.empty());
+        ASSERT_TRUE((s3_config.access_key.empty() && 
s3_config.secret_key.empty()) ||
+                    (!s3_config.access_key.empty() && 
!s3_config.secret_key.empty()));
+
+        obj_info->set_ak(s3_config.access_key);
+        obj_info->set_sk(s3_config.secret_key);
+        obj_info->set_endpoint(s3_config.endpoint);
+        obj_info->set_region(s3_config.region);
+        obj_info->set_bucket(s3_config.bucket);
+        obj_info->set_prefix(fmt::format("{}{}recycler_benchmark/{}", 
s3_config.prefix,
+                                         s3_config.prefix.ends_with('/') ? "" 
: "/",
+                                         butil::GenerateGUID()));
+        set_obj_store_provider(s3_config.provider, obj_info);
+        if (s3_config.access_key.empty()) {
+            obj_info->set_role_arn(s3_config.role_arn);
+            obj_info->set_external_id(s3_config.external_id);
+            
obj_info->set_cred_provider_type(CredProviderTypePB::INSTANCE_PROFILE);
+        }
+    }
+
+    void TearDown() override {
+        recycler_.reset();
+        txn_lazy_committer_.reset();
+        if (s3_accessor_) {
+            // The accessor is rooted at this run's GUID directory, never the 
external prefix.
+            const int ret = s3_accessor_->delete_all();
+            if (ret != 0) {
+                benchmark_failures_ += fmt::format("branch=s3_cleanup ret={} 
prefix={}\n", ret,
+                                                   
instance_.obj_info(0).prefix());
+            }
+            s3_accessor_.reset();
+        }
+
+        std::string report;
+        for (const auto& [operation_type, branches] : benchmark_results_) {
+            size_t branch_width = 24;
+            for (const auto& entry : branches) {
+                branch_width = std::max(branch_width, entry.first.size());
+            }
+            report += fmt::format("recycler benchmark: operation={}\n", 
operation_type);
+            report += fmt::format(
+                    "  {:<{}}  {:>12}  {:>13}  {:>16}  {:>12}  {:>12}  {:>12}"
+                    "  {:>12}\n",
+                    "branch", branch_width, "total_ms", "recycled_num", 
"recycled_bytes",
+                    "get_keys", "put_keys", "del_keys", "total_keys");
+            double total_elapsed_ms = 0;
+            RecycleMetrics total_recycle_metrics;
+            TxnKvCounts total_txn_kv_counts;
+            for (const auto& [branch, result] : branches) {
+                report += fmt::format(
+                        "  {:<{}}  {:>12.2f}  {:>13}  {:>16}  {:>12}  {:>12}  
{:>12}  {:>12}\n",
+                        branch, branch_width, result.elapsed_ms, 
result.metrics.num,
+                        result.metrics.bytes, result.txn_kv_counts.get, 
result.txn_kv_counts.put,
+                        result.txn_kv_counts.del, 
result.txn_kv_counts.total());
+                total_elapsed_ms += result.elapsed_ms;
+                total_recycle_metrics.num += result.metrics.num;
+                total_recycle_metrics.bytes += result.metrics.bytes;
+                total_txn_kv_counts += result.txn_kv_counts;
+            }
+            report += fmt::format(
+                    "  {:<{}}  {:>12.2f}  {:>13}  {:>16}  {:>12}  {:>12}  
{:>12}  {:>12}\n",
+                    "total", branch_width, total_elapsed_ms, 
total_recycle_metrics.num,
+                    total_recycle_metrics.bytes, total_txn_kv_counts.get, 
total_txn_kv_counts.put,
+                    total_txn_kv_counts.del, total_txn_kv_counts.total());
+        }
+        std::cout << report << std::flush;
+        if (!benchmark_failures_.empty()) {
+            ADD_FAILURE() << "recycler benchmark failures:\n" << 
benchmark_failures_;
+        }
+
+        thread_pool_group_ = {};
+        if (group_recycle_function_pool_) {
+            ASSERT_EQ(group_recycle_function_pool_->stop(), 0);
+        }
+        if (recycle_tablet_pool_) {
+            ASSERT_EQ(recycle_tablet_pool_->stop(), 0);
+        }
+        if (s3_producer_pool_) {
+            ASSERT_EQ(s3_producer_pool_->stop(), 0);
+        }
+        group_recycle_function_pool_.reset();
+        recycle_tablet_pool_.reset();
+        s3_producer_pool_.reset();
+        txn_kv_.reset();
+
+        config::force_immediate_recycle = old_force_immediate_recycle_;
+        config::retention_seconds = old_retention_seconds_;
+    }
+
+    RecycleMetrics read_recycle_metrics(const std::string& operation_type) 
const {
+        return {.num = 
g_bvar_recycler_instance_recycle_total_num_since_started.get(
+                        {kBenchmarkInstanceId, operation_type}),
+                .bytes = 
g_bvar_recycler_instance_recycle_total_bytes_since_started.get(
+                        {kBenchmarkInstanceId, operation_type})};
+    }
+
+    TxnKvCounts read_txn_kv_counts() const {
+        return {.get = g_bvar_txn_kv_get_count_normalized.get_value(),
+                .put = g_bvar_txn_kv_put.count() + 
g_bvar_txn_kv_atomic_set_ver_key.count() +
+                       g_bvar_txn_kv_atomic_set_ver_value.count() +
+                       g_bvar_txn_kv_atomic_add.count(),
+                .del = g_bvar_txn_kv_remove.count() + 
g_bvar_txn_kv_range_remove.count()};
+    }
+
+    void check_elapsed_ms(const std::string& branch, double actual_ms, double 
baseline_ms) {
+        const double limit_ms = baseline_ms * (1 + 
recycler_benchmark_duration_tolerance_ratio);
+        if (actual_ms > limit_ms) {
+            benchmark_failures_ +=
+                    fmt::format("branch={} actual_ms={:.2f} baseline_ms={:.2f} 
limit_ms={:.2f}\n",
+                                branch, actual_ms, baseline_ms, limit_ms);
+        }
+    }
+
+    void check_txn_kv_counts(const std::string& branch, const TxnKvCounts& 
actual,
+                             const TxnKvCounts& baseline) {
+        if (actual.get != baseline.get || actual.put != baseline.put ||
+            actual.del != baseline.del) {
+            benchmark_failures_ += fmt::format(
+                    "branch={} txn_kv_counts actual=[get={}, put={}, del={}] "
+                    "baseline=[get={}, put={}, del={}]\n",
+                    branch, actual.get, actual.put, actual.del, baseline.get, 
baseline.put,
+                    baseline.del);
+        }
+    }
+
+    // Time only the recycler call; seeding, assertions and metrics reads are 
excluded.
+    void measure(const std::string& operation_type, const std::string& branch,
+                 const std::string& phase) {
+        SCOPED_TRACE(phase);
+        const auto function_it = recycle_functions().find(operation_type);
+        ASSERT_NE(function_it, recycle_functions().end())
+                << "unknown recycler operation: " << operation_type;
+
+        const auto recycle_metrics_before = 
read_recycle_metrics(operation_type);
+        const auto txn_kv_counts_before = read_txn_kv_counts();
+        const auto start = std::chrono::steady_clock::now();
+        const int ret = (recycler_.get()->*(function_it->second))();
+        const auto elapsed =
+                std::chrono::duration<double, 
std::milli>(std::chrono::steady_clock::now() - start);
+        const auto recycle_metrics_after = 
read_recycle_metrics(operation_type);
+        const auto txn_kv_counts_after = read_txn_kv_counts();
+        const RecycleMetrics metrics {
+                .num = recycle_metrics_after.num - recycle_metrics_before.num,
+                .bytes = recycle_metrics_after.bytes - 
recycle_metrics_before.bytes};
+        const int64_t get_keys = txn_kv_counts_after.get - 
txn_kv_counts_before.get;
+        const int64_t put_keys = txn_kv_counts_after.put - 
txn_kv_counts_before.put;
+        const int64_t del_keys = txn_kv_counts_after.del - 
txn_kv_counts_before.del;
+        const TxnKvCounts txn_kv_counts {.get = get_keys, .put = put_keys, 
.del = del_keys};
+        auto& result = benchmark_results_[operation_type][branch];
+        result.elapsed_ms += elapsed.count();
+        result.metrics.num += metrics.num;
+        result.metrics.bytes += metrics.bytes;
+        result.txn_kv_counts += txn_kv_counts;
+        ASSERT_EQ(ret, 0) << "recycler operation failed: " << operation_type;
+    }
+
+    void seed_packed_recycle_rowsets(int64_t tablet_id_base, 
DeleteBitmapVersion bitmap_version) {
+        ASSERT_EQ(put_benchmark_schema(txn_kv_.get(), kBenchmarkInstanceId), 
0);
+        std::unique_ptr<Transaction> txn;
+        for (int64_t file_id = 0; file_id < kPackedFileCount; ++file_id) {
+            if (file_id % (kSeedCommitBatch / kRowsetsPerPackedFile) == 0) {
+                if (txn) {
+                    ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK);
+                }
+                ASSERT_EQ(txn_kv_->create_txn(&txn), TxnErrorCode::TXN_OK);
+            }
+            put_packed_recycle_rowsets(txn.get(), kBenchmarkInstanceId, 
tablet_id_base, file_id,
+                                       bitmap_version);
+        }
+        ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK);
+    }
+
+    int64_t count_recycle_rowsets() {
+        std::string begin = recycle_rowset_key({kBenchmarkInstanceId, 0, ""});
+        const auto end =
+                recycle_rowset_key({kBenchmarkInstanceId, 
std::numeric_limits<int64_t>::max(), ""});
+        std::unique_ptr<Transaction> txn;
+        if (txn_kv_->create_txn(&txn) != TxnErrorCode::TXN_OK) {
+            return -1;
+        }
+        int64_t count = 0;
+        std::unique_ptr<RangeGetIterator> it;
+        do {
+            if (txn->get(begin, end, &it) != TxnErrorCode::TXN_OK) {
+                return -1;
+            }
+            count += it->size();
+            begin = it->next_begin_key();
+        } while (it->more());
+        return count;
+    }
+
+    void check_recycled_bitmap_range(int64_t tablet_id_base) {
+        std::unique_ptr<Transaction> txn;
+        ASSERT_EQ(txn_kv_->create_txn(&txn), TxnErrorCode::TXN_OK);
+        std::unique_ptr<RangeGetIterator> it;
+        ASSERT_EQ(txn->get(versioned::meta_delete_bitmap_key(
+                                   {kBenchmarkInstanceId, tablet_id_base, ""}),
+                           versioned::meta_delete_bitmap_key(
+                                   {kBenchmarkInstanceId, tablet_id_base + 
kRowsetsPerBranch, ""}),
+                           &it),
+                  TxnErrorCode::TXN_OK);
+        ASSERT_FALSE(it->has_next());
+    }
+
+    std::shared_ptr<TxnKv> txn_kv_;
+    InstanceInfoPB instance_;
+    RecyclerThreadPoolGroup thread_pool_group_;
+    std::shared_ptr<SimpleThreadPool> s3_producer_pool_;
+    std::shared_ptr<SimpleThreadPool> recycle_tablet_pool_;
+    std::shared_ptr<SimpleThreadPool> group_recycle_function_pool_;
+    std::shared_ptr<TxnLazyCommitter> txn_lazy_committer_;
+    std::unique_ptr<InstanceRecycler> recycler_;
+    std::shared_ptr<S3Accessor> s3_accessor_;
+    std::map<std::string, std::map<std::string, BenchmarkResult>> 
benchmark_results_;
+    std::string benchmark_failures_;
+    bool reset_before_next_benchmark_ = false;
+    double test_elapsed_ms_ = 0;
+
+    bool old_force_immediate_recycle_ = false;
+    int64_t old_retention_seconds_ = 0;
+};
+
+// Non-overlapping tablet id ranges per branch so a single recycle_rowsets() 
run
+// can cover several branches at once without id collisions.
+constexpr int64_t kTabletBase = 1'000'000;
+constexpr int64_t kTabletStride = 1'000'000'000LL;
+
+int64_t tablet_base_for(RecycleRowsetBranch branch) {
+    return kTabletBase + static_cast<int64_t>(branch) * kTabletStride;
+}
+
+TEST_F(RecyclerBenchmarkTest, RecycleRowsets) {
+    RecycleRowsetConfigGuard config_guard;
+    const auto test_start = std::chrono::steady_clock::now();
+    DORIS_CLOUD_DEFER {
+        const auto elapsed = std::chrono::duration<double, std::milli>(
+                std::chrono::steady_clock::now() - test_start);
+        // Includes seeding and validation, unlike the sum of timed recycler 
calls.
+        test_elapsed_ms_ = elapsed.count();
+    };
+
+    run_benchmark("legacy_empty_resource", [&] {
+        ASSERT_EQ(seed_recycle_rowsets(txn_kv_.get(), kBenchmarkInstanceId,
+                                       
RecycleRowsetBranch::kLegacyEmptyResource, kRowsetsPerBranch,
+                                       
tablet_base_for(RecycleRowsetBranch::kLegacyEmptyResource)),
+                  0);
+        ASSERT_NO_FATAL_FAILURE(measure("recycle_rowsets", 
"legacy_empty_resource", "delete"));
+        ASSERT_EQ(count_recycle_rowsets(), 0);
+    });
+
+    run_benchmark("legacy_with_resource", [&] {
+        ASSERT_EQ(seed_recycle_rowsets(txn_kv_.get(), kBenchmarkInstanceId,
+                                       
RecycleRowsetBranch::kLegacyWithResource, kRowsetsPerBranch,
+                                       
tablet_base_for(RecycleRowsetBranch::kLegacyWithResource)),
+                  0);
+        ASSERT_NO_FATAL_FAILURE(measure("recycle_rowsets", 
"legacy_with_resource", "delete"));
+        ASSERT_EQ(count_recycle_rowsets(), 0);
+    });
+
+    run_benchmark("prepare_direct", [&] {
+        ASSERT_EQ(seed_recycle_rowsets(txn_kv_.get(), kBenchmarkInstanceId,
+                                       RecycleRowsetBranch::kPrepareDirect, 
kRowsetsPerBranch,
+                                       
tablet_base_for(RecycleRowsetBranch::kPrepareDirect)),
+                  0);
+        ASSERT_NO_FATAL_FAILURE(measure("recycle_rowsets", "prepare_direct", 
"delete"));
+        ASSERT_EQ(count_recycle_rowsets(), 0);
+    });
+
+    run_benchmark("prepare_mark", [&] {
+        const auto tablet_id_base = 
tablet_base_for(RecycleRowsetBranch::kPrepareMark);
+        ASSERT_EQ(seed_recycle_rowsets(txn_kv_.get(), kBenchmarkInstanceId,
+                                       RecycleRowsetBranch::kPrepareMark, 
kRowsetsPerBranch,
+                                       tablet_id_base),
+                  0);
+        ASSERT_NO_FATAL_FAILURE(measure("recycle_rowsets", "prepare_mark", 
"mark"));
+        ASSERT_EQ(count_recycle_rowsets(), kRowsetsPerBranch);
+        ASSERT_NO_FATAL_FAILURE(measure("recycle_rowsets", "prepare_mark", 
"delete"));
+        ASSERT_EQ(count_recycle_rowsets(), 0);
+    });
+
+    run_benchmark("prepare_abort", [&] {
+        const auto txn_id_base = 
tablet_base_for(RecycleRowsetBranch::kPrepareAbort);
+        ASSERT_EQ(seed_recycle_rowsets(txn_kv_.get(), kBenchmarkInstanceId,
+                                       RecycleRowsetBranch::kPrepareAbort, 
kRowsetsPerBranch,
+                                       txn_id_base),
+                  0);
+        ASSERT_NO_FATAL_FAILURE(measure("recycle_rowsets", "prepare_abort", 
"mark"));
+        ASSERT_EQ(count_recycle_rowsets(), kRowsetsPerBranch);
+        ASSERT_NO_FATAL_FAILURE(measure("recycle_rowsets", "prepare_abort", 
"abort_and_delete"));
+        ASSERT_EQ(count_recycle_rowsets(), 0);
+    });
+
+    // Keep the missing-schema case separate from the bitmap ablation below.
+    run_benchmark("compacted_without_schema", [&] {
+        const auto tablet_id_base = 
tablet_base_for(RecycleRowsetBranch::kCompactedWithData);
+        ASSERT_EQ(seed_recycle_rowsets(txn_kv_.get(), kBenchmarkInstanceId,
+                                       
RecycleRowsetBranch::kCompactedWithData, kRowsetsPerBranch,
+                                       tablet_id_base),
+                  0);
+        ASSERT_EQ(remove_benchmark_schema(txn_kv_.get(), 
kBenchmarkInstanceId), 0);
+        ASSERT_NO_FATAL_FAILURE(measure("recycle_rowsets", 
"compacted_without_schema", "delete"));
+        ASSERT_EQ(count_recycle_rowsets(), 0);
+    });
+
+    // Hold rowset count, schema and file sizes fixed; vary only data packing 
and
+    // bitmap version. V1 bitmap KVs survive rowset recycling, so isolate each 
case.
+    for (const auto& [bitmap_version, bitmap_name] :
+         {std::pair {DeleteBitmapVersion::kNone, "data_only"},
+          std::pair {DeleteBitmapVersion::kV1, "delete_bitmap_v1"},
+          std::pair {DeleteBitmapVersion::kV2, "delete_bitmap_v2"}}) {
+        SCOPED_TRACE(bitmap_name);
+        const auto tablet_id_base =
+                tablet_base_for(RecycleRowsetBranch::kCompactedWithData) +
+                (1 + 2 * static_cast<int64_t>(bitmap_version)) * 
kRowsetsPerBranch;
+        run_benchmark(fmt::format("compacted_with_data/{}", bitmap_name), [&] {
+            ASSERT_EQ(seed_recycle_rowsets(txn_kv_.get(), kBenchmarkInstanceId,
+                                           
RecycleRowsetBranch::kCompactedWithData,
+                                           kRowsetsPerBranch, tablet_id_base, 
true, bitmap_version),
+                      0);
+            ASSERT_EQ(count_recycle_rowsets(), kRowsetsPerBranch);
+            ASSERT_NO_FATAL_FAILURE(measure("recycle_rowsets",
+                                            
fmt::format("compacted_with_data/{}", bitmap_name),
+                                            "delete"));
+            ASSERT_EQ(count_recycle_rowsets(), 0);
+            
ASSERT_NO_FATAL_FAILURE(check_recycled_bitmap_range(tablet_id_base));
+        });
+
+        run_benchmark(fmt::format("compacted_with_packed_data/{}", 
bitmap_name), [&] {
+            ASSERT_NO_FATAL_FAILURE(seed_packed_recycle_rowsets(tablet_id_base 
+ kRowsetsPerBranch,
+                                                                
bitmap_version));
+
+            int64_t remaining = count_recycle_rowsets();
+            ASSERT_EQ(remaining, kRowsetsPerBranch);
+            for (int pass = 1; pass <= kMaxPackedRecyclePasses && remaining > 
0; ++pass) {
+                const int64_t previous_remaining = remaining;
+                // A worker may exhaust packed-file transaction retries while 
recycle_rowsets()
+                // still returns success. Include every pass in the reported 
total elapsed time.
+                ASSERT_NO_FATAL_FAILURE(
+                        measure("recycle_rowsets",
+                                fmt::format("compacted_with_packed_data/{}", 
bitmap_name),
+                                fmt::format("delete_pass_{}", pass)));
+                remaining = count_recycle_rowsets();
+                ASSERT_GE(remaining, 0);
+                ASSERT_LT(remaining, previous_remaining)
+                        << "packed rowset recycling made no progress, pass=" 
<< pass;
+            }
+            ASSERT_EQ(remaining, 0)
+                    << "packed rowset recycling exceeded " << 
kMaxPackedRecyclePasses << " passes";
+
+            // Check outside the timed calls that every packed file reached 
zero references.
+            std::unique_ptr<Transaction> txn;
+            ASSERT_EQ(txn_kv_->create_txn(&txn), TxnErrorCode::TXN_OK);
+            for (int64_t file_id = 0; file_id < kPackedFileCount; ++file_id) {
+                std::string value;
+                ASSERT_EQ(txn->get(packed_file_key({kBenchmarkInstanceId,
+                                                    
benchmark_packed_file_path(file_id)}),
+                                   &value),
+                          TxnErrorCode::TXN_KEY_NOT_FOUND);
+            }
+            ASSERT_NO_FATAL_FAILURE(
+                    check_recycled_bitmap_range(tablet_id_base + 
kRowsetsPerBranch));
+        });
+    }
+
+    run_benchmark("compacted_empty", [&] {
+        ASSERT_EQ(seed_recycle_rowsets(txn_kv_.get(), kBenchmarkInstanceId,
+                                       RecycleRowsetBranch::kCompactedEmpty, 
kRowsetsPerBranch,
+                                       
tablet_base_for(RecycleRowsetBranch::kCompactedEmpty)),
+                  0);
+        ASSERT_NO_FATAL_FAILURE(measure("recycle_rowsets", "compacted_empty", 
"delete"));
+        ASSERT_EQ(count_recycle_rowsets(), 0);
+    });
+
+    // Validate ordering outside the timed calls so callback KV reads do not 
skew timings.
+    run_benchmark("prepare_abort_before_delete", [&] {
+        const auto txn_id = 
tablet_base_for(RecycleRowsetBranch::kPrepareAbort) + kRowsetsPerBranch;
+        ASSERT_EQ(seed_recycle_rowsets(txn_kv_.get(), kBenchmarkInstanceId,
+                                       RecycleRowsetBranch::kPrepareAbort, 1, 
txn_id),
+                  0);
+        ASSERT_EQ(recycler_->recycle_rowsets(), 0);
+        ASSERT_EQ(count_recycle_rowsets(), 1);
+
+        auto* sp = SyncPoint::get_instance();
+        DORIS_CLOUD_DEFER {
+            sp->clear_all_call_backs();
+            sp->disable_processing();
+        };
+        sp->enable_processing();
+
+        ASSERT_EQ(recycler_->recycle_rowsets(), 0);
+        ASSERT_EQ(count_recycle_rowsets(), 0);
+    });
+
+    // Keep mixed-workload IDs separate from the transactions aborted above.
+    // Already-marked direct rowsets are deleted in the first pass.
+    run_benchmark("mixed", [&] {
+        constexpr int64_t mixed_tablet_offset = 7 * kTabletStride;
+        ASSERT_EQ(seed_recycle_rowsets(
+                          txn_kv_.get(), kBenchmarkInstanceId,
+                          RecycleRowsetBranch::kLegacyEmptyResource, 
kRowsetsPerBranch,
+                          mixed_tablet_offset +
+                                  
tablet_base_for(RecycleRowsetBranch::kLegacyEmptyResource)),
+                  0);
+        ASSERT_EQ(seed_recycle_rowsets(
+                          txn_kv_.get(), kBenchmarkInstanceId,
+                          RecycleRowsetBranch::kLegacyWithResource, 
kRowsetsPerBranch,
+                          mixed_tablet_offset +
+                                  
tablet_base_for(RecycleRowsetBranch::kLegacyWithResource)),
+                  0);
+        ASSERT_EQ(
+                seed_recycle_rowsets(
+                        txn_kv_.get(), kBenchmarkInstanceId, 
RecycleRowsetBranch::kPrepareDirect,
+                        kRowsetsPerBranch,
+                        mixed_tablet_offset + 
tablet_base_for(RecycleRowsetBranch::kPrepareDirect)),
+                0);
+        const auto mark_tablet_id_base =
+                mixed_tablet_offset + 
tablet_base_for(RecycleRowsetBranch::kPrepareMark);
+        ASSERT_EQ(seed_recycle_rowsets(txn_kv_.get(), kBenchmarkInstanceId,
+                                       RecycleRowsetBranch::kPrepareMark, 
kRowsetsPerBranch,
+                                       mark_tablet_id_base),
+                  0);
+        const auto txn_id_base =
+                mixed_tablet_offset + 
tablet_base_for(RecycleRowsetBranch::kPrepareAbort);
+        ASSERT_EQ(seed_recycle_rowsets(txn_kv_.get(), kBenchmarkInstanceId,
+                                       RecycleRowsetBranch::kPrepareAbort, 
kRowsetsPerBranch,
+                                       txn_id_base),
+                  0);
+        ASSERT_EQ(seed_recycle_rowsets(
+                          txn_kv_.get(), kBenchmarkInstanceId,
+                          RecycleRowsetBranch::kCompactedWithData, 
kRowsetsPerBranch,
+                          mixed_tablet_offset +
+                                  
tablet_base_for(RecycleRowsetBranch::kCompactedWithData)),
+                  0);
+        ASSERT_EQ(
+                seed_recycle_rowsets(txn_kv_.get(), kBenchmarkInstanceId,
+                                     RecycleRowsetBranch::kCompactedEmpty, 
kRowsetsPerBranch,
+                                     mixed_tablet_offset +
+                                             
tablet_base_for(RecycleRowsetBranch::kCompactedEmpty)),
+                0);
+        ASSERT_NO_FATAL_FAILURE(measure("recycle_rowsets", "mixed", 
"mark_and_delete_ready"));
+        ASSERT_EQ(count_recycle_rowsets(), 2 * kRowsetsPerBranch);
+        ASSERT_NO_FATAL_FAILURE(measure("recycle_rowsets", "mixed", 
"abort_and_delete_prepare"));
+        ASSERT_EQ(count_recycle_rowsets(), 0);
+    });
+
+    // Recorded with 10,000 rowsets per branch. Recalibrate if the workload 
changes.
+    static_assert(kRowsetsPerBranch == 10000);
+    const std::map<std::string, double> baseline_elapsed_ms = {
+            {"compacted_empty", 794.18},
+            {"compacted_with_data/data_only", 2285.80},
+            {"compacted_with_data/delete_bitmap_v1", 2461.84},
+            {"compacted_with_data/delete_bitmap_v2", 2726.75},
+            {"compacted_with_packed_data/data_only", 8583.32},
+            {"compacted_with_packed_data/delete_bitmap_v1", 9315.59},
+            {"compacted_with_packed_data/delete_bitmap_v2", 8948.54},
+            {"compacted_without_schema", 3924.04},
+            {"legacy_empty_resource", 945.75},
+            {"legacy_with_resource", 6718.37},
+            {"mixed", 35178.52},
+            {"prepare_abort", 16465.44},
+            {"prepare_direct", 5360.28},
+            {"prepare_mark", 7039.65},
+    };
+    const std::map<std::string, TxnKvCounts> baseline_txn_kv_counts = {
+            {"compacted_empty", {.get = 10000, .put = 0, .del = 10000}},
+            {"compacted_with_data/data_only", {.get = 10001, .put = 0, .del = 
10000}},
+            {"compacted_with_data/delete_bitmap_v1", {.get = 10000, .put = 0, 
.del = 10000}},
+            {"compacted_with_data/delete_bitmap_v2", {.get = 20000, .put = 0, 
.del = 20000}},
+            {"compacted_with_packed_data/data_only", {.get = 25000, .put = 
10000, .del = 15000}},
+            {"compacted_with_packed_data/delete_bitmap_v1",
+             {.get = 25000, .put = 10000, .del = 15000}},
+            {"compacted_with_packed_data/delete_bitmap_v2",
+             {.get = 35000, .put = 10000, .del = 25000}},
+            {"compacted_without_schema", {.get = 10000, .put = 0, .del = 
10000}},
+            {"legacy_empty_resource", {.get = 10000, .put = 0, .del = 10000}},
+            {"legacy_with_resource", {.get = 20000, .put = 0, .del = 20000}},
+            {"mixed", {.get = 200000, .put = 40000, .del = 120000}},
+            {"prepare_abort", {.get = 90000, .put = 30000, .del = 30000}},
+            {"prepare_direct", {.get = 20000, .put = 0, .del = 20000}},
+            {"prepare_mark", {.get = 40000, .put = 10000, .del = 20000}},
+    };
+    const auto& results = benchmark_results_["recycle_rowsets"];
+    double total_elapsed_ms = 0;
+    for (const auto& [branch, baseline_ms] : baseline_elapsed_ms) {
+        const auto result = results.find(branch);
+        if (result == results.end()) {
+            benchmark_failures_ +=
+                    fmt::format("branch={} did not produce a timing result\n", 
branch);
+            continue;
+        }
+        check_elapsed_ms(branch, result->second.elapsed_ms, baseline_ms);
+        check_txn_kv_counts(branch, result->second.txn_kv_counts,
+                            baseline_txn_kv_counts.at(branch));
+        total_elapsed_ms += result->second.elapsed_ms;
+    }
+    check_elapsed_ms("total_elapsed_ms", total_elapsed_ms, 110748.07);
+}
+
+} // namespace
+} // namespace doris::cloud


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

Reply via email to