From 5bb24755a66dfc9973f5556b49cbf283ffd78710 Mon Sep 17 00:00:00 2001 From: laihui <1353307710@qq.com> Date: Mon, 21 Sep 2026 16:40:39 +0800 Subject: [PATCH 1/5] [improvement](be) Isolate load control RPCs from heavy load work ### What problem does this PR solve? When load writes, flush waits, or closes occupy the shared BRPC heavy pool, `tablet_writer_open`, `tablet_writer_cancel`, and `open_load_stream` queue behind that work. In particular, a cancellation intended to release load resources cannot be dispatched promptly under heavy-pool saturation. Add two dedicated load pools: - `brpc_load_light`: writer open/cancel and stream open. - `brpc_load_heavy`: add-block RPCs (including the HTTP forwarding path) and `LoadStreamMgr` flush/pre-close/close tasks. Both pools have independent thread/queue settings and queue-size, active-thread, effective thread-limit, and effective queue-limit metrics. The four new `brpc_load_{heavy,light}_work_pool_{threads,max_queue_size}` settings require a restart and accept `-1` or a positive value. By default, load-heavy inherits the existing heavy-pool settings; load-light uses `max(32, CPU cores)` threads and `max(1024, CPU cores * 32)` queued requests. The split reserves execution capacity for load control requests. It does not remove locks or storage operations inside open/cancel, and it creates additional fixed worker threads. Existing generic heavy/light pools, RPC response contracts, and storage/transaction semantics are retained. No throughput or latency benchmark is claimed. ### Release note Isolate load open and cancellation RPCs from load write/close queues to reduce control-request queueing under load. Add independently configurable load-heavy and load-light BRPC worker pools and metrics. ### Check List (For Author) - Test - [x] Unit Test: added deterministic queue-routing and queue-full callback/status coverage for writer open/cancel, stream open, and add-block, plus streaming close-pool wiring. **Not executed**, as requested by the author. - Validation completed: clang-format 16.0.5 on all five changed C++ files and `git diff --check`. - Build/UT execution skipped at the author's request. An initial build-environment probe was stopped during third-party dependency setup; no BE build or test result is claimed. - Behavior changed: - [x] Yes. Load control and data requests use separate pools; existing generic RPC pools retain their configuration names. - Does this need documentation? - [x] No separate documentation change in this PR. New tuning defaults and restart requirements are documented alongside the BE configuration declarations. ### Review notes - Concurrency: only dispatch destinations change; handlers retain their existing locks and memory-tracking contexts. The new pools introduce no new lock order. - Lifetime: pools remain owned by `PInternalService`; `LoadStreamMgr` borrows its load-heavy pool using the existing ownership model. New metric hooks are deregistered on destruction. - Compatibility: no protobuf, persistent format, FE variable, visibility, or transaction protocol change. - Parallel paths: HTTP add-block delegates to the modified RPC; streaming open and flush/close are included; shared service dispatch covers cloud and local storage modes. - Failure handling: queue rejection reports the destination pool and preserves exactly-once completion behavior, including the empty cancellation response. - Coverage limit: the added tests cover dispatch/backpressure and wiring; end-to-end load behavior and performance were not exercised. --- be/src/common/config.cpp | 15 ++ be/src/common/config.h | 11 ++ be/src/service/internal_service.cpp | 129 ++++++++++--- be/src/service/internal_service.h | 3 + .../internal_service_load_work_pool_test.cpp | 177 ++++++++++++++++++ 5 files changed, 307 insertions(+), 28 deletions(-) create mode 100644 be/test/service/internal_service_load_work_pool_test.cpp diff --git a/be/src/common/config.cpp b/be/src/common/config.cpp index a0b479a4eeaefa..d156d488d2992e 100644 --- a/be/src/common/config.cpp +++ b/be/src/common/config.cpp @@ -656,6 +656,21 @@ DEFINE_mInt64(load_error_log_reserve_hours, "48"); // error log size limit, default 200MB DEFINE_mInt64(load_error_log_limit_bytes, "209715200"); +// Dedicated load pools. -1 keeps the existing heavy-pool capacity for load writes and +// uses CPU-scaled defaults for load control requests. These settings require a restart. +DEFINE_Int32(brpc_load_heavy_work_pool_threads, "-1"); +DEFINE_Int32(brpc_load_heavy_work_pool_max_queue_size, "-1"); +DEFINE_Int32(brpc_load_light_work_pool_threads, "-1"); +DEFINE_Int32(brpc_load_light_work_pool_max_queue_size, "-1"); +DEFINE_Validator(brpc_load_heavy_work_pool_threads, + [](const int config) -> bool { return config == -1 || config > 0; }); +DEFINE_Validator(brpc_load_heavy_work_pool_max_queue_size, + [](const int config) -> bool { return config == -1 || config > 0; }); +DEFINE_Validator(brpc_load_light_work_pool_threads, + [](const int config) -> bool { return config == -1 || config > 0; }); +DEFINE_Validator(brpc_load_light_work_pool_max_queue_size, + [](const int config) -> bool { return config == -1 || config > 0; }); + DEFINE_Int32(brpc_heavy_work_pool_threads, "-1"); DEFINE_Int32(brpc_peer_fetch_pool_threads, "-1"); DEFINE_Int32(brpc_light_work_pool_threads, "-1"); diff --git a/be/src/common/config.h b/be/src/common/config.h index 6bda64a14553a2..ca45230e448365 100644 --- a/be/src/common/config.h +++ b/be/src/common/config.h @@ -739,6 +739,17 @@ DECLARE_mInt64(load_error_log_reserve_hours); // error log size limit, default 200MB DECLARE_mInt64(load_error_log_limit_bytes); +// Dedicated load data pool for add-block and streaming flush/close work. +// -1 inherits brpc_heavy_work_pool_threads/max_queue_size (including CPU-scaled defaults). +DECLARE_Int32(brpc_load_heavy_work_pool_threads); +DECLARE_Int32(brpc_load_heavy_work_pool_max_queue_size); +// Dedicated load control pool for writer open/cancel and stream open. These handlers may +// acquire locks or access storage, so they must not use the query light pool. +// -1 selects max(32, CPU cores) threads and max(1024, CPU cores * 32) queued requests. +// All four load pool settings require a restart. +DECLARE_Int32(brpc_load_light_work_pool_threads); +DECLARE_Int32(brpc_load_light_work_pool_max_queue_size); + // be brpc interface is classified into two categories: light and heavy // each category has diffrent thread number // threads to handle heavy api interface, such as transmit_block etc diff --git a/be/src/service/internal_service.cpp b/be/src/service/internal_service.cpp index a5954695d9c00f..1a6d087e802e29 100644 --- a/be/src/service/internal_service.cpp +++ b/be/src/service/internal_service.cpp @@ -137,6 +137,15 @@ namespace doris { #include "common/compile_check_avoid_begin.h" using namespace ErrorCode; +DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(load_heavy_work_pool_queue_size, MetricUnit::NOUNIT); +DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(load_heavy_work_active_threads, MetricUnit::NOUNIT); +DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(load_heavy_work_pool_max_queue_size, MetricUnit::NOUNIT); +DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(load_heavy_work_max_threads, MetricUnit::NOUNIT); +DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(load_light_work_pool_queue_size, MetricUnit::NOUNIT); +DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(load_light_work_active_threads, MetricUnit::NOUNIT); +DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(load_light_work_pool_max_queue_size, MetricUnit::NOUNIT); +DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(load_light_work_max_threads, MetricUnit::NOUNIT); + DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(heavy_work_pool_queue_size, MetricUnit::NOUNIT); DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(peer_fetch_work_pool_queue_size, MetricUnit::NOUNIT); DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(light_work_pool_queue_size, MetricUnit::NOUNIT); @@ -158,6 +167,35 @@ DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(arrow_flight_work_max_threads, MetricUnit::NO static bvar::LatencyRecorder g_process_remote_fetch_rowsets_latency("process_remote_fetch_rowsets"); +static int32_t resolved_brpc_load_heavy_work_pool_threads() { + if (config::brpc_load_heavy_work_pool_threads != -1) { + return config::brpc_load_heavy_work_pool_threads; + } + return config::brpc_heavy_work_pool_threads != -1 ? config::brpc_heavy_work_pool_threads + : std::max(128, CpuInfo::num_cores() * 4); +} + +static int32_t resolved_brpc_load_heavy_work_pool_max_queue_size() { + if (config::brpc_load_heavy_work_pool_max_queue_size != -1) { + return config::brpc_load_heavy_work_pool_max_queue_size; + } + return config::brpc_heavy_work_pool_max_queue_size != -1 + ? config::brpc_heavy_work_pool_max_queue_size + : std::max(10240, CpuInfo::num_cores() * 320); +} + +static int32_t resolved_brpc_load_light_work_pool_threads() { + return config::brpc_load_light_work_pool_threads != -1 + ? config::brpc_load_light_work_pool_threads + : std::max(32, CpuInfo::num_cores()); +} + +static int32_t resolved_brpc_load_light_work_pool_max_queue_size() { + return config::brpc_load_light_work_pool_max_queue_size != -1 + ? config::brpc_load_light_work_pool_max_queue_size + : std::max(1024, CpuInfo::num_cores() * 32); +} + static int32_t resolved_brpc_peer_fetch_pool_threads() { return config::brpc_peer_fetch_pool_threads != -1 ? config::brpc_peer_fetch_pool_threads : std::max(64, CpuInfo::num_cores() * 2); @@ -214,7 +252,7 @@ class NewHttpClosure : public ::google::protobuf::Closure { PInternalService::PInternalService(ExecEnv* exec_env) : _exec_env(exec_env), - // heavy threadpool is used for load process and other process that will read disk or access network. + // General RPCs that read disk or access the network. _heavy_work_pool(config::brpc_heavy_work_pool_threads != -1 ? config::brpc_heavy_work_pool_threads : std::max(128, CpuInfo::num_cores() * 4), @@ -222,6 +260,13 @@ PInternalService::PInternalService(ExecEnv* exec_env) ? config::brpc_heavy_work_pool_max_queue_size : std::max(10240, CpuInfo::num_cores() * 320), "brpc_heavy"), + _load_heavy_work_pool(resolved_brpc_load_heavy_work_pool_threads(), + resolved_brpc_load_heavy_work_pool_max_queue_size(), + "brpc_load_heavy"), + // Open/cancel may block on storage or locks, but must not queue behind load writes. + _load_light_work_pool(resolved_brpc_load_light_work_pool_threads(), + resolved_brpc_load_light_work_pool_max_queue_size(), + "brpc_load_light"), // peer fetch threadpool isolates fetch_peer_data from heavy load traffic to avoid peer reads starving imports. _peer_fetch_pool(resolved_brpc_peer_fetch_pool_threads(), resolved_brpc_peer_fetch_pool_max_queue_size(), "brpc_peer_fetch"), @@ -241,6 +286,23 @@ PInternalService::PInternalService(ExecEnv* exec_env) ? config::brpc_arrow_flight_work_pool_max_queue_size : std::max(20480, CpuInfo::num_cores() * 640), "brpc_arrow_flight") { + REGISTER_HOOK_METRIC(load_heavy_work_pool_queue_size, + [this]() { return _load_heavy_work_pool.get_queue_size(); }); + REGISTER_HOOK_METRIC(load_heavy_work_active_threads, + [this]() { return _load_heavy_work_pool.get_active_threads(); }); + REGISTER_HOOK_METRIC(load_heavy_work_pool_max_queue_size, + []() { return resolved_brpc_load_heavy_work_pool_max_queue_size(); }); + REGISTER_HOOK_METRIC(load_heavy_work_max_threads, + []() { return resolved_brpc_load_heavy_work_pool_threads(); }); + REGISTER_HOOK_METRIC(load_light_work_pool_queue_size, + [this]() { return _load_light_work_pool.get_queue_size(); }); + REGISTER_HOOK_METRIC(load_light_work_active_threads, + [this]() { return _load_light_work_pool.get_active_threads(); }); + REGISTER_HOOK_METRIC(load_light_work_pool_max_queue_size, + []() { return resolved_brpc_load_light_work_pool_max_queue_size(); }); + REGISTER_HOOK_METRIC(load_light_work_max_threads, + []() { return resolved_brpc_load_light_work_pool_threads(); }); + REGISTER_HOOK_METRIC(heavy_work_pool_queue_size, [this]() { return _heavy_work_pool.get_queue_size(); }); REGISTER_HOOK_METRIC(peer_fetch_work_pool_queue_size, @@ -276,7 +338,7 @@ PInternalService::PInternalService(ExecEnv* exec_env) REGISTER_HOOK_METRIC(arrow_flight_work_max_threads, []() { return config::brpc_arrow_flight_work_pool_threads; }); - _exec_env->load_stream_mgr()->set_heavy_work_pool(&_heavy_work_pool); + _exec_env->load_stream_mgr()->set_heavy_work_pool(&_load_heavy_work_pool); CHECK_EQ(0, bthread_key_create(&AsyncIO::btls_io_ctx_key, AsyncIO::io_ctx_key_deleter)); } @@ -287,6 +349,15 @@ PInternalServiceImpl::PInternalServiceImpl(StorageEngine& engine, ExecEnv* exec_ PInternalServiceImpl::~PInternalServiceImpl() = default; PInternalService::~PInternalService() { + DEREGISTER_HOOK_METRIC(load_heavy_work_pool_queue_size); + DEREGISTER_HOOK_METRIC(load_heavy_work_active_threads); + DEREGISTER_HOOK_METRIC(load_heavy_work_pool_max_queue_size); + DEREGISTER_HOOK_METRIC(load_heavy_work_max_threads); + DEREGISTER_HOOK_METRIC(load_light_work_pool_queue_size); + DEREGISTER_HOOK_METRIC(load_light_work_active_threads); + DEREGISTER_HOOK_METRIC(load_light_work_pool_max_queue_size); + DEREGISTER_HOOK_METRIC(load_light_work_max_threads); + DEREGISTER_HOOK_METRIC(heavy_work_pool_queue_size); DEREGISTER_HOOK_METRIC(peer_fetch_work_pool_queue_size); DEREGISTER_HOOK_METRIC(light_work_pool_queue_size); @@ -313,7 +384,7 @@ void PInternalService::tablet_writer_open(google::protobuf::RpcController* contr const PTabletWriterOpenRequest* request, PTabletWriterOpenResult* response, google::protobuf::Closure* done) { - bool ret = _heavy_work_pool.try_offer([this, request, response, done]() { + bool ret = _load_light_work_pool.try_offer([this, request, response, done]() { VLOG_RPC << "tablet writer open, id=" << request->id() << ", index_id=" << request->index_id() << ", txn_id=" << request->txn_id(); signal::SignalTaskIdKeeper keeper(request->id()); @@ -327,7 +398,7 @@ void PInternalService::tablet_writer_open(google::protobuf::RpcController* contr st.to_protobuf(response->mutable_status()); }); if (!ret) { - offer_failed(response, done, _heavy_work_pool); + offer_failed(response, done, _load_light_work_pool); return; } } @@ -423,7 +494,7 @@ void PInternalService::open_load_stream(google::protobuf::RpcController* control const POpenLoadStreamRequest* request, POpenLoadStreamResponse* response, google::protobuf::Closure* done) { - bool ret = _heavy_work_pool.try_offer([this, controller, request, response, done]() { + bool ret = _load_light_work_pool.try_offer([this, controller, request, response, done]() { signal::SignalTaskIdKeeper keeper(request->load_id()); brpc::ClosureGuard done_guard(done); brpc::Controller* cntl = static_cast(controller); @@ -481,7 +552,7 @@ void PInternalService::open_load_stream(google::protobuf::RpcController* control st.to_protobuf(response->mutable_status()); }); if (!ret) { - offer_failed(response, done, _heavy_work_pool); + offer_failed(response, done, _load_light_work_pool); } } @@ -507,27 +578,29 @@ void PInternalService::tablet_writer_add_block(google::protobuf::RpcController* PTabletWriterAddBlockResult* response, google::protobuf::Closure* done) { int64_t submit_task_time_ns = MonotonicNanos(); - bool ret = _heavy_work_pool.try_offer([request, response, done, submit_task_time_ns, this]() { - int64_t wait_execution_time_ns = MonotonicNanos() - submit_task_time_ns; - brpc::ClosureGuard closure_guard(done); - int64_t execution_time_ns = 0; - { - SCOPED_RAW_TIMER(&execution_time_ns); - signal::SignalTaskIdKeeper keeper(request->id()); - auto st = _exec_env->load_channel_mgr()->add_batch(*request, response); - if (!st.ok()) { - LOG(WARNING) << "tablet writer add block failed, message=" << st - << ", id=" << request->id() << ", index_id=" << request->index_id() - << ", sender_id=" << request->sender_id() - << ", backend id=" << request->backend_id(); - } - st.to_protobuf(response->mutable_status()); - } - response->set_execution_time_us(execution_time_ns / NANOS_PER_MICRO); - response->set_wait_execution_time_us(wait_execution_time_ns / NANOS_PER_MICRO); - }); + bool ret = + _load_heavy_work_pool.try_offer([request, response, done, submit_task_time_ns, this]() { + int64_t wait_execution_time_ns = MonotonicNanos() - submit_task_time_ns; + brpc::ClosureGuard closure_guard(done); + int64_t execution_time_ns = 0; + { + SCOPED_RAW_TIMER(&execution_time_ns); + signal::SignalTaskIdKeeper keeper(request->id()); + auto st = _exec_env->load_channel_mgr()->add_batch(*request, response); + if (!st.ok()) { + LOG(WARNING) + << "tablet writer add block failed, message=" << st + << ", id=" << request->id() << ", index_id=" << request->index_id() + << ", sender_id=" << request->sender_id() + << ", backend id=" << request->backend_id(); + } + st.to_protobuf(response->mutable_status()); + } + response->set_execution_time_us(execution_time_ns / NANOS_PER_MICRO); + response->set_wait_execution_time_us(wait_execution_time_ns / NANOS_PER_MICRO); + }); if (!ret) { - offer_failed(response, done, _heavy_work_pool); + offer_failed(response, done, _load_heavy_work_pool); return; } } @@ -536,7 +609,7 @@ void PInternalService::tablet_writer_cancel(google::protobuf::RpcController* con const PTabletWriterCancelRequest* request, PTabletWriterCancelResult* response, google::protobuf::Closure* done) { - bool ret = _heavy_work_pool.try_offer([this, request, done]() { + bool ret = _load_light_work_pool.try_offer([this, request, done]() { VLOG_RPC << "tablet writer cancel, id=" << request->id() << ", index_id=" << request->index_id() << ", sender_id=" << request->sender_id(); signal::SignalTaskIdKeeper keeper(request->id()); @@ -549,7 +622,7 @@ void PInternalService::tablet_writer_cancel(google::protobuf::RpcController* con } }); if (!ret) { - offer_failed(response, done, _heavy_work_pool); + offer_failed(response, done, _load_light_work_pool); return; } } diff --git a/be/src/service/internal_service.h b/be/src/service/internal_service.h index 550f2af8637686..4422c5d5225abe 100644 --- a/be/src/service/internal_service.h +++ b/be/src/service/internal_service.h @@ -278,6 +278,9 @@ class PInternalService : public PBackendService { // define the interface for reading and writing data as heavy interface // otherwise as light interface FifoThreadPool _heavy_work_pool; + // Keep load control requests runnable when load writes or closes saturate their pool. + FifoThreadPool _load_heavy_work_pool; + FifoThreadPool _load_light_work_pool; FifoThreadPool _peer_fetch_pool; FifoThreadPool _light_work_pool; FifoThreadPool _arrow_flight_work_pool; diff --git a/be/test/service/internal_service_load_work_pool_test.cpp b/be/test/service/internal_service_load_work_pool_test.cpp new file mode 100644 index 00000000000000..c4e9d5fb23a02e --- /dev/null +++ b/be/test/service/internal_service_load_work_pool_test.cpp @@ -0,0 +1,177 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +#include + +#include +#include +#include +#include +#include + +#include "common/config.h" +#include "load/channel/load_stream_mgr.h" +#include "runtime/exec_env.h" +#include "service/internal_service.h" + +namespace doris { +namespace { + +// Hold the sole worker so routing and queue rejection can be checked without running +// storage handlers. Discard queued RPCs before releasing the worker at teardown. +class PausedLoadRpcPool { +public: + explicit PausedLoadRpcPool(FifoThreadPool& pool) : _pool(pool) { + auto resume = _resume.get_future().share(); + auto started = std::make_shared>(); + auto ready = started->get_future(); + CHECK(_pool.try_offer([started, resume]() { + started->set_value(); + resume.wait(); + })); + ready.wait(); + } + + ~PausedLoadRpcPool() { + _pool.shutdown(); + _resume.set_value(); + _pool.join(); + } + +private: + FifoThreadPool& _pool; + std::promise _resume; +}; + +class LoadRpcCountingClosure : public google::protobuf::Closure { +public: + void Run() override { ++calls; } + std::atomic calls {0}; +}; + +} // namespace + +class InternalServiceLoadWorkPoolTest : public testing::TestWithParam { +protected: + void SetUp() override { + // Keep all pools small and restore global configuration after each test. + for (auto* setting : + {&config::brpc_heavy_work_pool_threads, &config::brpc_heavy_work_pool_max_queue_size, + &config::brpc_light_work_pool_threads, &config::brpc_light_work_pool_max_queue_size, + &config::brpc_peer_fetch_pool_threads, &config::brpc_peer_fetch_pool_max_queue_size, + &config::brpc_arrow_flight_work_pool_threads, + &config::brpc_arrow_flight_work_pool_max_queue_size, + &config::brpc_load_heavy_work_pool_threads, + &config::brpc_load_heavy_work_pool_max_queue_size, + &config::brpc_load_light_work_pool_threads, + &config::brpc_load_light_work_pool_max_queue_size}) { + _saved_config.emplace_back(setting, *setting); + *setting = 1; + } + _exec_env._load_stream_mgr = std::make_unique(1); + _service = std::make_unique(&_exec_env); + for (auto* pool : {&_service->_heavy_work_pool, &_service->_light_work_pool, + &_service->_load_heavy_work_pool, &_service->_load_light_work_pool}) { + _paused_pools.push_back(std::make_unique(*pool)); + } + } + + void TearDown() override { + _paused_pools.clear(); + _exec_env.load_stream_mgr()->set_heavy_work_pool(nullptr); + _service.reset(); + _exec_env._load_stream_mgr.reset(); + for (const auto& [setting, value] : _saved_config) { + *setting = value; + } + } + + ExecEnv _exec_env; + std::unique_ptr _service; + std::vector> _paused_pools; + std::vector> _saved_config; +}; + +TEST_P(InternalServiceLoadWorkPoolTest, ControlRequestsBypassFullHeavyPools) { + ASSERT_TRUE(_service->_heavy_work_pool.try_offer([] {})); + ASSERT_TRUE(_service->_load_heavy_work_pool.try_offer([] {})); + + PTabletWriterOpenRequest open_request; + PTabletWriterOpenResult open_response; + PTabletWriterCancelRequest cancel_request; + PTabletWriterCancelResult cancel_response; + POpenLoadStreamRequest stream_request; + POpenLoadStreamResponse stream_response; + LoadRpcCountingClosure done; + auto submit = [&]() { + switch (GetParam()) { + case 0: + _service->tablet_writer_open(nullptr, &open_request, &open_response, &done); + break; + case 1: + _service->tablet_writer_cancel(nullptr, &cancel_request, &cancel_response, &done); + break; + case 2: + _service->open_load_stream(nullptr, &stream_request, &stream_response, &done); + break; + } + }; + + submit(); + EXPECT_EQ(done.calls.load(), 0); + EXPECT_EQ(_service->_load_light_work_pool.get_queue_size(), 1); + EXPECT_EQ(_service->_light_work_pool.get_queue_size(), 0); + + // A full control queue must invoke the closure exactly once and identify the right + // pool in responses that have a status (cancel's protobuf response is empty). + submit(); + EXPECT_EQ(done.calls.load(), 1); + if (GetParam() != 1) { + const auto& status = GetParam() == 0 ? open_response.status() : stream_response.status(); + EXPECT_EQ(status.status_code(), TStatusCode::CANCELLED); + ASSERT_EQ(status.error_msgs_size(), 1); + EXPECT_NE(status.error_msgs(0).find("brpc_load_light"), std::string::npos); + } +} + +INSTANTIATE_TEST_SUITE_P(LoadControl, InternalServiceLoadWorkPoolTest, testing::Values(0, 1, 2)); + +TEST_F(InternalServiceLoadWorkPoolTest, AddBlockUsesLoadHeavyPool) { + ASSERT_TRUE(_service->_heavy_work_pool.try_offer([] {})); + ASSERT_TRUE(_service->_load_light_work_pool.try_offer([] {})); + + PTabletWriterAddBlockRequest request; + PTabletWriterAddBlockResult response; + LoadRpcCountingClosure done; + _service->tablet_writer_add_block(nullptr, &request, &response, &done); + EXPECT_EQ(done.calls.load(), 0); + EXPECT_EQ(_service->_load_heavy_work_pool.get_queue_size(), 1); + EXPECT_EQ(_service->_light_work_pool.get_queue_size(), 0); + + _service->tablet_writer_add_block(nullptr, &request, &response, &done); + EXPECT_EQ(done.calls.load(), 1); + EXPECT_EQ(response.status().status_code(), TStatusCode::CANCELLED); + ASSERT_EQ(response.status().error_msgs_size(), 1); + EXPECT_NE(response.status().error_msgs(0).find("brpc_load_heavy"), std::string::npos); +} + +TEST_F(InternalServiceLoadWorkPoolTest, StreamingCloseUsesLoadHeavyPool) { + EXPECT_EQ(_exec_env.load_stream_mgr()->heavy_work_pool(), &_service->_load_heavy_work_pool); + EXPECT_NE(_exec_env.load_stream_mgr()->heavy_work_pool(), &_service->_heavy_work_pool); +} + +} // namespace doris From d2101811d59bbd63a3cfa3f0614a9efe5f7247e2 Mon Sep 17 00:00:00 2001 From: laihui <1353307710@qq.com> Date: Mon, 21 Sep 2026 16:52:14 +0800 Subject: [PATCH 2/5] [refactor](be) Reuse the existing heavy pool for load data work ### What problem does this PR solve? Problem Summary: Keep load open/cancel isolated with one dedicated load-light pool, while retaining the existing heavy pool for writes and streaming flush/close. Remove the redundant load-heavy pool and its settings and metrics, reducing the additional fixed worker count. Update dispatch and backpressure test coverage. ### Release note Only load control RPCs use the new load-light pool. Existing heavy-pool settings and metrics continue to apply to load data work. ### Check List (For Author) - Test: Updated unit tests; compilation and test execution skipped at user request. Formatting, build header hygiene, and git diff --check passed. - Behavior changed: Yes; load data work uses the original heavy pool. - Does this need documentation: No; configuration comments describe the new pool. --- be/src/common/config.cpp | 9 +- be/src/common/config.h | 6 +- be/src/service/internal_service.cpp | 82 +++++-------------- be/src/service/internal_service.h | 3 +- .../internal_service_load_work_pool_test.cpp | 20 ++--- 5 files changed, 33 insertions(+), 87 deletions(-) diff --git a/be/src/common/config.cpp b/be/src/common/config.cpp index d156d488d2992e..dd802c488ff7ec 100644 --- a/be/src/common/config.cpp +++ b/be/src/common/config.cpp @@ -656,16 +656,9 @@ DEFINE_mInt64(load_error_log_reserve_hours, "48"); // error log size limit, default 200MB DEFINE_mInt64(load_error_log_limit_bytes, "209715200"); -// Dedicated load pools. -1 keeps the existing heavy-pool capacity for load writes and -// uses CPU-scaled defaults for load control requests. These settings require a restart. -DEFINE_Int32(brpc_load_heavy_work_pool_threads, "-1"); -DEFINE_Int32(brpc_load_heavy_work_pool_max_queue_size, "-1"); +// Dedicated load control pool. -1 selects CPU-scaled defaults. Requires a restart. DEFINE_Int32(brpc_load_light_work_pool_threads, "-1"); DEFINE_Int32(brpc_load_light_work_pool_max_queue_size, "-1"); -DEFINE_Validator(brpc_load_heavy_work_pool_threads, - [](const int config) -> bool { return config == -1 || config > 0; }); -DEFINE_Validator(brpc_load_heavy_work_pool_max_queue_size, - [](const int config) -> bool { return config == -1 || config > 0; }); DEFINE_Validator(brpc_load_light_work_pool_threads, [](const int config) -> bool { return config == -1 || config > 0; }); DEFINE_Validator(brpc_load_light_work_pool_max_queue_size, diff --git a/be/src/common/config.h b/be/src/common/config.h index ca45230e448365..6dce07c9958153 100644 --- a/be/src/common/config.h +++ b/be/src/common/config.h @@ -739,14 +739,10 @@ DECLARE_mInt64(load_error_log_reserve_hours); // error log size limit, default 200MB DECLARE_mInt64(load_error_log_limit_bytes); -// Dedicated load data pool for add-block and streaming flush/close work. -// -1 inherits brpc_heavy_work_pool_threads/max_queue_size (including CPU-scaled defaults). -DECLARE_Int32(brpc_load_heavy_work_pool_threads); -DECLARE_Int32(brpc_load_heavy_work_pool_max_queue_size); // Dedicated load control pool for writer open/cancel and stream open. These handlers may // acquire locks or access storage, so they must not use the query light pool. // -1 selects max(32, CPU cores) threads and max(1024, CPU cores * 32) queued requests. -// All four load pool settings require a restart. +// Both load control pool settings require a restart. DECLARE_Int32(brpc_load_light_work_pool_threads); DECLARE_Int32(brpc_load_light_work_pool_max_queue_size); diff --git a/be/src/service/internal_service.cpp b/be/src/service/internal_service.cpp index 1a6d087e802e29..5eb6ffadc1c619 100644 --- a/be/src/service/internal_service.cpp +++ b/be/src/service/internal_service.cpp @@ -137,10 +137,6 @@ namespace doris { #include "common/compile_check_avoid_begin.h" using namespace ErrorCode; -DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(load_heavy_work_pool_queue_size, MetricUnit::NOUNIT); -DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(load_heavy_work_active_threads, MetricUnit::NOUNIT); -DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(load_heavy_work_pool_max_queue_size, MetricUnit::NOUNIT); -DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(load_heavy_work_max_threads, MetricUnit::NOUNIT); DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(load_light_work_pool_queue_size, MetricUnit::NOUNIT); DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(load_light_work_active_threads, MetricUnit::NOUNIT); DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(load_light_work_pool_max_queue_size, MetricUnit::NOUNIT); @@ -167,23 +163,6 @@ DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(arrow_flight_work_max_threads, MetricUnit::NO static bvar::LatencyRecorder g_process_remote_fetch_rowsets_latency("process_remote_fetch_rowsets"); -static int32_t resolved_brpc_load_heavy_work_pool_threads() { - if (config::brpc_load_heavy_work_pool_threads != -1) { - return config::brpc_load_heavy_work_pool_threads; - } - return config::brpc_heavy_work_pool_threads != -1 ? config::brpc_heavy_work_pool_threads - : std::max(128, CpuInfo::num_cores() * 4); -} - -static int32_t resolved_brpc_load_heavy_work_pool_max_queue_size() { - if (config::brpc_load_heavy_work_pool_max_queue_size != -1) { - return config::brpc_load_heavy_work_pool_max_queue_size; - } - return config::brpc_heavy_work_pool_max_queue_size != -1 - ? config::brpc_heavy_work_pool_max_queue_size - : std::max(10240, CpuInfo::num_cores() * 320); -} - static int32_t resolved_brpc_load_light_work_pool_threads() { return config::brpc_load_light_work_pool_threads != -1 ? config::brpc_load_light_work_pool_threads @@ -252,7 +231,7 @@ class NewHttpClosure : public ::google::protobuf::Closure { PInternalService::PInternalService(ExecEnv* exec_env) : _exec_env(exec_env), - // General RPCs that read disk or access the network. + // heavy threadpool is used for load process and other process that will read disk or access network. _heavy_work_pool(config::brpc_heavy_work_pool_threads != -1 ? config::brpc_heavy_work_pool_threads : std::max(128, CpuInfo::num_cores() * 4), @@ -260,9 +239,6 @@ PInternalService::PInternalService(ExecEnv* exec_env) ? config::brpc_heavy_work_pool_max_queue_size : std::max(10240, CpuInfo::num_cores() * 320), "brpc_heavy"), - _load_heavy_work_pool(resolved_brpc_load_heavy_work_pool_threads(), - resolved_brpc_load_heavy_work_pool_max_queue_size(), - "brpc_load_heavy"), // Open/cancel may block on storage or locks, but must not queue behind load writes. _load_light_work_pool(resolved_brpc_load_light_work_pool_threads(), resolved_brpc_load_light_work_pool_max_queue_size(), @@ -286,14 +262,6 @@ PInternalService::PInternalService(ExecEnv* exec_env) ? config::brpc_arrow_flight_work_pool_max_queue_size : std::max(20480, CpuInfo::num_cores() * 640), "brpc_arrow_flight") { - REGISTER_HOOK_METRIC(load_heavy_work_pool_queue_size, - [this]() { return _load_heavy_work_pool.get_queue_size(); }); - REGISTER_HOOK_METRIC(load_heavy_work_active_threads, - [this]() { return _load_heavy_work_pool.get_active_threads(); }); - REGISTER_HOOK_METRIC(load_heavy_work_pool_max_queue_size, - []() { return resolved_brpc_load_heavy_work_pool_max_queue_size(); }); - REGISTER_HOOK_METRIC(load_heavy_work_max_threads, - []() { return resolved_brpc_load_heavy_work_pool_threads(); }); REGISTER_HOOK_METRIC(load_light_work_pool_queue_size, [this]() { return _load_light_work_pool.get_queue_size(); }); REGISTER_HOOK_METRIC(load_light_work_active_threads, @@ -338,7 +306,7 @@ PInternalService::PInternalService(ExecEnv* exec_env) REGISTER_HOOK_METRIC(arrow_flight_work_max_threads, []() { return config::brpc_arrow_flight_work_pool_threads; }); - _exec_env->load_stream_mgr()->set_heavy_work_pool(&_load_heavy_work_pool); + _exec_env->load_stream_mgr()->set_heavy_work_pool(&_heavy_work_pool); CHECK_EQ(0, bthread_key_create(&AsyncIO::btls_io_ctx_key, AsyncIO::io_ctx_key_deleter)); } @@ -349,10 +317,6 @@ PInternalServiceImpl::PInternalServiceImpl(StorageEngine& engine, ExecEnv* exec_ PInternalServiceImpl::~PInternalServiceImpl() = default; PInternalService::~PInternalService() { - DEREGISTER_HOOK_METRIC(load_heavy_work_pool_queue_size); - DEREGISTER_HOOK_METRIC(load_heavy_work_active_threads); - DEREGISTER_HOOK_METRIC(load_heavy_work_pool_max_queue_size); - DEREGISTER_HOOK_METRIC(load_heavy_work_max_threads); DEREGISTER_HOOK_METRIC(load_light_work_pool_queue_size); DEREGISTER_HOOK_METRIC(load_light_work_active_threads); DEREGISTER_HOOK_METRIC(load_light_work_pool_max_queue_size); @@ -578,29 +542,27 @@ void PInternalService::tablet_writer_add_block(google::protobuf::RpcController* PTabletWriterAddBlockResult* response, google::protobuf::Closure* done) { int64_t submit_task_time_ns = MonotonicNanos(); - bool ret = - _load_heavy_work_pool.try_offer([request, response, done, submit_task_time_ns, this]() { - int64_t wait_execution_time_ns = MonotonicNanos() - submit_task_time_ns; - brpc::ClosureGuard closure_guard(done); - int64_t execution_time_ns = 0; - { - SCOPED_RAW_TIMER(&execution_time_ns); - signal::SignalTaskIdKeeper keeper(request->id()); - auto st = _exec_env->load_channel_mgr()->add_batch(*request, response); - if (!st.ok()) { - LOG(WARNING) - << "tablet writer add block failed, message=" << st - << ", id=" << request->id() << ", index_id=" << request->index_id() - << ", sender_id=" << request->sender_id() - << ", backend id=" << request->backend_id(); - } - st.to_protobuf(response->mutable_status()); - } - response->set_execution_time_us(execution_time_ns / NANOS_PER_MICRO); - response->set_wait_execution_time_us(wait_execution_time_ns / NANOS_PER_MICRO); - }); + bool ret = _heavy_work_pool.try_offer([request, response, done, submit_task_time_ns, this]() { + int64_t wait_execution_time_ns = MonotonicNanos() - submit_task_time_ns; + brpc::ClosureGuard closure_guard(done); + int64_t execution_time_ns = 0; + { + SCOPED_RAW_TIMER(&execution_time_ns); + signal::SignalTaskIdKeeper keeper(request->id()); + auto st = _exec_env->load_channel_mgr()->add_batch(*request, response); + if (!st.ok()) { + LOG(WARNING) << "tablet writer add block failed, message=" << st + << ", id=" << request->id() << ", index_id=" << request->index_id() + << ", sender_id=" << request->sender_id() + << ", backend id=" << request->backend_id(); + } + st.to_protobuf(response->mutable_status()); + } + response->set_execution_time_us(execution_time_ns / NANOS_PER_MICRO); + response->set_wait_execution_time_us(wait_execution_time_ns / NANOS_PER_MICRO); + }); if (!ret) { - offer_failed(response, done, _load_heavy_work_pool); + offer_failed(response, done, _heavy_work_pool); return; } } diff --git a/be/src/service/internal_service.h b/be/src/service/internal_service.h index 4422c5d5225abe..414e84b31273ee 100644 --- a/be/src/service/internal_service.h +++ b/be/src/service/internal_service.h @@ -278,8 +278,7 @@ class PInternalService : public PBackendService { // define the interface for reading and writing data as heavy interface // otherwise as light interface FifoThreadPool _heavy_work_pool; - // Keep load control requests runnable when load writes or closes saturate their pool. - FifoThreadPool _load_heavy_work_pool; + // Keep load control requests runnable when writes or closes saturate the heavy pool. FifoThreadPool _load_light_work_pool; FifoThreadPool _peer_fetch_pool; FifoThreadPool _light_work_pool; diff --git a/be/test/service/internal_service_load_work_pool_test.cpp b/be/test/service/internal_service_load_work_pool_test.cpp index c4e9d5fb23a02e..445ad1b73c3238 100644 --- a/be/test/service/internal_service_load_work_pool_test.cpp +++ b/be/test/service/internal_service_load_work_pool_test.cpp @@ -75,8 +75,6 @@ class InternalServiceLoadWorkPoolTest : public testing::TestWithParam { &config::brpc_peer_fetch_pool_threads, &config::brpc_peer_fetch_pool_max_queue_size, &config::brpc_arrow_flight_work_pool_threads, &config::brpc_arrow_flight_work_pool_max_queue_size, - &config::brpc_load_heavy_work_pool_threads, - &config::brpc_load_heavy_work_pool_max_queue_size, &config::brpc_load_light_work_pool_threads, &config::brpc_load_light_work_pool_max_queue_size}) { _saved_config.emplace_back(setting, *setting); @@ -85,7 +83,7 @@ class InternalServiceLoadWorkPoolTest : public testing::TestWithParam { _exec_env._load_stream_mgr = std::make_unique(1); _service = std::make_unique(&_exec_env); for (auto* pool : {&_service->_heavy_work_pool, &_service->_light_work_pool, - &_service->_load_heavy_work_pool, &_service->_load_light_work_pool}) { + &_service->_load_light_work_pool}) { _paused_pools.push_back(std::make_unique(*pool)); } } @@ -106,9 +104,8 @@ class InternalServiceLoadWorkPoolTest : public testing::TestWithParam { std::vector> _saved_config; }; -TEST_P(InternalServiceLoadWorkPoolTest, ControlRequestsBypassFullHeavyPools) { +TEST_P(InternalServiceLoadWorkPoolTest, ControlRequestsBypassFullHeavyPool) { ASSERT_TRUE(_service->_heavy_work_pool.try_offer([] {})); - ASSERT_TRUE(_service->_load_heavy_work_pool.try_offer([] {})); PTabletWriterOpenRequest open_request; PTabletWriterOpenResult open_response; @@ -150,8 +147,7 @@ TEST_P(InternalServiceLoadWorkPoolTest, ControlRequestsBypassFullHeavyPools) { INSTANTIATE_TEST_SUITE_P(LoadControl, InternalServiceLoadWorkPoolTest, testing::Values(0, 1, 2)); -TEST_F(InternalServiceLoadWorkPoolTest, AddBlockUsesLoadHeavyPool) { - ASSERT_TRUE(_service->_heavy_work_pool.try_offer([] {})); +TEST_F(InternalServiceLoadWorkPoolTest, AddBlockKeepsUsingHeavyPool) { ASSERT_TRUE(_service->_load_light_work_pool.try_offer([] {})); PTabletWriterAddBlockRequest request; @@ -159,19 +155,19 @@ TEST_F(InternalServiceLoadWorkPoolTest, AddBlockUsesLoadHeavyPool) { LoadRpcCountingClosure done; _service->tablet_writer_add_block(nullptr, &request, &response, &done); EXPECT_EQ(done.calls.load(), 0); - EXPECT_EQ(_service->_load_heavy_work_pool.get_queue_size(), 1); + EXPECT_EQ(_service->_heavy_work_pool.get_queue_size(), 1); EXPECT_EQ(_service->_light_work_pool.get_queue_size(), 0); _service->tablet_writer_add_block(nullptr, &request, &response, &done); EXPECT_EQ(done.calls.load(), 1); EXPECT_EQ(response.status().status_code(), TStatusCode::CANCELLED); ASSERT_EQ(response.status().error_msgs_size(), 1); - EXPECT_NE(response.status().error_msgs(0).find("brpc_load_heavy"), std::string::npos); + EXPECT_NE(response.status().error_msgs(0).find("brpc_heavy"), std::string::npos); } -TEST_F(InternalServiceLoadWorkPoolTest, StreamingCloseUsesLoadHeavyPool) { - EXPECT_EQ(_exec_env.load_stream_mgr()->heavy_work_pool(), &_service->_load_heavy_work_pool); - EXPECT_NE(_exec_env.load_stream_mgr()->heavy_work_pool(), &_service->_heavy_work_pool); +TEST_F(InternalServiceLoadWorkPoolTest, StreamingCloseKeepsUsingHeavyPool) { + EXPECT_EQ(_exec_env.load_stream_mgr()->heavy_work_pool(), &_service->_heavy_work_pool); + EXPECT_NE(_exec_env.load_stream_mgr()->heavy_work_pool(), &_service->_load_light_work_pool); } } // namespace doris From e9bd049f7f6a5a30133675ed90865bf696ce6ce3 Mon Sep 17 00:00:00 2001 From: laihui <1353307710@qq.com> Date: Mon, 21 Sep 2026 16:59:39 +0800 Subject: [PATCH 3/5] [refactor](be) Restrict the load-light pool to 32 cancellation workers ### What problem does this PR solve? Problem Summary: Open requests can block on locks or metadata RPCs and exhaust workers shared with cancellation. Route only tablet_writer_cancel to load-light, restore both open handlers to the existing heavy pool, and fix the cancellation worker count at 32. Keep queue capacity configurable and update metrics and tests. ### Release note Only tablet writer cancellation uses the dedicated 32-thread pool. Load open, write, and close continue using the existing heavy pool. ### Check List (For Author) - Test: Updated unit tests; compilation and execution skipped at user request. Clang-format 16, build header hygiene, and git diff --check passed. - Behavior changed: Yes; only cancellation moves to a fixed 32-thread pool. - Does this need documentation: No; code/config comments describe the settings. --- be/src/common/config.cpp | 6 +- be/src/common/config.h | 7 +- be/src/service/internal_service.cpp | 20 ++--- be/src/service/internal_service.h | 2 +- .../internal_service_load_work_pool_test.cpp | 85 +++++++++---------- 5 files changed, 55 insertions(+), 65 deletions(-) diff --git a/be/src/common/config.cpp b/be/src/common/config.cpp index dd802c488ff7ec..f4e1b551f40ace 100644 --- a/be/src/common/config.cpp +++ b/be/src/common/config.cpp @@ -656,11 +656,9 @@ DEFINE_mInt64(load_error_log_reserve_hours, "48"); // error log size limit, default 200MB DEFINE_mInt64(load_error_log_limit_bytes, "209715200"); -// Dedicated load control pool. -1 selects CPU-scaled defaults. Requires a restart. -DEFINE_Int32(brpc_load_light_work_pool_threads, "-1"); +// Queue for the fixed 32-thread load cancellation pool. -1 selects a CPU-scaled default. +// Requires a restart. DEFINE_Int32(brpc_load_light_work_pool_max_queue_size, "-1"); -DEFINE_Validator(brpc_load_light_work_pool_threads, - [](const int config) -> bool { return config == -1 || config > 0; }); DEFINE_Validator(brpc_load_light_work_pool_max_queue_size, [](const int config) -> bool { return config == -1 || config > 0; }); diff --git a/be/src/common/config.h b/be/src/common/config.h index 6dce07c9958153..2c79e01f84f537 100644 --- a/be/src/common/config.h +++ b/be/src/common/config.h @@ -739,11 +739,8 @@ DECLARE_mInt64(load_error_log_reserve_hours); // error log size limit, default 200MB DECLARE_mInt64(load_error_log_limit_bytes); -// Dedicated load control pool for writer open/cancel and stream open. These handlers may -// acquire locks or access storage, so they must not use the query light pool. -// -1 selects max(32, CPU cores) threads and max(1024, CPU cores * 32) queued requests. -// Both load control pool settings require a restart. -DECLARE_Int32(brpc_load_light_work_pool_threads); +// Queue for the dedicated load cancellation pool, which has a fixed 32 threads. +// -1 selects max(1024, CPU cores * 32) queued requests. Requires a restart. DECLARE_Int32(brpc_load_light_work_pool_max_queue_size); // be brpc interface is classified into two categories: light and heavy diff --git a/be/src/service/internal_service.cpp b/be/src/service/internal_service.cpp index 5eb6ffadc1c619..8f4fe61b7cf8ec 100644 --- a/be/src/service/internal_service.cpp +++ b/be/src/service/internal_service.cpp @@ -163,11 +163,7 @@ DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(arrow_flight_work_max_threads, MetricUnit::NO static bvar::LatencyRecorder g_process_remote_fetch_rowsets_latency("process_remote_fetch_rowsets"); -static int32_t resolved_brpc_load_light_work_pool_threads() { - return config::brpc_load_light_work_pool_threads != -1 - ? config::brpc_load_light_work_pool_threads - : std::max(32, CpuInfo::num_cores()); -} +static constexpr int32_t LOAD_LIGHT_WORK_POOL_THREADS = 32; static int32_t resolved_brpc_load_light_work_pool_max_queue_size() { return config::brpc_load_light_work_pool_max_queue_size != -1 @@ -239,8 +235,8 @@ PInternalService::PInternalService(ExecEnv* exec_env) ? config::brpc_heavy_work_pool_max_queue_size : std::max(10240, CpuInfo::num_cores() * 320), "brpc_heavy"), - // Open/cancel may block on storage or locks, but must not queue behind load writes. - _load_light_work_pool(resolved_brpc_load_light_work_pool_threads(), + // Keep cancellation dispatch independent of potentially blocking opens and writes. + _load_light_work_pool(LOAD_LIGHT_WORK_POOL_THREADS, resolved_brpc_load_light_work_pool_max_queue_size(), "brpc_load_light"), // peer fetch threadpool isolates fetch_peer_data from heavy load traffic to avoid peer reads starving imports. @@ -269,7 +265,7 @@ PInternalService::PInternalService(ExecEnv* exec_env) REGISTER_HOOK_METRIC(load_light_work_pool_max_queue_size, []() { return resolved_brpc_load_light_work_pool_max_queue_size(); }); REGISTER_HOOK_METRIC(load_light_work_max_threads, - []() { return resolved_brpc_load_light_work_pool_threads(); }); + []() { return LOAD_LIGHT_WORK_POOL_THREADS; }); REGISTER_HOOK_METRIC(heavy_work_pool_queue_size, [this]() { return _heavy_work_pool.get_queue_size(); }); @@ -348,7 +344,7 @@ void PInternalService::tablet_writer_open(google::protobuf::RpcController* contr const PTabletWriterOpenRequest* request, PTabletWriterOpenResult* response, google::protobuf::Closure* done) { - bool ret = _load_light_work_pool.try_offer([this, request, response, done]() { + bool ret = _heavy_work_pool.try_offer([this, request, response, done]() { VLOG_RPC << "tablet writer open, id=" << request->id() << ", index_id=" << request->index_id() << ", txn_id=" << request->txn_id(); signal::SignalTaskIdKeeper keeper(request->id()); @@ -362,7 +358,7 @@ void PInternalService::tablet_writer_open(google::protobuf::RpcController* contr st.to_protobuf(response->mutable_status()); }); if (!ret) { - offer_failed(response, done, _load_light_work_pool); + offer_failed(response, done, _heavy_work_pool); return; } } @@ -458,7 +454,7 @@ void PInternalService::open_load_stream(google::protobuf::RpcController* control const POpenLoadStreamRequest* request, POpenLoadStreamResponse* response, google::protobuf::Closure* done) { - bool ret = _load_light_work_pool.try_offer([this, controller, request, response, done]() { + bool ret = _heavy_work_pool.try_offer([this, controller, request, response, done]() { signal::SignalTaskIdKeeper keeper(request->load_id()); brpc::ClosureGuard done_guard(done); brpc::Controller* cntl = static_cast(controller); @@ -516,7 +512,7 @@ void PInternalService::open_load_stream(google::protobuf::RpcController* control st.to_protobuf(response->mutable_status()); }); if (!ret) { - offer_failed(response, done, _load_light_work_pool); + offer_failed(response, done, _heavy_work_pool); } } diff --git a/be/src/service/internal_service.h b/be/src/service/internal_service.h index 414e84b31273ee..db5cd4f31163ff 100644 --- a/be/src/service/internal_service.h +++ b/be/src/service/internal_service.h @@ -278,7 +278,7 @@ class PInternalService : public PBackendService { // define the interface for reading and writing data as heavy interface // otherwise as light interface FifoThreadPool _heavy_work_pool; - // Keep load control requests runnable when writes or closes saturate the heavy pool. + // Dedicated 32-thread pool for cancellation; open/write/close use the heavy pool. FifoThreadPool _load_light_work_pool; FifoThreadPool _peer_fetch_pool; FifoThreadPool _light_work_pool; diff --git a/be/test/service/internal_service_load_work_pool_test.cpp b/be/test/service/internal_service_load_work_pool_test.cpp index 445ad1b73c3238..9fdf230e8b3158 100644 --- a/be/test/service/internal_service_load_work_pool_test.cpp +++ b/be/test/service/internal_service_load_work_pool_test.cpp @@ -31,19 +31,21 @@ namespace doris { namespace { -// Hold the sole worker so routing and queue rejection can be checked without running +// Hold every worker so routing and queue rejection can be checked without running // storage handlers. Discard queued RPCs before releasing the worker at teardown. class PausedLoadRpcPool { public: explicit PausedLoadRpcPool(FifoThreadPool& pool) : _pool(pool) { auto resume = _resume.get_future().share(); - auto started = std::make_shared>(); - auto ready = started->get_future(); - CHECK(_pool.try_offer([started, resume]() { - started->set_value(); - resume.wait(); - })); - ready.wait(); + for (size_t i = 0; i < _pool._threads.size(); ++i) { + auto started = std::make_shared>(); + auto ready = started->get_future(); + CHECK(_pool.try_offer([started, resume]() { + started->set_value(); + resume.wait(); + })); + ready.wait(); + } } ~PausedLoadRpcPool() { @@ -68,14 +70,13 @@ class LoadRpcCountingClosure : public google::protobuf::Closure { class InternalServiceLoadWorkPoolTest : public testing::TestWithParam { protected: void SetUp() override { - // Keep all pools small and restore global configuration after each test. + // Keep configurable pools and queues small; cancellation retains its fixed 32 workers. for (auto* setting : {&config::brpc_heavy_work_pool_threads, &config::brpc_heavy_work_pool_max_queue_size, &config::brpc_light_work_pool_threads, &config::brpc_light_work_pool_max_queue_size, &config::brpc_peer_fetch_pool_threads, &config::brpc_peer_fetch_pool_max_queue_size, &config::brpc_arrow_flight_work_pool_threads, &config::brpc_arrow_flight_work_pool_max_queue_size, - &config::brpc_load_light_work_pool_threads, &config::brpc_load_light_work_pool_max_queue_size}) { _saved_config.emplace_back(setting, *setting); *setting = 1; @@ -104,15 +105,32 @@ class InternalServiceLoadWorkPoolTest : public testing::TestWithParam { std::vector> _saved_config; }; -TEST_P(InternalServiceLoadWorkPoolTest, ControlRequestsBypassFullHeavyPool) { +TEST_F(InternalServiceLoadWorkPoolTest, CancelBypassesFullHeavyPool) { + EXPECT_EQ(_service->_load_light_work_pool.get_active_threads(), 32); ASSERT_TRUE(_service->_heavy_work_pool.try_offer([] {})); + PTabletWriterCancelRequest request; + PTabletWriterCancelResult response; + LoadRpcCountingClosure done; + _service->tablet_writer_cancel(nullptr, &request, &response, &done); + EXPECT_EQ(done.calls.load(), 0); + EXPECT_EQ(_service->_load_light_work_pool.get_queue_size(), 1); + EXPECT_EQ(_service->_light_work_pool.get_queue_size(), 0); + + // Cancel's protobuf response is empty; queue rejection must still run the closure once. + _service->tablet_writer_cancel(nullptr, &request, &response, &done); + EXPECT_EQ(done.calls.load(), 1); +} + +TEST_P(InternalServiceLoadWorkPoolTest, OpenAndAddBlockKeepUsingHeavyPool) { + ASSERT_TRUE(_service->_load_light_work_pool.try_offer([] {})); + PTabletWriterOpenRequest open_request; PTabletWriterOpenResult open_response; - PTabletWriterCancelRequest cancel_request; - PTabletWriterCancelResult cancel_response; POpenLoadStreamRequest stream_request; POpenLoadStreamResponse stream_response; + PTabletWriterAddBlockRequest block_request; + PTabletWriterAddBlockResult block_response; LoadRpcCountingClosure done; auto submit = [&]() { switch (GetParam()) { @@ -120,50 +138,31 @@ TEST_P(InternalServiceLoadWorkPoolTest, ControlRequestsBypassFullHeavyPool) { _service->tablet_writer_open(nullptr, &open_request, &open_response, &done); break; case 1: - _service->tablet_writer_cancel(nullptr, &cancel_request, &cancel_response, &done); + _service->open_load_stream(nullptr, &stream_request, &stream_response, &done); break; case 2: - _service->open_load_stream(nullptr, &stream_request, &stream_response, &done); + _service->tablet_writer_add_block(nullptr, &block_request, &block_response, &done); break; } }; submit(); EXPECT_EQ(done.calls.load(), 0); - EXPECT_EQ(_service->_load_light_work_pool.get_queue_size(), 1); + EXPECT_EQ(_service->_heavy_work_pool.get_queue_size(), 1); EXPECT_EQ(_service->_light_work_pool.get_queue_size(), 0); - // A full control queue must invoke the closure exactly once and identify the right - // pool in responses that have a status (cancel's protobuf response is empty). submit(); EXPECT_EQ(done.calls.load(), 1); - if (GetParam() != 1) { - const auto& status = GetParam() == 0 ? open_response.status() : stream_response.status(); - EXPECT_EQ(status.status_code(), TStatusCode::CANCELLED); - ASSERT_EQ(status.error_msgs_size(), 1); - EXPECT_NE(status.error_msgs(0).find("brpc_load_light"), std::string::npos); - } + const auto& status = GetParam() == 0 ? open_response.status() + : GetParam() == 1 ? stream_response.status() + : block_response.status(); + EXPECT_EQ(status.status_code(), TStatusCode::CANCELLED); + ASSERT_EQ(status.error_msgs_size(), 1); + EXPECT_NE(status.error_msgs(0).find("brpc_heavy"), std::string::npos); } -INSTANTIATE_TEST_SUITE_P(LoadControl, InternalServiceLoadWorkPoolTest, testing::Values(0, 1, 2)); - -TEST_F(InternalServiceLoadWorkPoolTest, AddBlockKeepsUsingHeavyPool) { - ASSERT_TRUE(_service->_load_light_work_pool.try_offer([] {})); - - PTabletWriterAddBlockRequest request; - PTabletWriterAddBlockResult response; - LoadRpcCountingClosure done; - _service->tablet_writer_add_block(nullptr, &request, &response, &done); - EXPECT_EQ(done.calls.load(), 0); - EXPECT_EQ(_service->_heavy_work_pool.get_queue_size(), 1); - EXPECT_EQ(_service->_light_work_pool.get_queue_size(), 0); - - _service->tablet_writer_add_block(nullptr, &request, &response, &done); - EXPECT_EQ(done.calls.load(), 1); - EXPECT_EQ(response.status().status_code(), TStatusCode::CANCELLED); - ASSERT_EQ(response.status().error_msgs_size(), 1); - EXPECT_NE(response.status().error_msgs(0).find("brpc_heavy"), std::string::npos); -} +INSTANTIATE_TEST_SUITE_P(HeavyLoadRequests, InternalServiceLoadWorkPoolTest, + testing::Values(0, 1, 2)); TEST_F(InternalServiceLoadWorkPoolTest, StreamingCloseKeepsUsingHeavyPool) { EXPECT_EQ(_exec_env.load_stream_mgr()->heavy_work_pool(), &_service->_heavy_work_pool); From bd3fb0821efa35dc4c247bbe901b6279382767d8 Mon Sep 17 00:00:00 2001 From: laihui <1353307710@qq.com> Date: Mon, 21 Sep 2026 17:17:32 +0800 Subject: [PATCH 4/5] [fix](be) Make load cancellation worker count configurable ### What problem does this PR solve? Problem Summary: Replace the hard-coded cancellation worker count with brpc_load_light_work_pool_threads, defaulting to 32. Validate that the value is positive, use it when constructing the pool, and expose it in the capacity metric. Update unit test coverage to use a non-default worker count. ### Release note Allow configuring the cancellation pool worker count, default 32. A BE restart is required for changes to take effect. ### Check List (For Author) - Test: Unit tests updated but not run; compilation skipped at user request. Clang-format 16, header hygiene, and git diff --check passed. - Behavior changed: Yes; cancellation worker count is configurable. - Does this need documentation: No separate change; defaults and restart requirements are documented in configuration comments. --- be/src/common/config.cpp | 7 +++++-- be/src/common/config.h | 4 +++- be/src/service/internal_service.cpp | 6 ++---- be/src/service/internal_service.h | 2 +- be/test/service/internal_service_load_work_pool_test.cpp | 7 +++++-- 5 files changed, 16 insertions(+), 10 deletions(-) diff --git a/be/src/common/config.cpp b/be/src/common/config.cpp index f4e1b551f40ace..be97fb721a8678 100644 --- a/be/src/common/config.cpp +++ b/be/src/common/config.cpp @@ -656,8 +656,11 @@ DEFINE_mInt64(load_error_log_reserve_hours, "48"); // error log size limit, default 200MB DEFINE_mInt64(load_error_log_limit_bytes, "209715200"); -// Queue for the fixed 32-thread load cancellation pool. -1 selects a CPU-scaled default. -// Requires a restart. +// Dedicated load cancellation workers. Requires a restart. +DEFINE_Int32(brpc_load_light_work_pool_threads, "32"); +DEFINE_Validator(brpc_load_light_work_pool_threads, + [](const int config) -> bool { return config > 0; }); +// Queue capacity: -1 selects a CPU-scaled default. Requires a restart. DEFINE_Int32(brpc_load_light_work_pool_max_queue_size, "-1"); DEFINE_Validator(brpc_load_light_work_pool_max_queue_size, [](const int config) -> bool { return config == -1 || config > 0; }); diff --git a/be/src/common/config.h b/be/src/common/config.h index 2c79e01f84f537..fb25af4e86c47b 100644 --- a/be/src/common/config.h +++ b/be/src/common/config.h @@ -739,7 +739,9 @@ DECLARE_mInt64(load_error_log_reserve_hours); // error log size limit, default 200MB DECLARE_mInt64(load_error_log_limit_bytes); -// Queue for the dedicated load cancellation pool, which has a fixed 32 threads. +// Dedicated load cancellation workers, default 32. Must be positive; requires a restart. +DECLARE_Int32(brpc_load_light_work_pool_threads); +// Queue capacity for the dedicated load cancellation pool. // -1 selects max(1024, CPU cores * 32) queued requests. Requires a restart. DECLARE_Int32(brpc_load_light_work_pool_max_queue_size); diff --git a/be/src/service/internal_service.cpp b/be/src/service/internal_service.cpp index 8f4fe61b7cf8ec..eeed9f230d164a 100644 --- a/be/src/service/internal_service.cpp +++ b/be/src/service/internal_service.cpp @@ -163,8 +163,6 @@ DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(arrow_flight_work_max_threads, MetricUnit::NO static bvar::LatencyRecorder g_process_remote_fetch_rowsets_latency("process_remote_fetch_rowsets"); -static constexpr int32_t LOAD_LIGHT_WORK_POOL_THREADS = 32; - static int32_t resolved_brpc_load_light_work_pool_max_queue_size() { return config::brpc_load_light_work_pool_max_queue_size != -1 ? config::brpc_load_light_work_pool_max_queue_size @@ -236,7 +234,7 @@ PInternalService::PInternalService(ExecEnv* exec_env) : std::max(10240, CpuInfo::num_cores() * 320), "brpc_heavy"), // Keep cancellation dispatch independent of potentially blocking opens and writes. - _load_light_work_pool(LOAD_LIGHT_WORK_POOL_THREADS, + _load_light_work_pool(config::brpc_load_light_work_pool_threads, resolved_brpc_load_light_work_pool_max_queue_size(), "brpc_load_light"), // peer fetch threadpool isolates fetch_peer_data from heavy load traffic to avoid peer reads starving imports. @@ -265,7 +263,7 @@ PInternalService::PInternalService(ExecEnv* exec_env) REGISTER_HOOK_METRIC(load_light_work_pool_max_queue_size, []() { return resolved_brpc_load_light_work_pool_max_queue_size(); }); REGISTER_HOOK_METRIC(load_light_work_max_threads, - []() { return LOAD_LIGHT_WORK_POOL_THREADS; }); + []() { return config::brpc_load_light_work_pool_threads; }); REGISTER_HOOK_METRIC(heavy_work_pool_queue_size, [this]() { return _heavy_work_pool.get_queue_size(); }); diff --git a/be/src/service/internal_service.h b/be/src/service/internal_service.h index db5cd4f31163ff..8be21ea2a4067c 100644 --- a/be/src/service/internal_service.h +++ b/be/src/service/internal_service.h @@ -278,7 +278,7 @@ class PInternalService : public PBackendService { // define the interface for reading and writing data as heavy interface // otherwise as light interface FifoThreadPool _heavy_work_pool; - // Dedicated 32-thread pool for cancellation; open/write/close use the heavy pool. + // Dedicated pool for cancellation; open/write/close use the heavy pool. FifoThreadPool _load_light_work_pool; FifoThreadPool _peer_fetch_pool; FifoThreadPool _light_work_pool; diff --git a/be/test/service/internal_service_load_work_pool_test.cpp b/be/test/service/internal_service_load_work_pool_test.cpp index 9fdf230e8b3158..429cfccbd57e3d 100644 --- a/be/test/service/internal_service_load_work_pool_test.cpp +++ b/be/test/service/internal_service_load_work_pool_test.cpp @@ -70,17 +70,20 @@ class LoadRpcCountingClosure : public google::protobuf::Closure { class InternalServiceLoadWorkPoolTest : public testing::TestWithParam { protected: void SetUp() override { - // Keep configurable pools and queues small; cancellation retains its fixed 32 workers. + // Keep pools and queues small and restore configuration after each test. for (auto* setting : {&config::brpc_heavy_work_pool_threads, &config::brpc_heavy_work_pool_max_queue_size, &config::brpc_light_work_pool_threads, &config::brpc_light_work_pool_max_queue_size, &config::brpc_peer_fetch_pool_threads, &config::brpc_peer_fetch_pool_max_queue_size, &config::brpc_arrow_flight_work_pool_threads, &config::brpc_arrow_flight_work_pool_max_queue_size, + &config::brpc_load_light_work_pool_threads, &config::brpc_load_light_work_pool_max_queue_size}) { _saved_config.emplace_back(setting, *setting); *setting = 1; } + // Use a non-default value to verify that the cancellation pool honors configuration. + config::brpc_load_light_work_pool_threads = 3; _exec_env._load_stream_mgr = std::make_unique(1); _service = std::make_unique(&_exec_env); for (auto* pool : {&_service->_heavy_work_pool, &_service->_light_work_pool, @@ -106,7 +109,7 @@ class InternalServiceLoadWorkPoolTest : public testing::TestWithParam { }; TEST_F(InternalServiceLoadWorkPoolTest, CancelBypassesFullHeavyPool) { - EXPECT_EQ(_service->_load_light_work_pool.get_active_threads(), 32); + EXPECT_EQ(_service->_load_light_work_pool.get_active_threads(), 3); ASSERT_TRUE(_service->_heavy_work_pool.try_offer([] {})); PTabletWriterCancelRequest request; From 78c35feec7578bbb0064facdf53255dbcc66068f Mon Sep 17 00:00:00 2001 From: laihui <1353307710@qq.com> Date: Mon, 21 Sep 2026 21:01:42 +0800 Subject: [PATCH 5/5] [fix](be) Declare load cancellation pool metrics ### What problem does this PR solve? Problem Summary: REGISTER_HOOK_METRIC accesses a corresponding UIntGauge member on DorisMetrics. The load-light pool registered four new metrics without adding these members, causing internal_service.cpp to fail to compile with eight missing-member diagnostics in TeamCity build 1053775. Declare the four gauge pointers alongside the existing BRPC pool metrics. ### Release note None ### Check List (For Author) - Test: Matched all four missing members from the CI log to their declarations, metric prototypes, registrations and deregistrations. Clang-format 16, header hygiene and git diff --check passed. Local build and test execution skipped at the user's request; a passing CI rebuild has not yet been verified. - Behavior changed: No; complete the declarations required by existing hooks. - Does this need documentation: No. --- be/src/common/metrics/doris_metrics.h | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/be/src/common/metrics/doris_metrics.h b/be/src/common/metrics/doris_metrics.h index 852cdb28753e6f..971a55051b4b39 100644 --- a/be/src/common/metrics/doris_metrics.h +++ b/be/src/common/metrics/doris_metrics.h @@ -249,6 +249,11 @@ class DorisMetrics { IntCounter* upload_rowset_count = nullptr; IntCounter* upload_fail_count = nullptr; + UIntGauge* load_light_work_pool_queue_size = nullptr; + UIntGauge* load_light_work_active_threads = nullptr; + UIntGauge* load_light_work_pool_max_queue_size = nullptr; + UIntGauge* load_light_work_max_threads = nullptr; + UIntGauge* light_work_pool_queue_size = nullptr; UIntGauge* heavy_work_pool_queue_size = nullptr; UIntGauge* peer_fetch_work_pool_queue_size = nullptr;