From c228b296701b778f01ba0ea0279c98b5201f8306 Mon Sep 17 00:00:00 2001 From: yongman Date: Wed, 22 Jul 2026 16:54:55 +0800 Subject: [PATCH] LAC: smoothing tokens request and keep in high level Signed-off-by: yongman --- .../LocalAdmissionController.cpp | 61 ++++++++---- .../LocalAdmissionController.h | 5 + dbms/src/Flash/ResourceControl/TokenBucket.h | 2 + .../gtest_local_admission_controller.cpp | 98 +++++++++++++++++++ 4 files changed, 149 insertions(+), 17 deletions(-) diff --git a/dbms/src/Flash/ResourceControl/LocalAdmissionController.cpp b/dbms/src/Flash/ResourceControl/LocalAdmissionController.cpp index ed32425a591..77c69a0c812 100644 --- a/dbms/src/Flash/ResourceControl/LocalAdmissionController.cpp +++ b/dbms/src/Flash/ResourceControl/LocalAdmissionController.cpp @@ -85,7 +85,7 @@ std::optional ResourceGroup::buildRequestInfoIfNecessary(const S { acquire_tokens = getAcquireRUNumWithoutLock( consumption_delta_info.speed, - LocalAdmissionController::DEFAULT_TARGET_PERIOD.count(), + REFILL_TOKEN_INTERVAL.count(), LocalAdmissionController::ACQUIRE_RU_AMPLIFICATION); assert(acquire_tokens >= 0.0); @@ -138,26 +138,50 @@ bool ResourceGroup::shouldReportRUConsumption(const SteadyClock::time_point & no return false; } -double ResourceGroup::getAcquireRUNumWithoutLock(double speed, uint32_t n_sec, double amplification) const +bool ResourceGroup::shouldRefillToken(const SteadyClock::time_point & now) const { - assert(amplification > 1.0); + std::lock_guard lock(mu); + if (burstable || bucket_mode != normal_mode || request_in_progress) + return false; + + const auto elapsed = now - last_request_gac_timepoint; + RUNTIME_CHECK(elapsed.count() >= 0, elapsed.count()); + if (elapsed < REFILL_TOKEN_INTERVAL) + return false; - double remaining_ru = 0.0; - remaining_ru = bucket->peek(); + const auto refill_threshold = getTokenHighWatermarkWithoutLock() * REFILL_TOKEN_THRESHOLD_RATE; + return bucket->peek() <= refill_threshold; +} - // Appropriate amplification is necessary to prevent situation that GAC has sufficient RU, - // but user query speed is limited due to LAC requests too few RU. - double acquire_num = speed * n_sec * amplification; +double ResourceGroup::getTokenHighWatermarkWithoutLock() const +{ + // The resource group definition contains the global burst limit. Only use capacity as a local high watermark after + // GAC has returned the capacity assigned to this client. Before that, keep the startup fill rate as the watermark. + const auto high_watermark = has_gac_capacity ? bucket->getCapacity() : static_cast(user_ru_per_sec); + if unlikely (high_watermark <= 0.0 && !burstable) + return DEFAULT_BUFFER_TOKENS; + return high_watermark; +} - // This should not happen, but still add this to avoid stuck. - if unlikely (acquire_num == 0.0 && remaining_ru == 0.0) - acquire_num = DEFAULT_BUFFER_TOKENS; +double ResourceGroup::getAcquireRUNumWithoutLock(double speed, uint32_t n_sec, double amplification) const +{ + assert(amplification > 1.0); - // The purpose of subtracting remaining_ru is try to ensure that the number of local tokens - // always stays same with the amount consumed. - acquire_num -= remaining_ru; - acquire_num = (acquire_num > 0.0 ? acquire_num : 0.0); - return acquire_num; + const auto remaining_ru = bucket->peek(); + const auto high_watermark = getTokenHighWatermarkWithoutLock(); + auto acquire_num = high_watermark - remaining_ru; + if (acquire_num <= 0.0) + return 0.0; + + if (bucket->lowToken()) + return acquire_num; + + // Refill at most one second of predicted consumption in normal mode. The fallback batch gradually raises an idle + // bucket without transferring the whole capacity from GAC in one request. Low-token mode bypasses this limit above. + const auto refill_window = high_watermark * (1.0 - REFILL_TOKEN_THRESHOLD_RATE); + const auto fallback_batch = std::min(static_cast(DEFAULT_BUFFER_TOKENS), refill_window); + const auto incremental_batch = std::max(speed * n_sec * amplification, fallback_batch); + return std::min(acquire_num, incremental_batch); } void ResourceGroup::updateNormalMode(double add_tokens, double new_capacity, const SteadyClock::time_point & now) @@ -173,6 +197,7 @@ void ResourceGroup::updateNormalMode(double add_tokens, double new_capacity, con burstable = true; return; } + has_gac_capacity = true; auto config = bucket->getConfig(); std::string ori_bucket_info = bucket->toString(); @@ -209,6 +234,7 @@ void ResourceGroup::updateTrickleMode( burstable = true; return; } + has_gac_capacity = true; bucket_mode = TokenBucketMode::trickle_mode; double new_fill_rate = add_tokens / (static_cast(trickle_ms) / 1000); @@ -476,7 +502,8 @@ std::optional LocalAdmissionController::b for (const auto & ele : local_keyspace_resource_groups) { - const bool need_fetch_token = local_keyspace_low_token_resource_groups.contains(ele.first); + const bool need_fetch_token = local_keyspace_low_token_resource_groups.contains(ele.first) + || ele.second->shouldRefillToken(current_tick); const bool need_report = ele.second->shouldReportRUConsumption(current_tick); if (need_fetch_token || need_report) diff --git a/dbms/src/Flash/ResourceControl/LocalAdmissionController.h b/dbms/src/Flash/ResourceControl/LocalAdmissionController.h index c420268bea9..09d248ecccf 100644 --- a/dbms/src/Flash/ResourceControl/LocalAdmissionController.h +++ b/dbms/src/Flash/ResourceControl/LocalAdmissionController.h @@ -142,6 +142,8 @@ class ResourceGroup final : private boost::noncopyable static constexpr auto REPORT_RU_CONSUMPTION_DELTA_THRESHOLD = 100; static constexpr auto EXTENDING_REPORT_RU_CONSUMPTION_FACTOR = 4; static constexpr auto DEFAULT_BUFFER_TOKENS = 5000; + static constexpr auto REFILL_TOKEN_INTERVAL = std::chrono::seconds(1); + static constexpr double REFILL_TOKEN_THRESHOLD_RATE = 0.8; // Indicate the round trip time of gac request. static constexpr auto GAC_RTT_ANTICIPATION = std::chrono::seconds(1); @@ -204,9 +206,11 @@ class ResourceGroup final : private boost::noncopyable endRequestWithoutLock(); } bool shouldReportRUConsumption(const SteadyClock::time_point & now) const; + bool shouldRefillToken(const SteadyClock::time_point & now) const; std::optional buildRequestInfoIfNecessary(const SteadyClock::time_point & now); LACRUConsumptionDeltaInfo updateRUConsumptionDeltaInfoWithoutLock(); double getAcquireRUNumWithoutLock(double speed, uint32_t n_sec, double amplification) const; + double getTokenHighWatermarkWithoutLock() const; void updateRUConsumptionSpeedIfNecessary(const SteadyClock::time_point & now); // Called when user change config of resource group. @@ -310,6 +314,7 @@ class ResourceGroup final : private boost::noncopyable // Local token bucket. TokenBucketPtr bucket; TokenBucketMode bucket_mode = TokenBucketMode::normal_mode; + bool has_gac_capacity = false; // For compute priority. uint64_t cpu_time_in_ns = 0; diff --git a/dbms/src/Flash/ResourceControl/TokenBucket.h b/dbms/src/Flash/ResourceControl/TokenBucket.h index 61c052e0a63..eaa27376b0a 100644 --- a/dbms/src/Flash/ResourceControl/TokenBucket.h +++ b/dbms/src/Flash/ResourceControl/TokenBucket.h @@ -94,6 +94,8 @@ class TokenBucket final bool isStatic() const { return fill_rate == 0.0; } + double getCapacity() const { return capacity; } + std::string toString() const { return fmt::format("tokens: {}, fill_rate: {}, capacity: {}", tokens, fill_rate, capacity); diff --git a/dbms/src/Flash/ResourceControl/tests/gtest_local_admission_controller.cpp b/dbms/src/Flash/ResourceControl/tests/gtest_local_admission_controller.cpp index 0c74f3bac71..12e913c0a7d 100644 --- a/dbms/src/Flash/ResourceControl/tests/gtest_local_admission_controller.cpp +++ b/dbms/src/Flash/ResourceControl/tests/gtest_local_admission_controller.cpp @@ -94,6 +94,7 @@ TEST(LocalAdmissionControllerTest, LegacyBackendCachesSharedGroupAndOmitsKeyspac ASSERT_NE(default_group, nullptr); EXPECT_TRUE(lac.getPriority(keyspace_id, "default").has_value()); + default_group->consumeResource(1, 0); auto request_info = default_group->buildRequestInfoIfNecessary(SteadyClock::now()); ASSERT_TRUE(request_info.has_value()); EXPECT_EQ(request_info->keyspace_id, NullspaceID); @@ -126,6 +127,7 @@ TEST(LocalAdmissionControllerTest, KeyspaceScopedBackendCachesExactGroupAndKeeps EXPECT_EQ(lac.getCachedResourceGroupForTest(NullspaceID, "default"), nullptr); auto default_group = lac.getCachedResourceGroupForTest(keyspace_id, "default"); ASSERT_NE(default_group, nullptr); + default_group->consumeResource(1, 0); auto request_info = default_group->buildRequestInfoIfNecessary(SteadyClock::now()); ASSERT_TRUE(request_info.has_value()); EXPECT_EQ(request_info->keyspace_id, keyspace_id); @@ -178,4 +180,100 @@ TEST(LocalAdmissionControllerTest, UnexpectedPriorityFallsBackToMedium) EXPECT_EQ(analytics_group->user_priority_val, ResourceGroup::MediumPriorityValue); } +TEST(LocalAdmissionControllerTest, StartupRefillDoesNotUseGlobalBurstLimitAsHighWatermark) +{ + constexpr double fill_rate = 1000; + constexpr int64_t global_burst_limit = 10000; + constexpr double consumed_tokens = 500; + const auto start_time = SteadyClock::now() - 2 * ResourceGroup::REFILL_TOKEN_INTERVAL; + + resource_manager::ResourceGroup group_pb; + group_pb.set_name("startup"); + group_pb.set_mode(resource_manager::GroupMode::RUMode); + group_pb.set_priority(ResourceGroup::UserMediumPriority); + auto * settings = group_pb.mutable_r_u_settings()->mutable_r_u()->mutable_settings(); + settings->set_fill_rate(fill_rate); + settings->set_burst_limit(global_burst_limit); + + ResourceGroup group(NullspaceID, group_pb, start_time); + group.smooth_ru_consumption_speed = 0; + group.consumeResource(consumed_tokens, 0); + + EXPECT_FALSE(group.shouldRefillToken(start_time + ResourceGroup::REFILL_TOKEN_INTERVAL / 2)); + EXPECT_TRUE(group.shouldRefillToken(start_time + ResourceGroup::REFILL_TOKEN_INTERVAL)); + + const auto request_info = group.buildRequestInfoIfNecessary(start_time + ResourceGroup::REFILL_TOKEN_INTERVAL); + ASSERT_TRUE(request_info.has_value()); + EXPECT_DOUBLE_EQ(request_info->acquire_tokens, fill_rate * (1 - ResourceGroup::REFILL_TOKEN_THRESHOLD_RATE)); + EXPECT_DOUBLE_EQ(request_info->ru_consumption_delta, consumed_tokens); + + const auto response_time = start_time + ResourceGroup::REFILL_TOKEN_INTERVAL; + group.updateNormalMode(request_info->acquire_tokens, global_burst_limit, response_time); + EXPECT_FALSE(group.lowToken()); + EXPECT_TRUE(group.shouldRefillToken(response_time + ResourceGroup::REFILL_TOKEN_INTERVAL)); + + const auto next_request_info + = group.buildRequestInfoIfNecessary(response_time + ResourceGroup::REFILL_TOKEN_INTERVAL); + ASSERT_TRUE(next_request_info.has_value()); + EXPECT_DOUBLE_EQ( + next_request_info->acquire_tokens, + global_burst_limit * (1 - ResourceGroup::REFILL_TOKEN_THRESHOLD_RATE)); +} + +TEST(LocalAdmissionControllerTest, RefillTokensIncrementallyAboveLowWatermark) +{ + constexpr double capacity = 10000; + constexpr double consumed_tokens = 3000; + + ResourceGroup group( + "normal_refill", + ResourceGroup::UserMediumPriority, + capacity, + /*burstable_=*/false); + group.smooth_ru_consumption_speed = 0; + group.consumeResource(consumed_tokens, 0); + + ASSERT_TRUE(group.shouldRefillToken(SteadyClock::now())); + auto request_info = group.buildRequestInfoIfNecessary(SteadyClock::now()); + ASSERT_TRUE(request_info.has_value()); + EXPECT_DOUBLE_EQ(request_info->acquire_tokens, capacity * (1 - ResourceGroup::REFILL_TOKEN_THRESHOLD_RATE)); + EXPECT_DOUBLE_EQ(request_info->ru_consumption_delta, consumed_tokens); +} + +TEST(LocalAdmissionControllerTest, RefillTokensUsesPredictedConsumptionWhenHigher) +{ + constexpr double capacity = 10000; + constexpr double consumed_tokens = 3000; + + ResourceGroup group( + "high_speed", + ResourceGroup::UserMediumPriority, + capacity, + /*burstable_=*/false); + group.smooth_ru_consumption_speed = 4000; + group.consumeResource(consumed_tokens, 0); + + auto request_info = group.buildRequestInfoIfNecessary(SteadyClock::now()); + ASSERT_TRUE(request_info.has_value()); + EXPECT_DOUBLE_EQ(request_info->acquire_tokens, consumed_tokens); +} + +TEST(LocalAdmissionControllerTest, LowTokenRefillBypassesIncrementalLimit) +{ + constexpr double capacity = 10000; + constexpr double consumed_tokens = 7500; + + ResourceGroup group( + "emergency_refill", + ResourceGroup::UserMediumPriority, + capacity, + /*burstable_=*/false); + group.smooth_ru_consumption_speed = 0; + group.consumeResource(consumed_tokens, 0); + + auto request_info = group.buildRequestInfoIfNecessary(SteadyClock::now()); + ASSERT_TRUE(request_info.has_value()); + EXPECT_DOUBLE_EQ(request_info->acquire_tokens, consumed_tokens); +} + } // namespace DB::tests