This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.2
in repository https://gitbox.apache.org/repos/asf/doris.git
commit 658f5958b33ea293644a36178f66cbc7dd33d4a5
Author: github-actions[bot]
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Fri Oct 9 16:07:59 2026 +0800
branch-4.1: [fix](cloud) Delete versioned tablet indexes and metadata
atomically #68322 (#68817)
Cherry-picked from #68322
Co-authored-by: Yixuan Wang <[email protected]>
---
cloud/src/recycler/recycler.cpp | 68 +++++++++++++++------------
cloud/test/recycle_versioned_keys_test.cpp | 74 ++++++++++++++++++++++++++++++
2 files changed, 112 insertions(+), 30 deletions(-)
diff --git a/cloud/src/recycler/recycler.cpp b/cloud/src/recycler/recycler.cpp
index 76a1f048b59..7595d80317d 100644
--- a/cloud/src/recycler/recycler.cpp
+++ b/cloud/src/recycler/recycler.cpp
@@ -3361,9 +3361,6 @@ int InstanceRecycler::recycle_orphan_partitions() {
int InstanceRecycler::recycle_tablets(int64_t table_id, int64_t index_id,
RecyclerMetricsContext& metrics_context,
int64_t partition_id) {
- bool is_multi_version =
- instance_info_.has_multi_version_status() &&
- instance_info_.multi_version_status() !=
MultiVersionStatus::MULTI_VERSION_DISABLED;
int64_t num_scanned = 0;
std::atomic_long num_recycled = 0;
@@ -3500,7 +3497,43 @@ int InstanceRecycler::recycle_tablets(int64_t table_id,
int64_t index_id,
}
}
}
- if (is_multi_version) {
+ if (should_recycle_versioned_keys()) {
+ // Remove tablet indexes in the same transaction as tablet
metadata.
+ std::vector<std::string> versioned_idx_keys;
+ versioned_idx_keys.reserve(tablets_info.size());
+ for (const auto& tablet_info : tablets_info) {
+ versioned_idx_keys.push_back(
+ versioned::tablet_index_key({instance_id_,
tablet_info.tablet_id}));
+ }
+ std::vector<std::optional<std::string>> tablet_index_vals;
+ TxnErrorCode err = txn->batch_get(&tablet_index_vals,
versioned_idx_keys);
+ if (err != TxnErrorCode::TXN_OK) {
+ LOG_WARNING("failed to batch get tablet index kv")
+ .tag("instance_id", instance_id_)
+ .tag("num_tablets", tablets_info.size())
+ .tag("err", err);
+ return -1;
+ }
+ DCHECK_EQ(tablet_index_vals.size(), versioned_idx_keys.size());
+ for (size_t i = 0; i < tablets_info.size(); ++i) {
+ if (!tablet_index_vals[i].has_value()) {
+ continue;
+ }
+ const auto& tablet_info = tablets_info[i];
+ TabletIndexPB tablet_index_pb;
+ if
(!tablet_index_pb.ParseFromString(tablet_index_vals[i].value())) {
+ LOG_WARNING("failed to parse tablet index pb")
+ .tag("instance_id", instance_id_)
+ .tag("tablet_id", tablet_info.tablet_id);
+ return -1;
+ }
+ std::string versioned_inverted_idx_key =
versioned::tablet_inverted_index_key(
+ {instance_id_, tablet_index_pb.db_id(),
tablet_index_pb.table_id(),
+ tablet_index_pb.index_id(),
tablet_index_pb.partition_id(),
+ tablet_info.tablet_id});
+ txn->remove(versioned_inverted_idx_key);
+ txn->remove(versioned_idx_keys[i]);
+ }
for (auto& tablet_info : tablets_info) {
// Remove all versions of tablet compact stats for recycled
tablet
auto k = versioned::tablet_compact_stats_key({instance_id_,
tablet_info.tablet_id});
@@ -3534,6 +3567,7 @@ int InstanceRecycler::recycle_tablets(int64_t table_id,
int64_t index_id,
for (auto& k : init_rs_keys) {
txn->remove(k);
}
+
TEST_SYNC_POINT_CALLBACK("InstanceRecycler::recycle_tablets.before_commit",
txn.get());
if (TxnErrorCode err = txn->commit(); err != TxnErrorCode::TXN_OK) {
LOG(WARNING) << "failed to delete kvs related to tablets,
instance_id=" << instance_id_
<< ", err=" << err;
@@ -5470,32 +5504,6 @@ int InstanceRecycler::recycle_versioned_tablet(int64_t
tablet_id,
LOG(INFO) << "remove delete bitmap kv, tablet=" << tablet_id << ", begin="
<< hex(dbm_start_key)
<< " end=" << hex(dbm_end_key);
- std::string versioned_idx_key = versioned::tablet_index_key({instance_id_,
tablet_id});
- std::string tablet_index_val;
- err = txn->get(versioned_idx_key, &tablet_index_val);
- if (err != TxnErrorCode::TXN_KEY_NOT_FOUND && err != TxnErrorCode::TXN_OK)
{
- LOG_WARNING("failed to get tablet index kv")
- .tag("instance_id", instance_id_)
- .tag("tablet_id", tablet_id)
- .tag("err", err);
- ret = -1;
- } else if (err == TxnErrorCode::TXN_OK) {
- // If the tablet index kv exists, we need to delete it
- TabletIndexPB tablet_index_pb;
- if (!tablet_index_pb.ParseFromString(tablet_index_val)) {
- LOG_WARNING("failed to parse tablet index pb")
- .tag("instance_id", instance_id_)
- .tag("tablet_id", tablet_id);
- ret = -1;
- } else {
- std::string versioned_inverted_idx_key =
versioned::tablet_inverted_index_key(
- {instance_id_, tablet_index_pb.db_id(),
tablet_index_pb.table_id(),
- tablet_index_pb.index_id(),
tablet_index_pb.partition_id(), tablet_id});
- txn->remove(versioned_inverted_idx_key);
- txn->remove(versioned_idx_key);
- }
- }
-
err = txn->commit();
if (err != TxnErrorCode::TXN_OK) {
LOG(WARNING) << "failed to delete rowset kv of tablet " << tablet_id
<< ", err=" << err;
diff --git a/cloud/test/recycle_versioned_keys_test.cpp
b/cloud/test/recycle_versioned_keys_test.cpp
index 1551f895c30..763749dafb4 100644
--- a/cloud/test/recycle_versioned_keys_test.cpp
+++ b/cloud/test/recycle_versioned_keys_test.cpp
@@ -32,6 +32,7 @@
#include "common/defer.h"
#include "common/util.h"
+#include "cpp/sync_point.h"
#include "meta-service/meta_service.h"
#include "meta-store/codec.h"
#include "meta-store/document_message.h"
@@ -1537,6 +1538,79 @@ TEST(RecycleVersionedKeysTest,
RecycleTabletWithRowsetRefCountConcurrent) {
}
}
+TEST(RecycleVersionedKeysTest, RecycleTabletMetadataAndIndexesAtomically) {
+ auto meta_service = get_meta_service();
+ auto txn_kv = meta_service->txn_kv();
+ std::string instance_id = "recycle_tablet_metadata_and_indexes";
+ std::string cloud_unique_id = fmt::format("1:{}:0", instance_id);
+ ASSERT_NO_FATAL_FAILURE(create_and_refresh_instance(meta_service.get(),
instance_id));
+
+ int64_t db_id = 1, table_id = 2, index_id = 3, partition_id = 4, tablet_id
= 5;
+ ASSERT_NO_FATAL_FAILURE(prepare_and_commit_index(meta_service.get(),
cloud_unique_id, db_id,
+ table_id, index_id));
+ ASSERT_NO_FATAL_FAILURE(prepare_and_commit_partition(meta_service.get(),
cloud_unique_id, db_id,
+ table_id,
partition_id, index_id));
+ ASSERT_NO_FATAL_FAILURE(create_tablet(meta_service.get(), cloud_unique_id,
db_id, table_id,
+ index_id, partition_id, tablet_id));
+
+ auto check_tablet_keys = [&](bool exists) {
+ std::unique_ptr<Transaction> txn;
+ ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
+ for (const auto& key :
+ {meta_tablet_key({instance_id, table_id, index_id, partition_id,
tablet_id}),
+ meta_tablet_idx_key({instance_id, tablet_id}),
+ versioned::tablet_index_key({instance_id, tablet_id}),
+ versioned::tablet_inverted_index_key(
+ {instance_id, db_id, table_id, index_id, partition_id,
tablet_id})}) {
+ std::string value;
+ EXPECT_EQ(txn->get(key, &value),
+ exists ? TxnErrorCode::TXN_OK :
TxnErrorCode::TXN_KEY_NOT_FOUND)
+ << hex(key);
+ }
+ for (const auto& key : {versioned::meta_tablet_key({instance_id,
tablet_id}),
+ versioned::tablet_load_stats_key({instance_id,
tablet_id}),
+
versioned::tablet_compact_stats_key({instance_id, tablet_id})}) {
+ std::vector<std::pair<std::string, Versionstamp>> values;
+ ASSERT_NO_FATAL_FAILURE(versioned_get_all(txn_kv.get(), key,
values));
+ EXPECT_EQ(values.size(), exists ? 1 : 0) << hex(key);
+ }
+ };
+ ASSERT_NO_FATAL_FAILURE(check_tablet_keys(true));
+
+ InstanceInfoPB instance_info;
+ ASSERT_NO_FATAL_FAILURE(get_instance(meta_service.get(), cloud_unique_id,
instance_info));
+ auto recycler = get_instance_recycler(meta_service.get(), instance_info);
+ auto* sp = SyncPoint::get_instance();
+ DORIS_CLOUD_DEFER {
+ sp->clear_all_call_backs();
+ sp->disable_processing();
+ };
+ bool commit_attempted = false;
+ sp->set_call_back("InstanceRecycler::recycle_tablets.before_commit",
[&](auto&& args) {
+ commit_attempted = true;
+ // Force a real commit conflict after data recycling without changing
tablet keys.
+ auto* txn = try_any_cast<Transaction*>(args[0]);
+ std::string key = instance_key(instance_id);
+ std::string value;
+ ASSERT_EQ(txn->get(key, &value), TxnErrorCode::TXN_OK);
+ std::unique_ptr<Transaction> conflicting_txn;
+ ASSERT_EQ(txn_kv->create_txn(&conflicting_txn), TxnErrorCode::TXN_OK);
+ conflicting_txn->put(key, value);
+ ASSERT_EQ(conflicting_txn->commit(), TxnErrorCode::TXN_OK);
+ });
+ sp->enable_processing();
+
+ RecyclerMetricsContext ctx;
+ ASSERT_EQ(recycler->recycle_tablets(table_id, index_id, ctx), -1);
+ ASSERT_TRUE(commit_attempted);
+ ASSERT_NO_FATAL_FAILURE(check_tablet_keys(true));
+
+ sp->disable_processing();
+ recycler = get_instance_recycler(meta_service.get(), instance_info);
+ ASSERT_EQ(recycler->recycle_tablets(table_id, index_id, ctx), 0);
+ ASSERT_NO_FATAL_FAILURE(check_tablet_keys(false));
+}
+
// A test that simulates a drop index operation.
TEST(RecycleVersionedKeysTest, RecycleIndex) {
auto meta_service = get_meta_service();
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]