Coverage Report

Created: 2026-10-09 07:20

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