This is an automated email from the ASF dual-hosted git repository.

sollhui pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/master by this push:
     new ce233bc1f62 [enhance](load) Isolate load cancellation in a dedicated 
thread pool (#68328)
ce233bc1f62 is described below

commit ce233bc1f62b4d105530552e773fab792f92ae3a
Author: hui lai <[email protected]>
AuthorDate: Mon Sep 28 16:17:15 2026 +0800

    [enhance](load) Isolate load cancellation in a dedicated thread pool 
(#68328)
    
    ### What problem does this PR solve?
    
    `tablet_writer_cancel` shares the BRPC heavy work pool with load open,
    write, and close requests. When these operations saturate the pool,
    cancellation requests also queue behind them, delaying cancellation and
    resource release.
    
    Dispatch `tablet_writer_cancel` through a dedicated `brpc_load_light`
    pool. Writer open and stream open remain in the heavy pool because they
    can block on locks or metadata RPCs; writes and closes also keep their
    existing scheduling. This isolates cancellation dispatch from heavy-pool
    congestion, while cancellation still follows the existing locking and
    cleanup logic.
    
    The new pool adds queue-size, active-thread, and capacity metrics.
    `brpc_load_light_work_pool_threads` controls its worker count, defaults
    to `32`, and accepts positive values. Its queue capacity is controlled
    by `brpc_load_light_work_pool_max_queue_size`: `-1` selects `max(1024,
    CPU cores * 32)`, and a positive value sets an explicit limit. Both
    settings require a BE restart.
    
    ### Release note
    
    Isolate tablet writer cancellation RPCs in a dedicated thread pool to
    reduce cancellation delays caused by a saturated BRPC heavy work pool.
    
    ### Check List (For Author)
    
    - Test
    - [x] Unit Test: added coverage for a non-default worker count,
    cancellation routing, and queue-full completion, plus checks that open,
    add-block, and streaming close retain the heavy pool. Tests have not
    been run.
    - Passed: clang-format 16, build header hygiene checks, and `git diff
    --check`.
        - BE compilation and runtime tests have not been completed.
    - Behavior changed:
        - [x] Yes. Tablet writer cancellation uses a dedicated pool.
    - Does this need documentation?
    - [x] No separate documentation change. Pool sizing and queue
    configuration are documented in code comments.
---
 be/src/common/config.cpp                           |   9 ++
 be/src/common/config.h                             |   6 +
 be/src/common/metrics/doris_metrics.h              |   5 +
 be/src/service/internal_service.cpp                |  33 +++-
 be/src/service/internal_service.h                  |   2 +
 .../internal_service_load_work_pool_test.cpp       | 175 +++++++++++++++++++++
 6 files changed, 228 insertions(+), 2 deletions(-)

diff --git a/be/src/common/config.cpp b/be/src/common/config.cpp
index 7f857b812e6..85763ed7a92 100644
--- a/be/src/common/config.cpp
+++ b/be/src/common/config.cpp
@@ -658,6 +658,15 @@ DEFINE_mInt64(load_error_log_reserve_hours, "48");
 // error log size limit, default 200MB
 DEFINE_mInt64(load_error_log_limit_bytes, "209715200");
 
+// 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; });
+
 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 058e05f9293..63020e990a8 100644
--- a/be/src/common/config.h
+++ b/be/src/common/config.h
@@ -739,6 +739,12 @@ DECLARE_mInt64(load_error_log_reserve_hours);
 // error log size limit, default 200MB
 DECLARE_mInt64(load_error_log_limit_bytes);
 
+// 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);
+
 // 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/common/metrics/doris_metrics.h 
b/be/src/common/metrics/doris_metrics.h
index 852cdb28753..971a55051b4 100644
--- a/be/src/common/metrics/doris_metrics.h
+++ b/be/src/common/metrics/doris_metrics.h
@@ -249,6 +249,11 @@ public:
     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;
diff --git a/be/src/service/internal_service.cpp 
b/be/src/service/internal_service.cpp
index a5954695d9c..eeed9f230d1 100644
--- a/be/src/service/internal_service.cpp
+++ b/be/src/service/internal_service.cpp
@@ -137,6 +137,11 @@ namespace doris {
 #include "common/compile_check_avoid_begin.h"
 using namespace ErrorCode;
 
+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 +163,12 @@ 
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_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);
@@ -222,6 +233,10 @@ PInternalService::PInternalService(ExecEnv* exec_env)
                                    ? 
config::brpc_heavy_work_pool_max_queue_size
                                    : std::max(10240, CpuInfo::num_cores() * 
320),
                            "brpc_heavy"),
+          // Keep cancellation dispatch independent of potentially blocking 
opens and writes.
+          _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.
           _peer_fetch_pool(resolved_brpc_peer_fetch_pool_threads(),
                            resolved_brpc_peer_fetch_pool_max_queue_size(), 
"brpc_peer_fetch"),
@@ -241,6 +256,15 @@ 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_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 
config::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,
@@ -287,6 +311,11 @@ PInternalServiceImpl::PInternalServiceImpl(StorageEngine& 
engine, ExecEnv* exec_
 PInternalServiceImpl::~PInternalServiceImpl() = default;
 
 PInternalService::~PInternalService() {
+    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);
@@ -536,7 +565,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 +578,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 550f2af8637..8be21ea2a40 100644
--- a/be/src/service/internal_service.h
+++ b/be/src/service/internal_service.h
@@ -278,6 +278,8 @@ protected:
     // define the interface for reading and writing data as heavy interface
     // otherwise as light interface
     FifoThreadPool _heavy_work_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;
     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 00000000000..429cfccbd57
--- /dev/null
+++ b/be/test/service/internal_service_load_work_pool_test.cpp
@@ -0,0 +1,175 @@
+// 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 <gtest/gtest.h>
+
+#include <atomic>
+#include <future>
+#include <memory>
+#include <utility>
+#include <vector>
+
+#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 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();
+        for (size_t i = 0; i < _pool._threads.size(); ++i) {
+            auto started = std::make_shared<std::promise<void>>();
+            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<void> _resume;
+};
+
+class LoadRpcCountingClosure : public google::protobuf::Closure {
+public:
+    void Run() override { ++calls; }
+    std::atomic<int> calls {0};
+};
+
+} // namespace
+
+class InternalServiceLoadWorkPoolTest : public testing::TestWithParam<int> {
+protected:
+    void SetUp() override {
+        // 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<LoadStreamMgr>(1);
+        _service = std::make_unique<PInternalService>(&_exec_env);
+        for (auto* pool : {&_service->_heavy_work_pool, 
&_service->_light_work_pool,
+                           &_service->_load_light_work_pool}) {
+            
_paused_pools.push_back(std::make_unique<PausedLoadRpcPool>(*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<PInternalService> _service;
+    std::vector<std::unique_ptr<PausedLoadRpcPool>> _paused_pools;
+    std::vector<std::pair<int32_t*, int32_t>> _saved_config;
+};
+
+TEST_F(InternalServiceLoadWorkPoolTest, CancelBypassesFullHeavyPool) {
+    EXPECT_EQ(_service->_load_light_work_pool.get_active_threads(), 3);
+    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;
+    POpenLoadStreamRequest stream_request;
+    POpenLoadStreamResponse stream_response;
+    PTabletWriterAddBlockRequest block_request;
+    PTabletWriterAddBlockResult block_response;
+    LoadRpcCountingClosure done;
+    auto submit = [&]() {
+        switch (GetParam()) {
+        case 0:
+            _service->tablet_writer_open(nullptr, &open_request, 
&open_response, &done);
+            break;
+        case 1:
+            _service->open_load_stream(nullptr, &stream_request, 
&stream_response, &done);
+            break;
+        case 2:
+            _service->tablet_writer_add_block(nullptr, &block_request, 
&block_response, &done);
+            break;
+        }
+    };
+
+    submit();
+    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);
+
+    submit();
+    EXPECT_EQ(done.calls.load(), 1);
+    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(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);
+    EXPECT_NE(_exec_env.load_stream_mgr()->heavy_work_pool(), 
&_service->_load_light_work_pool);
+}
+
+} // namespace doris


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to