diff --git a/be/src/cloud/cloud_cumulative_compaction.cpp b/be/src/cloud/cloud_cumulative_compaction.cpp index d136944ae4b4ea..bf1adc81f18f27 100644 --- a/be/src/cloud/cloud_cumulative_compaction.cpp +++ b/be/src/cloud/cloud_cumulative_compaction.cpp @@ -409,7 +409,8 @@ Status CloudCumulativeCompaction::modify_rowsets() { } } // agg delete bitmap for pre rowsets - if (config::enable_agg_and_remove_pre_rowsets_delete_bitmap && + if (config::delete_bitmap_store_write_version != 2 && + config::enable_agg_and_remove_pre_rowsets_delete_bitmap && _tablet->keys_type() == KeysType::UNIQUE_KEYS && _tablet->enable_unique_key_merge_on_write() && _input_rowsets.size() != 1) { OlapStopWatch watch; @@ -425,17 +426,23 @@ Status CloudCumulativeCompaction::modify_rowsets() { std::sort(pre_rowsets.begin(), pre_rowsets.end(), Rowset::comparator); auto pre_rowsets_delete_bitmap = std::make_shared(_tablet->tablet_id()); std::map pre_rowset_to_versions; + std::unique_ptr pre_rowset_delete_bitmap_stats; + if (config::enable_remove_pre_rowsets_delete_bitmap_by_keys) { + pre_rowset_delete_bitmap_stats = + std::make_unique(); + } cloud_tablet()->agg_delete_bitmap_for_compaction( _output_rowset->start_version(), _output_rowset->end_version(), pre_rowsets, - pre_rowsets_delete_bitmap, pre_rowset_to_versions); + pre_rowsets_delete_bitmap, pre_rowset_to_versions, + pre_rowset_delete_bitmap_stats.get()); // update delete bitmap to ms DBUG_EXECUTE_IF( "CumulativeCompaction.modify_rowsets.cloud_update_delete_bitmap_without_lock.block", DBUG_BLOCK); auto status = _engine.meta_mgr().cloud_update_delete_bitmap_without_lock( *cloud_tablet(), pre_rowsets_delete_bitmap.get(), pre_rowset_to_versions, - cloud_tablet()->table_id(), _output_rowset->start_version(), - _output_rowset->end_version()); + pre_rowset_delete_bitmap_stats.get(), cloud_tablet()->table_id(), + _output_rowset->start_version(), _output_rowset->end_version()); if (!status.ok()) { LOG(WARNING) << "failed to agg pre rowsets delete bitmap to ms. tablet_id=" << _tablet->tablet_id() << ", pre rowset num=" << pre_rowsets.size() diff --git a/be/src/cloud/cloud_meta_mgr.cpp b/be/src/cloud/cloud_meta_mgr.cpp index dd1290f7024a54..fac5c8a73a7f59 100644 --- a/be/src/cloud/cloud_meta_mgr.cpp +++ b/be/src/cloud/cloud_meta_mgr.cpp @@ -2126,8 +2126,10 @@ Status CloudMetaMgr::update_delete_bitmap(const CloudTablet& tablet, int64_t loc Status CloudMetaMgr::cloud_update_delete_bitmap_without_lock( const CloudTablet& tablet, DeleteBitmap* delete_bitmap, - std::map& rowset_to_versions, int64_t table_id, - int64_t pre_rowset_agg_start_version, int64_t pre_rowset_agg_end_version) { + std::map& rowset_to_versions, + const CloudTablet::PreRowsetDeleteBitmapStats* pre_rowset_delete_bitmap_stats, + int64_t table_id, int64_t pre_rowset_agg_start_version, + int64_t pre_rowset_agg_end_version) { if (config::delete_bitmap_store_write_version == 2) { VLOG_DEBUG << "no need to agg delete bitmap v1 in ms because use v2"; return Status::OK(); @@ -2166,6 +2168,20 @@ Status CloudMetaMgr::cloud_update_delete_bitmap_without_lock( req.set_pre_rowset_agg_start_version(pre_rowset_agg_start_version); req.set_pre_rowset_agg_end_version(pre_rowset_agg_end_version); } + if (pre_rowset_delete_bitmap_stats != nullptr) { + for (const auto& [rowset_id, delete_bitmap_stats] : *pre_rowset_delete_bitmap_stats) { + auto* rowset_stats_pb = req.add_pre_rowset_delete_bitmap_stats(); + rowset_stats_pb->set_rowset_id(rowset_id); + for (const auto& [segment_id, version, delete_bitmap_size] : delete_bitmap_stats) { + auto* delete_bitmap_stat_pb = rowset_stats_pb->add_delete_bitmap_stats(); + delete_bitmap_stat_pb->set_segment_id(segment_id); + delete_bitmap_stat_pb->set_version(version); + delete_bitmap_stat_pb->set_delete_bitmap_size(delete_bitmap_size); + } + } + } + TEST_SYNC_POINT_RETURN_WITH_VALUE( + "CloudMetaMgr::cloud_update_delete_bitmap_without_lock.before_rpc", Status::OK(), &req); return retry_rpc(MetaServiceRPC::UPDATE_DELETE_BITMAP, req, &res, &MetaService_Stub::update_delete_bitmap, { diff --git a/be/src/cloud/cloud_meta_mgr.h b/be/src/cloud/cloud_meta_mgr.h index 9084df3d6ea42a..71144a0a7e40b9 100644 --- a/be/src/cloud/cloud_meta_mgr.h +++ b/be/src/cloud/cloud_meta_mgr.h @@ -153,8 +153,10 @@ class CloudMetaMgr { Status cloud_update_delete_bitmap_without_lock( const CloudTablet& tablet, DeleteBitmap* delete_bitmap, - std::map& rowset_to_versions, int64_t table_id, - int64_t pre_rowset_agg_start_version = 0, int64_t pre_rowset_agg_end_version = 0); + std::map& rowset_to_versions, + const CloudTablet::PreRowsetDeleteBitmapStats* pre_rowset_delete_bitmap_stats, + int64_t table_id, int64_t pre_rowset_agg_start_version = 0, + int64_t pre_rowset_agg_end_version = 0); Status get_delete_bitmap_update_lock(const CloudTablet& tablet, int64_t lock_id, int64_t initiator); diff --git a/be/src/cloud/cloud_tablet.cpp b/be/src/cloud/cloud_tablet.cpp index 7827c438bf9f54..ff883442bd182f 100644 --- a/be/src/cloud/cloud_tablet.cpp +++ b/be/src/cloud/cloud_tablet.cpp @@ -1432,12 +1432,34 @@ Status CloudTablet::calc_delete_bitmap_for_compaction( void CloudTablet::agg_delete_bitmap_for_compaction( int64_t start_version, int64_t end_version, const std::vector& pre_rowsets, - DeleteBitmapPtr& new_delete_bitmap, - std::map& pre_rowset_to_versions) { - for (auto& rowset : pre_rowsets) { + DeleteBitmapPtr& new_delete_bitmap, std::map& pre_rowset_to_versions, + PreRowsetDeleteBitmapStats* pre_rowset_delete_bitmap_stats) { + auto& delete_bitmap = tablet_meta()->delete_bitmap(); + for (const auto& rowset : pre_rowsets) { + if (pre_rowset_delete_bitmap_stats != nullptr) { + auto& rowset_delete_bitmap_stats = + (*pre_rowset_delete_bitmap_stats)[rowset->rowset_id().to_string()]; + std::shared_lock lock(delete_bitmap.lock); + const auto bitmap_start_version = static_cast(start_version); + const auto bitmap_end_version = static_cast(end_version); + for (uint32_t seg_id = 0; seg_id < rowset->num_segments(); ++seg_id) { + DeleteBitmap::BitmapKey segment_start {rowset->rowset_id(), seg_id, + bitmap_start_version}; + for (auto it = delete_bitmap.delete_bitmap.lower_bound(segment_start); + it != delete_bitmap.delete_bitmap.end(); ++it) { + const auto& [key, bitmap] = *it; + if (std::get<0>(key) != rowset->rowset_id() || std::get<1>(key) != seg_id || + std::get<2>(key) >= bitmap_end_version) { + break; + } + rowset_delete_bitmap_stats.emplace_back(seg_id, std::get<2>(key), + bitmap.getSizeInBytes()); + } + } + } for (uint32_t seg_id = 0; seg_id < rowset->num_segments(); ++seg_id) { - auto d = tablet_meta()->delete_bitmap().get_agg_without_cache( - {rowset->rowset_id(), seg_id, end_version}, start_version); + auto d = delete_bitmap.get_agg_without_cache({rowset->rowset_id(), seg_id, end_version}, + start_version); if (d->isEmpty()) { continue; } diff --git a/be/src/cloud/cloud_tablet.h b/be/src/cloud/cloud_tablet.h index 7fc3dec14512d2..3718184b915ac3 100644 --- a/be/src/cloud/cloud_tablet.h +++ b/be/src/cloud/cloud_tablet.h @@ -18,6 +18,7 @@ #pragma once #include +#include #include "storage/partial_update_info.h" #include "storage/rowset/rowset.h" @@ -76,6 +77,11 @@ struct RecycledRowsets { class CloudTablet final : public BaseTablet { public: + // rowset id -> [(segment id, version, serialized delete bitmap size)] + using PreRowsetDeleteBitmapStats = std::map< + std::string, + std::vector>>; + CloudTablet(CloudStorageEngine& engine, TabletMetaSharedPtr tablet_meta); ~CloudTablet() override; @@ -357,10 +363,11 @@ class CloudTablet final : public BaseTablet { // check that if the delete bitmap in delete bitmap cache has the same cardinality with the expected_delete_bitmap's Status check_delete_bitmap_cache(int64_t txn_id, DeleteBitmap* expected_delete_bitmap) override; - void agg_delete_bitmap_for_compaction(int64_t start_version, int64_t end_version, - const std::vector& pre_rowsets, - DeleteBitmapPtr& new_delete_bitmap, - std::map& pre_rowset_to_versions); + void agg_delete_bitmap_for_compaction( + int64_t start_version, int64_t end_version, + const std::vector& pre_rowsets, DeleteBitmapPtr& new_delete_bitmap, + std::map& pre_rowset_to_versions, + PreRowsetDeleteBitmapStats* pre_rowset_delete_bitmap_stats); bool need_remove_unused_rowsets(); diff --git a/be/src/common/config.cpp b/be/src/common/config.cpp index f365dab488a50e..0e2379862c6d2d 100644 --- a/be/src/common/config.cpp +++ b/be/src/common/config.cpp @@ -1431,6 +1431,8 @@ DEFINE_mBool(enable_mow_get_agg_by_cache, "true"); DEFINE_mBool(enable_mow_get_agg_correctness_check_core, "false"); DEFINE_mBool(enable_agg_and_remove_pre_rowsets_delete_bitmap, "true"); DEFINE_mBool(enable_check_agg_and_remove_pre_rowsets_delete_bitmap, "false"); +// Send pre-rowset delete bitmap stats so Meta Service removes the split keys individually. +DEFINE_mBool(enable_remove_pre_rowsets_delete_bitmap_by_keys, "true"); // The secure path with user files, used in the `local` table function. DEFINE_String(user_files_secure_path, "${DORIS_HOME}"); diff --git a/be/src/common/config.h b/be/src/common/config.h index 699bbdbee47a50..d7432a585fa1a4 100644 --- a/be/src/common/config.h +++ b/be/src/common/config.h @@ -1523,6 +1523,7 @@ DECLARE_mBool(enable_mow_get_agg_by_cache); DECLARE_mBool(enable_mow_get_agg_correctness_check_core); DECLARE_mBool(enable_agg_and_remove_pre_rowsets_delete_bitmap); DECLARE_mBool(enable_check_agg_and_remove_pre_rowsets_delete_bitmap); +DECLARE_mBool(enable_remove_pre_rowsets_delete_bitmap_by_keys); // The secure path with user files, used in the `local` table function. DECLARE_String(user_files_secure_path); diff --git a/be/test/cloud/cloud_meta_mgr_test.cpp b/be/test/cloud/cloud_meta_mgr_test.cpp index ff87378348ef10..c3da6ec9925d20 100644 --- a/be/test/cloud/cloud_meta_mgr_test.cpp +++ b/be/test/cloud/cloud_meta_mgr_test.cpp @@ -29,12 +29,14 @@ #include "cloud/cloud_storage_engine.h" #include "cloud/cloud_tablet.h" +#include "cloud/config.h" #include "cpp/sync_point.h" #include "load/stream_load/stream_load_context.h" #include "storage/olap_common.h" #include "storage/rowset/rowset_factory.h" #include "storage/rowset/rowset_meta.h" #include "storage/tablet/tablet_meta.h" +#include "util/defer_op.h" #include "util/uid_util.h" namespace doris { @@ -67,6 +69,79 @@ TEST_F(CloudMetaMgrTest, response_status_falls_back_for_invalid_actual_code) { EXPECT_EQ(get_response_code(status), MetaServiceCode::KV_TXN_CONFLICT); } +TEST_F(CloudMetaMgrTest, PreRowsetDeleteBitmapStatsRequestEncoding) { + CloudStorageEngine engine(EngineOptions {}); + CloudMetaMgr meta_mgr; + TabletMetaSharedPtr tablet_meta( + new TabletMeta(1001, 2, 15673, 15674, 4, 5, TTabletSchema(), 6, {{7, 8}}, + UniqueId(9, 10), TTabletType::TABLET_TYPE_DISK, TCompressionType::LZ4F)); + auto tablet = std::make_shared(engine, std::make_shared(*tablet_meta)); + DeleteBitmap delete_bitmap(tablet->tablet_id()); + std::map rowset_to_versions; + + int32_t old_write_version = config::delete_bitmap_store_write_version; + config::delete_bitmap_store_write_version = 1; + auto* sp = SyncPoint::get_instance(); + sp->clear_all_call_backs(); + sp->enable_processing(); + Defer cleanup {[&] { + config::delete_bitmap_store_write_version = old_write_version; + sp->disable_processing(); + sp->clear_all_call_backs(); + }}; + + auto capture_request = [&](const CloudTablet::PreRowsetDeleteBitmapStats* stats) { + bool called = false; + UpdateDeleteBitmapRequest captured_req; + SyncPoint::CallbackGuard guard; + sp->set_call_back( + "CloudMetaMgr::cloud_update_delete_bitmap_without_lock.before_rpc", + [&](auto&& args) { + auto* req = try_any_cast(args[0]); + captured_req.CopyFrom(*req); + called = true; + auto* ret = try_any_cast*>(args.back()); + ret->first = Status::OK(); + ret->second = true; + }, + &guard); + auto status = meta_mgr.cloud_update_delete_bitmap_without_lock( + *tablet, &delete_bitmap, rowset_to_versions, stats, tablet->table_id(), 1, 2); + EXPECT_TRUE(status.ok()) << status; + EXPECT_TRUE(called); + return captured_req; + }; + + auto config_disabled_req = capture_request(nullptr); + EXPECT_EQ(config_disabled_req.pre_rowset_delete_bitmap_stats_size(), 0); + + CloudTablet::PreRowsetDeleteBitmapStats empty_stats; + empty_stats.emplace( + "rowset_without_delete_bitmap", + std::vector> {}); + auto config_enabled_req = capture_request(&empty_stats); + ASSERT_EQ(config_enabled_req.pre_rowset_delete_bitmap_stats_size(), 1); + EXPECT_EQ(config_enabled_req.pre_rowset_delete_bitmap_stats(0).rowset_id(), + "rowset_without_delete_bitmap"); + EXPECT_EQ(config_enabled_req.pre_rowset_delete_bitmap_stats(0).delete_bitmap_stats_size(), 0); + + using DeleteBitmapStat = std::tuple; + CloudTablet::PreRowsetDeleteBitmapStats populated_stats; + populated_stats.emplace("rowset_with_delete_bitmap", + std::vector {{3, 7, 1024}, {8, 9, 2048}}); + auto populated_stats_req = capture_request(&populated_stats); + ASSERT_EQ(populated_stats_req.pre_rowset_delete_bitmap_stats_size(), 1); + const auto& rowset_stats = populated_stats_req.pre_rowset_delete_bitmap_stats(0); + EXPECT_EQ(rowset_stats.rowset_id(), "rowset_with_delete_bitmap"); + ASSERT_EQ(rowset_stats.delete_bitmap_stats_size(), 2); + EXPECT_EQ(rowset_stats.delete_bitmap_stats(0).segment_id(), 3); + EXPECT_EQ(rowset_stats.delete_bitmap_stats(0).version(), 7); + EXPECT_EQ(rowset_stats.delete_bitmap_stats(0).delete_bitmap_size(), 1024); + EXPECT_EQ(rowset_stats.delete_bitmap_stats(1).segment_id(), 8); + EXPECT_EQ(rowset_stats.delete_bitmap_stats(1).version(), 9); + EXPECT_EQ(rowset_stats.delete_bitmap_stats(1).delete_bitmap_size(), 2048); +} + static AbortTxnRequest get_abort_txn_request(CloudMetaMgr* meta_mgr, const StreamLoadContext& ctx) { auto* sp = SyncPoint::get_instance(); sp->clear_all_call_backs(); diff --git a/be/test/cloud/cloud_tablet_test.cpp b/be/test/cloud/cloud_tablet_test.cpp index 4c1bfbf79a3252..7b30b7541c0b15 100644 --- a/be/test/cloud/cloud_tablet_test.cpp +++ b/be/test/cloud/cloud_tablet_test.cpp @@ -74,6 +74,79 @@ class CloudTabletWarmUpStateTest : public testing::Test { CloudStorageEngine _engine; }; +class CloudTabletDeleteBitmapTest : public CloudTabletWarmUpStateTest {}; + +TEST_F(CloudTabletDeleteBitmapTest, AggDeleteBitmapForCompactionReturnsPreRowsetStats) { + auto rowset1 = create_rowset(Version(1, 1), 2); + auto rowset2 = create_rowset(Version(2, 2)); + auto rowset_without_delete_bitmap = create_rowset(Version(3, 3)); + ASSERT_NE(rowset1, nullptr); + ASSERT_NE(rowset2, nullptr); + ASSERT_NE(rowset_without_delete_bitmap, nullptr); + + roaring::Roaring before_range; + before_range.add(1); + roaring::Roaring at_start; + at_start.add(2); + at_start.add(3); + roaring::Roaring within_range; + within_range.add(4); + roaring::Roaring at_end; + at_end.add(5); + roaring::Roaring second_segment; + second_segment.add(6); + second_segment.add(7); + second_segment.add(8); + roaring::Roaring second_rowset; + second_rowset.add(9); + + auto& delete_bitmap = _tablet->tablet_meta()->delete_bitmap(); + delete_bitmap.set({rowset1->rowset_id(), 0, 4}, before_range); + delete_bitmap.set({rowset1->rowset_id(), 0, 5}, at_start); + delete_bitmap.set({rowset1->rowset_id(), 0, 6}, within_range); + delete_bitmap.set({rowset1->rowset_id(), 0, 7}, at_end); + delete_bitmap.set({rowset1->rowset_id(), 1, 6}, second_segment); + delete_bitmap.set({rowset2->rowset_id(), 0, 5}, second_rowset); + + auto aggregated_delete_bitmap = std::make_shared(_tablet->tablet_id()); + std::map pre_rowset_to_versions; + CloudTablet::PreRowsetDeleteBitmapStats pre_rowset_delete_bitmap_stats; + _tablet->agg_delete_bitmap_for_compaction( + 5, 7, {rowset1, rowset2, rowset_without_delete_bitmap}, aggregated_delete_bitmap, + pre_rowset_to_versions, &pre_rowset_delete_bitmap_stats); + + using DeleteBitmapStat = std::tuple; + EXPECT_EQ(pre_rowset_delete_bitmap_stats.at(rowset1->rowset_id().to_string()), + (std::vector { + {0, 5, at_start.getSizeInBytes()}, + {0, 6, within_range.getSizeInBytes()}, + {1, 6, second_segment.getSizeInBytes()}, + })); + EXPECT_EQ(pre_rowset_delete_bitmap_stats.at(rowset2->rowset_id().to_string()), + (std::vector {{0, 5, second_rowset.getSizeInBytes()}})); + EXPECT_TRUE( + pre_rowset_delete_bitmap_stats.at(rowset_without_delete_bitmap->rowset_id().to_string()) + .empty()); + + roaring::Roaring aggregated; + ASSERT_EQ(aggregated_delete_bitmap->get({rowset1->rowset_id(), 0, 7}, &aggregated), 0); + EXPECT_EQ(aggregated.cardinality(), 4); + ASSERT_EQ(aggregated_delete_bitmap->get({rowset1->rowset_id(), 1, 7}, &aggregated), 0); + EXPECT_EQ(aggregated.cardinality(), 3); + ASSERT_EQ(aggregated_delete_bitmap->get({rowset2->rowset_id(), 0, 7}, &aggregated), 0); + EXPECT_EQ(aggregated.cardinality(), 1); + EXPECT_EQ(pre_rowset_to_versions.at(rowset1->rowset_id().to_string()), 1); + EXPECT_EQ(pre_rowset_to_versions.at(rowset2->rowset_id().to_string()), 2); + + auto aggregated_without_stats = std::make_shared(_tablet->tablet_id()); + std::map rowset_versions_without_stats; + _tablet->agg_delete_bitmap_for_compaction( + 5, 7, {rowset1, rowset2, rowset_without_delete_bitmap}, aggregated_without_stats, + rowset_versions_without_stats, nullptr); + EXPECT_EQ(aggregated_without_stats->delete_bitmap.size(), 3); + EXPECT_EQ(rowset_versions_without_stats.size(), 2); +} + // Test get_rowset_warmup_state for non-existent rowset TEST_F(CloudTabletWarmUpStateTest, TestGetRowsetWarmupStateNonExistent) { auto rowset = create_rowset(Version(1, 1)); diff --git a/cloud/src/meta-service/meta_service.cpp b/cloud/src/meta-service/meta_service.cpp index 525e6391e0546c..7ca6c31aa4a63d 100644 --- a/cloud/src/meta-service/meta_service.cpp +++ b/cloud/src/meta-service/meta_service.cpp @@ -3941,6 +3941,136 @@ void _write_delete_bitmap_kvs(MetaServiceCode& code, std::string& msg, std::stri << " key_size: " << key.size() << " value_size: " << val.size(); } +static bool commit_pre_rowset_delete_bitmap_removal( + MetaServiceCode& code, std::string& msg, std::stringstream& ss, + const std::shared_ptr& txn_kv, std::unique_ptr& txn, KVStats& stats, + UpdateDeleteBitmapTxnStats& txn_stats, int64_t tablet_id, const std::string& rowset_id) { + auto txn_size = txn->approximate_bytes(); + LOG(INFO) << "commit delete bitmap point deletes before transaction size exceeds limit, " + "tablet_id=" + << tablet_id << ", rowset=" << rowset_id << ", txn_size=" << txn_size; + auto err = txn->commit(); + TEST_SYNC_POINT_CALLBACK("update_delete_bitmap:remove_pre_rowsets:commit", txn_size); + txn_stats.total_txn_put_keys += txn->num_put_keys(); + txn_stats.total_txn_put_bytes += txn->put_bytes(); + txn_stats.total_txn_size += txn_size; + txn_stats.total_txn_count++; + if (err != TxnErrorCode::TXN_OK) { + code = cast_as(err); + ss << "failed to remove pre rowsets delete bitmap, err=" << err + << " tablet_id=" << tablet_id << " rowset_id=" << rowset_id << " txn_size=" << txn_size; + msg = ss.str(); + g_bvar_update_delete_bitmap_fail_counter << 1; + return false; + } + stats.get_bytes += txn->get_bytes(); + stats.put_bytes += txn->put_bytes(); + stats.del_bytes += txn->delete_bytes(); + stats.get_counter += txn->num_get_keys(); + stats.put_counter += txn->num_put_keys(); + stats.del_counter += txn->num_del_keys(); + txn_stats.current_key_count = 0; + txn_stats.current_value_count = 0; + err = txn_kv->create_txn(&txn); + if (err != TxnErrorCode::TXN_OK) { + code = cast_as(err); + msg = "failed to init txn when removing pre rowsets delete bitmap"; + return false; + } + return true; +} + +static bool remove_pre_rowset_delete_bitmap( + MetaServiceCode& code, std::string& msg, std::stringstream& ss, + const std::shared_ptr& txn_kv, std::unique_ptr& txn, KVStats& stats, + const UpdateDeleteBitmapRequest* request, const std::string& instance_id, + const std::set& non_exist_rowset_ids, UpdateDeleteBitmapTxnStats& txn_stats) { + if (!request->has_pre_rowset_agg_start_version() || + !request->has_pre_rowset_agg_end_version() || + request->pre_rowset_agg_start_version() >= request->pre_rowset_agg_end_version()) { + return true; + } + + auto tablet_id = request->tablet_id(); + if (request->pre_rowset_delete_bitmap_stats_size() == 0) { + std::string pre_rowset_id; + for (size_t i = 0; i < request->rowset_ids_size(); ++i) { + if (request->rowset_ids(i) == pre_rowset_id) { + continue; + } + if (non_exist_rowset_ids.contains(request->rowset_ids(i))) { + LOG(INFO) << "skip remove pre rowsets delete bitmap, rowset_id=" + << request->rowset_ids(i) << " tablet_id=" << tablet_id + << " because the rowset does not exist"; + continue; + } + pre_rowset_id = request->rowset_ids(i); + auto delete_bitmap_start = + meta_delete_bitmap_key({instance_id, tablet_id, request->rowset_ids(i), + request->pre_rowset_agg_start_version(), 0}); + auto delete_bitmap_end = + meta_delete_bitmap_key({instance_id, tablet_id, request->rowset_ids(i), + request->pre_rowset_agg_end_version(), 0}); + txn->remove(delete_bitmap_start, delete_bitmap_end); + LOG(INFO) << "remove pre rowsets delete bitmap by range, tablet_id=" << tablet_id + << ", rowset=" << request->rowset_ids(i) + << ", start_version=" << request->pre_rowset_agg_start_version() + << ", end_version=" << request->pre_rowset_agg_end_version() + << ", start_key=" << hex(delete_bitmap_start) + << ", end_key=" << hex(delete_bitmap_end); + } + return true; + } + + for (const auto& rowset_stats : request->pre_rowset_delete_bitmap_stats()) { + if (non_exist_rowset_ids.contains(rowset_stats.rowset_id())) { + LOG(INFO) << "skip remove pre rowsets delete bitmap, rowset_id=" + << rowset_stats.rowset_id() << " tablet_id=" << tablet_id + << " because the rowset does not exist"; + continue; + } + uint64_t delete_key_count = 0; + for (const auto& delete_bitmap_stat : rowset_stats.delete_bitmap_stats()) { + DCHECK(delete_bitmap_stat.version() >= request->pre_rowset_agg_start_version() && + delete_bitmap_stat.version() < request->pre_rowset_agg_end_version()); + auto delete_bitmap_key = meta_delete_bitmap_key( + {instance_id, tablet_id, rowset_stats.rowset_id(), delete_bitmap_stat.version(), + delete_bitmap_stat.segment_id()}); + auto remove_key = [&](std::string_view key) { + auto txn_size = txn->approximate_bytes(); + if (txn_size > 0 && + txn_size + key.size() * 3 > static_cast(config::max_txn_commit_byte) && + !commit_pre_rowset_delete_bitmap_removal(code, msg, ss, txn_kv, txn, stats, + txn_stats, tablet_id, + rowset_stats.rowset_id())) { + return false; + } + txn->remove(key); + delete_key_count++; + return true; + }; + if (!remove_key(delete_bitmap_key)) { + return false; + } + auto split_key_count = + delete_bitmap_stat.delete_bitmap_size() / DEFAULT_BLOB_SPLIT_SIZE + + (delete_bitmap_stat.delete_bitmap_size() % DEFAULT_BLOB_SPLIT_SIZE != 0); + for (size_t i = 0; i < split_key_count; ++i) { + if (!remove_key(blob_key(delete_bitmap_key, i))) { + return false; + } + } + } + if (delete_key_count > 0) { + LOG(INFO) << "remove pre rowsets delete bitmap by keys, tablet_id=" << tablet_id + << ", rowset=" << rowset_stats.rowset_id() + << ", delete_bitmap_count=" << rowset_stats.delete_bitmap_stats_size() + << ", delete_key_count=" << delete_key_count; + } + } + return true; +} + void MetaServiceImpl::update_delete_bitmap(google::protobuf::RpcController* controller, const UpdateDeleteBitmapRequest* request, UpdateDeleteBitmapResponse* response, @@ -4162,35 +4292,9 @@ void MetaServiceImpl::update_delete_bitmap(google::protobuf::RpcController* cont if (code != MetaServiceCode::OK) return; } - // remove pre rowset delete bitmap - if (request->has_pre_rowset_agg_start_version() && request->has_pre_rowset_agg_end_version() && - request->pre_rowset_agg_start_version() < request->pre_rowset_agg_end_version()) { - std::string pre_rowset_id = ""; - for (size_t i = 0; i < request->rowset_ids_size(); ++i) { - if (request->rowset_ids(i) == pre_rowset_id) { - continue; - } - if (non_exist_rowset_ids.contains(request->rowset_ids(i))) { - LOG(INFO) << "skip remove pre rowsets delete bitmap, rowset_id=" - << request->rowset_ids(i) << " tablet_id=" << tablet_id - << " because the rowset does not exist"; - continue; - } - pre_rowset_id = request->rowset_ids(i); - auto delete_bitmap_start = - meta_delete_bitmap_key({instance_id, tablet_id, request->rowset_ids(i), - request->pre_rowset_agg_start_version(), 0}); - auto delete_bitmap_end = - meta_delete_bitmap_key({instance_id, tablet_id, request->rowset_ids(i), - request->pre_rowset_agg_end_version(), 0}); - txn->remove(delete_bitmap_start, delete_bitmap_end); - LOG(INFO) << "remove pre rowsets delete bitmap, tablet_id=" << tablet_id - << ", rowset=" << request->rowset_ids(i) - << ", start_version=" << request->pre_rowset_agg_start_version() - << ", end_version=" << request->pre_rowset_agg_end_version() - << ", start_key=" << hex(delete_bitmap_start) - << ", end_key=" << hex(delete_bitmap_end); - } + if (!remove_pre_rowset_delete_bitmap(code, msg, ss, txn_kv_, txn, stats, request, instance_id, + non_exist_rowset_ids, txn_stats)) { + return; } err = txn->commit(); txn_stats.total_txn_put_keys += txn->num_put_keys(); @@ -4327,6 +4431,9 @@ void MetaServiceImpl::get_delete_bitmap(google::protobuf::RpcController* control std::unique_ptr it; int64_t last_ver = -1; int64_t last_seg_id = -1; + std::string last_delete_bitmap_first_blob_key; + uint16_t next_blob_sequence = 0; + bool skip_last_delete_bitmap = false; int64_t round = 0; while (it == nullptr /* may be not init */ || it->more()) { if (test) { @@ -4348,6 +4455,31 @@ void MetaServiceImpl::get_delete_bitmap(google::protobuf::RpcController* control msg = ss.str(); return; } + if (!skip_last_delete_bitmap && !last_delete_bitmap_first_blob_key.empty()) { + std::string first_blob_value; + err = txn->get(last_delete_bitmap_first_blob_key, &first_blob_value); + if (err == TxnErrorCode::TXN_KEY_NOT_FOUND) { + delete_bitmap_byte -= + response->segment_delete_bitmaps( + response->segment_delete_bitmaps_size() - 1) + .size(); + delete_bitmap_num--; + response->mutable_rowset_ids()->RemoveLast(); + response->mutable_segment_ids()->RemoveLast(); + response->mutable_versions()->RemoveLast(); + response->mutable_segment_delete_bitmaps()->RemoveLast(); + skip_last_delete_bitmap = true; + LOG(WARNING) + << "skip incomplete delete bitmap whose first blob key " + "disappeared after transaction retry" + << ", tablet_id=" << tablet_id + << ", rowset_id=" << rowset_ids[i] << ", version=" << last_ver + << ", segment_id=" << last_seg_id; + } else if (err != TxnErrorCode::TXN_OK) { + retry++; + continue; + } + } if (test) { err = txn->get(start_key, end_key, &it, false, 2); } else { @@ -4373,7 +4505,15 @@ void MetaServiceImpl::get_delete_bitmap(google::protobuf::RpcController* control auto k1 = k; k1.remove_prefix(1); std::vector, int, int>> out; - decode_key(&k1, &out); + auto decode_ret = decode_key(&k1, &out); + if (decode_ret != 0 || (out.size() != 7 && out.size() != 8)) { + code = MetaServiceCode::KV_TXN_GET_ERR; + ss << "invalid delete bitmap key: " << hex(k) + << ", decode_ret=" << decode_ret << ", field_count=" << out.size(); + msg = ss.str(); + g_bvar_get_delete_bitmap_fail_counter << 1; + return; + } // 0x01 "meta" ${instance_id} "delete_bitmap" ${tablet_id} // ${rowset_id0} ${version1} ${segment_id0} -> DeleteBitmapPB auto ver = std::get(std::get<0>(out[5])); @@ -4382,15 +4522,58 @@ void MetaServiceImpl::get_delete_bitmap(google::protobuf::RpcController* control // FIXME: Don't expose the implementation details of splitting large value. // merge splitted large values (>90*1000) if (ver != last_ver || seg_id != last_seg_id) { + auto sequence = + out.size() == 8 + ? static_cast( + std::get(std::get<0>(out[7])) & 0xffff) + : uint16_t {0}; + last_ver = ver; + last_seg_id = seg_id; + last_delete_bitmap_first_blob_key.clear(); + next_blob_sequence = static_cast(sequence + 1); + // Key-based pre-rowset cleanup may leave an obsolete tail if its size estimate is + // too small. The bitmap is already aggregated, so the tail can be skipped and will + // be recycled with its rowset. + skip_last_delete_bitmap = sequence != 0; + if (skip_last_delete_bitmap) { + LOG(WARNING) + << "skip incomplete delete bitmap whose first blob " + "sequence is not zero" + << ", tablet_id=" << tablet_id + << ", rowset_id=" << rowset_ids[i] << ", version=" << ver + << ", segment_id=" << seg_id << ", first_sequence=" << sequence; + continue; + } + if (out.size() == 8) { + last_delete_bitmap_first_blob_key.assign(k.data(), k.size()); + } response->add_rowset_ids(rowset_ids[i]); response->add_segment_ids(seg_id); response->add_versions(ver); response->add_segment_delete_bitmaps(std::string(v)); - last_ver = ver; - last_seg_id = seg_id; delete_bitmap_num++; delete_bitmap_byte += v.length(); } else { + if (skip_last_delete_bitmap) { + continue; + } + auto sequence = + out.size() == 8 + ? static_cast( + std::get(std::get<0>(out[7])) & 0xffff) + : uint16_t {0}; + if (out.size() != 8 || sequence != next_blob_sequence) { + code = MetaServiceCode::KV_TXN_GET_ERR; + ss << "non-contiguous delete bitmap blob sequence, tablet_id=" + << tablet_id << ", rowset_id=" << rowset_ids[i] + << ", version=" << ver << ", segment_id=" << seg_id + << ", expected_sequence=" << next_blob_sequence + << ", actual_sequence=" << sequence; + msg = ss.str(); + g_bvar_get_delete_bitmap_fail_counter << 1; + return; + } + next_blob_sequence++; TEST_SYNC_POINT_CALLBACK("get_delete_bitmap_code", &code); if (code != MetaServiceCode::OK) { ss << "test get get_delete_bitmap fail, code=" diff --git a/cloud/src/meta-store/blob_message.cpp b/cloud/src/meta-store/blob_message.cpp index c7e3edef63dcdb..75c62151928bdd 100644 --- a/cloud/src/meta-store/blob_message.cpp +++ b/cloud/src/meta-store/blob_message.cpp @@ -40,6 +40,15 @@ static std::vector split_string(const std::string_view& str, i return substrings; } +std::string blob_key(std::string_view key, size_t sequence, uint8_t ver) { + std::string split_key(key); + int64_t suffix = ver; + suffix <<= 56; + suffix += static_cast(sequence); + encode_int64(suffix, &split_key); + return split_key; +} + bool ValueBuf::to_pb(google::protobuf::Message* pb) const { butil::IOBuf merge; for (auto&& it : iters) { @@ -151,12 +160,8 @@ void blob_put(Transaction* txn, std::string_view key, const google::protobuf::Me void blob_put(Transaction* txn, std::string_view key, std::string_view value, uint8_t ver, size_t split_size) { auto split_vec = split_string(value, split_size); - int64_t suffix_base = ver; - suffix_base <<= 56; for (size_t i = 0; i < split_vec.size(); ++i) { - std::string k(key); - encode_int64(suffix_base + i, &k); - txn->put(k, split_vec[i]); + txn->put(blob_key(key, i, ver), split_vec[i]); } } @@ -301,12 +306,8 @@ void blob_put(Transaction* txn, std::string_view key, std::string_view value, ui encode_versionstamp_end(&encoded_key); auto split_vec = split_string(value, split_size); - int64_t suffix_base = ver; - suffix_base <<= 56; for (size_t i = 0; i < split_vec.size(); ++i) { - std::string k(encoded_key); - encode_int64(suffix_base + i, &k); - txn->atomic_set_ver_key(k, offset, split_vec[i]); + txn->atomic_set_ver_key(blob_key(encoded_key, i, ver), offset, split_vec[i]); } } diff --git a/cloud/src/meta-store/blob_message.h b/cloud/src/meta-store/blob_message.h index 25d5ad30366dfd..525bc26957ba02 100644 --- a/cloud/src/meta-store/blob_message.h +++ b/cloud/src/meta-store/blob_message.h @@ -17,8 +17,10 @@ #pragma once +#include #include #include +#include #include "meta-store/txn_kv.h" #include "meta-store/txn_kv_error.h" @@ -29,6 +31,11 @@ class Message; namespace doris::cloud { +inline constexpr size_t DEFAULT_BLOB_SPLIT_SIZE = 90 * 1000; + +// Return a split blob key with the encoded version and sequence suffix. +std::string blob_key(std::string_view key, size_t sequence, uint8_t ver = 0); + /** * Supports splitting large values (>100KB) into multiple KVs, with a logical value size of up to the fdb transaction limit (<10MB). * Supports multi version format parsing of values (which can be any byte sequence format), and can recognize the version of values @@ -82,7 +89,7 @@ TxnErrorCode blob_get(Transaction* txn, std::string_view key, ValueBuf* val, boo * @param split_size how many byte sized fragments are the value split into */ void blob_put(Transaction* txn, std::string_view key, const google::protobuf::Message& pb, - uint8_t ver, size_t split_size = 90 * 1000); + uint8_t ver, size_t split_size = DEFAULT_BLOB_SPLIT_SIZE); /** * Put a KV, it's value may be bigger than 100k @@ -93,7 +100,7 @@ void blob_put(Transaction* txn, std::string_view key, const google::protobuf::Me * @param split_size how many byte sized fragments are the value split into */ void blob_put(Transaction* txn, std::string_view key, std::string_view value, uint8_t ver, - size_t split_size = 90 * 1000); + size_t split_size = DEFAULT_BLOB_SPLIT_SIZE); // Iterator for blob key-value pairs. // @@ -189,19 +196,20 @@ namespace versioned { // Put a blob message with a auto generated versionstamp. void blob_put(Transaction* txn, std::string_view key, const google::protobuf::Message& pb, - uint8_t ver = 0, size_t split_size = 90 * 1000); + uint8_t ver = 0, size_t split_size = DEFAULT_BLOB_SPLIT_SIZE); // Put a blob message with a auto generated versionstamp. void blob_put(Transaction* txn, std::string_view key, std::string_view value, uint8_t ver = 0, - size_t split_size = 90 * 1000); + size_t split_size = DEFAULT_BLOB_SPLIT_SIZE); // Put a blob message with a specified versionstamp. void blob_put(Transaction* txn, std::string_view key, Versionstamp v, - const google::protobuf::Message& pb, uint8_t ver = 0, size_t split_size = 90 * 1000); + const google::protobuf::Message& pb, uint8_t ver = 0, + size_t split_size = DEFAULT_BLOB_SPLIT_SIZE); // Put a blob message with a specified versionstamp. void blob_put(Transaction* txn, std::string_view key, Versionstamp v, std::string_view value, - uint8_t ver = 0, size_t split_size = 90 * 1000); + uint8_t ver = 0, size_t split_size = DEFAULT_BLOB_SPLIT_SIZE); } // namespace versioned diff --git a/cloud/test/meta_service_test.cpp b/cloud/test/meta_service_test.cpp index fc7dc03c91e578..28f353cff96bf5 100644 --- a/cloud/test/meta_service_test.cpp +++ b/cloud/test/meta_service_test.cpp @@ -40,6 +40,7 @@ #include "cpp/sync_point.h" #include "meta-service/meta_service_helper.h" #include "meta-store/blob_message.h" +#include "meta-store/codec.h" #include "meta-store/document_message.h" #include "meta-store/keys.h" #include "meta-store/mem_txn_kv.h" @@ -6911,9 +6912,25 @@ TEST(MetaServiceTest, UpdateDeleteBitmapWithException) { } } +static void put_delete_bitmap_test_rowset(MetaServiceProxy* meta_service, + const std::string& instance_id, int64_t tablet_id, + int64_t version, const std::string& rowset_id) { + std::unique_ptr txn; + ASSERT_EQ(meta_service->txn_kv()->create_txn(&txn), TxnErrorCode::TXN_OK); + std::string rowset_key = meta_rowset_key({instance_id, tablet_id, version}); + doris::RowsetMetaCloudPB rowset_meta; + rowset_meta.set_rowset_id(0); + rowset_meta.set_rowset_id_v2(rowset_id); + std::string rowset_value; + ASSERT_TRUE(rowset_meta.SerializeToString(&rowset_value)); + txn->put(rowset_key, rowset_value); + ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK); +} + void update_delete_bitmap_with_remove_pre(MetaServiceProxy* meta_service, int64_t table_id, int64_t tablet_id, bool inject = false, - bool rowset_non_exist = false) { + bool rowset_non_exist = false, + bool use_delete_bitmap_stats = true) { // create rowset, if `rowset_non_exist` enabled, only r4 exists { std::unique_ptr txn; @@ -7041,6 +7058,21 @@ void update_delete_bitmap_with_remove_pre(MetaServiceProxy* meta_service, int64_ update_delete_bitmap_req.add_segment_delete_bitmaps(new_large_value); update_delete_bitmap_req.add_pre_rowset_versions(rowset_version); } + if (use_delete_bitmap_stats) { + for (const auto& rowset : rowset_vector) { + auto* rowset_stats = update_delete_bitmap_req.add_pre_rowset_delete_bitmap_stats(); + rowset_stats->set_rowset_id(rowset); + for (const auto& [delete_bitmap_rowset, segment, version] : + rowset_segment_version_vector) { + if (delete_bitmap_rowset == rowset && version >= 4 && version < 6) { + auto* delete_bitmap_stat = rowset_stats->add_delete_bitmap_stats(); + delete_bitmap_stat->set_segment_id(segment); + delete_bitmap_stat->set_version(version); + delete_bitmap_stat->set_delete_bitmap_size(large_value.size()); + } + } + } + } meta_service->update_delete_bitmap(reinterpret_cast(&cntl), &update_delete_bitmap_req, &update_delete_bitmap_res, nullptr); @@ -7111,11 +7143,18 @@ TEST(MetaServiceTest, UpdateDeleteBitmapWithRemovePreDeleteBitmap) { SyncPoint::get_instance()->clear_all_call_backs(); }; - update_delete_bitmap_with_remove_pre(meta_service.get(), 200, 202); + update_delete_bitmap_with_remove_pre(meta_service.get(), 200, 202, false, false, false); int64_t max_txn_commit_byte = config::max_txn_commit_byte; config::max_txn_commit_byte = 1000; + int remove_pre_rowsets_commit_count = 0; + sp->set_call_back("update_delete_bitmap:remove_pre_rowsets:commit", + [&](auto&&) { remove_pre_rowsets_commit_count++; }); + sp->enable_processing(); update_delete_bitmap_with_remove_pre(meta_service.get(), 300, 302); + EXPECT_GT(remove_pre_rowsets_commit_count, 0); + sp->clear_all_call_backs(); + sp->disable_processing(); sp->set_call_back("update_delete_bitmap:commit:err", [&](auto&& args) { auto initiator = try_any_cast(args[0]); @@ -7132,6 +7171,286 @@ TEST(MetaServiceTest, UpdateDeleteBitmapWithRemovePreDeleteBitmap) { config::max_txn_commit_byte = max_txn_commit_byte; update_delete_bitmap_with_remove_pre(meta_service.get(), 500, 502, false, true); + + update_delete_bitmap_with_remove_pre(meta_service.get(), 600, 602); +} + +TEST(MetaServiceTest, EmptyPreRowsetStatsUsesKeyRemoval) { + auto meta_service = get_meta_service(); + extern std::string get_instance_id(const std::shared_ptr& rc_mgr, + const std::string& cloud_unique_id); + auto instance_id = get_instance_id(meta_service->resource_mgr(), "test_cloud_unique_id"); + constexpr int64_t tablet_id = 622; + const std::string rowset_id = "empty_pre_rowset_stats"; + put_delete_bitmap_test_rowset(meta_service.get(), instance_id, tablet_id, 1, rowset_id); + + auto old_delete_bitmap_key = meta_delete_bitmap_key({instance_id, tablet_id, rowset_id, 1, 0}); + const std::string old_delete_bitmap = "old_delete_bitmap"; + std::unique_ptr txn; + ASSERT_EQ(meta_service->txn_kv()->create_txn(&txn), TxnErrorCode::TXN_OK); + blob_put(txn.get(), old_delete_bitmap_key, old_delete_bitmap, 0); + ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK); + + brpc::Controller cntl; + UpdateDeleteBitmapRequest req; + UpdateDeleteBitmapResponse res; + req.set_cloud_unique_id("test_cloud_unique_id"); + req.set_table_id(620); + req.set_partition_id(621); + req.set_tablet_id(tablet_id); + req.set_lock_id(-3); + req.set_without_lock(true); + req.set_initiator(tablet_id); + req.set_pre_rowset_agg_start_version(1); + req.set_pre_rowset_agg_end_version(2); + req.add_rowset_ids(rowset_id); + req.add_segment_ids(0); + req.add_versions(2); + req.add_segment_delete_bitmaps("aggregated_delete_bitmap"); + req.add_pre_rowset_versions(1); + req.add_pre_rowset_delete_bitmap_stats()->set_rowset_id(rowset_id); + meta_service->update_delete_bitmap(reinterpret_cast(&cntl), + &req, &res, nullptr); + + ASSERT_EQ(res.status().code(), MetaServiceCode::OK); + ASSERT_EQ(meta_service->txn_kv()->create_txn(&txn), TxnErrorCode::TXN_OK); + ValueBuf val; + ASSERT_EQ(blob_get(txn.get(), old_delete_bitmap_key, &val), TxnErrorCode::TXN_OK); + EXPECT_EQ(val.value(), old_delete_bitmap); +} + +TEST(MetaServiceTest, UnderestimatedPreRowsetDeleteBitmapSizeLeavesSkippedTail) { + auto meta_service = get_meta_service(); + extern std::string get_instance_id(const std::shared_ptr& rc_mgr, + const std::string& cloud_unique_id); + auto instance_id = get_instance_id(meta_service->resource_mgr(), "test_cloud_unique_id"); + constexpr int64_t tablet_id = 632; + const std::string rowset_id = "underestimated_delete_bitmap"; + put_delete_bitmap_test_rowset(meta_service.get(), instance_id, tablet_id, 1, rowset_id); + + auto old_delete_bitmap_key = meta_delete_bitmap_key({instance_id, tablet_id, rowset_id, 1, 0}); + const std::string old_delete_bitmap(DEFAULT_BLOB_SPLIT_SIZE * 3 + 1, 'a'); + std::unique_ptr txn; + ASSERT_EQ(meta_service->txn_kv()->create_txn(&txn), TxnErrorCode::TXN_OK); + blob_put(txn.get(), old_delete_bitmap_key, old_delete_bitmap, 0); + ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK); + + const std::string aggregated_delete_bitmap = "aggregated_delete_bitmap"; + brpc::Controller update_cntl; + UpdateDeleteBitmapRequest update_req; + UpdateDeleteBitmapResponse update_res; + update_req.set_cloud_unique_id("test_cloud_unique_id"); + update_req.set_table_id(630); + update_req.set_partition_id(631); + update_req.set_tablet_id(tablet_id); + update_req.set_lock_id(-3); + update_req.set_without_lock(true); + update_req.set_initiator(tablet_id); + update_req.set_pre_rowset_agg_start_version(1); + update_req.set_pre_rowset_agg_end_version(2); + update_req.add_rowset_ids(rowset_id); + update_req.add_segment_ids(0); + update_req.add_versions(2); + update_req.add_segment_delete_bitmaps(aggregated_delete_bitmap); + update_req.add_pre_rowset_versions(1); + auto* rowset_stats = update_req.add_pre_rowset_delete_bitmap_stats(); + rowset_stats->set_rowset_id(rowset_id); + auto* delete_bitmap_stat = rowset_stats->add_delete_bitmap_stats(); + delete_bitmap_stat->set_segment_id(0); + delete_bitmap_stat->set_version(1); + delete_bitmap_stat->set_delete_bitmap_size(DEFAULT_BLOB_SPLIT_SIZE); + meta_service->update_delete_bitmap( + reinterpret_cast(&update_cntl), &update_req, + &update_res, nullptr); + ASSERT_EQ(update_res.status().code(), MetaServiceCode::OK); + + ASSERT_EQ(meta_service->txn_kv()->create_txn(&txn), TxnErrorCode::TXN_OK); + std::string value; + EXPECT_EQ(txn->get(blob_key(old_delete_bitmap_key, 0), &value), + TxnErrorCode::TXN_KEY_NOT_FOUND); + EXPECT_EQ(txn->get(blob_key(old_delete_bitmap_key, 1), &value), TxnErrorCode::TXN_OK); + + brpc::Controller get_cntl; + GetDeleteBitmapRequest get_req; + GetDeleteBitmapResponse get_res; + get_req.set_cloud_unique_id("test_cloud_unique_id"); + get_req.set_tablet_id(tablet_id); + get_req.add_rowset_ids(rowset_id); + get_req.add_begin_versions(0); + get_req.add_end_versions(3); + meta_service->get_delete_bitmap(reinterpret_cast(&get_cntl), + &get_req, &get_res, nullptr); + + ASSERT_EQ(get_res.status().code(), MetaServiceCode::OK); + ASSERT_EQ(get_res.rowset_ids_size(), 1); + ASSERT_EQ(get_res.segment_ids_size(), 1); + ASSERT_EQ(get_res.versions_size(), 1); + ASSERT_EQ(get_res.segment_delete_bitmaps_size(), 1); + EXPECT_EQ(get_res.rowset_ids(0), rowset_id); + EXPECT_EQ(get_res.segment_ids(0), 0); + EXPECT_EQ(get_res.versions(0), 2); + EXPECT_EQ(get_res.segment_delete_bitmaps(0), aggregated_delete_bitmap); +} + +static void test_get_delete_bitmap_during_point_cleanup(size_t delete_bitmap_size) { + auto meta_service = get_meta_service(); + auto sp = SyncPoint::get_instance(); + DORIS_CLOUD_DEFER { + sp->clear_all_call_backs(); + sp->disable_processing(); + }; + + extern std::string get_instance_id(const std::shared_ptr& rc_mgr, + const std::string& cloud_unique_id); + auto instance_id = get_instance_id(meta_service->resource_mgr(), "test_cloud_unique_id"); + constexpr int64_t table_id = 640; + constexpr int64_t partition_id = 641; + constexpr int64_t tablet_id = 642; + const std::string rowset_id = "point_cleanup_during_paginated_read"; + put_delete_bitmap_test_rowset(meta_service.get(), instance_id, tablet_id, 1, rowset_id); + + auto old_delete_bitmap_key = meta_delete_bitmap_key({instance_id, tablet_id, rowset_id, 1, 0}); + std::string old_delete_bitmap(DEFAULT_BLOB_SPLIT_SIZE, 'a'); + old_delete_bitmap.append(DEFAULT_BLOB_SPLIT_SIZE, 'b'); + old_delete_bitmap.append(DEFAULT_BLOB_SPLIT_SIZE, 'c'); + old_delete_bitmap.push_back('d'); + std::unique_ptr txn; + ASSERT_EQ(meta_service->txn_kv()->create_txn(&txn), TxnErrorCode::TXN_OK); + blob_put(txn.get(), old_delete_bitmap_key, old_delete_bitmap, 0); + ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK); + + const std::string aggregated_delete_bitmap = "aggregated_delete_bitmap"; + UpdateDeleteBitmapRequest update_req; + update_req.set_cloud_unique_id("test_cloud_unique_id"); + update_req.set_table_id(table_id); + update_req.set_partition_id(partition_id); + update_req.set_tablet_id(tablet_id); + update_req.set_lock_id(-3); + update_req.set_without_lock(true); + update_req.set_initiator(tablet_id); + update_req.set_pre_rowset_agg_start_version(1); + update_req.set_pre_rowset_agg_end_version(2); + update_req.add_rowset_ids(rowset_id); + update_req.add_segment_ids(0); + update_req.add_versions(2); + update_req.add_segment_delete_bitmaps(aggregated_delete_bitmap); + update_req.add_pre_rowset_versions(1); + auto* rowset_stats = update_req.add_pre_rowset_delete_bitmap_stats(); + rowset_stats->set_rowset_id(rowset_id); + auto* delete_bitmap_stat = rowset_stats->add_delete_bitmap_stats(); + delete_bitmap_stat->set_segment_id(0); + delete_bitmap_stat->set_version(1); + delete_bitmap_stat->set_delete_bitmap_size(delete_bitmap_size); + + bool cleanup_done = false; + sp->set_call_back("get_delete_bitmap_test", + [&](auto&& args) { *try_any_cast(args[0]) = true; }); + sp->set_call_back("get_delete_bitmap_err", [&](auto&& args) { + auto round = *try_any_cast(args[0]); + if (round != 1 || cleanup_done) { + return; + } + cleanup_done = true; + brpc::Controller update_cntl; + UpdateDeleteBitmapResponse update_res; + meta_service->update_delete_bitmap( + reinterpret_cast(&update_cntl), &update_req, + &update_res, nullptr); + ASSERT_EQ(update_res.status().code(), MetaServiceCode::OK); + *try_any_cast(args[1]) = TxnErrorCode::TXN_TOO_OLD; + }); + sp->enable_processing(); + + brpc::Controller get_cntl; + GetDeleteBitmapRequest get_req; + GetDeleteBitmapResponse get_res; + get_req.set_cloud_unique_id("test_cloud_unique_id"); + get_req.set_tablet_id(tablet_id); + get_req.add_rowset_ids(rowset_id); + get_req.add_begin_versions(0); + get_req.add_end_versions(3); + meta_service->get_delete_bitmap(reinterpret_cast(&get_cntl), + &get_req, &get_res, nullptr); + + ASSERT_TRUE(cleanup_done); + ASSERT_EQ(get_res.status().code(), MetaServiceCode::OK); + ASSERT_EQ(get_res.rowset_ids_size(), 1); + ASSERT_EQ(get_res.segment_ids_size(), 1); + ASSERT_EQ(get_res.versions_size(), 1); + ASSERT_EQ(get_res.segment_delete_bitmaps_size(), 1); + EXPECT_EQ(get_res.rowset_ids(0), rowset_id); + EXPECT_EQ(get_res.segment_ids(0), 0); + EXPECT_EQ(get_res.versions(0), 2); + EXPECT_EQ(get_res.segment_delete_bitmaps(0), aggregated_delete_bitmap); +} + +TEST(MetaServiceTest, GetDeleteBitmapSkipsPointCleanupTailAfterTxnTooOld) { + test_get_delete_bitmap_during_point_cleanup(DEFAULT_BLOB_SPLIT_SIZE * 3); +} + +TEST(MetaServiceTest, GetDeleteBitmapDropsPartialValueAfterAllBlobKeysAreRemoved) { + test_get_delete_bitmap_during_point_cleanup(DEFAULT_BLOB_SPLIT_SIZE * 3 + 1); +} + +TEST(MetaServiceTest, RemovePreDeleteBitmapBatchesEachBlobKey) { + auto meta_service = get_meta_service(); + auto sp = SyncPoint::get_instance(); + int64_t old_max_txn_commit_byte = config::max_txn_commit_byte; + DORIS_CLOUD_DEFER { + config::max_txn_commit_byte = old_max_txn_commit_byte; + sp->clear_all_call_backs(); + sp->disable_processing(); + }; + + constexpr size_t max_txn_commit_byte = 300; + config::max_txn_commit_byte = max_txn_commit_byte; + std::vector committed_txn_sizes; + sp->set_call_back("update_delete_bitmap:remove_pre_rowsets:commit", [&](auto&& args) { + committed_txn_sizes.push_back(try_any_cast(args[0])); + }); + sp->enable_processing(); + + extern std::string get_instance_id(const std::shared_ptr& rc_mgr, + const std::string& cloud_unique_id); + auto instance_id = get_instance_id(meta_service->resource_mgr(), "test_cloud_unique_id"); + constexpr int64_t tablet_id = 612; + const std::string rowset_id = "batch_blob_keys_rowset"; + const std::string value(DEFAULT_BLOB_SPLIT_SIZE * 20, 'a'); + auto delete_bitmap_key = meta_delete_bitmap_key({instance_id, tablet_id, rowset_id, 1, 0}); + std::unique_ptr txn; + ASSERT_EQ(meta_service->txn_kv()->create_txn(&txn), TxnErrorCode::TXN_OK); + blob_put(txn.get(), delete_bitmap_key, value, 0); + ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK); + + brpc::Controller cntl; + UpdateDeleteBitmapRequest req; + UpdateDeleteBitmapResponse res; + req.set_cloud_unique_id("test_cloud_unique_id"); + req.set_table_id(610); + req.set_partition_id(611); + req.set_tablet_id(tablet_id); + req.set_lock_id(-3); + req.set_without_lock(true); + req.set_initiator(tablet_id); + req.set_pre_rowset_agg_start_version(1); + req.set_pre_rowset_agg_end_version(2); + auto* rowset_stats = req.add_pre_rowset_delete_bitmap_stats(); + rowset_stats->set_rowset_id(rowset_id); + auto* delete_bitmap_stat = rowset_stats->add_delete_bitmap_stats(); + delete_bitmap_stat->set_segment_id(0); + delete_bitmap_stat->set_version(1); + delete_bitmap_stat->set_delete_bitmap_size(value.size()); + meta_service->update_delete_bitmap(reinterpret_cast(&cntl), + &req, &res, nullptr); + + ASSERT_EQ(res.status().code(), MetaServiceCode::OK); + ASSERT_FALSE(committed_txn_sizes.empty()); + for (auto txn_size : committed_txn_sizes) { + EXPECT_LE(txn_size, max_txn_commit_byte); + } + ASSERT_EQ(meta_service->txn_kv()->create_txn(&txn), TxnErrorCode::TXN_OK); + ValueBuf val; + EXPECT_EQ(blob_get(txn.get(), delete_bitmap_key, &val), TxnErrorCode::TXN_KEY_NOT_FOUND); } TEST(MetaServiceTest, GetDeleteBitmapWithIdx) { diff --git a/gensrc/proto/cloud.proto b/gensrc/proto/cloud.proto index 51c522b6272cee..00250beb89b3be 100644 --- a/gensrc/proto/cloud.proto +++ b/gensrc/proto/cloud.proto @@ -1903,6 +1903,17 @@ message DeleteBitmapStoragePB { optional PackedSliceLocationPB packed_slice_location = 3; } +message PreRowsetDeleteBitmapStatsPB { + message DeleteBitmapStatPB { + optional uint32 segment_id = 1; + optional int64 version = 2; + optional uint64 delete_bitmap_size = 3; + } + + optional string rowset_id = 1; + repeated DeleteBitmapStatPB delete_bitmap_stats = 2; +} + message UpdateDeleteBitmapRequest { optional string cloud_unique_id = 1; // For auth optional int64 table_id = 2; @@ -1928,6 +1939,8 @@ message UpdateDeleteBitmapRequest { optional int64 pre_rowset_agg_end_version = 17; // when update delete_bitmap of pre rowsets, check the rowset exists repeated int64 pre_rowset_versions = 18; + // original delete bitmaps removed after pre rowset aggregation + repeated PreRowsetDeleteBitmapStatsPB pre_rowset_delete_bitmap_stats = 19; // 50-51: used for storage v2 repeated string delta_rowset_ids = 50; repeated DeleteBitmapStoragePB delete_bitmap_storages = 51;