From b24762b80c4016c45d0db0ef03bed184359d54d3 Mon Sep 17 00:00:00 2001 From: Yixuan Wang Date: Mon, 13 Jul 2026 19:18:25 +0800 Subject: [PATCH] 1 --- cloud/src/recycler/recycler.cpp | 210 +++++++++++------------------ cloud/src/recycler/recycler.h | 4 +- cloud/test/recycler_test.cpp | 232 +++++++++++++++++++++++++++++++- 3 files changed, 306 insertions(+), 140 deletions(-) diff --git a/cloud/src/recycler/recycler.cpp b/cloud/src/recycler/recycler.cpp index bd1b9bff878cee..9af6a45037ddc6 100644 --- a/cloud/src/recycler/recycler.cpp +++ b/cloud/src/recycler/recycler.cpp @@ -1872,22 +1872,32 @@ int InstanceRecycler::abort_job_for_related_rowset(const RowsetMetaCloudPB& rows return -1; } - std::string job_id {}; - if (!job_pb.compaction().empty()) { - for (const auto& c : job_pb.compaction()) { - if (c.id() == rowset_meta.job_id()) { - job_id = c.id(); - break; - } + const TabletCompactionJobPB* matched_compaction = nullptr; + for (const auto& compaction : job_pb.compaction()) { + if (compaction.id() == rowset_meta.job_id()) { + matched_compaction = &compaction; + break; } - } else if (job_pb.has_schema_change()) { - job_id = job_pb.schema_change().id(); + } + const TabletSchemaChangeJobPB* matched_schema_change = nullptr; + if (matched_compaction == nullptr && job_pb.has_schema_change() && + job_pb.schema_change().id() == rowset_meta.job_id()) { + matched_schema_change = &job_pb.schema_change(); } - if (!job_id.empty() && rowset_meta.job_id() == job_id) { + if (matched_compaction != nullptr || matched_schema_change != nullptr) { LOG(INFO) << "begin to abort job for related rowset, job_id=" << rowset_meta.job_id() << " instance_id=" << instance_id_ << " tablet_id=" << tablet_idx.tablet_id(); - req.mutable_job()->CopyFrom(job_pb); + // Compaction jobs belong to the rowset's tablet. A schema-change job is mirrored under + // the new tablet key, but its recorded index remains the base tablet so ABORT can remove + // both the base and mirrored schema-change records. + if (matched_compaction != nullptr) { + req.mutable_job()->mutable_idx()->CopyFrom(tablet_idx); + req.mutable_job()->add_compaction()->CopyFrom(*matched_compaction); + } else { + req.mutable_job()->mutable_idx()->CopyFrom(job_pb.idx()); + req.mutable_job()->mutable_schema_change()->CopyFrom(*matched_schema_change); + } req.set_action(FinishTabletJobRequest::ABORT); _finish_tablet_job(&req, &res, instance_id_, txn, txn_kv_.get(), delete_bitmap_lock_white_list_.get(), resource_mgr_.get(), code, msg, @@ -1906,7 +1916,7 @@ int InstanceRecycler::abort_job_for_related_rowset(const RowsetMetaCloudPB& rows LOG(INFO) << "there is no job for related rowset, directly recycle rowset data" << ", instance_id=" << instance_id_ << ", tablet_id=" << tablet_idx.tablet_id() - << ", job_id=" << job_id + << ", job_id=" << rowset_meta.job_id() << ", rowset_id=" << rowset_meta.rowset_id_v2(); // clang-format on } @@ -1954,15 +1964,9 @@ struct DeferredRecyclePrepareDeleteTask { int64_t tablet_id = 0; }; -template -std::optional make_deferred_abort_task(const T& rowset_meta_pb) { - if constexpr (std::is_same_v) { - if (rowset_meta_pb.type() != RecycleRowsetPB::PREPARE) { - return std::nullopt; - } - } - - const auto& rs_meta = rowset_meta(rowset_meta_pb); +std::optional make_deferred_abort_task( + const RowsetMetaCloudPB& rowset_meta_pb) { + const auto& rs_meta = rowset_meta_pb; DeferredRecycleAbortTask task; task.tablet_id = rs_meta.tablet_id(); task.start_version = rs_meta.start_version(); @@ -1981,10 +1985,8 @@ std::optional make_deferred_abort_task(const T& rowset return std::nullopt; } -template -bool need_mark_rowset_as_recycled(const T& rowset_meta_pb) { - const auto& rs_meta = rowset_meta(rowset_meta_pb); - return !rs_meta.has_is_recycled() || !rs_meta.is_recycled(); +bool need_mark_rowset_as_recycled(const RowsetMetaCloudPB& rowset_meta_pb) { + return !rowset_meta_pb.has_is_recycled() || !rowset_meta_pb.is_recycled(); } template @@ -2017,10 +2019,14 @@ int batch_mark_rowsets_as_recycled(TxnKv* txn_kv, const std::string& instance_id << " key=" << hex(key); return -1; } - if (!need_mark_rowset_as_recycled(rowset_meta_pb)) { + if (!need_mark_rowset_as_recycled(rowset_meta(rowset_meta_pb))) { continue; } mutable_rowset_meta(rowset_meta_pb)->set_is_recycled(true); + if constexpr (std::is_same_v) { + [[maybe_unused]] auto type = rowset_meta_pb.type(); + TEST_SYNC_POINT_CALLBACK("InstanceRecycler::batch_mark_rowsets_as_recycled", &type); + } val.clear(); rowset_meta_pb.SerializeToString(&val); txn->put(key, val); @@ -2034,11 +2040,9 @@ int batch_mark_rowsets_as_recycled(TxnKv* txn_kv, const std::string& instance_id return 0; } -template int collect_deferred_abort_tasks(TxnKv* txn_kv, const std::string& instance_id, const std::vector& keys, - std::vector* abort_tasks, - bool skip_base_version) { + std::vector* abort_tasks) { constexpr size_t kAbortCheckBatchSize = 256; for (size_t offset = 0; offset < keys.size(); offset += kAbortCheckBatchSize) { size_t limit = std::min(keys.size(), offset + kAbortCheckBatchSize); @@ -2061,16 +2065,13 @@ int collect_deferred_abort_tasks(TxnKv* txn_kv, const std::string& instance_id, << " key=" << hex(key); return -1; } - T rowset_meta_pb; + RowsetMetaCloudPB rowset_meta_pb; if (!rowset_meta_pb.ParseFromString(val)) { LOG(WARNING) << "failed to parse rowset meta, instance_id=" << instance_id << " key=" << hex(key); return -1; } - if (skip_base_version && rowset_meta(rowset_meta_pb).end_version() == 1) { - continue; - } - if (auto abort_task = make_deferred_abort_task(rowset_meta_pb); + if (auto abort_task = make_deferred_abort_task(rowset_meta(rowset_meta_pb)); abort_task.has_value()) { abort_tasks->emplace_back(std::move(*abort_task)); } @@ -2079,12 +2080,9 @@ int collect_deferred_abort_tasks(TxnKv* txn_kv, const std::string& instance_id, return 0; } -template -int InstanceRecycler::batch_abort_txn_or_job_for_recycle(const std::vector& keys, - bool skip_base_version) { +int InstanceRecycler::batch_abort_txn_or_job_for_recycle(const std::vector& keys) { std::vector abort_tasks; - if (collect_deferred_abort_tasks(txn_kv_.get(), instance_id_, keys, &abort_tasks, - skip_base_version) != 0) { + if (collect_deferred_abort_tasks(txn_kv_.get(), instance_id_, keys, &abort_tasks) != 0) { LOG(WARNING) << "failed to collect rowset abort tasks, instance_id=" << instance_id_; return -1; } @@ -2112,48 +2110,6 @@ int InstanceRecycler::batch_abort_txn_or_job_for_recycle(const std::vector& keys, - std::vector* delete_tasks) { - constexpr size_t kPrepareCheckBatchSize = 256; - for (size_t offset = 0; offset < keys.size(); offset += kPrepareCheckBatchSize) { - size_t limit = std::min(keys.size(), offset + kPrepareCheckBatchSize); - std::unique_ptr txn; - TxnErrorCode err = txn_kv->create_txn(&txn); - if (err != TxnErrorCode::TXN_OK) { - LOG(WARNING) << "failed to create txn, instance_id=" << instance_id; - return -1; - } - for (size_t idx = offset; idx < limit; ++idx) { - const std::string& key = keys[idx]; - std::string val; - err = txn->get(key, &val); - if (err == TxnErrorCode::TXN_KEY_NOT_FOUND) { - // has already been removed - continue; - } - if (err != TxnErrorCode::TXN_OK) { - LOG(WARNING) << "failed to get recycle rowset, instance_id=" << instance_id - << " key=" << hex(key); - return -1; - } - RecycleRowsetPB rowset; - if (!rowset.ParseFromString(val)) { - LOG(WARNING) << "failed to parse recycle rowset, instance_id=" << instance_id - << " key=" << hex(key); - return -1; - } - if (rowset.type() != RecycleRowsetPB::PREPARE) { - continue; - } - const auto& rs_meta = rowset.rowset_meta(); - delete_tasks->push_back( - {key, rs_meta.resource_id(), rs_meta.rowset_id_v2(), rs_meta.tablet_id()}); - } - } - return 0; -} - int InstanceRecycler::recycle_ref_rowsets(bool* has_unrecycled_rowsets) { const std::string task_name = "recycle_ref_rowsets"; *has_unrecycled_rowsets = false; @@ -5000,7 +4956,6 @@ int InstanceRecycler::recycle_rowsets() { }; std::vector rowset_keys_to_mark_recycled; - std::vector rowset_keys_to_abort; std::vector prepare_rowset_keys_to_delete; // Store keys of rowset recycled by background workers @@ -5010,9 +4965,6 @@ int InstanceRecycler::recycle_rowsets() { auto worker_pool = std::make_unique( config::instance_recycler_worker_pool_size, "recycle_rowsets"); worker_pool->start(); - auto mark_abort_worker_pool = std::make_unique( - config::instance_recycler_worker_pool_size, "recycle_rs_mark_abort"); - mark_abort_worker_pool->start(); // TODO bacth delete auto delete_versioned_delete_bitmap_kvs = [&](int64_t tablet_id, const std::string& rowset_id) { std::string dbm_start_key = @@ -5132,31 +5084,7 @@ int InstanceRecycler::recycle_rowsets() { segment_metrics_context_.total_recycled_num += rowset.rowset_meta().num_segments(); return 0; } - auto* rowset_meta = rowset.mutable_rowset_meta(); - if (config::enable_mark_delete_rowset_before_recycle) { - if (need_mark_rowset_as_recycled(rowset)) { - rowset_keys_to_mark_recycled.emplace_back(k); - LOG(INFO) << "rowset queued to mark as recycled, recycler will delete data and kv " - "at next turn, instance_id=" - << instance_id_ << " tablet_id=" << rowset_meta->tablet_id() - << " version=[" << rowset_meta->start_version() << '-' - << rowset_meta->end_version() << "]"; - return 0; - } - } - - if (config::enable_abort_txn_and_job_for_delete_rowset_before_recycle && - rowset_meta->end_version() != 1) { - if (make_deferred_abort_task(rowset).has_value()) { - LOG(INFO) << "rowset queued to abort related txn or job after current scan batch, " - "instance_id=" - << instance_id_ << " tablet_id=" << rowset_meta->tablet_id() - << " version=[" << rowset_meta->start_version() << '-' - << rowset_meta->end_version() << "]"; - rowset_keys_to_abort.emplace_back(k); - } - } // TODO(plat1ko): check rowset not referenced if (!rowset_meta->has_resource_id()) [[unlikely]] { // impossible @@ -5178,6 +5106,36 @@ int InstanceRecycler::recycle_rowsets() { << " creation_time=" << rowset_meta->creation_time() << " task_type=" << metrics_context.operation_type; if (rowset.type() == RecycleRowsetPB::PREPARE) { + if (config::enable_mark_delete_rowset_before_recycle) { + if (need_mark_rowset_as_recycled(rowset.rowset_meta())) { + rowset_keys_to_mark_recycled.emplace_back(k); + LOG(INFO) << "rowset queued to mark as recycled, recycler will delete data and " + "kv " + "at next turn, instance_id=" + << instance_id_ << " tablet_id=" << rowset_meta->tablet_id() + << " version=[" << rowset_meta->start_version() << '-' + << rowset_meta->end_version() << "]"; + return 0; + } + } + + if (config::enable_abort_txn_and_job_for_delete_rowset_before_recycle && + rowset_meta->end_version() != 1) { + int ret = 0; + if (rowset_meta->has_load_id()) { + DCHECK(rowset_meta->has_txn_id() && rowset_meta->txn_id() > 0); + ret = abort_txn_for_related_rowset(rowset_meta->txn_id()); + } else if (rowset_meta->has_job_id()) { + ret = abort_job_for_related_rowset(*rowset_meta); + } + if (ret != 0) { + LOG(WARNING) << "failed to abort txn or job for related rowset, instance_id=" + << instance_id_ << " tablet_id=" << rowset_meta->tablet_id() + << " version=[" << rowset_meta->start_version() << '-' + << rowset_meta->end_version() << "]"; + return -1; + } + } // unable to calculate file path, can only be deleted by rowset id prefix num_prepare += 1; if (delete_rowset_data_by_prefix(std::string(k), rowset_meta->resource_id(), @@ -5221,21 +5179,17 @@ int InstanceRecycler::recycle_rowsets() { }); }; - auto submit_mark_abort_rowset_job = [&](std::vector rowset_keys_to_mark, - std::vector rowset_keys_to_abort) { - if (rowset_keys_to_mark.empty() && rowset_keys_to_abort.empty()) { + auto submit_mark_rowset_job = [&](std::vector rowset_keys_to_mark) { + if (rowset_keys_to_mark.empty()) { return; } - mark_abort_worker_pool->submit([&, rowset_keys_to_mark = std::move(rowset_keys_to_mark), - rowset_keys_to_abort = - std::move(rowset_keys_to_abort)]() mutable { + worker_pool->submit([&, rowset_keys_to_mark = std::move(rowset_keys_to_mark)]() mutable { auto start = steady_clock::now(); DORIS_CLOUD_DEFER { auto cost = duration_cast(steady_clock::now() - start).count(); LOG(INFO) << "finish mark and abort rowset job, instance_id=" << instance_id_ << ' ' << "cost_ms=" << cost << ' ' - << "rowset_keys_to_mark.size()=" << rowset_keys_to_mark.size() << ' ' - << "rowset_keys_to_abort.size()=" << rowset_keys_to_abort.size(); + << "rowset_keys_to_mark.size()=" << rowset_keys_to_mark.size(); }; if (!rowset_keys_to_mark.empty() && batch_mark_rowsets_as_recycled(txn_kv_.get(), instance_id_, @@ -5245,26 +5199,14 @@ int InstanceRecycler::recycle_rowsets() { << "rowset_keys_to_mark.size()=" << rowset_keys_to_mark.size(); return; } - if (!rowset_keys_to_abort.empty() && - batch_abort_txn_or_job_for_recycle(rowset_keys_to_abort, true) != - 0) { - LOG(WARNING) << "failed to batch abort txn or job for related rowset, " - "instance_id=" - << instance_id_ << ' ' - << "rowset_keys_to_abort.size()=" << rowset_keys_to_abort.size(); - return; - } }); }; bool scan_finished = false; auto loop_done = [&]() -> int { std::vector mark_keys_to_process; - std::vector abort_keys_to_process; mark_keys_to_process.swap(rowset_keys_to_mark_recycled); - abort_keys_to_process.swap(rowset_keys_to_abort); - submit_mark_abort_rowset_job(std::move(mark_keys_to_process), - std::move(abort_keys_to_process)); + submit_mark_rowset_job(std::move(mark_keys_to_process)); if (!scan_finished && rowsets.size() < delete_rowset_batch_size) { return 0; } @@ -5332,7 +5274,6 @@ int InstanceRecycler::recycle_rowsets() { ret = -1; } - mark_abort_worker_pool->stop(); worker_pool->stop(); if (!async_recycled_rowset_keys.empty()) { @@ -6085,8 +6026,7 @@ int InstanceRecycler::recycle_tmp_rowsets() { return; } if (!abort_keys_to_process.empty() && - batch_abort_txn_or_job_for_recycle(abort_keys_to_process, - false) != 0) { + batch_abort_txn_or_job_for_recycle(abort_keys_to_process) != 0) { LOG(WARNING) << "failed to batch abort txn or job for releated rowset, instance_id=" << instance_id_; return; @@ -7351,7 +7291,9 @@ int InstanceRecycler::scan_and_statistics_rowsets() { return 0; } - if(!rowset_meta->has_is_recycled() || !rowset_meta->is_recycled()) { + if (config::enable_mark_delete_rowset_before_recycle && + rowset.type() == RecycleRowsetPB::PREPARE && + (!rowset_meta->has_is_recycled() || !rowset_meta->is_recycled())) { return 0; } diff --git a/cloud/src/recycler/recycler.h b/cloud/src/recycler/recycler.h index 2e8b6d57f7f3ae..b3ce7e90d2a99c 100644 --- a/cloud/src/recycler/recycler.h +++ b/cloud/src/recycler/recycler.h @@ -602,9 +602,7 @@ class InstanceRecycler { int abort_txn_for_related_rowset(int64_t txn_id); int abort_job_for_related_rowset(const RowsetMetaCloudPB& rowset_meta); - template - int batch_abort_txn_or_job_for_recycle(const std::vector& keys, - bool skip_base_version); + int batch_abort_txn_or_job_for_recycle(const std::vector& keys); private: std::atomic_bool stopped_ {false}; diff --git a/cloud/test/recycler_test.cpp b/cloud/test/recycler_test.cpp index ad8075cec9c55f..c480227e9c1b36 100644 --- a/cloud/test/recycler_test.cpp +++ b/cloud/test/recycler_test.cpp @@ -1430,6 +1430,69 @@ TEST(RecyclerTest, next_recycle_rowset_tablet_key_overwrites_existing_buffer) { EXPECT_TRUE(std::get(std::get<0>(out[4])).empty()); } +TEST(RecyclerTest, recycle_rowsets_only_marks_prepare_rowsets_as_recycled) { + config::retention_seconds = 0; + auto old_enable_mark = config::enable_mark_delete_rowset_before_recycle; + config::enable_mark_delete_rowset_before_recycle = true; + DORIS_CLOUD_DEFER { + config::enable_mark_delete_rowset_before_recycle = old_enable_mark; + }; + + auto txn_kv = std::make_shared(); + ASSERT_EQ(txn_kv->init(), 0); + + InstanceInfoPB instance; + instance.set_instance_id(instance_id); + auto obj_info = instance.add_obj_info(); + obj_info->set_id("recycle_rowsets_only_marks_prepare"); + obj_info->set_ak(config::test_s3_ak); + obj_info->set_sk(config::test_s3_sk); + obj_info->set_endpoint(config::test_s3_endpoint); + obj_info->set_region(config::test_s3_region); + obj_info->set_bucket(config::test_s3_bucket); + obj_info->set_prefix("recycle_rowsets_only_marks_prepare"); + + InstanceRecycler recycler(txn_kv, instance, thread_group, + std::make_shared(txn_kv)); + ASSERT_EQ(recycler.init(), 0); + auto accessor = recycler.accessor_map_.begin()->second; + + doris::TabletSchemaCloudPB schema; + schema.set_schema_version(1); + constexpr int64_t index_id = 10001; + constexpr int64_t tablet_id = 10002; + auto prepare_rowset = + create_rowset("recycle_rowsets_only_marks_prepare", tablet_id, index_id, 1, schema); + auto compact_rowset = + create_rowset("recycle_rowsets_only_marks_prepare", tablet_id, index_id, 1, schema); + ASSERT_EQ(create_recycle_rowset(txn_kv.get(), accessor.get(), prepare_rowset, + RecycleRowsetPB::PREPARE, true), + 0); + ASSERT_EQ(create_recycle_rowset(txn_kv.get(), accessor.get(), compact_rowset, + RecycleRowsetPB::COMPACT, true), + 0); + + std::atomic marked_prepare_rowset_count = 0; + std::atomic marked_non_prepare_rowset_count = 0; + auto sp = SyncPoint::get_instance(); + DORIS_CLOUD_DEFER { + sp->clear_all_call_backs(); + }; + sp->set_call_back("InstanceRecycler::batch_mark_rowsets_as_recycled", [&](auto&& args) { + auto* type = try_any_cast(args[0]); + if (*type == RecycleRowsetPB::PREPARE) { + ++marked_prepare_rowset_count; + } else { + ++marked_non_prepare_rowset_count; + } + }); + sp->enable_processing(); + + ASSERT_EQ(recycler.recycle_rowsets(), 0); + EXPECT_EQ(marked_prepare_rowset_count.load(), 1); + EXPECT_EQ(marked_non_prepare_rowset_count.load(), 0); +} + TEST(RecyclerTest, recycle_rowsets_tablet_batch_limit_recycles_remaining_in_next_round) { config::retention_seconds = 0; auto txn_kv = std::make_shared(); @@ -1648,17 +1711,14 @@ TEST(RecyclerTest, recycle_rowsets_limit_per_tablet_batch) { return it->size(); }; - ASSERT_EQ(recycler.recycle_rowsets(), 0); ASSERT_EQ(recycler.recycle_rowsets(), 0); EXPECT_EQ(count_recycle_rowsets(tablet_id0), 3); EXPECT_EQ(count_recycle_rowsets(tablet_id1), 3); - ASSERT_EQ(recycler.recycle_rowsets(), 0); ASSERT_EQ(recycler.recycle_rowsets(), 0); EXPECT_EQ(count_recycle_rowsets(tablet_id0), 1); EXPECT_EQ(count_recycle_rowsets(tablet_id1), 1); - ASSERT_EQ(recycler.recycle_rowsets(), 0); ASSERT_EQ(recycler.recycle_rowsets(), 0); EXPECT_EQ(count_recycle_rowsets(tablet_id0), 0); EXPECT_EQ(count_recycle_rowsets(tablet_id1), 0); @@ -9150,6 +9210,172 @@ TEST(RecyclerTest, abort_job_for_related_rowset_when_tablet_recycled) { << "Should return 0 when tablet is already recycled (parallel recycle scenario)"; } +TEST(RecyclerTest, abort_exact_job_and_reject_expired_job_for_related_rowset) { + auto txn_kv = std::make_shared(); + ASSERT_EQ(txn_kv->init(), 0); + + const std::string test_instance_id = "instance_id_recycle_test"; + InstanceInfoPB instance; + instance.set_instance_id(test_instance_id); + auto obj_info = instance.add_obj_info(); + obj_info->set_id("abort_exact_job_and_reject_expired_job_for_related_rowset"); + 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("abort_exact_job_and_reject_expired_job_for_related_rowset"); + + InstanceRecycler recycler(txn_kv, instance, thread_group, + std::make_shared(txn_kv)); + ASSERT_EQ(recycler.init(), 0); + + constexpr int64_t table_id = 22000; + constexpr int64_t index_id = 22001; + constexpr int64_t partition_id = 22002; + constexpr int64_t tablet_id = 22003; + ASSERT_EQ(create_tablet(txn_kv.get(), table_id, index_id, partition_id, tablet_id), 0); + + TabletIndexPB tablet_idx; + ASSERT_EQ(get_tablet_idx(txn_kv.get(), test_instance_id, tablet_id, tablet_idx), 0); + + TabletJobInfoPB job; + job.mutable_idx()->CopyFrom(tablet_idx); + auto* first_compaction = job.add_compaction(); + first_compaction->set_id("job_a"); + first_compaction->set_expiration(current_time + 1000); + auto* target_compaction = job.add_compaction(); + target_compaction->set_id("job_b"); + target_compaction->set_expiration(current_time + 1000); + + auto job_key = job_tablet_key({test_instance_id, tablet_idx.table_id(), tablet_idx.index_id(), + tablet_idx.partition_id(), tablet_idx.tablet_id()}); + std::unique_ptr txn; + ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK); + txn->put(job_key, job.SerializeAsString()); + ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK); + + RowsetMetaCloudPB rowset_meta; + rowset_meta.set_tablet_id(tablet_id); + rowset_meta.set_job_id("job_b"); + rowset_meta.set_rowset_id_v2("rowset_b"); + ASSERT_EQ(recycler.abort_job_for_related_rowset(rowset_meta), 0); + + ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK); + std::string job_val; + ASSERT_EQ(txn->get(job_key, &job_val), TxnErrorCode::TXN_OK); + TabletJobInfoPB remaining_job; + ASSERT_TRUE(remaining_job.ParseFromString(job_val)); + ASSERT_EQ(remaining_job.compaction_size(), 1); + EXPECT_EQ(remaining_job.compaction(0).id(), "job_a"); + + job.Clear(); + first_compaction = job.add_compaction(); + first_compaction->set_id("job_a"); + first_compaction->set_expiration(current_time + 1000); + target_compaction = job.add_compaction(); + target_compaction->set_id("job_b"); + target_compaction->set_expiration(current_time - 1000); + ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK); + txn->put(job_key, job.SerializeAsString()); + ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK); + + ASSERT_EQ(recycler.abort_job_for_related_rowset(rowset_meta), -1); + ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK); + ASSERT_EQ(txn->get(job_key, &job_val), TxnErrorCode::TXN_OK); + ASSERT_TRUE(remaining_job.ParseFromString(job_val)); + ASSERT_EQ(remaining_job.compaction_size(), 2); + EXPECT_EQ(remaining_job.compaction(0).id(), "job_a"); + EXPECT_EQ(remaining_job.compaction(1).id(), "job_b"); +} + +TEST(RecyclerTest, abort_schema_change_from_new_tablet_rowset) { + auto txn_kv = std::make_shared(); + ASSERT_EQ(txn_kv->init(), 0); + + const std::string test_instance_id = "instance_id_recycle_test"; + InstanceInfoPB instance; + instance.set_instance_id(test_instance_id); + auto* obj_info = instance.add_obj_info(); + obj_info->set_id("abort_schema_change_from_new_tablet_rowset"); + 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("abort_schema_change_from_new_tablet_rowset"); + + InstanceRecycler recycler(txn_kv, instance, thread_group, + std::make_shared(txn_kv)); + ASSERT_EQ(recycler.init(), 0); + + constexpr int64_t table_id = 23000; + constexpr int64_t base_index_id = 23001; + constexpr int64_t new_index_id = 23002; + constexpr int64_t partition_id = 23003; + constexpr int64_t base_tablet_id = 23004; + constexpr int64_t new_tablet_id = 23005; + + TabletIndexPB base_tablet_idx; + base_tablet_idx.set_table_id(table_id); + base_tablet_idx.set_index_id(base_index_id); + base_tablet_idx.set_partition_id(partition_id); + base_tablet_idx.set_tablet_id(base_tablet_id); + + TabletIndexPB new_tablet_idx; + new_tablet_idx.set_table_id(table_id); + new_tablet_idx.set_index_id(new_index_id); + new_tablet_idx.set_partition_id(partition_id); + new_tablet_idx.set_tablet_id(new_tablet_id); + + TabletJobInfoPB job; + job.mutable_idx()->CopyFrom(base_tablet_idx); + auto* schema_change = job.mutable_schema_change(); + schema_change->set_id("schema_change_job"); + schema_change->set_initiator("BE1"); + schema_change->set_expiration(current_time + 1000); + schema_change->mutable_new_tablet_idx()->CopyFrom(new_tablet_idx); + + auto base_job_key = job_tablet_key( + {test_instance_id, table_id, base_index_id, partition_id, base_tablet_id}); + auto new_job_key = + job_tablet_key({test_instance_id, table_id, new_index_id, partition_id, new_tablet_id}); + + doris::TabletMetaCloudPB new_tablet_meta; + new_tablet_meta.set_tablet_id(new_tablet_id); + new_tablet_meta.set_tablet_state(doris::TabletStatePB::PB_NOTREADY); + + std::unique_ptr txn; + ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK); + txn->put(meta_tablet_idx_key({test_instance_id, new_tablet_id}), + new_tablet_idx.SerializeAsString()); + txn->put(meta_tablet_key( + {test_instance_id, table_id, new_index_id, partition_id, new_tablet_id}), + new_tablet_meta.SerializeAsString()); + txn->put(base_job_key, job.SerializeAsString()); + txn->put(new_job_key, job.SerializeAsString()); + ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK); + + RowsetMetaCloudPB rowset_meta; + rowset_meta.set_tablet_id(new_tablet_id); + rowset_meta.set_job_id(schema_change->id()); + rowset_meta.set_rowset_id_v2("new_tablet_output_rowset"); + ASSERT_EQ(recycler.abort_job_for_related_rowset(rowset_meta), 0); + + ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK); + std::string base_job_val; + ASSERT_EQ(txn->get(base_job_key, &base_job_val), TxnErrorCode::TXN_OK); + TabletJobInfoPB base_job; + ASSERT_TRUE(base_job.ParseFromString(base_job_val)); + EXPECT_FALSE(base_job.has_schema_change()); + + std::string new_job_val; + ASSERT_EQ(txn->get(new_job_key, &new_job_val), TxnErrorCode::TXN_OK); + TabletJobInfoPB new_job; + ASSERT_TRUE(new_job.ParseFromString(new_job_val)); + EXPECT_FALSE(new_job.has_schema_change()); +} + TEST(RecyclerTest, recycle_tablet_with_delete_file_failure) { // If object deletion fails, recycle_tablets should report failure and keep the tablet KV.