Coverage Report

Created: 2026-09-28 11:58

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/service/internal_service.cpp
Line
Count
Source
1
// Licensed to the Apache Software Foundation (ASF) under one
2
// or more contributor license agreements.  See the NOTICE file
3
// distributed with this work for additional information
4
// regarding copyright ownership.  The ASF licenses this file
5
// to you under the Apache License, Version 2.0 (the
6
// "License"); you may not use this file except in compliance
7
// with the License.  You may obtain a copy of the License at
8
//
9
//   http://www.apache.org/licenses/LICENSE-2.0
10
//
11
// Unless required by applicable law or agreed to in writing,
12
// software distributed under the License is distributed on an
13
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14
// KIND, either express or implied.  See the License for the
15
// specific language governing permissions and limitations
16
// under the License.
17
18
#include "service/internal_service.h"
19
20
#include <assert.h>
21
#include <brpc/closure_guard.h>
22
#include <brpc/controller.h>
23
#include <bthread/bthread.h>
24
#include <bthread/types.h>
25
#include <butil/errno.h>
26
#include <butil/iobuf.h>
27
#include <fcntl.h>
28
#include <fmt/core.h>
29
#include <gen_cpp/DataSinks_types.h>
30
#include <gen_cpp/MasterService_types.h>
31
#include <gen_cpp/PaloInternalService_types.h>
32
#include <gen_cpp/PlanNodes_types.h>
33
#include <gen_cpp/Status_types.h>
34
#include <gen_cpp/Types_types.h>
35
#include <gen_cpp/internal_service.pb.h>
36
#include <gen_cpp/olap_file.pb.h>
37
#include <gen_cpp/segment_v2.pb.h>
38
#include <gen_cpp/types.pb.h>
39
#include <google/protobuf/stubs/callback.h>
40
#include <stddef.h>
41
#include <stdint.h>
42
#include <sys/stat.h>
43
44
#include <algorithm>
45
#include <exception>
46
#include <memory>
47
#include <set>
48
#include <sstream>
49
#include <string>
50
#include <utility>
51
#include <vector>
52
53
#include "cloud/cloud_storage_engine.h"
54
#include "cloud/cloud_tablet_mgr.h"
55
#include "cloud/config.h"
56
#include "common/config.h"
57
#include "common/exception.h"
58
#include "common/logging.h"
59
#include "common/metrics/doris_metrics.h"
60
#include "common/metrics/metrics.h"
61
#include "common/signal_handler.h"
62
#include "common/status.h"
63
#include "core/block/block.h"
64
#include "core/data_type/data_type.h"
65
#include "exec/common/variant_util.h"
66
#include "exec/exchange/vdata_stream_mgr.h"
67
#include "exec/rowid_fetcher.h"
68
#include "exec/runtime_filter/runtime_filter_mgr.h"
69
#include "exec/sink/writer/varrow_flight_result_writer.h"
70
#include "exec/sink/writer/vmysql_result_writer.h"
71
#include "exprs/function/dictionary_factory.h"
72
#include "format/arrow/arrow_row_batch.h"
73
#include "format/csv/csv_reader.h"
74
#include "format/generic_reader.h"
75
#include "format/jni/jni_reader.h"
76
#include "format/json/new_json_reader.h"
77
#include "format/orc/vorc_reader.h"
78
#include "format/parquet/vparquet_reader.h"
79
#include "format/text/text_reader.h"
80
#include "io/fs/local_file_system.h"
81
#include "io/fs/stream_load_pipe.h"
82
#include "io/io_common.h"
83
#include "load/channel/load_channel_mgr.h"
84
#include "load/channel/load_stream_mgr.h"
85
#include "load/delta_writer/delta_writer.h"
86
#include "load/group_commit/wal/wal_manager.h"
87
#include "load/routine_load/routine_load_task_executor.h"
88
#include "load/stream_load/new_load_stream_mgr.h"
89
#include "load/stream_load/stream_load_context.h"
90
#include "runtime/cache/result_cache.h"
91
#include "runtime/cdc_client_mgr.h"
92
#include "runtime/descriptors.h"
93
#include "runtime/exec_env.h"
94
#include "runtime/fold_constant_executor.h"
95
#include "runtime/fragment_mgr.h"
96
#include "runtime/query_context.h"
97
#include "runtime/result_block_buffer.h"
98
#include "runtime/result_buffer_mgr.h"
99
#include "runtime/runtime_profile.h"
100
#include "runtime/thread_context.h"
101
#include "runtime/workload_group/workload_group.h"
102
#include "runtime/workload_group/workload_group_manager.h"
103
#include "service/backend_options.h"
104
#include "service/point_query_executor.h"
105
#include "storage/data_dir.h"
106
#include "storage/olap_common.h"
107
#include "storage/olap_define.h"
108
#include "storage/rowset/rowset.h"
109
#include "storage/rowset/rowset_meta.h"
110
#include "storage/segment/column_reader.h"
111
#include "storage/storage_engine.h"
112
#include "storage/tablet/tablet_fwd.h"
113
#include "storage/tablet/tablet_manager.h"
114
#include "storage/tablet/tablet_schema.h"
115
#include "storage/txn/txn_manager.h"
116
#include "util/async_io.h"
117
#include "util/brpc_client_cache.h"
118
#include "util/brpc_closure.h"
119
#include "util/jdbc_utils.h"
120
#include "util/jsonb/serialize.h"
121
#include "util/md5.h"
122
#include "util/network_util.h"
123
#include "util/proto_util.h"
124
#include "util/stopwatch.hpp"
125
#include "util/string_util.h"
126
#include "util/thrift_util.h"
127
#include "util/time.h"
128
#include "util/uid_util.h"
129
130
namespace google {
131
namespace protobuf {
132
class RpcController;
133
} // namespace protobuf
134
} // namespace google
135
136
namespace doris {
137
#include "common/compile_check_avoid_begin.h"
138
using namespace ErrorCode;
139
140
DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(load_light_work_pool_queue_size, MetricUnit::NOUNIT);
141
DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(load_light_work_active_threads, MetricUnit::NOUNIT);
142
DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(load_light_work_pool_max_queue_size, MetricUnit::NOUNIT);
143
DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(load_light_work_max_threads, MetricUnit::NOUNIT);
144
145
DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(heavy_work_pool_queue_size, MetricUnit::NOUNIT);
146
DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(peer_fetch_work_pool_queue_size, MetricUnit::NOUNIT);
147
DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(light_work_pool_queue_size, MetricUnit::NOUNIT);
148
DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(heavy_work_active_threads, MetricUnit::NOUNIT);
149
DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(peer_fetch_work_active_threads, MetricUnit::NOUNIT);
150
DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(light_work_active_threads, MetricUnit::NOUNIT);
151
152
DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(heavy_work_pool_max_queue_size, MetricUnit::NOUNIT);
153
DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(peer_fetch_work_pool_max_queue_size, MetricUnit::NOUNIT);
154
DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(light_work_pool_max_queue_size, MetricUnit::NOUNIT);
155
DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(heavy_work_max_threads, MetricUnit::NOUNIT);
156
DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(peer_fetch_work_max_threads, MetricUnit::NOUNIT);
157
DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(light_work_max_threads, MetricUnit::NOUNIT);
158
159
DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(arrow_flight_work_pool_queue_size, MetricUnit::NOUNIT);
160
DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(arrow_flight_work_active_threads, MetricUnit::NOUNIT);
161
DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(arrow_flight_work_pool_max_queue_size, MetricUnit::NOUNIT);
162
DEFINE_GAUGE_METRIC_PROTOTYPE_2ARG(arrow_flight_work_max_threads, MetricUnit::NOUNIT);
163
164
static bvar::LatencyRecorder g_process_remote_fetch_rowsets_latency("process_remote_fetch_rowsets");
165
166
323
static int32_t resolved_brpc_load_light_work_pool_max_queue_size() {
167
323
    return config::brpc_load_light_work_pool_max_queue_size != -1
168
323
                   ? config::brpc_load_light_work_pool_max_queue_size
169
323
                   : std::max(1024, CpuInfo::num_cores() * 32);
170
323
}
171
172
323
static int32_t resolved_brpc_peer_fetch_pool_threads() {
173
323
    return config::brpc_peer_fetch_pool_threads != -1 ? config::brpc_peer_fetch_pool_threads
174
323
                                                      : std::max(64, CpuInfo::num_cores() * 2);
175
323
}
176
177
323
static int32_t resolved_brpc_peer_fetch_pool_max_queue_size() {
178
323
    return config::brpc_peer_fetch_pool_max_queue_size != -1
179
323
                   ? config::brpc_peer_fetch_pool_max_queue_size
180
323
                   : std::max(4096, CpuInfo::num_cores() * 128);
181
323
}
182
183
template <typename T>
184
concept CanCancel = requires(T* response) { response->mutable_status(); };
185
186
template <typename T>
187
1
void offer_failed(T* response, google::protobuf::Closure* done, const FifoThreadPool& pool) {
188
1
    brpc::ClosureGuard closure_guard(done);
189
1
    LOG(WARNING) << "fail to offer request to the work pool, pool=" << pool.get_info();
190
1
}
_ZN5doris12offer_failedINS_25PTabletWriterCancelResultEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Line
Count
Source
187
1
void offer_failed(T* response, google::protobuf::Closure* done, const FifoThreadPool& pool) {
188
1
    brpc::ClosureGuard closure_guard(done);
189
    LOG(WARNING) << "fail to offer request to the work pool, pool=" << pool.get_info();
190
1
}
Unexecuted instantiation: _ZN5doris12offer_failedINS_14PCacheResponseEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedINS_17PFetchCacheResultEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
191
192
template <CanCancel T>
193
3
void offer_failed(T* response, google::protobuf::Closure* done, const FifoThreadPool& pool) {
194
3
    brpc::ClosureGuard closure_guard(done);
195
    // Should use status to generate protobuf message, because it will encoding Backend Info
196
    // into the error message and then we could know which backend's pool is full.
197
3
    Status st = Status::Error<TStatusCode::CANCELLED>(
198
3
            "fail to offer request to the work pool, pool={}", pool.get_info());
199
3
    st.to_protobuf(response->mutable_status());
200
3
    LOG(WARNING) << "cancelled due to fail to offer request to the work pool, pool="
201
3
                 << pool.get_info();
202
3
}
_ZN5doris12offer_failedITkNS_9CanCancelENS_23PTabletWriterOpenResultEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Line
Count
Source
193
1
void offer_failed(T* response, google::protobuf::Closure* done, const FifoThreadPool& pool) {
194
1
    brpc::ClosureGuard closure_guard(done);
195
    // Should use status to generate protobuf message, because it will encoding Backend Info
196
    // into the error message and then we could know which backend's pool is full.
197
1
    Status st = Status::Error<TStatusCode::CANCELLED>(
198
1
            "fail to offer request to the work pool, pool={}", pool.get_info());
199
1
    st.to_protobuf(response->mutable_status());
200
    LOG(WARNING) << "cancelled due to fail to offer request to the work pool, pool="
201
1
                 << pool.get_info();
202
1
}
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_23PExecPlanFragmentResultEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
_ZN5doris12offer_failedITkNS_9CanCancelENS_23POpenLoadStreamResponseEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Line
Count
Source
193
1
void offer_failed(T* response, google::protobuf::Closure* done, const FifoThreadPool& pool) {
194
1
    brpc::ClosureGuard closure_guard(done);
195
    // Should use status to generate protobuf message, because it will encoding Backend Info
196
    // into the error message and then we could know which backend's pool is full.
197
1
    Status st = Status::Error<TStatusCode::CANCELLED>(
198
1
            "fail to offer request to the work pool, pool={}", pool.get_info());
199
1
    st.to_protobuf(response->mutable_status());
200
    LOG(WARNING) << "cancelled due to fail to offer request to the work pool, pool="
201
1
                 << pool.get_info();
202
1
}
_ZN5doris12offer_failedITkNS_9CanCancelENS_27PTabletWriterAddBlockResultEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Line
Count
Source
193
1
void offer_failed(T* response, google::protobuf::Closure* done, const FifoThreadPool& pool) {
194
1
    brpc::ClosureGuard closure_guard(done);
195
    // Should use status to generate protobuf message, because it will encoding Backend Info
196
    // into the error message and then we could know which backend's pool is full.
197
1
    Status st = Status::Error<TStatusCode::CANCELLED>(
198
1
            "fail to offer request to the work pool, pool={}", pool.get_info());
199
1
    st.to_protobuf(response->mutable_status());
200
    LOG(WARNING) << "cancelled due to fail to offer request to the work pool, pool="
201
1
                 << pool.get_info();
202
1
}
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_25PCancelPlanFragmentResultEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_21PFetchArrowDataResultEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_26POutfileWriteSuccessResultEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_23PFetchTableSchemaResultEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_29PFetchArrowFlightSchemaResultEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_24PTabletKeyLookupResponseEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_25PJdbcTestConnectionResultEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_20PFetchColIdsResponseEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_26PFetchRemoteSchemaResponseEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_12PProxyResultEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_20PMergeFilterResponseEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_23PSendFilterSizeResponseEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_23PSyncFilterSizeResponseEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_22PPublishFilterResponseEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_15PSendDataResultEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_13PCommitResultEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_15PRollbackResultEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_19PConstantExprResultEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_26PTransmitRecCTEBlockResultEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_20PRerunFragmentResultEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_20PResetGlobalRfResultEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_19PTransmitDataResultEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_24PCheckRPCChannelResponseEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_24PResetRPCChannelResponseEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_13PGlobResponseEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_26PGroupCommitInsertResponseEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_24PGetWalQueueSizeResponseEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_22PGetBeResourceResponseEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
Unexecuted instantiation: _ZN5doris12offer_failedITkNS_9CanCancelENS_23PRequestCdcClientResultEEEvPT_PN6google8protobuf7ClosureERKNS_14WorkThreadPoolILb0EEE
203
204
template <typename T>
205
class NewHttpClosure : public ::google::protobuf::Closure {
206
public:
207
    NewHttpClosure(google::protobuf::Closure* done) : _done(done) {}
208
0
    NewHttpClosure(T* request, google::protobuf::Closure* done) : _request(request), _done(done) {}
Unexecuted instantiation: _ZN5doris14NewHttpClosureINS_28PTabletWriterAddBlockRequestEEC2EPS1_PN6google8protobuf7ClosureE
Unexecuted instantiation: _ZN5doris14NewHttpClosureINS_19PTransmitDataParamsEEC2EPS1_PN6google8protobuf7ClosureE
209
210
0
    void Run() override {
211
0
        if (_request != nullptr) {
212
0
            delete _request;
213
0
            _request = nullptr;
214
0
        }
215
0
        if (_done != nullptr) {
216
0
            _done->Run();
217
0
        }
218
0
        delete this;
219
0
    }
Unexecuted instantiation: _ZN5doris14NewHttpClosureINS_28PTabletWriterAddBlockRequestEE3RunEv
Unexecuted instantiation: _ZN5doris14NewHttpClosureINS_19PTransmitDataParamsEE3RunEv
220
221
private:
222
    T* _request = nullptr;
223
    google::protobuf::Closure* _done = nullptr;
224
};
225
226
PInternalService::PInternalService(ExecEnv* exec_env)
227
12
        : _exec_env(exec_env),
228
          // heavy threadpool is used for load process and other process that will read disk or access network.
229
12
          _heavy_work_pool(config::brpc_heavy_work_pool_threads != -1
230
12
                                   ? config::brpc_heavy_work_pool_threads
231
12
                                   : std::max(128, CpuInfo::num_cores() * 4),
232
12
                           config::brpc_heavy_work_pool_max_queue_size != -1
233
12
                                   ? config::brpc_heavy_work_pool_max_queue_size
234
12
                                   : std::max(10240, CpuInfo::num_cores() * 320),
235
12
                           "brpc_heavy"),
236
          // Keep cancellation dispatch independent of potentially blocking opens and writes.
237
12
          _load_light_work_pool(config::brpc_load_light_work_pool_threads,
238
12
                                resolved_brpc_load_light_work_pool_max_queue_size(),
239
12
                                "brpc_load_light"),
240
          // peer fetch threadpool isolates fetch_peer_data from heavy load traffic to avoid peer reads starving imports.
241
12
          _peer_fetch_pool(resolved_brpc_peer_fetch_pool_threads(),
242
12
                           resolved_brpc_peer_fetch_pool_max_queue_size(), "brpc_peer_fetch"),
243
244
          // light threadpool should be only used in query processing logic. All hanlers should be very light, not locked, not access disk.
245
12
          _light_work_pool(config::brpc_light_work_pool_threads != -1
246
12
                                   ? config::brpc_light_work_pool_threads
247
12
                                   : std::max(128, CpuInfo::num_cores() * 4),
248
12
                           config::brpc_light_work_pool_max_queue_size != -1
249
12
                                   ? config::brpc_light_work_pool_max_queue_size
250
12
                                   : std::max(10240, CpuInfo::num_cores() * 320),
251
12
                           "brpc_light"),
252
12
          _arrow_flight_work_pool(config::brpc_arrow_flight_work_pool_threads != -1
253
12
                                          ? config::brpc_arrow_flight_work_pool_threads
254
12
                                          : std::max(512, CpuInfo::num_cores() * 2),
255
12
                                  config::brpc_arrow_flight_work_pool_max_queue_size != -1
256
12
                                          ? config::brpc_arrow_flight_work_pool_max_queue_size
257
12
                                          : std::max(20480, CpuInfo::num_cores() * 640),
258
12
                                  "brpc_arrow_flight") {
259
12
    REGISTER_HOOK_METRIC(load_light_work_pool_queue_size,
260
12
                         [this]() { return _load_light_work_pool.get_queue_size(); });
261
12
    REGISTER_HOOK_METRIC(load_light_work_active_threads,
262
12
                         [this]() { return _load_light_work_pool.get_active_threads(); });
263
12
    REGISTER_HOOK_METRIC(load_light_work_pool_max_queue_size,
264
12
                         []() { return resolved_brpc_load_light_work_pool_max_queue_size(); });
265
12
    REGISTER_HOOK_METRIC(load_light_work_max_threads,
266
12
                         []() { return config::brpc_load_light_work_pool_threads; });
267
268
12
    REGISTER_HOOK_METRIC(heavy_work_pool_queue_size,
269
12
                         [this]() { return _heavy_work_pool.get_queue_size(); });
270
12
    REGISTER_HOOK_METRIC(peer_fetch_work_pool_queue_size,
271
12
                         [this]() { return _peer_fetch_pool.get_queue_size(); });
272
12
    REGISTER_HOOK_METRIC(light_work_pool_queue_size,
273
12
                         [this]() { return _light_work_pool.get_queue_size(); });
274
12
    REGISTER_HOOK_METRIC(heavy_work_active_threads,
275
12
                         [this]() { return _heavy_work_pool.get_active_threads(); });
276
12
    REGISTER_HOOK_METRIC(peer_fetch_work_active_threads,
277
12
                         [this]() { return _peer_fetch_pool.get_active_threads(); });
278
12
    REGISTER_HOOK_METRIC(light_work_active_threads,
279
12
                         [this]() { return _light_work_pool.get_active_threads(); });
280
281
12
    REGISTER_HOOK_METRIC(heavy_work_pool_max_queue_size,
282
12
                         []() { return config::brpc_heavy_work_pool_max_queue_size; });
283
12
    REGISTER_HOOK_METRIC(peer_fetch_work_pool_max_queue_size,
284
12
                         []() { return resolved_brpc_peer_fetch_pool_max_queue_size(); });
285
12
    REGISTER_HOOK_METRIC(light_work_pool_max_queue_size,
286
12
                         []() { return config::brpc_light_work_pool_max_queue_size; });
287
12
    REGISTER_HOOK_METRIC(heavy_work_max_threads,
288
12
                         []() { return config::brpc_heavy_work_pool_threads; });
289
12
    REGISTER_HOOK_METRIC(peer_fetch_work_max_threads,
290
12
                         []() { return resolved_brpc_peer_fetch_pool_threads(); });
291
12
    REGISTER_HOOK_METRIC(light_work_max_threads,
292
12
                         []() { return config::brpc_light_work_pool_threads; });
293
294
12
    REGISTER_HOOK_METRIC(arrow_flight_work_pool_queue_size,
295
12
                         [this]() { return _arrow_flight_work_pool.get_queue_size(); });
296
12
    REGISTER_HOOK_METRIC(arrow_flight_work_active_threads,
297
12
                         [this]() { return _arrow_flight_work_pool.get_active_threads(); });
298
12
    REGISTER_HOOK_METRIC(arrow_flight_work_pool_max_queue_size,
299
12
                         []() { return config::brpc_arrow_flight_work_pool_max_queue_size; });
300
12
    REGISTER_HOOK_METRIC(arrow_flight_work_max_threads,
301
12
                         []() { return config::brpc_arrow_flight_work_pool_threads; });
302
303
12
    _exec_env->load_stream_mgr()->set_heavy_work_pool(&_heavy_work_pool);
304
305
12
    CHECK_EQ(0, bthread_key_create(&AsyncIO::btls_io_ctx_key, AsyncIO::io_ctx_key_deleter));
306
12
}
307
308
PInternalServiceImpl::PInternalServiceImpl(StorageEngine& engine, ExecEnv* exec_env)
309
6
        : PInternalService(exec_env), _engine(engine) {}
310
311
3
PInternalServiceImpl::~PInternalServiceImpl() = default;
312
313
8
PInternalService::~PInternalService() {
314
8
    DEREGISTER_HOOK_METRIC(load_light_work_pool_queue_size);
315
8
    DEREGISTER_HOOK_METRIC(load_light_work_active_threads);
316
8
    DEREGISTER_HOOK_METRIC(load_light_work_pool_max_queue_size);
317
8
    DEREGISTER_HOOK_METRIC(load_light_work_max_threads);
318
319
8
    DEREGISTER_HOOK_METRIC(heavy_work_pool_queue_size);
320
8
    DEREGISTER_HOOK_METRIC(peer_fetch_work_pool_queue_size);
321
8
    DEREGISTER_HOOK_METRIC(light_work_pool_queue_size);
322
8
    DEREGISTER_HOOK_METRIC(heavy_work_active_threads);
323
8
    DEREGISTER_HOOK_METRIC(peer_fetch_work_active_threads);
324
8
    DEREGISTER_HOOK_METRIC(light_work_active_threads);
325
326
8
    DEREGISTER_HOOK_METRIC(heavy_work_pool_max_queue_size);
327
8
    DEREGISTER_HOOK_METRIC(peer_fetch_work_pool_max_queue_size);
328
8
    DEREGISTER_HOOK_METRIC(light_work_pool_max_queue_size);
329
8
    DEREGISTER_HOOK_METRIC(heavy_work_max_threads);
330
8
    DEREGISTER_HOOK_METRIC(peer_fetch_work_max_threads);
331
8
    DEREGISTER_HOOK_METRIC(light_work_max_threads);
332
333
8
    DEREGISTER_HOOK_METRIC(arrow_flight_work_pool_queue_size);
334
8
    DEREGISTER_HOOK_METRIC(arrow_flight_work_active_threads);
335
8
    DEREGISTER_HOOK_METRIC(arrow_flight_work_pool_max_queue_size);
336
8
    DEREGISTER_HOOK_METRIC(arrow_flight_work_max_threads);
337
338
8
    CHECK_EQ(0, bthread_key_delete(AsyncIO::btls_io_ctx_key));
339
8
}
340
341
void PInternalService::tablet_writer_open(google::protobuf::RpcController* controller,
342
                                          const PTabletWriterOpenRequest* request,
343
                                          PTabletWriterOpenResult* response,
344
48.6k
                                          google::protobuf::Closure* done) {
345
48.7k
    bool ret = _heavy_work_pool.try_offer([this, request, response, done]() {
346
48.7k
        VLOG_RPC << "tablet writer open, id=" << request->id()
347
0
                 << ", index_id=" << request->index_id() << ", txn_id=" << request->txn_id();
348
48.7k
        signal::SignalTaskIdKeeper keeper(request->id());
349
48.7k
        brpc::ClosureGuard closure_guard(done);
350
48.7k
        auto st = _exec_env->load_channel_mgr()->open(*request);
351
48.7k
        if (!st.ok()) {
352
0
            LOG(WARNING) << "load channel open failed, message=" << st << ", id=" << request->id()
353
0
                         << ", index_id=" << request->index_id()
354
0
                         << ", txn_id=" << request->txn_id();
355
0
        }
356
48.7k
        st.to_protobuf(response->mutable_status());
357
48.7k
    });
358
48.6k
    if (!ret) {
359
1
        offer_failed(response, done, _heavy_work_pool);
360
1
        return;
361
1
    }
362
48.6k
}
363
364
void PInternalService::exec_plan_fragment(google::protobuf::RpcController* controller,
365
                                          const PExecPlanFragmentRequest* request,
366
                                          PExecPlanFragmentResult* response,
367
140k
                                          google::protobuf::Closure* done) {
368
140k
    timeval tv {};
369
140k
    gettimeofday(&tv, nullptr);
370
140k
    response->set_received_time(tv.tv_sec * 1000LL + tv.tv_usec / 1000);
371
140k
    bool ret = _light_work_pool.try_offer([this, controller, request, response, done]() {
372
140k
        _exec_plan_fragment_in_pthread(controller, request, response, done);
373
140k
    });
374
140k
    if (!ret) {
375
0
        offer_failed(response, done, _light_work_pool);
376
0
        return;
377
0
    }
378
140k
}
379
380
void PInternalService::_exec_plan_fragment_in_pthread(google::protobuf::RpcController* controller,
381
                                                      const PExecPlanFragmentRequest* request,
382
                                                      PExecPlanFragmentResult* response,
383
227k
                                                      google::protobuf::Closure* done) {
384
227k
    timeval tv1 {};
385
227k
    gettimeofday(&tv1, nullptr);
386
227k
    response->set_execution_time(tv1.tv_sec * 1000LL + tv1.tv_usec / 1000);
387
227k
    brpc::ClosureGuard closure_guard(done);
388
227k
    auto st = Status::OK();
389
227k
    bool compact = request->has_compact() ? request->compact() : false;
390
227k
    PFragmentRequestVersion version =
391
227k
            request->has_version() ? request->version() : PFragmentRequestVersion::VERSION_1;
392
227k
    try {
393
227k
        st = _exec_plan_fragment_impl(request->request(), version, compact);
394
227k
    } catch (const Exception& e) {
395
0
        st = e.to_status();
396
0
    } catch (const std::exception& e) {
397
0
        st = Status::Error(ErrorCode::INTERNAL_ERROR, e.what());
398
0
    } catch (...) {
399
0
        st = Status::Error(ErrorCode::INTERNAL_ERROR,
400
0
                           "_exec_plan_fragment_impl meet unknown error");
401
0
    }
402
227k
    if (!st.ok()) {
403
1.41k
        LOG(WARNING) << "exec plan fragment failed, errmsg=" << st;
404
1.41k
    }
405
226k
    st.to_protobuf(response->mutable_status());
406
226k
    timeval tv2 {};
407
226k
    gettimeofday(&tv2, nullptr);
408
226k
    response->set_execution_done_time(tv2.tv_sec * 1000LL + tv2.tv_usec / 1000);
409
226k
}
410
411
void PInternalService::exec_plan_fragment_prepare(google::protobuf::RpcController* controller,
412
                                                  const PExecPlanFragmentRequest* request,
413
                                                  PExecPlanFragmentResult* response,
414
86.4k
                                                  google::protobuf::Closure* done) {
415
86.4k
    timeval tv {};
416
86.4k
    gettimeofday(&tv, nullptr);
417
86.4k
    response->set_received_time(tv.tv_sec * 1000LL + tv.tv_usec / 1000);
418
86.5k
    bool ret = _light_work_pool.try_offer([this, controller, request, response, done]() {
419
86.5k
        _exec_plan_fragment_in_pthread(controller, request, response, done);
420
86.5k
    });
421
86.4k
    if (!ret) {
422
0
        offer_failed(response, done, _light_work_pool);
423
0
        return;
424
0
    }
425
86.4k
}
426
427
void PInternalService::exec_plan_fragment_start(google::protobuf::RpcController* /*controller*/,
428
                                                const PExecPlanFragmentStartRequest* request,
429
                                                PExecPlanFragmentResult* result,
430
86.4k
                                                google::protobuf::Closure* done) {
431
86.4k
    timeval tv {};
432
86.4k
    gettimeofday(&tv, nullptr);
433
86.4k
    result->set_received_time(tv.tv_sec * 1000LL + tv.tv_usec / 1000);
434
86.5k
    bool ret = _light_work_pool.try_offer([this, request, result, done]() {
435
86.5k
        timeval tv1 {};
436
86.5k
        gettimeofday(&tv1, nullptr);
437
86.5k
        result->set_execution_time(tv1.tv_sec * 1000LL + tv1.tv_usec / 1000);
438
86.5k
        brpc::ClosureGuard closure_guard(done);
439
86.5k
        auto st = _exec_env->fragment_mgr()->start_query_execution(request);
440
86.5k
        st.to_protobuf(result->mutable_status());
441
86.5k
        timeval tv2 {};
442
86.5k
        gettimeofday(&tv2, nullptr);
443
86.5k
        result->set_execution_done_time(tv2.tv_sec * 1000LL + tv2.tv_usec / 1000);
444
86.5k
    });
445
86.4k
    if (!ret) {
446
0
        offer_failed(result, done, _light_work_pool);
447
0
        return;
448
0
    }
449
86.4k
}
450
451
void PInternalService::open_load_stream(google::protobuf::RpcController* controller,
452
                                        const POpenLoadStreamRequest* request,
453
                                        POpenLoadStreamResponse* response,
454
42
                                        google::protobuf::Closure* done) {
455
42
    bool ret = _heavy_work_pool.try_offer([this, controller, request, response, done]() {
456
40
        signal::SignalTaskIdKeeper keeper(request->load_id());
457
40
        brpc::ClosureGuard done_guard(done);
458
40
        brpc::Controller* cntl = static_cast<brpc::Controller*>(controller);
459
40
        brpc::StreamOptions stream_options;
460
461
40
        LOG(INFO) << "open load stream, load_id=" << request->load_id()
462
40
                  << ", src_id=" << request->src_id();
463
464
40
        std::vector<BaseTabletSPtr> tablets;
465
40
        for (const auto& req : request->tablets()) {
466
20
            BaseTabletSPtr tablet;
467
20
            if (auto res = ExecEnv::get_tablet(req.tablet_id()); !res.has_value()) [[unlikely]] {
468
0
                auto st = std::move(res).error();
469
0
                st.to_protobuf(response->mutable_status());
470
0
                cntl->SetFailed(st.to_string());
471
0
                return;
472
20
            } else {
473
20
                tablet = std::move(res).value();
474
20
            }
475
20
            auto resp = response->add_tablet_schemas();
476
20
            resp->set_index_id(req.index_id());
477
20
            resp->set_enable_unique_key_merge_on_write(tablet->enable_unique_key_merge_on_write());
478
20
            tablet->tablet_schema()->to_schema_pb(resp->mutable_tablet_schema());
479
20
            tablets.push_back(tablet);
480
20
        }
481
40
        if (!tablets.empty()) {
482
20
            auto* tablet_load_infos = response->mutable_tablet_load_rowset_num_infos();
483
20
            for (const auto& tablet : tablets) {
484
20
                BaseDeltaWriter::collect_tablet_load_rowset_num_info(tablet.get(),
485
20
                                                                     tablet_load_infos);
486
20
            }
487
20
        }
488
489
40
        LoadStream* load_stream = nullptr;
490
40
        auto st = _exec_env->load_stream_mgr()->open_load_stream(request, load_stream);
491
40
        if (!st.ok()) {
492
0
            st.to_protobuf(response->mutable_status());
493
0
            return;
494
0
        }
495
496
40
        stream_options.handler = load_stream;
497
40
        stream_options.idle_timeout_ms = request->idle_timeout_ms();
498
40
        DBUG_EXECUTE_IF("PInternalServiceImpl.open_load_stream.set_idle_timeout",
499
40
                        { stream_options.idle_timeout_ms = 1; });
500
501
40
        StreamId streamid;
502
40
        if (brpc::StreamAccept(&streamid, *cntl, &stream_options) != 0) {
503
0
            st = Status::Cancelled("Fail to accept stream {}", streamid);
504
0
            st.to_protobuf(response->mutable_status());
505
0
            cntl->SetFailed(st.to_string());
506
0
            return;
507
0
        }
508
509
40
        VLOG_DEBUG << "get streamid =" << streamid;
510
40
        st.to_protobuf(response->mutable_status());
511
40
    });
512
42
    if (!ret) {
513
1
        offer_failed(response, done, _heavy_work_pool);
514
1
    }
515
42
}
516
517
void PInternalService::tablet_writer_add_block_by_http(google::protobuf::RpcController* controller,
518
                                                       const ::doris::PEmptyRequest* request,
519
                                                       PTabletWriterAddBlockResult* response,
520
0
                                                       google::protobuf::Closure* done) {
521
0
    PTabletWriterAddBlockRequest* new_request = new PTabletWriterAddBlockRequest();
522
0
    google::protobuf::Closure* new_done =
523
0
            new NewHttpClosure<PTabletWriterAddBlockRequest>(new_request, done);
524
0
    brpc::Controller* cntl = static_cast<brpc::Controller*>(controller);
525
0
    Status st = attachment_extract_request_contain_block<PTabletWriterAddBlockRequest>(new_request,
526
0
                                                                                       cntl);
527
0
    if (st.ok()) {
528
0
        tablet_writer_add_block(controller, new_request, response, new_done);
529
0
    } else {
530
0
        st.to_protobuf(response->mutable_status());
531
0
    }
532
0
}
533
534
void PInternalService::tablet_writer_add_block(google::protobuf::RpcController* controller,
535
                                               const PTabletWriterAddBlockRequest* request,
536
                                               PTabletWriterAddBlockResult* response,
537
54.6k
                                               google::protobuf::Closure* done) {
538
54.6k
    int64_t submit_task_time_ns = MonotonicNanos();
539
54.9k
    bool ret = _heavy_work_pool.try_offer([request, response, done, submit_task_time_ns, this]() {
540
54.9k
        int64_t wait_execution_time_ns = MonotonicNanos() - submit_task_time_ns;
541
54.9k
        brpc::ClosureGuard closure_guard(done);
542
54.9k
        int64_t execution_time_ns = 0;
543
54.9k
        {
544
54.9k
            SCOPED_RAW_TIMER(&execution_time_ns);
545
54.9k
            signal::SignalTaskIdKeeper keeper(request->id());
546
54.9k
            auto st = _exec_env->load_channel_mgr()->add_batch(*request, response);
547
54.9k
            if (!st.ok()) {
548
46
                LOG(WARNING) << "tablet writer add block failed, message=" << st
549
46
                             << ", id=" << request->id() << ", index_id=" << request->index_id()
550
46
                             << ", sender_id=" << request->sender_id()
551
46
                             << ", backend id=" << request->backend_id();
552
46
            }
553
54.9k
            st.to_protobuf(response->mutable_status());
554
54.9k
        }
555
54.9k
        response->set_execution_time_us(execution_time_ns / NANOS_PER_MICRO);
556
54.9k
        response->set_wait_execution_time_us(wait_execution_time_ns / NANOS_PER_MICRO);
557
54.9k
    });
558
54.6k
    if (!ret) {
559
1
        offer_failed(response, done, _heavy_work_pool);
560
1
        return;
561
1
    }
562
54.6k
}
563
564
void PInternalService::tablet_writer_cancel(google::protobuf::RpcController* controller,
565
                                            const PTabletWriterCancelRequest* request,
566
                                            PTabletWriterCancelResult* response,
567
103
                                            google::protobuf::Closure* done) {
568
103
    bool ret = _load_light_work_pool.try_offer([this, request, done]() {
569
101
        VLOG_RPC << "tablet writer cancel, id=" << request->id()
570
0
                 << ", index_id=" << request->index_id() << ", sender_id=" << request->sender_id();
571
101
        signal::SignalTaskIdKeeper keeper(request->id());
572
101
        brpc::ClosureGuard closure_guard(done);
573
101
        auto st = _exec_env->load_channel_mgr()->cancel(*request);
574
101
        if (!st.ok()) {
575
0
            LOG(WARNING) << "tablet writer cancel failed, id=" << request->id()
576
0
                         << ", index_id=" << request->index_id()
577
0
                         << ", sender_id=" << request->sender_id();
578
0
        }
579
101
    });
580
103
    if (!ret) {
581
1
        offer_failed(response, done, _load_light_work_pool);
582
1
        return;
583
1
    }
584
103
}
585
586
Status PInternalService::_exec_plan_fragment_impl(
587
        const std::string& ser_request, PFragmentRequestVersion version, bool compact,
588
227k
        const std::function<void(RuntimeState*, Status*)>& cb) {
589
    // Sometimes the BE do not receive the first heartbeat message and it receives request from FE
590
    // If BE execute this fragment, it will core when it wants to get some property from master info.
591
227k
    if (ExecEnv::GetInstance()->cluster_info() == nullptr) {
592
0
        return Status::InternalError(
593
0
                "Have not receive the first heartbeat message from master, not ready to provide "
594
0
                "service");
595
0
    }
596
227k
    CHECK(version == PFragmentRequestVersion::VERSION_3)
597
0
            << "only support version 3, received " << version;
598
227k
    if (version == PFragmentRequestVersion::VERSION_3) {
599
227k
        TPipelineFragmentParamsList t_request;
600
227k
        {
601
227k
            const uint8_t* buf = (const uint8_t*)ser_request.data();
602
227k
            uint32_t len = ser_request.size();
603
227k
            RETURN_IF_ERROR(deserialize_thrift_msg(buf, &len, compact, &t_request));
604
227k
        }
605
606
227k
        const auto& fragment_list = t_request.params_list;
607
227k
        if (fragment_list.empty()) {
608
0
            return Status::InternalError("Invalid TPipelineFragmentParamsList!");
609
0
        }
610
227k
        MonotonicStopWatch timer;
611
227k
        timer.start();
612
613
        // work for old version frontend
614
227k
        if (!t_request.__isset.runtime_filter_info) {
615
72.2k
            TRuntimeFilterInfo runtime_filter_info;
616
72.2k
            auto local_param = fragment_list[0].local_params[0];
617
72.2k
            if (local_param.__isset.runtime_filter_params) {
618
72.1k
                runtime_filter_info.__set_runtime_filter_params(local_param.runtime_filter_params);
619
72.1k
            }
620
72.2k
            if (local_param.__isset.topn_filter_descs) {
621
0
                runtime_filter_info.__set_topn_filter_descs(local_param.topn_filter_descs);
622
0
            }
623
72.2k
            t_request.__set_runtime_filter_info(runtime_filter_info);
624
72.2k
        }
625
626
332k
        for (const TPipelineFragmentParams& fragment : fragment_list) {
627
332k
            if (cb) {
628
29
                RETURN_IF_ERROR(_exec_env->fragment_mgr()->exec_plan_fragment(
629
29
                        fragment, QuerySource::INTERNAL_FRONTEND, cb, t_request));
630
332k
            } else {
631
332k
                RETURN_IF_ERROR(_exec_env->fragment_mgr()->exec_plan_fragment(
632
332k
                        fragment, QuerySource::INTERNAL_FRONTEND, t_request));
633
332k
            }
634
332k
        }
635
226k
        timer.stop();
636
226k
        double cost_secs = static_cast<double>(timer.elapsed_time()) / 1000000000ULL;
637
226k
        if (cost_secs > 5) {
638
15
            LOG_WARNING("Prepare {} fragments of query {} costs {} seconds, it costs too much",
639
15
                        fragment_list.size(), print_id(fragment_list.front().query_id), cost_secs);
640
15
        }
641
642
226k
        return Status::OK();
643
227k
    } else {
644
0
        return Status::InternalError("invalid version");
645
0
    }
646
227k
}
647
648
void PInternalService::cancel_plan_fragment(google::protobuf::RpcController* /*controller*/,
649
                                            const PCancelPlanFragmentRequest* request,
650
                                            PCancelPlanFragmentResult* result,
651
154k
                                            google::protobuf::Closure* done) {
652
155k
    bool ret = _light_work_pool.try_offer([this, request, result, done]() {
653
155k
        brpc::ClosureGuard closure_guard(done);
654
155k
        signal::SignalTaskIdKeeper keeper(request->finst_id());
655
155k
        Status st = Status::OK();
656
657
155k
        const bool has_cancel_reason = request->has_cancel_reason();
658
155k
        const bool has_cancel_status = request->has_cancel_status();
659
        // During upgrade only LIMIT_REACH is used, other reason is changed to internal error
660
155k
        Status actual_cancel_status = Status::OK();
661
        // Convert PPlanFragmentCancelReason to Status
662
155k
        if (has_cancel_status) {
663
            // If fe set cancel status, then it is new FE now, should use cancel status.
664
155k
            actual_cancel_status = Status::create<false>(request->cancel_status());
665
155k
        } else if (has_cancel_reason) {
666
            // If fe not set cancel status, but set cancel reason, should convert cancel reason
667
            // to cancel status here.
668
0
            if (request->cancel_reason() == PPlanFragmentCancelReason::LIMIT_REACH) {
669
0
                actual_cancel_status = Status::Error<ErrorCode::LIMIT_REACH>("limit reach");
670
0
            } else {
671
                // Use cancel reason as error message
672
0
                actual_cancel_status = Status::InternalError(
673
0
                        PPlanFragmentCancelReason_Name(request->cancel_reason()));
674
0
            }
675
0
        } else {
676
0
            actual_cancel_status = Status::InternalError("unknown error");
677
0
        }
678
679
155k
        TUniqueId query_id;
680
155k
        query_id.__set_hi(request->query_id().hi());
681
155k
        query_id.__set_lo(request->query_id().lo());
682
155k
        LOG(INFO) << fmt::format("Cancel query {}, reason: {}", print_id(query_id),
683
155k
                                 actual_cancel_status.to_string());
684
155k
        _exec_env->fragment_mgr()->cancel_query(query_id, actual_cancel_status);
685
686
        // TODO: the logic seems useless, cancel only return Status::OK. remove it
687
155k
        st.to_protobuf(result->mutable_status());
688
155k
    });
689
154k
    if (!ret) {
690
0
        offer_failed(result, done, _light_work_pool);
691
0
        return;
692
0
    }
693
154k
}
694
695
void PInternalService::fetch_data(google::protobuf::RpcController* controller,
696
                                  const PFetchDataRequest* request, PFetchDataResult* result,
697
329k
                                  google::protobuf::Closure* done) {
698
    // fetch_data is a light operation which will put a request rather than wait inplace when there's no data ready.
699
    // when there's data ready, use brpc to send. there's queue in brpc service. won't take it too long.
700
329k
    auto ctx = GetResultBatchCtx::create_shared(result, done);
701
329k
    TUniqueId unique_id = UniqueId(request->finst_id()).to_thrift(); // query_id or instance_id
702
329k
    std::shared_ptr<MySQLResultBlockBuffer> buffer;
703
329k
    Status st = ExecEnv::GetInstance()->result_mgr()->find_buffer(unique_id, buffer);
704
329k
    if (!st.ok()) {
705
0
        LOG(WARNING) << "Result buffer not found! finst ID: " << print_id(unique_id);
706
0
        return;
707
0
    }
708
329k
    if (st = buffer->get_batch(ctx); !st.ok()) {
709
19
        LOG(WARNING) << "fetch_data failed: " << st.to_string();
710
19
    }
711
329k
}
712
713
void PInternalService::fetch_arrow_data(google::protobuf::RpcController* controller,
714
                                        const PFetchArrowDataRequest* request,
715
                                        PFetchArrowDataResult* result,
716
0
                                        google::protobuf::Closure* done) {
717
0
    bool ret = _arrow_flight_work_pool.try_offer([request, result, done]() {
718
0
        auto ctx = GetArrowResultBatchCtx::create_shared(result, done);
719
0
        TUniqueId unique_id = UniqueId(request->finst_id()).to_thrift(); // query_id or instance_id
720
0
        std::shared_ptr<ArrowFlightResultBlockBuffer> arrow_buffer;
721
0
        auto st = ExecEnv::GetInstance()->result_mgr()->find_buffer(unique_id, arrow_buffer);
722
0
        if (!st.ok()) {
723
0
            LOG(WARNING) << "Result buffer not found! Query ID: " << print_id(unique_id);
724
0
            return;
725
0
        }
726
0
        if (st = arrow_buffer->get_batch(ctx); !st.ok()) {
727
0
            LOG(WARNING) << "fetch_arrow_data failed: " << st.to_string();
728
0
        }
729
0
    });
730
0
    if (!ret) {
731
0
        offer_failed(result, done, _arrow_flight_work_pool);
732
0
        return;
733
0
    }
734
0
}
735
736
void PInternalService::outfile_write_success(google::protobuf::RpcController* controller,
737
                                             const POutfileWriteSuccessRequest* request,
738
                                             POutfileWriteSuccessResult* result,
739
4
                                             google::protobuf::Closure* done) {
740
4
    bool ret = _heavy_work_pool.try_offer([request, result, done]() {
741
4
        VLOG_RPC << "outfile write success file";
742
4
        brpc::ClosureGuard closure_guard(done);
743
4
        TResultFileSink result_file_sink;
744
4
        Status st = Status::OK();
745
4
        {
746
4
            const uint8_t* buf = (const uint8_t*)(request->result_file_sink().data());
747
4
            uint32_t len = request->result_file_sink().size();
748
4
            st = deserialize_thrift_msg(buf, &len, false, &result_file_sink);
749
4
            if (!st.ok()) {
750
0
                LOG(WARNING) << "outfile write success file failed, errmsg = " << st;
751
0
                st.to_protobuf(result->mutable_status());
752
0
                return;
753
0
            }
754
4
        }
755
756
4
        TResultFileSinkOptions file_options = result_file_sink.file_options;
757
4
        std::stringstream ss;
758
4
        ss << file_options.file_path << file_options.success_file_name;
759
4
        std::string file_name = ss.str();
760
4
        if (result_file_sink.storage_backend_type == TStorageBackendType::LOCAL) {
761
            // For local file writer, the file_path is a local dir.
762
            // Here we do a simple security verification by checking whether the file exists.
763
            // Because the file path is currently arbitrarily specified by the user,
764
            // Doris is not responsible for ensuring the correctness of the path.
765
            // This is just to prevent overwriting the existing file.
766
4
            bool exists = true;
767
4
            st = io::global_local_filesystem()->exists(file_name, &exists);
768
4
            if (!st.ok()) {
769
0
                LOG(WARNING) << "outfile write success filefailed, errmsg = " << st;
770
0
                st.to_protobuf(result->mutable_status());
771
0
                return;
772
0
            }
773
4
            if (exists) {
774
0
                st = Status::InternalError("File already exists: {}", file_name);
775
0
            }
776
4
            if (!st.ok()) {
777
0
                LOG(WARNING) << "outfile write success file failed, errmsg = " << st;
778
0
                st.to_protobuf(result->mutable_status());
779
0
                return;
780
0
            }
781
4
        }
782
783
4
        auto file_type_res =
784
4
                FileFactory::convert_storage_type(result_file_sink.storage_backend_type);
785
4
        if (!file_type_res.has_value()) [[unlikely]] {
786
0
            st = std::move(file_type_res).error();
787
0
            st.to_protobuf(result->mutable_status());
788
0
            LOG(WARNING) << "encounter unkonw type=" << result_file_sink.storage_backend_type
789
0
                         << ", st=" << st;
790
0
            return;
791
0
        }
792
793
4
        auto&& res = FileFactory::create_file_writer(file_type_res.value(), ExecEnv::GetInstance(),
794
4
                                                     file_options.broker_addresses,
795
4
                                                     file_options.broker_properties, file_name,
796
4
                                                     {
797
4
                                                             .write_file_cache = false,
798
4
                                                             .sync_file_data = false,
799
4
                                                     });
800
4
        using T = std::decay_t<decltype(res)>;
801
4
        if (!res.has_value()) [[unlikely]] {
802
0
            st = std::forward<T>(res).error();
803
0
            st.to_protobuf(result->mutable_status());
804
0
            return;
805
0
        }
806
807
4
        std::unique_ptr<doris::io::FileWriter> _file_writer_impl = std::forward<T>(res).value();
808
        // must write somthing because s3 file writer can not writer empty file
809
4
        st = _file_writer_impl->append({"success"});
810
4
        if (!st.ok()) {
811
0
            LOG(WARNING) << "outfile write success filefailed, errmsg=" << st;
812
0
            st.to_protobuf(result->mutable_status());
813
0
            return;
814
0
        }
815
4
        st = _file_writer_impl->close();
816
4
        if (!st.ok()) {
817
0
            LOG(WARNING) << "outfile write success filefailed, errmsg=" << st;
818
0
            st.to_protobuf(result->mutable_status());
819
0
            return;
820
0
        }
821
4
    });
822
4
    if (!ret) {
823
0
        offer_failed(result, done, _heavy_work_pool);
824
0
        return;
825
0
    }
826
4
}
827
828
void PInternalService::fetch_table_schema(google::protobuf::RpcController* controller,
829
                                          const PFetchTableSchemaRequest* request,
830
                                          PFetchTableSchemaResult* result,
831
645
                                          google::protobuf::Closure* done) {
832
645
    bool ret = _heavy_work_pool.try_offer([request, result, done]() {
833
645
        VLOG_RPC << "fetch table schema";
834
645
        brpc::ClosureGuard closure_guard(done);
835
645
        TFileScanRange file_scan_range;
836
645
        Status st = Status::OK();
837
645
        {
838
645
            const uint8_t* buf = (const uint8_t*)(request->file_scan_range().data());
839
645
            uint32_t len = request->file_scan_range().size();
840
645
            st = deserialize_thrift_msg(buf, &len, false, &file_scan_range);
841
645
            if (!st.ok()) {
842
0
                LOG(WARNING) << "fetch table schema failed, errmsg=" << st;
843
0
                st.to_protobuf(result->mutable_status());
844
0
                return;
845
0
            }
846
645
        }
847
645
        if (file_scan_range.__isset.ranges == false) {
848
0
            st = Status::InternalError("can not get TFileRangeDesc.");
849
0
            st.to_protobuf(result->mutable_status());
850
0
            return;
851
0
        }
852
645
        if (file_scan_range.__isset.params == false) {
853
0
            st = Status::InternalError("can not get TFileScanRangeParams.");
854
0
            st.to_protobuf(result->mutable_status());
855
0
            return;
856
0
        }
857
645
        const TFileRangeDesc& range = file_scan_range.ranges.at(0);
858
645
        const TFileScanRangeParams& params = file_scan_range.params;
859
860
645
        std::shared_ptr<MemTrackerLimiter> mem_tracker = MemTrackerLimiter::create_shared(
861
645
                MemTrackerLimiter::Type::OTHER,
862
645
                fmt::format("InternalService::fetch_table_schema:{}#{}", params.format_type,
863
645
                            params.file_type));
864
645
        SCOPED_ATTACH_TASK(mem_tracker);
865
866
        // make sure profile is desctructed after reader cause PrefetchBufferedReader
867
        // might asynchronouslly access the profile
868
645
        std::unique_ptr<RuntimeProfile> profile =
869
645
                std::make_unique<RuntimeProfile>("FetchTableSchema");
870
645
        std::unique_ptr<GenericReader> reader(nullptr);
871
645
        auto io_ctx = std::make_shared<io::IOContext>();
872
645
        auto file_cache_statis = std::make_shared<io::FileCacheStatistics>();
873
645
        auto file_reader_stats = std::make_shared<io::FileReaderStats>();
874
645
        io_ctx->file_cache_stats = file_cache_statis.get();
875
645
        io_ctx->file_reader_stats = file_reader_stats.get();
876
645
        constexpr size_t fetch_schema_batch_size = 4064;
877
        // file_slots is no use, but the lifetime should be longer than reader
878
645
        std::vector<SlotDescriptor*> file_slots;
879
645
        switch (params.format_type) {
880
403
        case TFileFormatType::FORMAT_CSV_PLAIN:
881
403
        case TFileFormatType::FORMAT_CSV_GZ:
882
403
        case TFileFormatType::FORMAT_CSV_BZ2:
883
403
        case TFileFormatType::FORMAT_CSV_LZ4FRAME:
884
403
        case TFileFormatType::FORMAT_CSV_LZ4BLOCK:
885
403
        case TFileFormatType::FORMAT_CSV_SNAPPYBLOCK:
886
403
        case TFileFormatType::FORMAT_CSV_LZOP:
887
403
        case TFileFormatType::FORMAT_CSV_DEFLATE: {
888
403
            reader = CsvReader::create_unique(nullptr, profile.get(), nullptr, params, range,
889
403
                                              file_slots, fetch_schema_batch_size, io_ctx.get(),
890
403
                                              io_ctx);
891
403
            break;
892
403
        }
893
0
        case TFileFormatType::FORMAT_TEXT: {
894
0
            reader = TextReader::create_unique(nullptr, profile.get(), nullptr, params, range,
895
0
                                               file_slots, fetch_schema_batch_size, io_ctx.get());
896
0
            break;
897
403
        }
898
119
        case TFileFormatType::FORMAT_PARQUET: {
899
119
            reader = ParquetReader::create_unique(params, range, io_ctx, nullptr);
900
119
            break;
901
403
        }
902
102
        case TFileFormatType::FORMAT_ORC: {
903
102
            reader = OrcReader::create_unique(params, range, fetch_schema_batch_size, "", io_ctx);
904
102
            break;
905
403
        }
906
21
        case TFileFormatType::FORMAT_JSON: {
907
21
            reader = NewJsonReader::create_unique(profile.get(), params, range, file_slots,
908
21
                                                  fetch_schema_batch_size, io_ctx.get(), io_ctx);
909
21
            break;
910
403
        }
911
0
        default:
912
0
            st = Status::InternalError("Not supported file format in fetch table schema: {}",
913
0
                                       params.format_type);
914
0
            st.to_protobuf(result->mutable_status());
915
0
            return;
916
645
        }
917
645
        if (!st.ok()) {
918
0
            LOG(WARNING) << "failed to create reader, errmsg=" << st;
919
0
            st.to_protobuf(result->mutable_status());
920
0
            return;
921
0
        }
922
645
        st = reader->init_schema_reader();
923
645
        if (!st.ok()) {
924
1
            LOG(WARNING) << "failed to init reader, errmsg=" << st;
925
1
            st.to_protobuf(result->mutable_status());
926
1
            return;
927
1
        }
928
644
        std::vector<std::string> col_names;
929
644
        std::vector<DataTypePtr> col_types;
930
644
        st = reader->get_parsed_schema(&col_names, &col_types);
931
644
        if (!st.ok()) {
932
4
            LOG(WARNING) << "fetch table schema failed, errmsg=" << st;
933
4
            st.to_protobuf(result->mutable_status());
934
4
            return;
935
4
        }
936
640
        result->set_column_nums(col_names.size());
937
4.89k
        for (size_t idx = 0; idx < col_names.size(); ++idx) {
938
4.25k
            result->add_column_names(col_names[idx]);
939
4.25k
        }
940
4.89k
        for (size_t idx = 0; idx < col_types.size(); ++idx) {
941
4.25k
            PTypeDesc* type_desc = result->add_column_types();
942
4.25k
            col_types[idx]->to_protobuf(type_desc);
943
4.25k
        }
944
640
        st.to_protobuf(result->mutable_status());
945
640
    });
946
645
    if (!ret) {
947
0
        offer_failed(result, done, _heavy_work_pool);
948
0
        return;
949
0
    }
950
645
}
951
952
void PInternalService::fetch_arrow_flight_schema(google::protobuf::RpcController* controller,
953
                                                 const PFetchArrowFlightSchemaRequest* request,
954
                                                 PFetchArrowFlightSchemaResult* result,
955
37
                                                 google::protobuf::Closure* done) {
956
37
    bool ret = _arrow_flight_work_pool.try_offer([request, result, done]() {
957
37
        brpc::ClosureGuard closure_guard(done);
958
37
        std::shared_ptr<arrow::Schema> schema;
959
37
        std::shared_ptr<ArrowFlightResultBlockBuffer> buffer;
960
37
        auto st = ExecEnv::GetInstance()->result_mgr()->find_buffer(
961
37
                UniqueId(request->finst_id()).to_thrift(), buffer);
962
37
        if (!st.ok()) {
963
0
            LOG(WARNING) << "fetch arrow flight schema failed, errmsg=" << st;
964
0
            st.to_protobuf(result->mutable_status());
965
0
            return;
966
0
        }
967
37
        st = buffer->get_schema(&schema);
968
37
        if (!st.ok()) {
969
0
            LOG(WARNING) << "fetch arrow flight schema failed, errmsg=" << st;
970
0
            st.to_protobuf(result->mutable_status());
971
0
            return;
972
0
        }
973
974
37
        std::string schema_str;
975
37
        st = serialize_arrow_schema(&schema, &schema_str);
976
37
        if (st.ok()) {
977
37
            result->set_schema(std::move(schema_str));
978
37
            if (!config::public_host.empty()) {
979
0
                result->set_be_arrow_flight_ip(config::public_host);
980
0
            }
981
37
            if (config::arrow_flight_sql_proxy_port != -1) {
982
0
                result->set_be_arrow_flight_port(config::arrow_flight_sql_proxy_port);
983
0
            }
984
37
        }
985
37
        st.to_protobuf(result->mutable_status());
986
37
    });
987
37
    if (!ret) {
988
0
        offer_failed(result, done, _arrow_flight_work_pool);
989
0
        return;
990
0
    }
991
37
}
992
993
Status PInternalService::_tablet_fetch_data(const PTabletKeyLookupRequest* request,
994
237
                                            PTabletKeyLookupResponse* response) {
995
237
    PointQueryExecutor executor;
996
237
    RETURN_IF_ERROR(executor.init(request, response));
997
237
    if (response->has_need_resend_query_context() && response->need_resend_query_context()) {
998
1
        return Status::OK();
999
1
    }
1000
236
    RETURN_IF_ERROR(executor.lookup_up());
1001
233
    executor.print_profile();
1002
233
    return Status::OK();
1003
236
}
1004
1005
void PInternalService::tablet_fetch_data(google::protobuf::RpcController* controller,
1006
                                         const PTabletKeyLookupRequest* request,
1007
                                         PTabletKeyLookupResponse* response,
1008
237
                                         google::protobuf::Closure* done) {
1009
237
    bool ret = _light_work_pool.try_offer([this, controller, request, response, done]() {
1010
237
        [[maybe_unused]] auto* cntl = static_cast<brpc::Controller*>(controller);
1011
237
        brpc::ClosureGuard guard(done);
1012
237
        Status st = _tablet_fetch_data(request, response);
1013
237
        st.to_protobuf(response->mutable_status());
1014
237
    });
1015
237
    if (!ret) {
1016
0
        offer_failed(response, done, _light_work_pool);
1017
0
        return;
1018
0
    }
1019
237
}
1020
1021
void PInternalService::test_jdbc_connection(google::protobuf::RpcController* controller,
1022
                                            const PJdbcTestConnectionRequest* request,
1023
                                            PJdbcTestConnectionResult* result,
1024
3
                                            google::protobuf::Closure* done) {
1025
3
    if (!doris::config::enable_java_support) {
1026
0
        doris::Status status = doris::Status::InternalError(
1027
0
                "you can change be config enable_java_support to true and restart be.");
1028
0
        status.to_protobuf(result->mutable_status());
1029
0
        done->Run();
1030
0
        return;
1031
0
    }
1032
3
    bool ret = _heavy_work_pool.try_offer([request, result, done]() {
1033
3
        VLOG_RPC << "test jdbc connection";
1034
3
        brpc::ClosureGuard closure_guard(done);
1035
3
        std::shared_ptr<MemTrackerLimiter> mem_tracker = MemTrackerLimiter::create_shared(
1036
3
                MemTrackerLimiter::Type::OTHER,
1037
3
                fmt::format("InternalService::test_jdbc_connection"));
1038
3
        SCOPED_ATTACH_TASK(mem_tracker);
1039
3
        TTableDescriptor table_desc;
1040
3
        Status st = Status::OK();
1041
3
        {
1042
3
            const uint8_t* buf = (const uint8_t*)request->jdbc_table().data();
1043
3
            uint32_t len = request->jdbc_table().size();
1044
3
            st = deserialize_thrift_msg(buf, &len, false, &table_desc);
1045
3
            if (!st.ok()) {
1046
0
                LOG(WARNING) << "test jdbc connection failed, errmsg=" << st;
1047
0
                st.to_protobuf(result->mutable_status());
1048
0
                return;
1049
0
            }
1050
3
        }
1051
3
        TJdbcTable jdbc_table = (table_desc.jdbcTable);
1052
1053
        // Resolve driver URL to absolute file:// path
1054
3
        std::string driver_url;
1055
3
        st = JdbcUtils::resolve_driver_url(jdbc_table.jdbc_driver_url, &driver_url);
1056
3
        if (!st.ok()) {
1057
0
            st.to_protobuf(result->mutable_status());
1058
0
            return;
1059
0
        }
1060
1061
        // Build params for JdbcConnectionTester
1062
3
        std::map<std::string, std::string> params;
1063
3
        params["jdbc_url"] = jdbc_table.jdbc_url;
1064
3
        params["jdbc_user"] = jdbc_table.jdbc_user;
1065
3
        params["jdbc_password"] = jdbc_table.jdbc_password;
1066
3
        params["jdbc_driver_class"] = jdbc_table.jdbc_driver_class;
1067
3
        params["jdbc_driver_url"] = driver_url;
1068
        // The catalog's expected MD5. Without it JdbcConnectionTester reads "" and
1069
        // JdbcDriverUtils.checksumVerifier("") is a no-op, so the one request whose entire job is
1070
        // to validate a catalog definition accepted any jar at all behind the driver URL.
1071
3
        params["jdbc_driver_checksum"] = jdbc_table.jdbc_driver_checksum;
1072
3
        params["query_sql"] = request->query_str();
1073
3
        params["catalog_id"] = std::to_string(jdbc_table.catalog_id);
1074
3
        params["connection_pool_min_size"] = std::to_string(jdbc_table.connection_pool_min_size);
1075
3
        params["connection_pool_max_size"] = std::to_string(jdbc_table.connection_pool_max_size);
1076
3
        params["connection_pool_max_wait_time"] =
1077
3
                std::to_string(jdbc_table.connection_pool_max_wait_time);
1078
3
        params["connection_pool_max_life_time"] =
1079
3
                std::to_string(jdbc_table.connection_pool_max_life_time);
1080
3
        params["connection_pool_keep_alive"] =
1081
3
                jdbc_table.connection_pool_keep_alive ? "true" : "false";
1082
3
        params["clean_datasource"] = "true";
1083
        // Map jdbc_table_type (TOdbcTableType enum value) to string name
1084
        // for JdbcTypeHandlerFactory to select the correct type handler.
1085
        // This ensures the right validation query is used (e.g. Oracle: "SELECT 1 FROM dual").
1086
3
        if (request->has_jdbc_table_type()) {
1087
3
            std::string type_name;
1088
3
            switch (request->jdbc_table_type()) {
1089
3
            case 0:
1090
3
                type_name = "MYSQL";
1091
3
                break;
1092
0
            case 1:
1093
0
                type_name = "ORACLE";
1094
0
                break;
1095
0
            case 2:
1096
0
                type_name = "POSTGRESQL";
1097
0
                break;
1098
0
            case 3:
1099
0
                type_name = "SQLSERVER";
1100
0
                break;
1101
0
            case 6:
1102
0
                type_name = "CLICKHOUSE";
1103
0
                break;
1104
0
            case 7:
1105
0
                type_name = "SAP_HANA";
1106
0
                break;
1107
0
            case 8:
1108
0
                type_name = "TRINO";
1109
0
                break;
1110
0
            case 9:
1111
0
                type_name = "PRESTO";
1112
0
                break;
1113
0
            case 10:
1114
0
                type_name = "OCEANBASE";
1115
0
                break;
1116
0
            case 11:
1117
0
                type_name = "OCEANBASE_ORACLE";
1118
0
                break;
1119
0
            case 13:
1120
0
                type_name = "DB2";
1121
0
                break;
1122
0
            case 14:
1123
0
                type_name = "GBASE";
1124
0
                break;
1125
0
            default:
1126
0
                break;
1127
3
            }
1128
3
            if (!type_name.empty()) {
1129
3
                params["table_type"] = type_name;
1130
3
            }
1131
3
        }
1132
        // required_fields and columns_types are required by JniReader
1133
3
        params["required_fields"] = "result";
1134
3
        params["columns_types"] = "int";
1135
1136
        // Use JniReader to create JdbcConnectionTester, which tests
1137
        // the connection in its open() method.
1138
3
        auto jni_reader = std::make_unique<JniReader>(Jni::plugin::JDBC_CONNECTION_TESTER, params);
1139
3
        st = jni_reader->open(nullptr, nullptr);
1140
3
        st.to_protobuf(result->mutable_status());
1141
1142
3
        Status close_st = jni_reader->close();
1143
3
        if (!close_st.ok()) {
1144
0
            LOG(WARNING) << "Failed to close JDBC connection tester: " << close_st.msg();
1145
0
        }
1146
3
    });
1147
1148
3
    if (!ret) {
1149
0
        offer_failed(result, done, _heavy_work_pool);
1150
0
        return;
1151
0
    }
1152
3
}
1153
1154
void PInternalServiceImpl::get_column_ids_by_tablet_ids(google::protobuf::RpcController* controller,
1155
                                                        const PFetchColIdsRequest* request,
1156
                                                        PFetchColIdsResponse* response,
1157
0
                                                        google::protobuf::Closure* done) {
1158
0
    bool ret = _light_work_pool.try_offer([this, controller, request, response, done]() {
1159
0
        _get_column_ids_by_tablet_ids(controller, request, response, done);
1160
0
    });
1161
0
    if (!ret) {
1162
0
        offer_failed(response, done, _light_work_pool);
1163
0
        return;
1164
0
    }
1165
0
}
1166
1167
void PInternalServiceImpl::_get_column_ids_by_tablet_ids(
1168
        google::protobuf::RpcController* controller, const PFetchColIdsRequest* request,
1169
0
        PFetchColIdsResponse* response, google::protobuf::Closure* done) {
1170
0
    brpc::ClosureGuard guard(done);
1171
0
    [[maybe_unused]] auto* cntl = static_cast<brpc::Controller*>(controller);
1172
0
    TabletManager* tablet_mgr = _engine.tablet_manager();
1173
0
    const auto& params = request->params();
1174
0
    for (const auto& param : params) {
1175
0
        int64_t index_id = param.indexid();
1176
0
        const auto& tablet_ids = param.tablet_ids();
1177
0
        std::set<std::set<int32_t>> filter_set;
1178
0
        std::map<int32_t, const TabletColumn*> id_to_column;
1179
0
        for (const int64_t tablet_id : tablet_ids) {
1180
0
            TabletSharedPtr tablet = tablet_mgr->get_tablet(tablet_id);
1181
0
            if (tablet == nullptr) {
1182
0
                std::stringstream ss;
1183
0
                ss << "cannot get tablet by id:" << tablet_id;
1184
0
                LOG(WARNING) << ss.str();
1185
0
                response->mutable_status()->set_status_code(TStatusCode::ILLEGAL_STATE);
1186
0
                response->mutable_status()->add_error_msgs(ss.str());
1187
0
                return;
1188
0
            }
1189
            // check schema consistency, column ids should be the same
1190
0
            const auto& columns = tablet->tablet_schema()->columns();
1191
1192
0
            std::set<int32_t> column_ids;
1193
0
            for (const auto& col : columns) {
1194
0
                column_ids.insert(col->unique_id());
1195
0
            }
1196
0
            filter_set.insert(std::move(column_ids));
1197
1198
0
            if (id_to_column.empty()) {
1199
0
                for (const auto& col : columns) {
1200
0
                    id_to_column.insert(std::pair {col->unique_id(), col.get()});
1201
0
                }
1202
0
            } else {
1203
0
                for (const auto& col : columns) {
1204
0
                    auto it = id_to_column.find(col->unique_id());
1205
0
                    if (it == id_to_column.end() || *(it->second) != *col) {
1206
0
                        ColumnPB prev_col_pb;
1207
0
                        ColumnPB curr_col_pb;
1208
0
                        if (it != id_to_column.end()) {
1209
0
                            it->second->to_schema_pb(&prev_col_pb);
1210
0
                        }
1211
0
                        col->to_schema_pb(&curr_col_pb);
1212
0
                        std::stringstream ss;
1213
0
                        ss << "consistency check failed: index{ " << index_id << " }"
1214
0
                           << " got inconsistent schema, prev column: " << prev_col_pb.DebugString()
1215
0
                           << " current column: " << curr_col_pb.DebugString();
1216
0
                        LOG(WARNING) << ss.str();
1217
0
                        response->mutable_status()->set_status_code(TStatusCode::ILLEGAL_STATE);
1218
0
                        response->mutable_status()->add_error_msgs(ss.str());
1219
0
                        return;
1220
0
                    }
1221
0
                }
1222
0
            }
1223
0
        }
1224
1225
0
        if (filter_set.size() > 1) {
1226
            // consistecy check failed
1227
0
            std::stringstream ss;
1228
0
            ss << "consistency check failed: index{" << index_id << "}"
1229
0
               << "got inconsistent schema";
1230
0
            LOG(WARNING) << ss.str();
1231
0
            response->mutable_status()->set_status_code(TStatusCode::ILLEGAL_STATE);
1232
0
            response->mutable_status()->add_error_msgs(ss.str());
1233
0
            return;
1234
0
        }
1235
        // consistency check passed, use the first tablet to be the representative
1236
0
        TabletSharedPtr tablet = tablet_mgr->get_tablet(tablet_ids[0]);
1237
0
        const auto& columns = tablet->tablet_schema()->columns();
1238
0
        auto entry = response->add_entries();
1239
0
        entry->set_index_id(index_id);
1240
0
        auto col_name_to_id = entry->mutable_col_name_to_id();
1241
0
        for (const auto& column : columns) {
1242
0
            (*col_name_to_id)[column->name()] = column->unique_id();
1243
0
        }
1244
0
    }
1245
0
    response->mutable_status()->set_status_code(TStatusCode::OK);
1246
0
}
1247
1248
template <class RPCResponse>
1249
struct AsyncRPCContext {
1250
    RPCResponse response;
1251
    brpc::Controller cntl;
1252
    brpc::CallId cid;
1253
};
1254
1255
void PInternalService::fetch_remote_tablet_schema(google::protobuf::RpcController* controller,
1256
                                                  const PFetchRemoteSchemaRequest* request,
1257
                                                  PFetchRemoteSchemaResponse* response,
1258
52
                                                  google::protobuf::Closure* done) {
1259
52
    bool ret = _heavy_work_pool.try_offer([request, response, done]() {
1260
52
        brpc::ClosureGuard closure_guard(done);
1261
52
        Status st = Status::OK();
1262
52
        std::shared_ptr<MemTrackerLimiter> mem_tracker = MemTrackerLimiter::create_shared(
1263
52
                MemTrackerLimiter::Type::OTHER,
1264
52
                fmt::format("InternalService::fetch_remote_tablet_schema"));
1265
52
        SCOPED_ATTACH_TASK(mem_tracker);
1266
52
        if (request->is_coordinator()) {
1267
            // Spawn rpc request to none coordinator nodes, and finally merge them all
1268
26
            PFetchRemoteSchemaRequest remote_request(*request);
1269
            // set it none coordinator to get merged schema
1270
26
            remote_request.set_is_coordinator(false);
1271
26
            using PFetchRemoteTabletSchemaRpcContext = AsyncRPCContext<PFetchRemoteSchemaResponse>;
1272
26
            std::vector<PFetchRemoteTabletSchemaRpcContext> rpc_contexts(
1273
26
                    request->tablet_location_size());
1274
52
            for (int i = 0; i < request->tablet_location_size(); ++i) {
1275
26
                std::string host = request->tablet_location(i).host();
1276
26
                int32_t brpc_port = request->tablet_location(i).brpc_port();
1277
26
                std::shared_ptr<PBackendService_Stub> stub(
1278
26
                        ExecEnv::GetInstance()->brpc_internal_client_cache()->get_client(
1279
26
                                host, brpc_port));
1280
26
                if (stub == nullptr) {
1281
0
                    LOG(WARNING) << "Failed to init rpc to " << host << ":" << brpc_port;
1282
0
                    st = Status::InternalError("Failed to init rpc to {}:{}", host, brpc_port);
1283
0
                    continue;
1284
0
                }
1285
26
                rpc_contexts[i].cid = rpc_contexts[i].cntl.call_id();
1286
26
                rpc_contexts[i].cntl.set_timeout_ms(config::fetch_remote_schema_rpc_timeout_ms);
1287
26
                stub->fetch_remote_tablet_schema(&rpc_contexts[i].cntl, &remote_request,
1288
26
                                                 &rpc_contexts[i].response, brpc::DoNothing());
1289
26
            }
1290
26
            std::vector<TabletSchemaSPtr> schemas;
1291
26
            for (auto& rpc_context : rpc_contexts) {
1292
26
                brpc::Join(rpc_context.cid);
1293
26
                if (!st.ok()) {
1294
                    // make sure all flying rpc request is joined
1295
0
                    continue;
1296
0
                }
1297
26
                if (rpc_context.cntl.Failed()) {
1298
0
                    LOG(WARNING) << "fetch_remote_tablet_schema rpc err:"
1299
0
                                 << rpc_context.cntl.ErrorText();
1300
0
                    ExecEnv::GetInstance()->brpc_internal_client_cache()->erase(
1301
0
                            rpc_context.cntl.remote_side());
1302
0
                    st = Status::InternalError("fetch_remote_tablet_schema rpc err: {}",
1303
0
                                               rpc_context.cntl.ErrorText());
1304
0
                }
1305
26
                if (rpc_context.response.status().status_code() != 0) {
1306
0
                    st = Status::create(rpc_context.response.status());
1307
0
                }
1308
26
                if (rpc_context.response.has_merged_schema()) {
1309
26
                    TabletSchemaSPtr schema = std::make_shared<TabletSchema>();
1310
26
                    schema->init_from_pb(rpc_context.response.merged_schema());
1311
26
                    schemas.push_back(schema);
1312
26
                }
1313
26
            }
1314
26
            if (!schemas.empty() && st.ok()) {
1315
                // merge all
1316
26
                TabletSchemaSPtr merged_schema;
1317
26
                st = variant_util::get_least_common_schema(schemas, nullptr, merged_schema);
1318
26
                if (!st.ok()) {
1319
0
                    LOG(WARNING) << "Failed to get least common schema: " << st.to_string();
1320
0
                    st = Status::InternalError("Failed to get least common schema: {}",
1321
0
                                               st.to_string());
1322
0
                }
1323
26
                VLOG_DEBUG << "dump schema:" << merged_schema->dump_structure();
1324
26
                merged_schema->reserve_extracted_columns();
1325
26
                merged_schema->to_schema_pb(response->mutable_merged_schema());
1326
26
            }
1327
26
            st.to_protobuf(response->mutable_status());
1328
26
            return;
1329
26
        } else {
1330
            // This is not a coordinator, get it's tablet and merge schema
1331
26
            std::vector<int64_t> target_tablets;
1332
26
            for (int i = 0; i < request->tablet_location_size(); ++i) {
1333
26
                const auto& location = request->tablet_location(i);
1334
26
                auto backend = BackendOptions::get_local_backend();
1335
                // If this is the target backend
1336
26
                if (backend.host == location.host() && config::brpc_port == location.brpc_port()) {
1337
26
                    target_tablets.assign(location.tablet_id().begin(), location.tablet_id().end());
1338
26
                    break;
1339
26
                }
1340
26
            }
1341
26
            if (!target_tablets.empty()) {
1342
26
                std::vector<TabletSchemaSPtr> tablet_schemas;
1343
1.33k
                for (int64_t tablet_id : target_tablets) {
1344
1.33k
                    auto res = ExecEnv::get_tablet(tablet_id);
1345
1.33k
                    if (!res.has_value()) {
1346
                        // just ignore
1347
0
                        LOG(WARNING) << "tablet does not exist, tablet id is " << tablet_id;
1348
0
                        continue;
1349
0
                    }
1350
1.33k
                    auto tablet = res.value();
1351
1.33k
                    auto rowsets = tablet->get_snapshot_rowset();
1352
1.33k
                    auto schema =
1353
1.33k
                            variant_util::VariantCompactionUtil::calculate_variant_extended_schema(
1354
1.33k
                                    rowsets, tablet->tablet_schema());
1355
1.33k
                    tablet_schemas.push_back(schema);
1356
1.33k
                }
1357
26
                if (!tablet_schemas.empty()) {
1358
                    // merge all
1359
26
                    TabletSchemaSPtr merged_schema;
1360
26
                    st = variant_util::get_least_common_schema(tablet_schemas, nullptr,
1361
26
                                                               merged_schema);
1362
26
                    if (!st.ok()) {
1363
0
                        LOG(WARNING) << "Failed to get least common schema: " << st.to_string();
1364
0
                        st = Status::InternalError("Failed to get least common schema: {}",
1365
0
                                                   st.to_string());
1366
0
                    }
1367
26
                    merged_schema->to_schema_pb(response->mutable_merged_schema());
1368
26
                    VLOG_DEBUG << "dump schema:" << merged_schema->dump_structure();
1369
26
                }
1370
26
            }
1371
26
            st.to_protobuf(response->mutable_status());
1372
26
        }
1373
52
    });
1374
52
    if (!ret) {
1375
0
        offer_failed(response, done, _heavy_work_pool);
1376
0
    }
1377
52
}
1378
1379
void PInternalService::report_stream_load_status(google::protobuf::RpcController* controller,
1380
                                                 const PReportStreamLoadStatusRequest* request,
1381
                                                 PReportStreamLoadStatusResponse* response,
1382
0
                                                 google::protobuf::Closure* done) {
1383
0
    TUniqueId load_id;
1384
0
    load_id.__set_hi(request->load_id().hi());
1385
0
    load_id.__set_lo(request->load_id().lo());
1386
0
    Status st = Status::OK();
1387
0
    auto stream_load_ctx = _exec_env->new_load_stream_mgr()->get(load_id);
1388
0
    if (!stream_load_ctx) {
1389
0
        st = Status::InternalError("unknown stream load id: {}", UniqueId(load_id).to_string());
1390
0
    }
1391
0
    stream_load_ctx->load_status_promise.set_value(st);
1392
0
    st.to_protobuf(response->mutable_status());
1393
0
}
1394
1395
void PInternalService::get_info(google::protobuf::RpcController* controller,
1396
                                const PProxyRequest* request, PProxyResult* response,
1397
3.11k
                                google::protobuf::Closure* done) {
1398
3.11k
    bool ret = _exec_env->routine_load_task_executor()->get_thread_pool().submit_func([this,
1399
3.11k
                                                                                       request,
1400
3.11k
                                                                                       response,
1401
3.11k
                                                                                       done]() {
1402
3.11k
        brpc::ClosureGuard closure_guard(done);
1403
        // PProxyRequest is defined in gensrc/proto/internal_service.proto
1404
        // Currently it supports 2 kinds of requests:
1405
        // 1. get all kafka partition ids for given topic
1406
        // 2. get all kafka partition offsets for given topic and timestamp.
1407
3.11k
        int timeout_ms = request->has_timeout_secs() ? request->timeout_secs() * 1000 : 60 * 1000;
1408
3.11k
        if (request->has_kafka_meta_request()) {
1409
3.11k
            const PKafkaMetaProxyRequest& kafka_request = request->kafka_meta_request();
1410
3.11k
            if (!kafka_request.offset_flags().empty()) {
1411
67
                std::vector<PIntegerPair> partition_offsets;
1412
67
                Status st = _exec_env->routine_load_task_executor()
1413
67
                                    ->get_kafka_real_offsets_for_partitions(
1414
67
                                            request->kafka_meta_request(), &partition_offsets,
1415
67
                                            timeout_ms);
1416
67
                if (st.ok()) {
1417
67
                    PKafkaPartitionOffsets* part_offsets = response->mutable_partition_offsets();
1418
67
                    for (const auto& entry : partition_offsets) {
1419
67
                        PIntegerPair* res = part_offsets->add_offset_times();
1420
67
                        res->set_key(entry.key());
1421
67
                        res->set_val(entry.val());
1422
67
                    }
1423
67
                }
1424
67
                st.to_protobuf(response->mutable_status());
1425
67
                return;
1426
3.04k
            } else if (!kafka_request.partition_id_for_latest_offsets().empty()) {
1427
                // get latest offsets for specified partition ids
1428
2.97k
                std::vector<PIntegerPair> partition_offsets;
1429
2.97k
                Status st = _exec_env->routine_load_task_executor()
1430
2.97k
                                    ->get_kafka_latest_offsets_for_partitions(
1431
2.97k
                                            request->kafka_meta_request(), &partition_offsets,
1432
2.97k
                                            timeout_ms);
1433
2.97k
                if (st.ok()) {
1434
2.97k
                    PKafkaPartitionOffsets* part_offsets = response->mutable_partition_offsets();
1435
2.97k
                    for (const auto& entry : partition_offsets) {
1436
2.97k
                        PIntegerPair* res = part_offsets->add_offset_times();
1437
2.97k
                        res->set_key(entry.key());
1438
2.97k
                        res->set_val(entry.val());
1439
2.97k
                    }
1440
2.97k
                }
1441
2.97k
                st.to_protobuf(response->mutable_status());
1442
2.97k
                return;
1443
2.97k
            } else if (!kafka_request.offset_times().empty()) {
1444
                // if offset_times() has elements, which means this request is to get offset by timestamp.
1445
1
                std::vector<PIntegerPair> partition_offsets;
1446
1
                Status st = _exec_env->routine_load_task_executor()
1447
1
                                    ->get_kafka_partition_offsets_for_times(
1448
1
                                            request->kafka_meta_request(), &partition_offsets,
1449
1
                                            timeout_ms);
1450
1
                if (st.ok()) {
1451
1
                    PKafkaPartitionOffsets* part_offsets = response->mutable_partition_offsets();
1452
1
                    for (const auto& entry : partition_offsets) {
1453
1
                        PIntegerPair* res = part_offsets->add_offset_times();
1454
1
                        res->set_key(entry.key());
1455
1
                        res->set_val(entry.val());
1456
1
                    }
1457
1
                }
1458
1
                st.to_protobuf(response->mutable_status());
1459
1
                return;
1460
76
            } else {
1461
                // get partition ids of topic
1462
76
                std::vector<int32_t> partition_ids;
1463
76
                Status st = _exec_env->routine_load_task_executor()->get_kafka_partition_meta(
1464
76
                        request->kafka_meta_request(), &partition_ids);
1465
76
                if (st.ok()) {
1466
73
                    PKafkaMetaProxyResult* kafka_result = response->mutable_kafka_meta_result();
1467
73
                    for (int32_t id : partition_ids) {
1468
73
                        kafka_result->add_partition_ids(id);
1469
73
                    }
1470
73
                }
1471
76
                st.to_protobuf(response->mutable_status());
1472
76
                return;
1473
76
            }
1474
3.11k
        }
1475
0
        if (request->has_kinesis_meta_request()) {
1476
0
            std::vector<std::string> shard_ids;
1477
0
            Status st = _exec_env->routine_load_task_executor()->get_kinesis_shard_meta(
1478
0
                    request->kinesis_meta_request(), &shard_ids);
1479
0
            if (st.ok()) {
1480
0
                PKinesisMetaProxyResult* kinesis_result = response->mutable_kinesis_meta_result();
1481
0
                for (const auto& shard_id : shard_ids) {
1482
0
                    kinesis_result->add_shard_ids(shard_id);
1483
0
                }
1484
0
            }
1485
0
            st.to_protobuf(response->mutable_status());
1486
0
            return;
1487
0
        }
1488
0
        Status::OK().to_protobuf(response->mutable_status());
1489
0
    });
1490
3.11k
    if (!ret) {
1491
0
        offer_failed(response, done, _heavy_work_pool);
1492
0
        return;
1493
0
    }
1494
3.11k
}
1495
1496
void PInternalService::update_cache(google::protobuf::RpcController* controller,
1497
                                    const PUpdateCacheRequest* request, PCacheResponse* response,
1498
67.3k
                                    google::protobuf::Closure* done) {
1499
67.4k
    bool ret = _light_work_pool.try_offer([this, request, response, done]() {
1500
67.4k
        brpc::ClosureGuard closure_guard(done);
1501
67.4k
        _exec_env->result_cache()->update(request, response);
1502
67.4k
    });
1503
67.3k
    if (!ret) {
1504
0
        offer_failed(response, done, _light_work_pool);
1505
0
        return;
1506
0
    }
1507
67.3k
}
1508
1509
void PInternalService::fetch_cache(google::protobuf::RpcController* controller,
1510
                                   const PFetchCacheRequest* request, PFetchCacheResult* result,
1511
3.22k
                                   google::protobuf::Closure* done) {
1512
3.22k
    bool ret = _light_work_pool.try_offer([this, request, result, done]() {
1513
3.22k
        brpc::ClosureGuard closure_guard(done);
1514
3.22k
        _exec_env->result_cache()->fetch(request, result);
1515
3.22k
    });
1516
3.22k
    if (!ret) {
1517
0
        offer_failed(result, done, _light_work_pool);
1518
0
        return;
1519
0
    }
1520
3.22k
}
1521
1522
void PInternalService::clear_cache(google::protobuf::RpcController* controller,
1523
                                   const PClearCacheRequest* request, PCacheResponse* response,
1524
0
                                   google::protobuf::Closure* done) {
1525
0
    bool ret = _light_work_pool.try_offer([this, request, response, done]() {
1526
0
        brpc::ClosureGuard closure_guard(done);
1527
0
        _exec_env->result_cache()->clear(request, response);
1528
0
    });
1529
0
    if (!ret) {
1530
0
        offer_failed(response, done, _light_work_pool);
1531
0
        return;
1532
0
    }
1533
0
}
1534
1535
void PInternalService::merge_filter(::google::protobuf::RpcController* controller,
1536
                                    const ::doris::PMergeFilterRequest* request,
1537
                                    ::doris::PMergeFilterResponse* response,
1538
3.25k
                                    ::google::protobuf::Closure* done) {
1539
3.26k
    bool ret = _light_work_pool.try_offer([this, controller, request, response, done]() {
1540
3.26k
        signal::SignalTaskIdKeeper keeper(request->query_id());
1541
3.26k
        brpc::ClosureGuard closure_guard(done);
1542
3.26k
        auto attachment = static_cast<brpc::Controller*>(controller)->request_attachment();
1543
3.26k
        butil::IOBufAsZeroCopyInputStream zero_copy_input_stream(attachment);
1544
3.26k
        Status st;
1545
3.26k
        try {
1546
3.26k
            st = _exec_env->fragment_mgr()->merge_filter(request, &zero_copy_input_stream);
1547
3.26k
        } catch (Exception& e) {
1548
0
            st = e.to_status();
1549
0
        }
1550
3.26k
        st.to_protobuf(response->mutable_status());
1551
3.26k
    });
1552
3.25k
    if (!ret) {
1553
0
        offer_failed(response, done, _light_work_pool);
1554
0
        return;
1555
0
    }
1556
3.25k
}
1557
1558
void PInternalService::send_filter_size(::google::protobuf::RpcController* controller,
1559
                                        const ::doris::PSendFilterSizeRequest* request,
1560
                                        ::doris::PSendFilterSizeResponse* response,
1561
100
                                        ::google::protobuf::Closure* done) {
1562
100
    bool ret = _light_work_pool.try_offer([this, request, response, done]() {
1563
100
        signal::SignalTaskIdKeeper keeper(request->query_id());
1564
100
        brpc::ClosureGuard closure_guard(done);
1565
100
        Status st;
1566
100
        try {
1567
100
            st = _exec_env->fragment_mgr()->send_filter_size(request);
1568
100
        } catch (Exception& e) {
1569
0
            st = e.to_status();
1570
0
        }
1571
100
        st.to_protobuf(response->mutable_status());
1572
100
    });
1573
100
    if (!ret) {
1574
0
        offer_failed(response, done, _light_work_pool);
1575
0
        return;
1576
0
    }
1577
100
}
1578
1579
void PInternalService::sync_filter_size(::google::protobuf::RpcController* controller,
1580
                                        const ::doris::PSyncFilterSizeRequest* request,
1581
                                        ::doris::PSyncFilterSizeResponse* response,
1582
100
                                        ::google::protobuf::Closure* done) {
1583
100
    bool ret = _light_work_pool.try_offer([this, request, response, done]() {
1584
100
        signal::SignalTaskIdKeeper keeper(request->query_id());
1585
100
        brpc::ClosureGuard closure_guard(done);
1586
100
        Status st;
1587
100
        try {
1588
100
            st = _exec_env->fragment_mgr()->sync_filter_size(request);
1589
100
        } catch (Exception& e) {
1590
0
            st = e.to_status();
1591
0
        }
1592
100
        st.to_protobuf(response->mutable_status());
1593
100
    });
1594
100
    if (!ret) {
1595
0
        offer_failed(response, done, _light_work_pool);
1596
0
        return;
1597
0
    }
1598
100
}
1599
1600
void PInternalService::apply_filterv2(::google::protobuf::RpcController* controller,
1601
                                      const ::doris::PPublishFilterRequestV2* request,
1602
                                      ::doris::PPublishFilterResponse* response,
1603
2.16k
                                      ::google::protobuf::Closure* done) {
1604
2.16k
    bool ret = _light_work_pool.try_offer([this, controller, request, response, done]() {
1605
2.16k
        signal::SignalTaskIdKeeper keeper(request->query_id());
1606
2.16k
        brpc::ClosureGuard closure_guard(done);
1607
2.16k
        const butil::IOBuf& request_attachment =
1608
2.16k
                static_cast<brpc::Controller*>(controller)->request_attachment();
1609
2.16k
        butil::IOBuf apply_attachment = request_attachment;
1610
2.16k
        butil::IOBuf forward_attachment = request_attachment;
1611
2.16k
        butil::IOBufAsZeroCopyInputStream zero_copy_input_stream(apply_attachment);
1612
2.16k
        VLOG_NOTICE << "rpc apply_filterv2 recv";
1613
2.16k
        Status st;
1614
2.16k
        try {
1615
2.16k
            st = _exec_env->fragment_mgr()->apply_filterv2(request, &zero_copy_input_stream);
1616
2.16k
        } catch (Exception& e) {
1617
0
            st = e.to_status();
1618
0
        }
1619
2.16k
        if (!st.ok()) {
1620
0
            LOG(WARNING) << "apply filter meet error: " << st.to_string();
1621
0
        }
1622
2.16k
        std::weak_ptr<QueryContext> forward_ctx;
1623
2.16k
        if (auto query_ctx = _exec_env->fragment_mgr()->get_query_ctx(
1624
2.16k
                    UniqueId(request->query_id()).to_thrift())) {
1625
2.16k
            if (!query_ctx->ignore_runtime_filter_error()) {
1626
2.16k
                forward_ctx = query_ctx;
1627
2.16k
            }
1628
2.16k
        }
1629
2.16k
        Status forward_st = forward_runtime_filter(*request, forward_attachment, forward_ctx);
1630
2.16k
        if (!forward_st.ok()) {
1631
0
            LOG(WARNING) << "forward runtime filter meet error: " << forward_st.to_string();
1632
0
            if (st.ok()) {
1633
0
                st = std::move(forward_st);
1634
0
            }
1635
0
        }
1636
2.16k
        st.to_protobuf(response->mutable_status());
1637
2.16k
    });
1638
2.16k
    if (!ret) {
1639
0
        offer_failed(response, done, _light_work_pool);
1640
0
        return;
1641
0
    }
1642
2.16k
}
1643
1644
void PInternalService::send_data(google::protobuf::RpcController* controller,
1645
                                 const PSendDataRequest* request, PSendDataResult* response,
1646
44
                                 google::protobuf::Closure* done) {
1647
44
    bool ret = _heavy_work_pool.try_offer([this, request, response, done]() {
1648
44
        brpc::ClosureGuard closure_guard(done);
1649
44
        TUniqueId load_id;
1650
44
        load_id.hi = request->load_id().hi();
1651
44
        load_id.lo = request->load_id().lo();
1652
        // On 1.2.3 we add load id to send data request and using load id to get pipe
1653
44
        auto stream_load_ctx = _exec_env->new_load_stream_mgr()->get(load_id);
1654
44
        if (stream_load_ctx == nullptr) {
1655
0
            response->mutable_status()->set_status_code(1);
1656
0
            response->mutable_status()->add_error_msgs("could not find stream load context");
1657
44
        } else {
1658
44
            auto pipe = stream_load_ctx->pipe;
1659
157
            for (int i = 0; i < request->data_size(); ++i) {
1660
113
                std::unique_ptr<PDataRow> row(new PDataRow());
1661
113
                row->CopyFrom(request->data(i));
1662
113
                Status s = pipe->append(std::move(row));
1663
113
                if (!s.ok()) {
1664
0
                    response->mutable_status()->set_status_code(1);
1665
0
                    response->mutable_status()->add_error_msgs(s.to_string());
1666
0
                    return;
1667
0
                }
1668
113
            }
1669
44
            response->mutable_status()->set_status_code(0);
1670
44
        }
1671
44
    });
1672
44
    if (!ret) {
1673
0
        offer_failed(response, done, _heavy_work_pool);
1674
0
        return;
1675
0
    }
1676
44
}
1677
1678
void PInternalService::commit(google::protobuf::RpcController* controller,
1679
                              const PCommitRequest* request, PCommitResult* response,
1680
44
                              google::protobuf::Closure* done) {
1681
44
    bool ret = _heavy_work_pool.try_offer([this, request, response, done]() {
1682
44
        brpc::ClosureGuard closure_guard(done);
1683
44
        TUniqueId load_id;
1684
44
        load_id.hi = request->load_id().hi();
1685
44
        load_id.lo = request->load_id().lo();
1686
1687
44
        auto stream_load_ctx = _exec_env->new_load_stream_mgr()->get(load_id);
1688
44
        if (stream_load_ctx == nullptr) {
1689
0
            response->mutable_status()->set_status_code(1);
1690
0
            response->mutable_status()->add_error_msgs("could not find stream load context");
1691
44
        } else {
1692
44
            static_cast<void>(stream_load_ctx->pipe->finish());
1693
44
            response->mutable_status()->set_status_code(0);
1694
44
        }
1695
44
    });
1696
44
    if (!ret) {
1697
0
        offer_failed(response, done, _heavy_work_pool);
1698
0
        return;
1699
0
    }
1700
44
}
1701
1702
void PInternalService::rollback(google::protobuf::RpcController* controller,
1703
                                const PRollbackRequest* request, PRollbackResult* response,
1704
5
                                google::protobuf::Closure* done) {
1705
5
    bool ret = _heavy_work_pool.try_offer([this, request, response, done]() {
1706
5
        brpc::ClosureGuard closure_guard(done);
1707
5
        TUniqueId load_id;
1708
5
        load_id.hi = request->load_id().hi();
1709
5
        load_id.lo = request->load_id().lo();
1710
5
        auto stream_load_ctx = _exec_env->new_load_stream_mgr()->get(load_id);
1711
5
        if (stream_load_ctx == nullptr) {
1712
0
            response->mutable_status()->set_status_code(1);
1713
0
            response->mutable_status()->add_error_msgs("could not find stream load context");
1714
5
        } else {
1715
5
            stream_load_ctx->pipe->cancel("rollback");
1716
5
            response->mutable_status()->set_status_code(0);
1717
5
        }
1718
5
    });
1719
5
    if (!ret) {
1720
0
        offer_failed(response, done, _heavy_work_pool);
1721
0
        return;
1722
0
    }
1723
5
}
1724
1725
void PInternalService::fold_constant_expr(google::protobuf::RpcController* controller,
1726
                                          const PConstantExprRequest* request,
1727
                                          PConstantExprResult* response,
1728
574
                                          google::protobuf::Closure* done) {
1729
574
    bool ret = _light_work_pool.try_offer([request, response, done]() {
1730
574
        brpc::ClosureGuard closure_guard(done);
1731
574
        TFoldConstantParams t_request;
1732
574
        Status st = Status::OK();
1733
574
        {
1734
574
            const uint8_t* buf = (const uint8_t*)request->request().data();
1735
574
            uint32_t len = request->request().size();
1736
574
            st = deserialize_thrift_msg(buf, &len, false, &t_request);
1737
574
        }
1738
574
        if (!st.ok()) {
1739
0
            LOG(WARNING) << "exec fold constant expr failed, errmsg=" << st
1740
0
                         << " .and query_id_is: " << t_request.query_id;
1741
0
            st.to_protobuf(response->mutable_status());
1742
0
            return;
1743
0
        }
1744
574
        auto fold_func = [&]() -> Status {
1745
574
            std::unique_ptr<FoldConstantExecutor> fold_executor =
1746
574
                    std::make_unique<FoldConstantExecutor>();
1747
574
            RETURN_IF_ERROR_OR_CATCH_EXCEPTION(
1748
574
                    fold_executor->fold_constant_vexpr(t_request, response));
1749
529
            return Status::OK();
1750
574
        };
1751
574
        st = fold_func();
1752
574
        if (!st.ok()) {
1753
45
            LOG(WARNING) << "exec fold constant expr failed, errmsg=" << st
1754
45
                         << " .and query_id_is: " << t_request.query_id;
1755
45
        }
1756
574
        st.to_protobuf(response->mutable_status());
1757
574
    });
1758
574
    if (!ret) {
1759
0
        offer_failed(response, done, _light_work_pool);
1760
0
        return;
1761
0
    }
1762
574
}
1763
1764
void PInternalService::transmit_rec_cte_block(google::protobuf::RpcController* controller,
1765
                                              const PTransmitRecCTEBlockParams* request,
1766
                                              PTransmitRecCTEBlockResult* response,
1767
4.26k
                                              google::protobuf::Closure* done) {
1768
4.26k
    bool ret = _light_work_pool.try_offer([this, request, response, done]() {
1769
4.26k
        brpc::ClosureGuard closure_guard(done);
1770
4.26k
        auto st = _exec_env->fragment_mgr()->transmit_rec_cte_block(
1771
4.26k
                UniqueId(request->query_id()).to_thrift(),
1772
4.26k
                UniqueId(request->fragment_instance_id()).to_thrift(), request->node_id(),
1773
4.26k
                request->blocks(), request->eos());
1774
4.26k
        st.to_protobuf(response->mutable_status());
1775
4.26k
    });
1776
4.26k
    if (!ret) {
1777
0
        offer_failed(response, done, _light_work_pool);
1778
0
        return;
1779
0
    }
1780
4.26k
}
1781
1782
void PInternalService::rerun_fragment(google::protobuf::RpcController* controller,
1783
                                      const PRerunFragmentParams* request,
1784
                                      PRerunFragmentResult* response,
1785
10.8k
                                      google::protobuf::Closure* done) {
1786
10.8k
    bool ret = _light_work_pool.try_offer([this, request, response, done]() {
1787
        // Use shared_ptr<ClosureGuard> so we can transfer ownership to the PFC.
1788
        // For wait_for_destroy/final_close, the guard is stored in the PFC and the RPC
1789
        // response is deferred until the PFC is fully destroyed. For rebuild/submit,
1790
        // the guard fires immediately when this lambda returns.
1791
10.8k
        std::shared_ptr<brpc::ClosureGuard> closure_guard =
1792
10.8k
                std::make_shared<brpc::ClosureGuard>(done);
1793
10.8k
        auto st = _exec_env->fragment_mgr()->rerun_fragment(
1794
10.8k
                closure_guard, UniqueId(request->query_id()).to_thrift(), request->fragment_id(),
1795
10.8k
                request->stage());
1796
10.8k
        st.to_protobuf(response->mutable_status());
1797
10.8k
    });
1798
10.8k
    if (!ret) {
1799
0
        offer_failed(response, done, _light_work_pool);
1800
0
        return;
1801
0
    }
1802
10.8k
}
1803
1804
void PInternalService::reset_global_rf(google::protobuf::RpcController* controller,
1805
                                       const PResetGlobalRfParams* request,
1806
                                       PResetGlobalRfResult* response,
1807
2.06k
                                       google::protobuf::Closure* done) {
1808
2.06k
    bool ret = _light_work_pool.try_offer([this, request, response, done]() {
1809
2.06k
        brpc::ClosureGuard closure_guard(done);
1810
2.06k
        auto st = _exec_env->fragment_mgr()->reset_global_rf(
1811
2.06k
                UniqueId(request->query_id()).to_thrift(), request->filter_ids());
1812
2.06k
        st.to_protobuf(response->mutable_status());
1813
2.06k
    });
1814
2.06k
    if (!ret) {
1815
0
        offer_failed(response, done, _light_work_pool);
1816
0
        return;
1817
0
    }
1818
2.06k
}
1819
1820
void PInternalService::transmit_block(google::protobuf::RpcController* controller,
1821
                                      const PTransmitDataParams* request,
1822
                                      PTransmitDataResult* response,
1823
1.84M
                                      google::protobuf::Closure* done) {
1824
1.84M
    int64_t receive_time = GetCurrentTimeNanos();
1825
1.86M
    if (config::enable_bthread_transmit_block) {
1826
1.86M
        response->set_receive_time(receive_time);
1827
        // under high concurrency, thread pool will have a lot of lock contention.
1828
        // May offer failed to the thread pool, so that we should avoid using thread
1829
        // pool here.
1830
1.86M
        _transmit_block(controller, request, response, done, Status::OK(), 0);
1831
18.4E
    } else {
1832
18.4E
        bool ret = _light_work_pool.try_offer([this, controller, request, response, done,
1833
18.4E
                                               receive_time]() {
1834
0
            response->set_receive_time(receive_time);
1835
            // Sometimes transmit block function is the last owner of PlanFragmentExecutor
1836
            // It will release the object. And the object maybe a JNIContext.
1837
            // JNIContext will hold some TLS object. It could not work correctly under bthread
1838
            // Context. So that put the logic into pthread.
1839
            // But this is rarely happens, so this config is disabled by default.
1840
0
            _transmit_block(controller, request, response, done, Status::OK(),
1841
0
                            GetCurrentTimeNanos() - receive_time);
1842
0
        });
1843
18.4E
        if (!ret) {
1844
0
            offer_failed(response, done, _light_work_pool);
1845
0
            return;
1846
0
        }
1847
18.4E
    }
1848
1.84M
}
1849
1850
void PInternalService::transmit_block_by_http(google::protobuf::RpcController* controller,
1851
                                              const PEmptyRequest* request,
1852
                                              PTransmitDataResult* response,
1853
0
                                              google::protobuf::Closure* done) {
1854
0
    int64_t receive_time = GetCurrentTimeNanos();
1855
0
    bool ret = _heavy_work_pool.try_offer([this, controller, response, done, receive_time]() {
1856
0
        PTransmitDataParams* new_request = new PTransmitDataParams();
1857
0
        google::protobuf::Closure* new_done =
1858
0
                new NewHttpClosure<PTransmitDataParams>(new_request, done);
1859
0
        brpc::Controller* cntl = static_cast<brpc::Controller*>(controller);
1860
0
        Status st =
1861
0
                attachment_extract_request_contain_block<PTransmitDataParams>(new_request, cntl);
1862
0
        _transmit_block(controller, new_request, response, new_done, st,
1863
0
                        GetCurrentTimeNanos() - receive_time);
1864
0
    });
1865
0
    if (!ret) {
1866
0
        offer_failed(response, done, _heavy_work_pool);
1867
0
        return;
1868
0
    }
1869
0
}
1870
1871
void PInternalService::_transmit_block(google::protobuf::RpcController* controller,
1872
                                       const PTransmitDataParams* request,
1873
                                       PTransmitDataResult* response,
1874
                                       google::protobuf::Closure* done, const Status& extract_st,
1875
1.84M
                                       const int64_t wait_for_worker) {
1876
1.85M
    if (request->has_query_id()) {
1877
18.4E
        VLOG_ROW << "transmit block: fragment_instance_id=" << print_id(request->finst_id())
1878
18.4E
                 << " query_id=" << print_id(request->query_id()) << " node=" << request->node_id();
1879
1.85M
    }
1880
1881
    // The response is accessed when done->Run is called in transmit_block(),
1882
    // give response a default value to avoid null pointers in high concurrency.
1883
1.84M
    Status st;
1884
1.84M
    if (extract_st.ok()) {
1885
1.84M
        st = _exec_env->vstream_mgr()->transmit_block(request, &done, wait_for_worker);
1886
1.84M
        if (!st.ok() && !st.is<END_OF_FILE>()) {
1887
0
            LOG(WARNING) << "transmit_block failed, message=" << st
1888
0
                         << ", fragment_instance_id=" << print_id(request->finst_id())
1889
0
                         << ", node=" << request->node_id()
1890
0
                         << ", from sender_id: " << request->sender_id()
1891
0
                         << ", be_number: " << request->be_number()
1892
0
                         << ", packet_seq: " << request->packet_seq();
1893
0
        }
1894
1.84M
    } else {
1895
1.17k
        st = extract_st;
1896
1.17k
    }
1897
1.85M
    if (done != nullptr) {
1898
1.85M
        st.to_protobuf(response->mutable_status());
1899
1.85M
        done->Run();
1900
1.85M
    }
1901
1.84M
}
1902
1903
void PInternalService::check_rpc_channel(google::protobuf::RpcController* controller,
1904
                                         const PCheckRPCChannelRequest* request,
1905
                                         PCheckRPCChannelResponse* response,
1906
0
                                         google::protobuf::Closure* done) {
1907
0
    bool ret = _light_work_pool.try_offer([request, response, done]() {
1908
0
        brpc::ClosureGuard closure_guard(done);
1909
0
        response->mutable_status()->set_status_code(0);
1910
0
        if (request->data().size() != request->size()) {
1911
0
            std::stringstream ss;
1912
0
            ss << "data size not same, expected: " << request->size()
1913
0
               << ", actual: " << request->data().size();
1914
0
            response->mutable_status()->add_error_msgs(ss.str());
1915
0
            response->mutable_status()->set_status_code(1);
1916
1917
0
        } else {
1918
0
            Md5Digest digest;
1919
0
            digest.update(static_cast<const void*>(request->data().c_str()),
1920
0
                          request->data().size());
1921
0
            digest.digest();
1922
0
            if (!iequal(digest.hex(), request->md5())) {
1923
0
                std::stringstream ss;
1924
0
                ss << "md5 not same, expected: " << request->md5() << ", actual: " << digest.hex();
1925
0
                response->mutable_status()->add_error_msgs(ss.str());
1926
0
                response->mutable_status()->set_status_code(1);
1927
0
            }
1928
0
        }
1929
0
    });
1930
0
    if (!ret) {
1931
0
        offer_failed(response, done, _light_work_pool);
1932
0
        return;
1933
0
    }
1934
0
}
1935
1936
void PInternalService::reset_rpc_channel(google::protobuf::RpcController* controller,
1937
                                         const PResetRPCChannelRequest* request,
1938
                                         PResetRPCChannelResponse* response,
1939
0
                                         google::protobuf::Closure* done) {
1940
0
    bool ret = _light_work_pool.try_offer([request, response, done]() {
1941
0
        brpc::ClosureGuard closure_guard(done);
1942
0
        response->mutable_status()->set_status_code(0);
1943
0
        if (request->all()) {
1944
0
            int size = ExecEnv::GetInstance()->brpc_internal_client_cache()->size();
1945
0
            if (size > 0) {
1946
0
                std::vector<std::string> endpoints;
1947
0
                ExecEnv::GetInstance()->brpc_internal_client_cache()->get_all(&endpoints);
1948
0
                ExecEnv::GetInstance()->brpc_internal_client_cache()->clear();
1949
0
                *response->mutable_channels() = {endpoints.begin(), endpoints.end()};
1950
0
            }
1951
0
        } else {
1952
0
            for (const std::string& endpoint : request->endpoints()) {
1953
0
                if (!ExecEnv::GetInstance()->brpc_internal_client_cache()->exist(endpoint)) {
1954
0
                    response->mutable_status()->add_error_msgs(endpoint + ": not found.");
1955
0
                    continue;
1956
0
                }
1957
1958
0
                if (ExecEnv::GetInstance()->brpc_internal_client_cache()->erase(endpoint)) {
1959
0
                    response->add_channels(endpoint);
1960
0
                } else {
1961
0
                    response->mutable_status()->add_error_msgs(endpoint + ": reset failed.");
1962
0
                }
1963
0
            }
1964
0
            if (request->endpoints_size() != response->channels_size()) {
1965
0
                response->mutable_status()->set_status_code(1);
1966
0
            }
1967
0
        }
1968
0
    });
1969
0
    if (!ret) {
1970
0
        offer_failed(response, done, _light_work_pool);
1971
0
        return;
1972
0
    }
1973
0
}
1974
1975
void PInternalService::hand_shake(google::protobuf::RpcController* controller,
1976
                                  const PHandShakeRequest* request, PHandShakeResponse* response,
1977
1.78k
                                  google::protobuf::Closure* done) {
1978
    // The light pool may be full. Handshake is used to check the connection state of brpc.
1979
    // Should not be interfered by the thread pool logic.
1980
1.78k
    brpc::ClosureGuard closure_guard(done);
1981
1.78k
    if (request->has_hello()) {
1982
1.78k
        response->set_hello(request->hello());
1983
1.78k
    }
1984
1.78k
    response->mutable_status()->set_status_code(0);
1985
1.78k
}
1986
1987
void PInternalServiceImpl::request_slave_tablet_pull_rowset(
1988
        google::protobuf::RpcController* controller, const PTabletWriteSlaveRequest* request,
1989
0
        PTabletWriteSlaveResult* response, google::protobuf::Closure* done) {
1990
0
    brpc::ClosureGuard closure_guard(done);
1991
0
    Status::NotSupported("single replica load has been removed")
1992
0
            .to_protobuf(response->mutable_status());
1993
0
}
1994
1995
void PInternalServiceImpl::response_slave_tablet_pull_rowset(
1996
        google::protobuf::RpcController* controller, const PTabletWriteSlaveDoneRequest* request,
1997
0
        PTabletWriteSlaveDoneResult* response, google::protobuf::Closure* done) {
1998
0
    brpc::ClosureGuard closure_guard(done);
1999
0
    Status::NotSupported("single replica load has been removed")
2000
0
            .to_protobuf(response->mutable_status());
2001
0
}
2002
2003
void PInternalService::multiget_data(google::protobuf::RpcController* controller,
2004
                                     const PMultiGetRequest* request, PMultiGetResponse* response,
2005
0
                                     google::protobuf::Closure* done) {
2006
0
    brpc::ClosureGuard closure_guard(done);
2007
0
#pragma GCC diagnostic push
2008
0
#pragma GCC diagnostic ignored "-Wdeprecated-declarations"
2009
    // The deprecated response field is retained only to reject legacy callers explicitly.
2010
0
    Status::NotSupported("multiget_data is deprecated; use multiget_data_v2")
2011
0
            .to_protobuf(response->mutable_status());
2012
0
#pragma GCC diagnostic pop
2013
0
}
2014
2015
void PInternalService::multiget_data_v2(google::protobuf::RpcController* controller,
2016
                                        const PMultiGetRequestV2* request,
2017
                                        PMultiGetResponseV2* response,
2018
1.17k
                                        google::protobuf::Closure* done) {
2019
1.17k
    std::vector<uint64_t> id_set;
2020
1.17k
    id_set.push_back(request->wg_id());
2021
1.17k
    auto wg = ExecEnv::GetInstance()->workload_group_mgr()->get_group(id_set);
2022
1.17k
    Status st = Status::OK();
2023
2024
1.17k
    if (!wg) [[unlikely]] {
2025
0
        brpc::ClosureGuard closure_guard(done);
2026
0
        st = Status::Error<TStatusCode::CANCELLED>("fail to find wg: wg id:" +
2027
0
                                                   std::to_string(request->wg_id()));
2028
0
        st.to_protobuf(response->mutable_status());
2029
0
        return;
2030
0
    }
2031
2032
1.17k
    doris::TaskScheduler* exec_sched = nullptr;
2033
1.17k
    ScannerScheduler* scan_sched = nullptr;
2034
1.17k
    ScannerScheduler* remote_scan_sched = nullptr;
2035
1.17k
    wg->get_query_scheduler(&exec_sched, &scan_sched, &remote_scan_sched);
2036
1.17k
    DCHECK(remote_scan_sched);
2037
2038
1.17k
    st = remote_scan_sched->submit_scan_task(
2039
1.17k
            SimplifiedScanTask(
2040
1.17k
                    [request, response, done]() {
2041
1.17k
                        SCOPED_ATTACH_TASK(ExecEnv::GetInstance()->rowid_storage_reader_tracker());
2042
1.17k
                        signal::set_signal_task_id(request->query_id());
2043
                        // multi get data by rowid
2044
1.17k
                        MonotonicStopWatch watch;
2045
1.17k
                        watch.start();
2046
1.17k
                        brpc::ClosureGuard closure_guard(done);
2047
1.17k
                        response->mutable_status()->set_status_code(0);
2048
1.17k
                        Status st = RowIdStorageReader::read_by_rowids(*request, response);
2049
1.17k
                        st.to_protobuf(response->mutable_status());
2050
1.17k
                        LOG(INFO) << "multiget_data finished, cost(us):"
2051
1.17k
                                  << watch.elapsed_time() / 1000;
2052
1.17k
                        return true;
2053
1.17k
                    },
2054
1.17k
                    nullptr, nullptr),
2055
1.17k
            fmt::format("{}-multiget_data_v2", print_id(request->query_id())));
2056
2057
1.17k
    if (!st.ok()) {
2058
0
        brpc::ClosureGuard closure_guard(done);
2059
0
        st.to_protobuf(response->mutable_status());
2060
0
    }
2061
1.17k
}
2062
2063
void PInternalServiceImpl::get_tablet_rowset_versions(google::protobuf::RpcController* cntl_base,
2064
                                                      const PGetTabletVersionsRequest* request,
2065
                                                      PGetTabletVersionsResponse* response,
2066
0
                                                      google::protobuf::Closure* done) {
2067
0
    brpc::ClosureGuard closure_guard(done);
2068
0
    VLOG_DEBUG << "receive get tablet versions request: " << request->DebugString();
2069
0
    _engine.get_tablet_rowset_versions(request, response);
2070
0
}
2071
2072
void PInternalService::glob(google::protobuf::RpcController* controller,
2073
                            const PGlobRequest* request, PGlobResponse* response,
2074
27
                            google::protobuf::Closure* done) {
2075
27
    bool ret = _heavy_work_pool.try_offer([request, response, done]() {
2076
27
        brpc::ClosureGuard closure_guard(done);
2077
27
        std::vector<io::FileInfo> files;
2078
27
        Status st = io::global_local_filesystem()->safe_glob(request->pattern(), &files);
2079
27
        if (st.ok()) {
2080
24
            for (auto& file : files) {
2081
24
                PGlobResponse_PFileInfo* pfile = response->add_files();
2082
24
                pfile->set_file(file.file_name);
2083
24
                pfile->set_size(file.file_size);
2084
24
            }
2085
22
        }
2086
27
        st.to_protobuf(response->mutable_status());
2087
27
    });
2088
27
    if (!ret) {
2089
0
        offer_failed(response, done, _heavy_work_pool);
2090
0
        return;
2091
0
    }
2092
27
}
2093
2094
void PInternalService::group_commit_insert(google::protobuf::RpcController* controller,
2095
                                           const PGroupCommitInsertRequest* request,
2096
                                           PGroupCommitInsertResponse* response,
2097
29
                                           google::protobuf::Closure* done) {
2098
29
    TUniqueId load_id;
2099
29
    load_id.__set_hi(request->load_id().hi());
2100
29
    load_id.__set_lo(request->load_id().lo());
2101
29
    std::shared_ptr<std::mutex> lock = std::make_shared<std::mutex>();
2102
29
    std::shared_ptr<bool> is_done = std::make_shared<bool>(false);
2103
29
    bool ret = _heavy_work_pool.try_offer([this, request, response, done, load_id, lock,
2104
29
                                           is_done]() {
2105
29
        brpc::ClosureGuard closure_guard(done);
2106
29
        std::shared_ptr<StreamLoadContext> ctx = std::make_shared<StreamLoadContext>(_exec_env);
2107
29
        auto pipe = std::make_shared<io::StreamLoadPipe>(
2108
29
                io::kMaxPipeBufferedBytes /* max_buffered_bytes */, 64 * 1024 /* min_chunk_size */,
2109
29
                -1 /* total_length */, true /* use_proto */);
2110
29
        ctx->pipe = pipe;
2111
29
        Status st = _exec_env->new_load_stream_mgr()->put(load_id, ctx);
2112
29
        if (st.ok()) {
2113
29
            try {
2114
29
                st = _exec_plan_fragment_impl(
2115
29
                        request->exec_plan_fragment_request().request(),
2116
29
                        request->exec_plan_fragment_request().version(),
2117
29
                        request->exec_plan_fragment_request().compact(),
2118
29
                        [&, response, done, load_id, lock, is_done](RuntimeState* state,
2119
29
                                                                    Status* status) {
2120
29
                            std::lock_guard<std::mutex> lock1(*lock);
2121
29
                            if (*is_done) {
2122
0
                                return;
2123
0
                            }
2124
29
                            *is_done = true;
2125
29
                            brpc::ClosureGuard cb_closure_guard(done);
2126
29
                            response->set_label(state->import_label());
2127
29
                            response->set_txn_id(state->wal_id());
2128
29
                            response->set_loaded_rows(state->num_rows_load_success());
2129
29
                            response->set_filtered_rows(state->num_rows_load_filtered());
2130
29
                            status->to_protobuf(response->mutable_status());
2131
29
                            if (!state->get_error_log_file_path().empty()) {
2132
0
                                response->set_error_url(
2133
0
                                        to_load_error_http_path(state->get_error_log_file_path()));
2134
0
                            }
2135
29
                            if (!state->get_first_error_msg().empty()) {
2136
0
                                response->set_first_error_msg(state->get_first_error_msg());
2137
0
                            }
2138
29
                            _exec_env->new_load_stream_mgr()->remove(load_id);
2139
29
                        });
2140
29
            } catch (const Exception& e) {
2141
0
                st = e.to_status();
2142
0
            } catch (const std::exception& e) {
2143
0
                st = Status::Error(ErrorCode::INTERNAL_ERROR, e.what());
2144
0
            } catch (...) {
2145
0
                st = Status::Error(ErrorCode::INTERNAL_ERROR,
2146
0
                                   "_exec_plan_fragment_impl meet unknown error");
2147
0
            }
2148
29
            if (!st.ok()) {
2149
0
                LOG(WARNING) << "exec plan fragment failed, load_id=" << print_id(load_id)
2150
0
                             << ", errmsg=" << st;
2151
0
                std::lock_guard<std::mutex> lock1(*lock);
2152
0
                if (*is_done) {
2153
0
                    closure_guard.release();
2154
0
                } else {
2155
0
                    *is_done = true;
2156
0
                    st.to_protobuf(response->mutable_status());
2157
0
                    _exec_env->new_load_stream_mgr()->remove(load_id);
2158
0
                }
2159
29
            } else {
2160
29
                closure_guard.release();
2161
66
                for (int i = 0; i < request->data().size(); ++i) {
2162
37
                    std::unique_ptr<PDataRow> row(new PDataRow());
2163
37
                    row->CopyFrom(request->data(i));
2164
37
                    st = pipe->append(std::move(row));
2165
37
                    if (!st.ok()) {
2166
0
                        break;
2167
0
                    }
2168
37
                }
2169
29
                if (st.ok()) {
2170
29
                    static_cast<void>(pipe->finish());
2171
29
                }
2172
29
            }
2173
29
        }
2174
29
    });
2175
29
    if (!ret) {
2176
0
        _exec_env->new_load_stream_mgr()->remove(load_id);
2177
0
        offer_failed(response, done, _heavy_work_pool);
2178
0
        return;
2179
0
    }
2180
29
};
2181
2182
void PInternalService::get_wal_queue_size(google::protobuf::RpcController* controller,
2183
                                          const PGetWalQueueSizeRequest* request,
2184
                                          PGetWalQueueSizeResponse* response,
2185
1.14k
                                          google::protobuf::Closure* done) {
2186
1.16k
    bool ret = _heavy_work_pool.try_offer([this, request, response, done]() {
2187
1.16k
        brpc::ClosureGuard closure_guard(done);
2188
1.16k
        Status st = Status::OK();
2189
1.16k
        auto table_id = request->table_id();
2190
1.16k
        auto count = _exec_env->wal_mgr()->get_wal_queue_size(table_id);
2191
1.16k
        response->set_size(count);
2192
1.16k
        response->mutable_status()->set_status_code(st.code());
2193
1.16k
    });
2194
1.14k
    if (!ret) {
2195
0
        offer_failed(response, done, _heavy_work_pool);
2196
0
    }
2197
1.14k
}
2198
2199
void PInternalService::get_be_resource(google::protobuf::RpcController* controller,
2200
                                       const PGetBeResourceRequest* request,
2201
                                       PGetBeResourceResponse* response,
2202
0
                                       google::protobuf::Closure* done) {
2203
0
    bool ret = _heavy_work_pool.try_offer([response, done]() {
2204
0
        brpc::ClosureGuard closure_guard(done);
2205
0
        int64_t mem_limit = MemInfo::mem_limit();
2206
0
        int64_t mem_usage = PerfCounters::get_vm_rss();
2207
2208
0
        PGlobalResourceUsage* global_resource_usage = response->mutable_global_be_resource_usage();
2209
0
        global_resource_usage->set_mem_limit(mem_limit);
2210
0
        global_resource_usage->set_mem_usage(mem_usage);
2211
2212
0
        Status st = Status::OK();
2213
0
        response->mutable_status()->set_status_code(st.code());
2214
0
    });
2215
0
    if (!ret) {
2216
0
        offer_failed(response, done, _heavy_work_pool);
2217
0
    }
2218
0
}
2219
2220
void PInternalService::delete_dictionary(google::protobuf::RpcController* controller,
2221
                                         const PDeleteDictionaryRequest* request,
2222
                                         PDeleteDictionaryResponse* response,
2223
3
                                         google::protobuf::Closure* done) {
2224
3
    brpc::ClosureGuard closure_guard(done);
2225
3
    Status st = ExecEnv::GetInstance()->dict_factory()->delete_dict(request->dictionary_id());
2226
3
    st.to_protobuf(response->mutable_status());
2227
3
}
2228
2229
void PInternalService::commit_refresh_dictionary(google::protobuf::RpcController* controller,
2230
                                                 const PCommitRefreshDictionaryRequest* request,
2231
                                                 PCommitRefreshDictionaryResponse* response,
2232
89
                                                 google::protobuf::Closure* done) {
2233
89
    brpc::ClosureGuard closure_guard(done);
2234
89
    Status st = ExecEnv::GetInstance()->dict_factory()->commit_refresh_dict(
2235
89
            request->dictionary_id(), request->version_id());
2236
89
    st.to_protobuf(response->mutable_status());
2237
89
}
2238
2239
void PInternalService::abort_refresh_dictionary(google::protobuf::RpcController* controller,
2240
                                                const PAbortRefreshDictionaryRequest* request,
2241
                                                PAbortRefreshDictionaryResponse* response,
2242
2
                                                google::protobuf::Closure* done) {
2243
2
    brpc::ClosureGuard closure_guard(done);
2244
2
    Status st = ExecEnv::GetInstance()->dict_factory()->abort_refresh_dict(request->dictionary_id(),
2245
2
                                                                           request->version_id());
2246
2
    st.to_protobuf(response->mutable_status());
2247
2
}
2248
2249
void PInternalService::get_tablet_rowsets(google::protobuf::RpcController* controller,
2250
                                          const PGetTabletRowsetsRequest* request,
2251
                                          PGetTabletRowsetsResponse* response,
2252
0
                                          google::protobuf::Closure* done) {
2253
0
    DCHECK(config::is_cloud_mode());
2254
0
    auto start_time = GetMonoTimeMicros();
2255
0
    Defer defer {
2256
0
            [&]() { g_process_remote_fetch_rowsets_latency << GetMonoTimeMicros() - start_time; }};
2257
0
    brpc::ClosureGuard closure_guard(done);
2258
0
    LOG(INFO) << "process get tablet rowsets, request=" << request->ShortDebugString();
2259
0
    if (!request->has_tablet_id() || !request->has_version_start() || !request->has_version_end()) {
2260
0
        Status::InvalidArgument("missing params tablet/version_start/version_end")
2261
0
                .to_protobuf(response->mutable_status());
2262
0
        return;
2263
0
    }
2264
0
    CloudStorageEngine& storage = ExecEnv::GetInstance()->storage_engine().to_cloud();
2265
2266
0
    auto maybe_tablet =
2267
0
            storage.tablet_mgr().get_tablet(request->tablet_id(), /*warmup data*/ false,
2268
0
                                            /*syn_delete_bitmap*/ false, /*delete_bitmap*/ nullptr,
2269
0
                                            /*local_only*/ true);
2270
0
    if (!maybe_tablet) {
2271
0
        maybe_tablet.error().to_protobuf(response->mutable_status());
2272
0
        return;
2273
0
    }
2274
0
    auto tablet = maybe_tablet.value();
2275
0
    Result<CaptureRowsetResult> ret;
2276
0
    {
2277
0
        std::shared_lock l(tablet->get_header_lock());
2278
0
        ret = tablet->capture_consistent_rowsets_unlocked(
2279
0
                {request->version_start(), request->version_end()},
2280
0
                CaptureRowsetOps {.enable_fetch_rowsets_from_peers = false});
2281
0
    }
2282
0
    if (!ret) {
2283
0
        ret.error().to_protobuf(response->mutable_status());
2284
0
        return;
2285
0
    }
2286
0
    auto rowsets = std::move(ret.value().rowsets);
2287
0
    for (const auto& rs : rowsets) {
2288
0
        RowsetMetaPB meta;
2289
0
        rs->rowset_meta()->to_rowset_pb(&meta);
2290
0
        response->mutable_rowsets()->Add(std::move(meta));
2291
0
    }
2292
0
    if (request->has_delete_bitmap_keys()) {
2293
0
        DCHECK(tablet->need_read_delete_bitmap());
2294
0
        auto delete_bitmap = std::move(ret.value().delete_bitmap);
2295
0
        auto keys_pb = request->delete_bitmap_keys();
2296
0
        size_t len = keys_pb.rowset_ids().size();
2297
0
        DCHECK_EQ(len, keys_pb.segment_ids().size());
2298
0
        DCHECK_EQ(len, keys_pb.versions().size());
2299
0
        std::set<DeleteBitmap::BitmapKey> keys;
2300
0
        for (size_t i = 0; i < len; ++i) {
2301
0
            RowsetId rs_id;
2302
0
            rs_id.init(keys_pb.rowset_ids(i));
2303
0
            keys.emplace(rs_id, keys_pb.segment_ids(i), keys_pb.versions(i));
2304
0
        }
2305
0
        auto diffset = delete_bitmap->diffset(keys).to_pb();
2306
0
        *response->mutable_delete_bitmap() = std::move(diffset);
2307
0
    }
2308
0
    Status::OK().to_protobuf(response->mutable_status());
2309
0
}
2310
2311
void PInternalService::request_cdc_client(google::protobuf::RpcController* controller,
2312
                                          const PRequestCdcClientRequest* request,
2313
                                          PRequestCdcClientResult* result,
2314
0
                                          google::protobuf::Closure* done) {
2315
0
    bool ret = _heavy_work_pool.try_offer([this, request, result, done]() {
2316
0
        _exec_env->cdc_client_mgr()->request_cdc_client_impl(request, result, done);
2317
0
    });
2318
2319
0
    if (!ret) {
2320
0
        offer_failed(result, done, _heavy_work_pool);
2321
0
        return;
2322
0
    }
2323
0
}
2324
2325
void PInternalService::sync_tablet_meta(google::protobuf::RpcController* controller,
2326
                                        const PSyncTabletMetaRequest* request,
2327
                                        PSyncTabletMetaResponse* response,
2328
0
                                        google::protobuf::Closure* done) {
2329
0
    brpc::ClosureGuard closure_guard(done);
2330
0
    Status::NotSupported("sync_tablet_meta only supports cloud mode")
2331
0
            .to_protobuf(response->mutable_status());
2332
0
}
2333
2334
#include "common/compile_check_avoid_end.h"
2335
} // namespace doris