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 1378d72082f [improvement](cloud) Point delete versioned delete bitmap 
keys when recycling rowsets (#67827)
1378d72082f is described below

commit 1378d72082fc6d2600088c63e74bfb932bf93b0d
Author: Yixuan Wang <[email protected]>
AuthorDate: Fri Sep 11 15:18:02 2026 +0800

    [improvement](cloud) Point delete versioned delete bitmap keys when 
recycling rowsets (#67827)
---
 cloud/src/recycler/recycler.cpp |  90 +++++++++++++++---
 cloud/src/recycler/recycler.h   |  10 +-
 cloud/test/recycler_test.cpp    | 196 ++++++++++++++++++++++++++++++++++++++++
 3 files changed, 280 insertions(+), 16 deletions(-)

diff --git a/cloud/src/recycler/recycler.cpp b/cloud/src/recycler/recycler.cpp
index 73c2c6a1b63..680dfd09a76 100644
--- a/cloud/src/recycler/recycler.cpp
+++ b/cloud/src/recycler/recycler.cpp
@@ -221,6 +221,57 @@ static int txn_remove(TxnKv* txn_kv, 
std::vector<std::string> keys) {
     }
 }
 
+// Remove versioned delete bitmap keys grouped by rowset.
+// Each inner vector represents all DBM shard keys for one rowset that MUST be 
deleted
+// atomically in the same txn. Transaction splitting only occurs between 
rowsets, never
+// within a rowset's shard keys.
+//
+// This ensures that a DBM blob with multiple shards is either fully deleted 
or not
+// deleted at all, preventing partial deletion that would leave the blob 
unrecoverable.
+//
+// return 0 for success otherwise error
+static int delete_versioned_delete_bitmap_by_rowset(
+        TxnKv* txn_kv, const std::vector<std::vector<std::string>>& 
rowset_dbm_key_groups) {
+    if (rowset_dbm_key_groups.empty()) {
+        return 0;
+    }
+    size_t idx = 0;
+    while (idx < rowset_dbm_key_groups.size()) {
+        std::unique_ptr<Transaction> txn;
+        TxnErrorCode err = txn_kv->create_txn(&txn);
+        if (err != TxnErrorCode::TXN_OK) {
+            return -1;
+        }
+
+        bool has_keys = false;
+        while (idx < rowset_dbm_key_groups.size()) {
+            const auto& keys = rowset_dbm_key_groups[idx];
+            const auto keys_bytes =
+                    std::accumulate(keys.begin(), keys.end(), size_t {0},
+                                    [](size_t sum, const auto& key) { return 
sum + key.size(); });
+
+            if (has_keys && txn->approximate_bytes() + keys_bytes >= 
config::max_txn_commit_byte) {
+                break;
+            }
+
+            for (const auto& key : keys) {
+                txn->remove(key);
+            }
+
+            has_keys = true;
+            ++idx;
+        }
+
+        
TEST_SYNC_POINT_CALLBACK("delete_versioned_delete_bitmap_by_rowset::commit_result");
+        err = txn->commit();
+        if (err != TxnErrorCode::TXN_OK) {
+            LOG(WARNING) << "failed to remove delete bitmap keys, err=" << err;
+            return -1;
+        }
+    }
+    return 0;
+}
+
 void scan_restore_job_rowset(
         Transaction* txn, const std::string& instance_id, int64_t tablet_id, 
MetaServiceCode& code,
         std::string& msg,
@@ -4111,8 +4162,8 @@ int 
InstanceRecycler::decrement_packed_file_ref_counts(const doris::RowsetMetaCl
 }
 
 int InstanceRecycler::decrement_delete_bitmap_packed_file_ref_counts(
-        int64_t tablet_id, const std::string& rowset_id,
-        DeleteBitmapStorageType* out_storage_type) {
+        int64_t tablet_id, const std::string& rowset_id, 
DeleteBitmapStorageType* out_storage_type,
+        std::vector<std::string>* keys) {
     if (out_storage_type) {
         *out_storage_type = DeleteBitmapStorageType::NOT_FOUND;
     }
@@ -4149,6 +4200,12 @@ int 
InstanceRecycler::decrement_delete_bitmap_packed_file_ref_counts(
         return -1;
     }
 
+    if (keys) {
+        for (auto& key : dbm_val.keys()) {
+            keys->push_back(std::move(key));
+        }
+    }
+
     DeleteBitmapStoragePB storage;
     if (!dbm_val.to_pb(&storage)) {
         LOG_WARNING("failed to parse delete bitmap storage")
@@ -4458,7 +4515,8 @@ int InstanceRecycler::delete_packed_file_and_kv(const 
std::string& packed_file_p
 
 int InstanceRecycler::delete_rowset_data(
         const std::map<std::string, doris::RowsetMetaCloudPB>& rowsets, 
RowsetRecyclingState type,
-        RecyclerMetricsContext& metrics_context) {
+        RecyclerMetricsContext& metrics_context,
+        std::vector<std::vector<std::string>>* delete_bitmap_key_groups) {
     int ret = 0;
     // resource_id -> file_paths
     std::map<std::string, std::vector<std::string>> resource_file_paths;
@@ -4513,10 +4571,11 @@ int InstanceRecycler::delete_rowset_data(
             continue;
         }
 
-        // Process delete bitmap - check where it's stored.
         DeleteBitmapStorageType delete_bitmap_storage_type = 
DeleteBitmapStorageType::NOT_FOUND;
-        if (decrement_delete_bitmap_packed_file_ref_counts(tablet_id, 
rowset_id,
-                                                           
&delete_bitmap_storage_type) != 0) {
+        std::vector<std::string> rowset_dbm_keys;
+        if (decrement_delete_bitmap_packed_file_ref_counts(
+                    tablet_id, rowset_id, &delete_bitmap_storage_type,
+                    delete_bitmap_key_groups ? &rowset_dbm_keys : nullptr) != 
0) {
             LOG_WARNING("failed to decrement delete bitmap packed file ref 
count")
                     .tag("instance_id", instance_id_)
                     .tag("tablet_id", tablet_id)
@@ -4524,6 +4583,9 @@ int InstanceRecycler::delete_rowset_data(
             ret = -1;
             continue;
         }
+        if (!rowset_dbm_keys.empty() && delete_bitmap_key_groups) {
+            delete_bitmap_key_groups->emplace_back(std::move(rowset_dbm_keys));
+        }
         if (delete_bitmap_storage_type == 
DeleteBitmapStorageType::STANDALONE_FILE) {
             file_paths.push_back(delete_bitmap_path(tablet_id, rowset_id));
         }
@@ -5962,17 +6024,19 @@ int InstanceRecycler::recycle_rowsets() {
                                              std::map<std::string, 
RowsetMetaCloudPB> rowsets) {
         worker_pool->submit([&, rowset_keys_to_delete = std::move(rowset_keys),
                              rowsets_to_delete = std::move(rowsets)]() {
+            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,
-                                   metrics_context) != 0) {
+                                   metrics_context, 
&versioned_delete_bitmap_key_groups) != 0) {
                 LOG(WARNING) << "failed to delete rowset data, instance_id=" 
<< instance_id_;
                 return;
             }
-            for (const auto& [_, rs] : rowsets_to_delete) {
-                if (delete_versioned_delete_bitmap_kvs(rs.partition_id(), 
rs.tablet_id(),
-                                                       rs.rowset_id_v2()) != 
0) {
-                    return;
-                }
+            if (!versioned_delete_bitmap_key_groups.empty() &&
+                delete_versioned_delete_bitmap_by_rowset(txn_kv_.get(),
+                                                         
versioned_delete_bitmap_key_groups) != 0) {
+                LOG(WARNING) << "failed to delete versioned delete bitmap kv, 
instance_id="
+                             << instance_id_;
+                return;
             }
             if (txn_remove(txn_kv_.get(), rowset_keys_to_delete) != 0) {
                 LOG(WARNING) << "failed to delete recycle rowset kv, 
instance_id=" << instance_id_;
@@ -6016,6 +6080,7 @@ int InstanceRecycler::recycle_rowsets() {
         // 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) {
@@ -6045,7 +6110,6 @@ int InstanceRecycler::recycle_rowsets() {
                                                            
std::make_move_iterator(end));
             submit_delete_rowset_data_job(std::move(rowset_keys_to_remove), 
{});
         }
-
         return 0;
     };
 
diff --git a/cloud/src/recycler/recycler.h b/cloud/src/recycler/recycler.h
index 23a2e1563f6..77b017d3e40 100644
--- a/cloud/src/recycler/recycler.h
+++ b/cloud/src/recycler/recycler.h
@@ -525,8 +525,10 @@ private:
     int delete_delete_bitmap_kvs(int64_t tablet_id, const std::string& 
rowset_id);
 
     // return 0 for success otherwise error
-    int delete_rowset_data(const std::map<std::string, 
doris::RowsetMetaCloudPB>& rowsets,
-                           RowsetRecyclingState type, RecyclerMetricsContext& 
metrics_context);
+    int delete_rowset_data(
+            const std::map<std::string, doris::RowsetMetaCloudPB>& rowsets,
+            RowsetRecyclingState type, RecyclerMetricsContext& metrics_context,
+            std::vector<std::vector<std::string>>* delete_bitmap_key_groups = 
nullptr);
 
     // Decrement packed file ref counts for rowset segments.
     // Returns 0 for success, -1 for error.
@@ -542,9 +544,11 @@ private:
     // Process delete bitmap storage and decrement packed file ref count when 
needed.
     // Returns 0 for success, -1 for error.
     // out_storage_type: if not null, will be set to the delete bitmap storage 
type.
+    // keys: if not null, will collect all versioned delete bitmap keys for 
batch deletion.
     int decrement_delete_bitmap_packed_file_ref_counts(int64_t tablet_id,
                                                        const std::string& 
rowset_id,
-                                                       
DeleteBitmapStorageType* out_storage_type);
+                                                       
DeleteBitmapStorageType* out_storage_type,
+                                                       
std::vector<std::string>* keys = nullptr);
 
     int delete_packed_file_and_kv(const std::string& packed_file_path,
                                   const std::string& packed_key,
diff --git a/cloud/test/recycler_test.cpp b/cloud/test/recycler_test.cpp
index 880ac077dc8..a93359785d3 100644
--- a/cloud/test/recycler_test.cpp
+++ b/cloud/test/recycler_test.cpp
@@ -3686,6 +3686,78 @@ TEST(RecyclerTest, recycle_tablet_packed_file_ref_count) 
{
     ASSERT_FALSE(list_iter->has_next());
 }
 
+TEST(RecyclerTest, delete_bitmap_collect_shard_keys) {
+    auto txn_kv = std::make_shared<MemTxnKv>();
+    ASSERT_EQ(txn_kv->init(), 0);
+
+    constexpr std::string_view kResourceId = "delete_bitmap_shard_keys";
+    InstanceInfoPB instance;
+    instance.set_instance_id(instance_id);
+    auto obj_info = instance.add_obj_info();
+    obj_info->set_id(std::string(kResourceId));
+    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(std::string(kResourceId));
+
+    InstanceRecycler recycler(txn_kv, instance, thread_group,
+                              std::make_shared<TxnLazyCommitter>(txn_kv));
+    ASSERT_EQ(recycler.init(), 0);
+
+    constexpr int64_t tablet_id = 50001;
+    const std::string rowset_id = "rowset_for_delete_bitmap";
+
+    // Build a delete bitmap storage stored in FDB whose serialized value is
+    // large enough to be split into multiple KVs (shards).
+    DeleteBitmapStoragePB storage;
+    storage.set_store_in_fdb(true);
+    auto* delete_bitmap = storage.mutable_delete_bitmap();
+    delete_bitmap->add_rowset_ids(rowset_id);
+    delete_bitmap->add_segment_ids(0);
+    delete_bitmap->add_versions(1);
+    // The serialized value exceeds the default blob split size (90KB), so it 
is
+    // stored as multiple shard KVs.
+    delete_bitmap->add_segment_delete_bitmaps(std::string(200 * 1000, 'x'));
+
+    std::string dbm_key = versioned::meta_delete_bitmap_key({instance_id, 
tablet_id, rowset_id});
+    std::unique_ptr<Transaction> txn;
+    ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
+    cloud::blob_put(txn.get(), dbm_key, storage, 0);
+    ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK);
+
+    // Collect the actual shard keys to compare against.
+    std::vector<std::string> expected_keys;
+    ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
+    std::unique_ptr<RangeGetIterator> it;
+    std::string begin_key = dbm_key;
+    std::string end_key = dbm_key;
+    encode_int64(INT64_MAX, &end_key);
+    while (it == nullptr || it->more()) {
+        ASSERT_EQ(txn->get(begin_key, end_key, &it), TxnErrorCode::TXN_OK);
+        while (it->has_next()) {
+            auto [k, _] = it->next();
+            expected_keys.emplace_back(k.data(), k.size());
+        }
+        begin_key = it->next_begin_key();
+    }
+    ASSERT_GT(expected_keys.size(), 1);
+
+    // Collect shard keys and verify to_pb still works afterwards.
+    std::vector<std::string> keys;
+    InstanceRecycler::DeleteBitmapStorageType storage_type =
+            InstanceRecycler::DeleteBitmapStorageType::NOT_FOUND;
+    ASSERT_EQ(0, 
recycler.decrement_delete_bitmap_packed_file_ref_counts(tablet_id, rowset_id,
+                                                                         
&storage_type, &keys));
+
+    // The value is split into several shard keys and all of them are returned.
+    EXPECT_EQ(keys, expected_keys);
+    EXPECT_GT(keys.size(), 1);
+    // to_pb succeeded after collecting keys: function returned 0 and resolved 
IN_FDB.
+    EXPECT_EQ(storage_type, InstanceRecycler::DeleteBitmapStorageType::IN_FDB);
+}
+
 TEST(RecyclerTest, recycle_indexes) {
     config::retention_seconds = 0;
     auto txn_kv = std::make_shared<MemTxnKv>();
@@ -11871,4 +11943,128 @@ TEST(RecyclerTest, 
RecycleInstanceFilterReadsConfigDynamically) {
     EXPECT_FALSE(filter_out_instance("instance1"));
     EXPECT_TRUE(filter_out_instance("instance2"));
 }
+
+TEST(RecyclerTest, delete_versioned_delete_bitmap_by_rowset) {
+    auto txn_kv = std::make_shared<MemTxnKv>();
+    ASSERT_EQ(txn_kv->init(), 0);
+
+    // Test 1: Empty input - should succeed
+    {
+        std::vector<std::vector<std::string>> empty_groups;
+        EXPECT_EQ(delete_versioned_delete_bitmap_by_rowset(txn_kv.get(), 
empty_groups), 0);
+    }
+
+    // Test 2: Single rowset with multiple shard keys - should commit 
atomically
+    {
+        std::vector<std::vector<std::string>> single_rowset_groups;
+        std::vector<std::string> rowset1_keys = {"key1_shard0", "key1_shard1", 
"key1_shard2"};
+        single_rowset_groups.push_back(rowset1_keys);
+
+        // Pre-populate keys
+        for (const auto& key : rowset1_keys) {
+            std::unique_ptr<Transaction> txn;
+            ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
+            txn->put(key, "value");
+            ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK);
+        }
+
+        // Delete via function
+        EXPECT_EQ(delete_versioned_delete_bitmap_by_rowset(txn_kv.get(), 
single_rowset_groups), 0);
+
+        // Verify all keys are deleted
+        for (const auto& key : rowset1_keys) {
+            std::unique_ptr<Transaction> txn;
+            ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
+            std::string val;
+            EXPECT_EQ(txn->get(key, &val), TxnErrorCode::TXN_KEY_NOT_FOUND);
+        }
+    }
+
+    // Test 3: Multiple rowsets
+    {
+        std::vector<std::vector<std::string>> multi_rowset_groups;
+        std::vector<std::string> rowset1_keys = {"r1_shard0", "r1_shard1", 
"r1_shard2"};
+        multi_rowset_groups.push_back(rowset1_keys);
+        std::vector<std::string> rowset2_keys = {"r2_shard0", "r2_shard1", 
"r2_shard2",
+                                                 "r2_shard3"};
+        multi_rowset_groups.push_back(rowset2_keys);
+        std::vector<std::string> rowset3_keys = {"r3_shard0", "r3_shard1"};
+        multi_rowset_groups.push_back(rowset3_keys);
+
+        // Pre-populate all keys
+        for (const auto& group : multi_rowset_groups) {
+            for (const auto& key : group) {
+                std::unique_ptr<Transaction> txn;
+                ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
+                txn->put(key, "value");
+                ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK);
+            }
+        }
+
+        // Delete via function
+        EXPECT_EQ(delete_versioned_delete_bitmap_by_rowset(txn_kv.get(), 
multi_rowset_groups), 0);
+
+        // Verify all keys are deleted
+        for (const auto& group : multi_rowset_groups) {
+            for (const auto& key : group) {
+                std::unique_ptr<Transaction> txn;
+                ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
+                std::string val;
+                EXPECT_EQ(txn->get(key, &val), 
TxnErrorCode::TXN_KEY_NOT_FOUND);
+            }
+        }
+    }
+
+    // Test 4: Batching with byte limit
+    {
+        config::max_txn_commit_byte = 800;
+        DORIS_CLOUD_DEFER {
+            config::max_txn_commit_byte = 10000000;
+        };
+
+        std::vector<std::vector<std::string>> large_groups;
+        for (int i = 0; i < 15; ++i) {
+            std::vector<std::string> rowset_keys;
+            for (int j = 0; j < 3; ++j) {
+                // Each key is ~30 bytes: "large_rowset_0000_shard_0000"
+                
rowset_keys.push_back(fmt::format("large_rowset_{:04d}_shard_{:04d}", i, j));
+            }
+            large_groups.push_back(rowset_keys);
+        }
+
+        // Pre-populate - each group is ~90 bytes (3 keys * 30 bytes), so 800 
limit allows ~8 groups per txn
+        for (const auto& group : large_groups) {
+            for (const auto& key : group) {
+                std::unique_ptr<Transaction> txn;
+                ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
+                txn->put(key, "v");
+                ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK);
+            }
+        }
+
+        // Track commits
+        std::atomic<int> commit_count {0};
+        auto sp = SyncPoint::get_instance();
+        DORIS_CLOUD_DEFER {
+            sp->clear_all_call_backs();
+        };
+        
sp->set_call_back("delete_versioned_delete_bitmap_by_rowset::commit_result",
+                          [&](auto&&) { commit_count++; });
+        sp->enable_processing();
+
+        EXPECT_EQ(delete_versioned_delete_bitmap_by_rowset(txn_kv.get(), 
large_groups), 0);
+        EXPECT_GT(commit_count.load(), 1) << "Should batch into multiple 
transactions";
+
+        // Verify all deleted
+        for (const auto& group : large_groups) {
+            for (const auto& key : group) {
+                std::unique_ptr<Transaction> txn;
+                ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
+                std::string val;
+                EXPECT_EQ(txn->get(key, &val), 
TxnErrorCode::TXN_KEY_NOT_FOUND);
+            }
+        }
+    }
+}
+
 } // namespace doris::cloud


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

Reply via email to