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]