be/src/service/doris_main.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 <arrow/flight/client.h> |
19 | | #include <arrow/flight/sql/client.h> |
20 | | #include <arrow/scalar.h> |
21 | | #include <arrow/status.h> |
22 | | #include <arrow/table.h> |
23 | | #include <butil/macros.h> |
24 | | // IWYU pragma: no_include <bthread/errno.h> |
25 | | #include <errno.h> // IWYU pragma: keep |
26 | | #include <fcntl.h> |
27 | | #include <fmt/core.h> |
28 | | #if !defined(__SANITIZE_ADDRESS__) && !defined(ADDRESS_SANITIZER) && !defined(LEAK_SANITIZER) && \ |
29 | | !defined(THREAD_SANITIZER) && !defined(USE_JEMALLOC) |
30 | | #include <gperftools/malloc_extension.h> // IWYU pragma: keep |
31 | | #endif |
32 | | #include <libgen.h> |
33 | | #include <setjmp.h> |
34 | | #include <signal.h> |
35 | | #include <stdint.h> |
36 | | #include <stdio.h> |
37 | | #include <stdlib.h> |
38 | | #include <unistd.h> |
39 | | |
40 | | #include <cstring> |
41 | | #include <functional> |
42 | | #include <memory> |
43 | | #include <ostream> |
44 | | #include <string> |
45 | | #include <string_view> |
46 | | #include <thread> |
47 | | #include <tuple> |
48 | | #include <vector> |
49 | | |
50 | | #include "cloud/cloud_backend_service.h" |
51 | | #include "cloud/config.h" |
52 | | #include "common/phdr_cache.h" |
53 | | #include "common/stack_trace.h" |
54 | | #if defined(__ELF__) && !defined(__FreeBSD__) |
55 | | #include "common/symbol_index.h" |
56 | | #endif |
57 | | #include "runtime/memory/mem_tracker_limiter.h" |
58 | | #include "storage/tablet/tablet_schema_cache.h" |
59 | | #include "storage/utils.h" |
60 | | #include "util/concurrency_stats.h" |
61 | | |
62 | | #if defined(LEAK_SANITIZER) |
63 | | #include <sanitizer/lsan_interface.h> |
64 | | #endif |
65 | | |
66 | | #include <curl/curl.h> |
67 | | #include <thrift/TOutput.h> |
68 | | |
69 | | #include "agent/heartbeat_server.h" |
70 | | #include "common/config.h" |
71 | | #include "common/daemon.h" |
72 | | #include "common/logging.h" |
73 | | #include "common/signal_handler.h" |
74 | | #include "common/status.h" |
75 | | #include "io/cache/block_file_cache_factory.h" |
76 | | #include "load/stream_load/stream_load_recorder_manager.h" |
77 | | #include "runtime/exec_env.h" |
78 | | #include "runtime/user_function_cache.h" |
79 | | #include "service/arrow_flight/flight_sql_service.h" |
80 | | #include "service/backend_options.h" |
81 | | #include "service/backend_service.h" |
82 | | #include "service/http_service.h" |
83 | | #include "service/server/be_server_starter_factory.h" |
84 | | #include "storage/options.h" |
85 | | #include "storage/storage_engine.h" |
86 | | #include "udf/python/python_env.h" |
87 | | #include "util/debug_util.h" |
88 | | #include "util/disk_info.h" |
89 | | #include "util/jni_plugin_registry.h" |
90 | | #include "util/mem_info.h" |
91 | | #include "util/string_util.h" |
92 | | #include "util/thread.h" |
93 | | #include "util/thrift_rpc_helper.h" |
94 | | #include "util/thrift_server.h" |
95 | | #include "util/uid_util.h" |
96 | | |
97 | | namespace doris {} // namespace doris |
98 | | |
99 | | static void help(const char*); |
100 | | |
101 | | extern "C" { |
102 | | void __lsan_do_leak_check(); |
103 | | int __llvm_profile_write_file(); |
104 | | } |
105 | | |
106 | | namespace doris { |
107 | | |
108 | 7 | void signal_handler(int signal) { |
109 | 7 | if (signal == SIGINT || signal == SIGTERM) { |
110 | 7 | k_doris_exit = true; |
111 | 7 | } |
112 | | // SIGQUIT deliberately does nothing here; see init_signals(). |
113 | 7 | } |
114 | | |
115 | 21 | int install_signal(int signo, void (*handler)(int)) { |
116 | 21 | struct sigaction sa; |
117 | 21 | memset(&sa, 0, sizeof(struct sigaction)); |
118 | 21 | sa.sa_handler = handler; |
119 | | // Restartable syscalls stay restartable. It matters most for SIGQUIT: unlike the two |
120 | | // shutdown signals, that one is sent to a HEALTHY BE - `kill -3` is the operator's habit |
121 | | // for asking a running process for a thread dump, and since the handler now produces |
122 | | // nothing they send it again. Without SA_RESTART each of those turns whatever syscall the |
123 | | // receiving thread happened to be in into EINTR, in the middle of normal serving. |
124 | 21 | sa.sa_flags = SA_RESTART; |
125 | 21 | sigemptyset(&sa.sa_mask); |
126 | 21 | auto ret = sigaction(signo, &sa, nullptr); |
127 | 21 | if (ret != 0) { |
128 | 0 | char buf[64]; |
129 | 0 | LOG(ERROR) << "install signal failed, signo=" << signo << ", errno=" << errno |
130 | 0 | << ", errmsg=" << strerror_r(errno, buf, sizeof(buf)); |
131 | 0 | } |
132 | 21 | return ret; |
133 | 21 | } |
134 | | |
135 | 7 | void init_signals() { |
136 | 7 | auto ret = install_signal(SIGINT, signal_handler); |
137 | 7 | if (ret < 0) { |
138 | 0 | exit(-1); |
139 | 0 | } |
140 | 7 | ret = install_signal(SIGTERM, signal_handler); |
141 | 7 | if (ret < 0) { |
142 | 0 | exit(-1); |
143 | 0 | } |
144 | | // SIGQUIT is taken over even though the BE does nothing with it, because its default |
145 | | // action is not "nothing": it terminates the process and dumps core. `kill -3 <pid>` is |
146 | | // what an operator reaches for to get a thread dump out of a process that looks stuck, |
147 | | // and until the JVM started running with -Xrs it got one - the JVM installed a handler |
148 | | // for this signal along with the shutdown ones. -Xrs stops it from doing that (see |
149 | | // JvmLauncher::_build_options), which would leave SIGQUIT at SIG_DFL and turn that |
150 | | // habitual command into a crash. Handling it and ignoring it is the pre-change |
151 | | // behaviour minus the thread dump; jcmd and jstack, which attach rather than signal, |
152 | | // are how to get one now. |
153 | | // |
154 | | // Installed with a handler rather than SIG_IGN so that the disposition is inheritable |
155 | | // by nothing and visible to JvmLauncher's BeOwnedSignalGuard, which saves and restores |
156 | | // it around any VM creation it does not control. |
157 | 7 | ret = install_signal(SIGQUIT, signal_handler); |
158 | 7 | if (ret < 0) { |
159 | 0 | exit(-1); |
160 | 0 | } |
161 | 7 | } |
162 | | |
163 | 8 | static void thrift_output(const char* x) { |
164 | 8 | LOG(WARNING) << "thrift internal message: " << x; |
165 | 8 | } |
166 | | |
167 | | } // namespace doris |
168 | | |
169 | | // These code is referenced from clickhouse |
170 | | // It is used to check the SIMD instructions |
171 | | enum class InstructionFail { |
172 | | NONE = 0, |
173 | | SSE3 = 1, |
174 | | SSSE3 = 2, |
175 | | SSE4_1 = 3, |
176 | | SSE4_2 = 4, |
177 | | POPCNT = 5, |
178 | | AVX = 6, |
179 | | AVX2 = 7, |
180 | | AVX512 = 8, |
181 | | ARM_NEON = 9 |
182 | | }; |
183 | | |
184 | 0 | auto instruction_fail_to_string(InstructionFail fail) { |
185 | 0 | switch (fail) { |
186 | 0 | #define ret(x) return std::make_tuple(STDERR_FILENO, x, ARRAY_SIZE(x) - 1) |
187 | 0 | case InstructionFail::NONE: |
188 | 0 | ret("NONE"); |
189 | 0 | case InstructionFail::SSE3: |
190 | 0 | ret("SSE3"); |
191 | 0 | case InstructionFail::SSSE3: |
192 | 0 | ret("SSSE3"); |
193 | 0 | case InstructionFail::SSE4_1: |
194 | 0 | ret("SSE4.1"); |
195 | 0 | case InstructionFail::SSE4_2: |
196 | 0 | ret("SSE4.2"); |
197 | 0 | case InstructionFail::POPCNT: |
198 | 0 | ret("POPCNT"); |
199 | 0 | case InstructionFail::AVX: |
200 | 0 | ret("AVX"); |
201 | 0 | case InstructionFail::AVX2: |
202 | 0 | ret("AVX2"); |
203 | 0 | case InstructionFail::AVX512: |
204 | 0 | ret("AVX512"); |
205 | 0 | case InstructionFail::ARM_NEON: |
206 | 0 | ret("ARM_NEON"); |
207 | 0 | } |
208 | | |
209 | 0 | LOG(ERROR) << "Unrecognized instruction fail value." << std::endl; |
210 | 0 | exit(-1); |
211 | 0 | } |
212 | | |
213 | | sigjmp_buf jmpbuf; |
214 | | |
215 | 0 | void sig_ill_check_handler(int, siginfo_t*, void*) { |
216 | 0 | siglongjmp(jmpbuf, 1); |
217 | 0 | } |
218 | | |
219 | | /// Check if necessary SSE extensions are available by trying to execute some sse instructions. |
220 | | /// If instruction is unavailable, SIGILL will be sent by kernel. |
221 | 8 | void check_required_instructions_impl(volatile InstructionFail& fail) { |
222 | 8 | #if defined(__SSE3__) |
223 | 8 | fail = InstructionFail::SSE3; |
224 | 8 | __asm__ volatile("addsubpd %%xmm0, %%xmm0" : : : "xmm0"); |
225 | 8 | #endif |
226 | | |
227 | 8 | #if defined(__SSSE3__) |
228 | 8 | fail = InstructionFail::SSSE3; |
229 | 8 | __asm__ volatile("pabsw %%xmm0, %%xmm0" : : : "xmm0"); |
230 | | |
231 | 8 | #endif |
232 | | |
233 | 8 | #if defined(__SSE4_1__) |
234 | 8 | fail = InstructionFail::SSE4_1; |
235 | 8 | __asm__ volatile("pmaxud %%xmm0, %%xmm0" : : : "xmm0"); |
236 | 8 | #endif |
237 | | |
238 | 8 | #if defined(__SSE4_2__) |
239 | 8 | fail = InstructionFail::SSE4_2; |
240 | 8 | __asm__ volatile("pcmpgtq %%xmm0, %%xmm0" : : : "xmm0"); |
241 | 8 | #endif |
242 | | |
243 | | /// Defined by -msse4.2 |
244 | 8 | #if defined(__POPCNT__) |
245 | 8 | fail = InstructionFail::POPCNT; |
246 | 8 | { |
247 | 8 | uint64_t a = 0; |
248 | 8 | uint64_t b = 0; |
249 | 8 | __asm__ volatile("popcnt %1, %0" : "=r"(a) : "r"(b) :); |
250 | 8 | } |
251 | 8 | #endif |
252 | | |
253 | 8 | #if defined(__AVX__) |
254 | 8 | fail = InstructionFail::AVX; |
255 | 8 | __asm__ volatile("vaddpd %%ymm0, %%ymm0, %%ymm0" : : : "ymm0"); |
256 | 8 | #endif |
257 | | |
258 | 8 | #if defined(__AVX2__) |
259 | 8 | fail = InstructionFail::AVX2; |
260 | 8 | __asm__ volatile("vpabsw %%ymm0, %%ymm0" : : : "ymm0"); |
261 | 8 | #endif |
262 | | |
263 | | #if defined(__AVX512__) |
264 | | fail = InstructionFail::AVX512; |
265 | | __asm__ volatile("vpabsw %%zmm0, %%zmm0" : : : "zmm0"); |
266 | | #endif |
267 | | |
268 | | #if defined(__ARM_NEON__) |
269 | | fail = InstructionFail::ARM_NEON; |
270 | | #ifndef __APPLE__ |
271 | | __asm__ volatile("vadd.i32 q8, q8, q8" : : : "q8"); |
272 | | #endif |
273 | | #endif |
274 | | |
275 | 8 | fail = InstructionFail::NONE; |
276 | 8 | } |
277 | | |
278 | 0 | bool write_retry(int fd, const char* data, size_t size) { |
279 | 0 | if (!size) size = strlen(data); |
280 | |
|
281 | 0 | while (size != 0) { |
282 | 0 | ssize_t res = ::write(fd, data, size); |
283 | |
|
284 | 0 | if ((-1 == res || 0 == res) && errno != EINTR) return false; |
285 | | |
286 | 0 | if (res > 0) { |
287 | 0 | data += res; |
288 | 0 | size -= res; |
289 | 0 | } |
290 | 0 | } |
291 | | |
292 | 0 | return true; |
293 | 0 | } |
294 | | |
295 | | /// Macros to avoid using strlen(), since it may fail if SSE is not supported. |
296 | | #define WRITE_ERROR(data) \ |
297 | 0 | do { \ |
298 | 0 | static_assert(__builtin_constant_p(data)); \ |
299 | 0 | if (!write_retry(STDERR_FILENO, data, ARRAY_SIZE(data) - 1)) _Exit(1); \ |
300 | 0 | } while (false) |
301 | | |
302 | | /// Check SSE and others instructions availability. Calls exit on fail. |
303 | | /// This function must be called as early as possible, even before main, because static initializers may use unavailable instructions. |
304 | 8 | void check_required_instructions() { |
305 | 8 | struct sigaction sa {}; |
306 | 8 | struct sigaction sa_old {}; |
307 | 8 | sa.sa_sigaction = sig_ill_check_handler; |
308 | 8 | sa.sa_flags = SA_SIGINFO; |
309 | 8 | auto signal = SIGILL; |
310 | 8 | if (sigemptyset(&sa.sa_mask) != 0 || sigaddset(&sa.sa_mask, signal) != 0 || |
311 | 8 | sigaction(signal, &sa, &sa_old) != 0) { |
312 | | /// You may wonder about strlen. |
313 | | /// Typical implementation of strlen is using SSE4.2 or AVX2. |
314 | | /// But this is not the case because it's compiler builtin and is executed at compile time. |
315 | |
|
316 | 0 | WRITE_ERROR("Can not set signal handler\n"); |
317 | 0 | _Exit(1); |
318 | 0 | } |
319 | | |
320 | 8 | volatile InstructionFail fail = InstructionFail::NONE; |
321 | | |
322 | 8 | if (sigsetjmp(jmpbuf, 1)) { |
323 | 0 | WRITE_ERROR("Instruction check fail. The CPU does not support "); |
324 | 0 | if (!std::apply(write_retry, instruction_fail_to_string(fail))) _Exit(1); |
325 | 0 | WRITE_ERROR(" instruction set.\n"); |
326 | 0 | WRITE_ERROR( |
327 | 0 | "For example, if your CPU does not support AVX2, you need to rebuild the Doris BE " |
328 | 0 | "with: USE_AVX2=0 sh build.sh --be"); |
329 | 0 | _Exit(1); |
330 | 0 | } |
331 | | |
332 | 8 | check_required_instructions_impl(fail); |
333 | | |
334 | 8 | if (sigaction(signal, &sa_old, nullptr)) { |
335 | 0 | WRITE_ERROR("Can not set signal handler\n"); |
336 | 0 | _Exit(1); |
337 | 0 | } |
338 | 8 | } |
339 | | |
340 | | struct Checker { |
341 | 8 | Checker() { check_required_instructions(); } |
342 | | } checker |
343 | | #ifndef __APPLE__ |
344 | | __attribute__((init_priority(101))) /// Run before other static initializers. |
345 | | #endif |
346 | | ; |
347 | | |
348 | | // A startup failure that happens after ExecEnv::init() has run must terminate the |
349 | | // process the same way normal shutdown does (see the _exit(0) at the end of main): |
350 | | // in the default mode we _exit() immediately, skipping global destructors and the |
351 | | // LeakSanitizer atexit check. Init-time singletons (e.g. the internal workload |
352 | | // group's task scheduler) intentionally live for the whole process lifetime, so |
353 | | // running the leak check on this abnormal-exit path reports them as false-positive |
354 | | // leaks. enable_graceful_exit_check is honored so memleak-check mode still runs LSAN. |
355 | 0 | [[noreturn]] static void exit_on_startup_failure() { |
356 | 0 | if (!doris::config::enable_graceful_exit_check) { |
357 | 0 | google::FlushLogFiles(google::GLOG_INFO); |
358 | 0 | _exit(1); |
359 | 0 | } |
360 | | |
361 | | // exit() starts destroying function-local statics while background threads are still |
362 | | // running. Tear down ExecEnv first so StorageEngine workers are stopped and joined before |
363 | | // they can race with those destructors. |
364 | 0 | if (doris::ExecEnv::ready()) { |
365 | 0 | doris::ExecEnv::GetInstance()->destroy(); |
366 | 0 | } |
367 | 0 | google::FlushLogFiles(google::GLOG_INFO); |
368 | 0 | exit(1); |
369 | 0 | } |
370 | | |
371 | 8 | int main(int argc, char** argv) { |
372 | 8 | doris::signal::InstallFailureSignalHandler(); |
373 | | // create StackTraceCache Instance, at the beginning, other static destructors may use. |
374 | 8 | StackTrace::createCache(); |
375 | | // extern doris::ErrorCode::ErrorCodeInitializer error_code_init; |
376 | | // Some developers will modify status.h and we use a very ticky logic to init error_states |
377 | | // and it maybe not inited. So add a check here. |
378 | 8 | doris::ErrorCode::error_code_init.check_init(); |
379 | | // check if print version or help |
380 | 8 | if (argc > 1) { |
381 | 1 | if (strcmp(argv[1], "--version") == 0 || strcmp(argv[1], "-v") == 0) { |
382 | 1 | puts(doris::get_build_version(false).c_str()); |
383 | 1 | exit(0); |
384 | 1 | } else if (strcmp(argv[1], "--help") == 0 || strcmp(argv[1], "-?") == 0) { |
385 | 0 | help(basename(argv[0])); |
386 | 0 | exit(0); |
387 | 0 | } |
388 | 1 | } |
389 | | |
390 | 7 | if (getenv("DORIS_HOME") == nullptr) { |
391 | 0 | fprintf(stderr, "you need set DORIS_HOME environment variable.\n"); |
392 | 0 | exit(-1); |
393 | 0 | } |
394 | 7 | if (getenv("PID_DIR") == nullptr) { |
395 | 0 | fprintf(stderr, "you need set PID_DIR environment variable.\n"); |
396 | 0 | exit(-1); |
397 | 0 | } |
398 | | |
399 | 7 | SCOPED_INIT_THREAD_CONTEXT(); |
400 | | |
401 | 7 | using doris::Status; |
402 | 7 | using std::string; |
403 | | |
404 | | // open pid file, obtain file lock and save pid |
405 | 7 | string pid_file = string(getenv("PID_DIR")) + "/be.pid"; |
406 | 7 | int fd = open(pid_file.c_str(), O_RDWR | O_CREAT, |
407 | 7 | S_IRUSR | S_IWUSR | S_IRGRP | S_IWGRP | S_IROTH); |
408 | 7 | if (fd < 0) { |
409 | 0 | fprintf(stderr, "fail to create pid file."); |
410 | 0 | exit(-1); |
411 | 0 | } |
412 | | |
413 | 7 | string pid = std::to_string((long)getpid()); |
414 | 7 | pid += "\n"; |
415 | 7 | size_t length = write(fd, pid.c_str(), pid.size()); |
416 | 7 | if (length != pid.size()) { |
417 | 0 | fprintf(stderr, "fail to save pid into pid file."); |
418 | 0 | exit(-1); |
419 | 0 | } |
420 | | |
421 | | // descriptor will be leaked when failing to close fd |
422 | 7 | if (::close(fd) < 0) { |
423 | 0 | fprintf(stderr, "failed to close fd of pidfile."); |
424 | 0 | exit(-1); |
425 | 0 | } |
426 | | |
427 | | // init config. |
428 | | // the config in be_custom.conf will overwrite the config in be.conf |
429 | | // Must init custom config after init config, separately. |
430 | | // Because the path of custom config file is defined in be.conf |
431 | 7 | string conffile = string(getenv("DORIS_HOME")) + "/conf/be.conf"; |
432 | 7 | if (!doris::config::init(conffile.c_str(), true, true, true)) { |
433 | 0 | fprintf(stderr, "error read config file. \n"); |
434 | 0 | return -1; |
435 | 0 | } |
436 | | |
437 | 7 | string custom_conffile = doris::config::custom_config_dir + "/be_custom.conf"; |
438 | 7 | if (!doris::config::init(custom_conffile.c_str(), true, false, false)) { |
439 | 0 | fprintf(stderr, "error read custom config file. \n"); |
440 | 0 | return -1; |
441 | 0 | } |
442 | | |
443 | 7 | doris::init_be_version_metrics(); |
444 | | |
445 | | // ATTN: Callers that want to override default gflags variables should do so before calling this method |
446 | 7 | google::ParseCommandLineFlags(&argc, &argv, true); |
447 | | // ATTN: MUST init before LOG |
448 | 7 | doris::init_glog("be"); |
449 | | |
450 | 7 | LOG(INFO) << doris::get_version_string(false); |
451 | | |
452 | 7 | doris::init_thrift_logging(); |
453 | | |
454 | 7 | if (doris::config::enable_fuzzy_mode) { |
455 | 7 | Status status = doris::config::set_fuzzy_configs(); |
456 | 7 | if (!status.ok()) { |
457 | 0 | LOG(WARNING) << "Failed to initialize fuzzy config: " << status; |
458 | 0 | exit(1); |
459 | 0 | } |
460 | 7 | } |
461 | | |
462 | | #if !defined(__SANITIZE_ADDRESS__) && !defined(ADDRESS_SANITIZER) && !defined(LEAK_SANITIZER) && \ |
463 | | !defined(THREAD_SANITIZER) && !defined(USE_JEMALLOC) |
464 | | // Change the total TCMalloc thread cache size if necessary. |
465 | | const size_t kDefaultTotalThreadCacheBytes = 1024 * 1024 * 1024; |
466 | | if (!MallocExtension::instance()->SetNumericProperty("tcmalloc.max_total_thread_cache_bytes", |
467 | | kDefaultTotalThreadCacheBytes)) { |
468 | | fprintf(stderr, "Failed to change TCMalloc total thread cache size.\n"); |
469 | | return -1; |
470 | | } |
471 | | #endif |
472 | | |
473 | 7 | std::vector<doris::StorePath> paths; |
474 | 7 | auto olap_res = doris::parse_conf_store_paths(doris::config::storage_root_path, &paths); |
475 | 7 | if (!olap_res) { |
476 | 0 | LOG(ERROR) << "parse config storage path failed, path=" << doris::config::storage_root_path; |
477 | 0 | exit(-1); |
478 | 0 | } |
479 | | |
480 | 7 | std::vector<doris::StorePath> spill_paths; |
481 | 7 | if (doris::config::spill_storage_root_path.empty()) { |
482 | 7 | doris::config::spill_storage_root_path = doris::config::storage_root_path; |
483 | 7 | } |
484 | 7 | olap_res = doris::parse_conf_store_paths(doris::config::spill_storage_root_path, &spill_paths); |
485 | 7 | if (!olap_res) { |
486 | 0 | LOG(ERROR) << "parse config spill storage path failed, path=" |
487 | 0 | << doris::config::spill_storage_root_path; |
488 | 0 | exit(-1); |
489 | 0 | } |
490 | 7 | std::set<std::string> broken_paths; |
491 | 7 | doris::parse_conf_broken_store_paths(doris::config::broken_storage_path, &broken_paths); |
492 | | |
493 | 7 | auto it = paths.begin(); |
494 | 18 | for (; it != paths.end();) { |
495 | 11 | if (broken_paths.count(it->path) > 0) { |
496 | 0 | if (doris::config::ignore_broken_disk) { |
497 | 0 | LOG(WARNING) << "ignore broken disk, path = " << it->path; |
498 | 0 | it = paths.erase(it); |
499 | 0 | } else { |
500 | 0 | LOG(ERROR) << "a broken disk is found " << it->path; |
501 | 0 | exit(-1); |
502 | 0 | } |
503 | 11 | } else if (!doris::check_datapath_rw(it->path)) { |
504 | 0 | if (doris::config::ignore_broken_disk) { |
505 | 0 | LOG(WARNING) << "read write test file failed, path=" << it->path; |
506 | 0 | it = paths.erase(it); |
507 | 0 | } else { |
508 | 0 | LOG(ERROR) << "read write test file failed, path=" << it->path; |
509 | | // if only one disk and the disk is full, also need exit because rocksdb will open failed |
510 | 0 | exit(-1); |
511 | 0 | } |
512 | 11 | } else { |
513 | 11 | ++it; |
514 | 11 | } |
515 | 11 | } |
516 | | |
517 | 7 | if (paths.empty()) { |
518 | 0 | LOG(ERROR) << "All disks are broken, exit."; |
519 | 0 | exit(-1); |
520 | 0 | } |
521 | | |
522 | 7 | it = spill_paths.begin(); |
523 | 18 | for (; it != spill_paths.end();) { |
524 | 11 | if (!doris::check_datapath_rw(it->path)) { |
525 | 0 | if (doris::config::ignore_broken_disk) { |
526 | 0 | LOG(WARNING) << "read write test file failed, path=" << it->path; |
527 | 0 | it = spill_paths.erase(it); |
528 | 0 | } else { |
529 | 0 | LOG(ERROR) << "read write test file failed, path=" << it->path; |
530 | 0 | exit(-1); |
531 | 0 | } |
532 | 11 | } else { |
533 | 11 | ++it; |
534 | 11 | } |
535 | 11 | } |
536 | 7 | if (spill_paths.empty()) { |
537 | 0 | LOG(ERROR) << "All spill disks are broken, exit."; |
538 | 0 | exit(-1); |
539 | 0 | } |
540 | | |
541 | | // initialize libcurl here to avoid concurrent initialization |
542 | 7 | auto curl_ret = curl_global_init(CURL_GLOBAL_ALL); |
543 | 7 | if (curl_ret != 0) { |
544 | 0 | LOG(ERROR) << "fail to initialize libcurl, curl_ret=" << curl_ret; |
545 | 0 | exit(-1); |
546 | 0 | } |
547 | | // add logger for thrift internal |
548 | 7 | apache::thrift::TOutput::instance().setOutputFunction(doris::thrift_output); |
549 | | |
550 | 7 | Status status = Status::OK(); |
551 | | // No JVM is started here on purpose. It is created by the first Java feature that asks |
552 | | // for it - a JNI table format, a Java UDF, an hdfs access - and a BE that uses none of |
553 | | // them runs without one. See Jni::JvmLauncher. |
554 | | |
555 | 7 | if (doris::config::enable_python_udf_support) { |
556 | 7 | if (std::string python_udf_root_path = |
557 | 7 | fmt::format("{}/lib/udf/python", std::getenv("DORIS_HOME")); |
558 | 7 | !std::filesystem::exists(python_udf_root_path)) { |
559 | 1 | std::filesystem::create_directories(python_udf_root_path); |
560 | 1 | } |
561 | | |
562 | | // Normalize and trim all Python-related config parameters |
563 | 7 | std::string python_env_mode = |
564 | 7 | std::string(doris::trim(doris::to_lower(doris::config::python_env_mode))); |
565 | 7 | std::string python_conda_root_path = |
566 | 7 | std::string(doris::trim(doris::config::python_conda_root_path)); |
567 | 7 | std::string python_venv_root_path = |
568 | 7 | std::string(doris::trim(doris::config::python_venv_root_path)); |
569 | 7 | std::string python_venv_interpreter_paths = |
570 | 7 | std::string(doris::trim(doris::config::python_venv_interpreter_paths)); |
571 | | |
572 | 7 | if (python_env_mode == "conda") { |
573 | 1 | if (python_conda_root_path.empty()) { |
574 | 0 | LOG(ERROR) |
575 | 0 | << "Python conda root path is empty, please set `python_conda_root_path` " |
576 | 0 | "or set `enable_python_udf_support` to `false`"; |
577 | 0 | exit(1); |
578 | 0 | } |
579 | 1 | LOG(INFO) << "Doris backend python version manager is initialized. Python conda " |
580 | 1 | "root path: " |
581 | 1 | << python_conda_root_path; |
582 | 1 | status = doris::PythonVersionManager::instance().init(doris::PythonEnvType::CONDA, |
583 | 1 | python_conda_root_path, ""); |
584 | 6 | } else if (python_env_mode == "venv") { |
585 | 6 | if (python_venv_root_path.empty()) { |
586 | 0 | LOG(ERROR) |
587 | 0 | << "Python venv root path is empty, please set `python_venv_root_path` or " |
588 | 0 | "set `enable_python_udf_support` to `false`"; |
589 | 0 | exit(1); |
590 | 0 | } |
591 | 6 | if (python_venv_interpreter_paths.empty()) { |
592 | 0 | LOG(ERROR) |
593 | 0 | << "Python interpreter paths is empty, please set " |
594 | 0 | "`python_venv_interpreter_paths` or set `enable_python_udf_support` to " |
595 | 0 | "`false`"; |
596 | 0 | exit(1); |
597 | 0 | } |
598 | 6 | LOG(INFO) << "Doris backend python version manager is initialized. Python venv " |
599 | 6 | "root path: " |
600 | 6 | << python_venv_root_path |
601 | 6 | << ", python interpreter paths: " << python_venv_interpreter_paths; |
602 | 6 | status = doris::PythonVersionManager::instance().init(doris::PythonEnvType::VENV, |
603 | 6 | python_venv_root_path, |
604 | 6 | python_venv_interpreter_paths); |
605 | 6 | } else { |
606 | 0 | status = Status::InvalidArgument( |
607 | 0 | "Python env mode is invalid, should be `conda` or `venv`. If you don't want to " |
608 | 0 | "enable the Python UDF function, please set `enable_python_udf_support` to " |
609 | 0 | "`false`"); |
610 | 0 | } |
611 | | |
612 | 7 | if (!status.ok()) { |
613 | 0 | LOG(ERROR) << "Failed to initialize python version manager: " << status; |
614 | 0 | exit(1); |
615 | 0 | } |
616 | 7 | LOG(INFO) << doris::PythonVersionManager::instance().to_string(); |
617 | 7 | } |
618 | | |
619 | | // SIGINT and SIGTERM are how the BE is asked to shut down, and the handler installed |
620 | | // here does nothing but raise the flag the loop at the end of main() waits on, so the |
621 | | // shutdown stays orderly. SIGQUIT is claimed here as well, so that it does nothing at |
622 | | // all rather than killing the BE with a core dump. A JVM would rather turn the first |
623 | | // two into a Java Shutdown.exit() and answer the third with a thread dump, and it |
624 | | // installs handlers of its own for all three when it starts. The JVM used to be created |
625 | | // a few lines above this call, which is what left these handlers on top; now that it is |
626 | | // created on demand, Jni::JvmLauncher::_bootstrap() is what puts them back once the JVM |
627 | | // has had its way with them. |
628 | | // https://www.oracle.com/java/technologies/javase/signals.html |
629 | 7 | doris::init_signals(); |
630 | | // ATTN: MUST init before `ExecEnv`, `StorageEngine` and other daemon services |
631 | | // |
632 | | // Daemon ───┬──► StorageEngine ──► ExecEnv ──► Disk/Mem/CpuInfo |
633 | | // │ |
634 | | // │ |
635 | | // BackendService ─┘ |
636 | 7 | doris::CpuInfo::init(); |
637 | 7 | doris::DiskInfo::init(); |
638 | 7 | doris::MemInfo::init(); |
639 | | |
640 | 7 | LOG(INFO) << doris::CpuInfo::debug_string(); |
641 | 7 | LOG(INFO) << doris::DiskInfo::debug_string(); |
642 | 7 | LOG(INFO) << doris::MemInfo::debug_string(); |
643 | | |
644 | | // Doris-patched GNU libunwind reads PHDR metadata from our lock-free snapshot instead of |
645 | | // entering glibc dl_iterate_phdr while jemalloc profiling or signal-context unwinding may |
646 | | // already be involved in loader-lock-sensitive code. Configure libunwind before daemon threads |
647 | | // start so all later heap-profile and stack-trace unwinds use the same lock-safe policy. |
648 | 7 | configureLibunwindPHDRCache(); |
649 | 7 | updatePHDRCache(); |
650 | 7 | LOG(INFO) << "PHDR cache enabled: " << hasPHDRCache(); |
651 | 7 | #if defined(__ELF__) && !defined(__FreeBSD__) |
652 | 7 | auto symbol_index = doris::SymbolIndex::instance(); |
653 | 7 | LOG(INFO) << "SymbolIndex preloaded: objects=" << symbol_index->objects().size() |
654 | 7 | << " symbols=" << symbol_index->symbols().size(); |
655 | 7 | #endif |
656 | 7 | if (!doris::BackendOptions::init()) { |
657 | 0 | exit(-1); |
658 | 0 | } |
659 | | |
660 | | // init exec env |
661 | 7 | auto* exec_env(doris::ExecEnv::GetInstance()); |
662 | 7 | status = doris::ExecEnv::init(doris::ExecEnv::GetInstance(), paths, spill_paths, broken_paths); |
663 | 7 | if (status != Status::OK()) { |
664 | 0 | std::cerr << "failed to init doris storage engine, res=" << status; |
665 | 0 | exit_on_startup_failure(); |
666 | 0 | } |
667 | | |
668 | | // Start concurrency stats manager |
669 | 7 | doris::ConcurrencyStatsManager::instance().start(); |
670 | | |
671 | | // begin to start services |
672 | 7 | doris::ThriftRpcHelper::setup(exec_env); |
673 | | // 1. thrift server with be_port |
674 | 7 | std::shared_ptr<doris::BaseBackendService> service; |
675 | 7 | std::unique_ptr<doris::server::IServerStarter> backend_thrift_starter; |
676 | 7 | std::unique_ptr<doris::server::IServerStarter> brpc_starter; |
677 | 7 | std::unique_ptr<doris::server::IServerStarter> http_starter; |
678 | 7 | std::unique_ptr<doris::server::IServerStarter> heartbeat_thrift_starter; |
679 | 7 | std::unique_ptr<doris::server::IServerStarter> flight_starter; |
680 | 7 | bool backend_thrift_started = false; |
681 | 7 | bool brpc_started = false; |
682 | 7 | bool http_started = false; |
683 | 7 | bool heartbeat_thrift_started = false; |
684 | 7 | bool flight_started = false; |
685 | 7 | auto stop_and_join_server = [](std::unique_ptr<doris::server::IServerStarter>& starter, |
686 | 7 | bool started) { |
687 | 0 | if (starter != nullptr) { |
688 | 0 | if (started) { |
689 | 0 | starter->stop(); |
690 | 0 | } |
691 | 0 | starter->join(); |
692 | 0 | } |
693 | 0 | }; |
694 | 7 | std::function<void(Status&, std::string_view)> stop_work_if_error = [&](Status& status, |
695 | 70 | std::string_view msg) { |
696 | 70 | if (!status.ok()) { |
697 | 0 | std::cerr << msg << '\n'; |
698 | 0 | service->stop_works(); |
699 | 0 | if (doris::config::enable_graceful_exit_check) { |
700 | 0 | stop_and_join_server(flight_starter, flight_started); |
701 | 0 | stop_and_join_server(heartbeat_thrift_starter, heartbeat_thrift_started); |
702 | 0 | stop_and_join_server(http_starter, http_started); |
703 | 0 | stop_and_join_server(backend_thrift_starter, backend_thrift_started); |
704 | 0 | stop_and_join_server(brpc_starter, brpc_started); |
705 | 0 | service.reset(); |
706 | 0 | } |
707 | 0 | exit_on_startup_failure(); |
708 | 0 | } |
709 | 70 | }; |
710 | | |
711 | 7 | if (doris::config::is_cloud_mode()) { |
712 | 1 | service = std::make_shared<doris::CloudBackendService>( |
713 | 1 | exec_env->storage_engine().to_cloud(), exec_env); |
714 | 6 | } else { |
715 | 6 | service = std::make_shared<doris::BackendService>(exec_env->storage_engine().to_local(), |
716 | 6 | exec_env); |
717 | 6 | } |
718 | | |
719 | 7 | status = doris::server::create_backend_thrift_starter(exec_env, doris::config::be_port, service, |
720 | 7 | &backend_thrift_starter); |
721 | 7 | stop_work_if_error(status, "Failed to create BE server, exiting"); |
722 | 7 | status = backend_thrift_starter->start(); |
723 | 7 | backend_thrift_started = status.ok(); |
724 | 7 | stop_work_if_error(status, "Doris BE server did not start correctly, exiting"); |
725 | | |
726 | | // 2. brpc service |
727 | 7 | status = doris::server::create_brpc_starter(exec_env, doris::config::brpc_port, |
728 | 7 | doris::config::brpc_num_threads, &brpc_starter); |
729 | 7 | stop_work_if_error(status, "Failed to create BRPC service, exiting"); |
730 | 7 | status = brpc_starter->start(); |
731 | 7 | brpc_started = status.ok(); |
732 | 7 | stop_work_if_error(status, "BRPC service did not start correctly, exiting"); |
733 | | |
734 | | // 3. http service |
735 | 7 | status = |
736 | 7 | doris::server::create_http_starter(exec_env, doris::config::webserver_port, |
737 | 7 | doris::config::webserver_num_workers, &http_starter); |
738 | 7 | stop_work_if_error(status, "Failed to create BE HTTP service, exiting"); |
739 | 7 | status = http_starter->start(); |
740 | 7 | http_started = status.ok(); |
741 | 7 | stop_work_if_error(status, "Doris Be http service did not start correctly, exiting"); |
742 | | |
743 | | // 4. heart beat server |
744 | 7 | doris::ClusterInfo* cluster_info = exec_env->cluster_info(); |
745 | 7 | status = doris::server::create_heartbeat_thrift_starter( |
746 | 7 | exec_env, doris::config::heartbeat_service_port, |
747 | 7 | doris::config::heartbeat_service_thread_count, cluster_info, &heartbeat_thrift_starter); |
748 | 7 | stop_work_if_error(status, "Heartbeat services did not start correctly, exiting"); |
749 | | |
750 | 7 | status = heartbeat_thrift_starter->start(); |
751 | 7 | heartbeat_thrift_started = status.ok(); |
752 | 7 | stop_work_if_error(status, "Doris BE HeartBeat Service did not start correctly, exiting: " + |
753 | 7 | status.to_string()); |
754 | | |
755 | | // 5. arrow flight service |
756 | 7 | status = doris::server::create_flight_starter(doris::config::arrow_flight_sql_port, |
757 | 7 | &flight_starter); |
758 | 7 | stop_work_if_error(status, "Failed to create Arrow Flight service, exiting"); |
759 | 7 | status = flight_starter->start(); |
760 | 7 | flight_started = status.ok(); |
761 | 7 | stop_work_if_error( |
762 | 7 | status, "Arrow Flight Service did not start correctly, exiting, " + status.to_string()); |
763 | | |
764 | | // 6. start daemon thread to do clean or gc jobs |
765 | 7 | doris::Daemon daemon; |
766 | 7 | daemon.start(); |
767 | | |
768 | 7 | exec_env->storage_engine().notify_listeners(); |
769 | | |
770 | 7 | doris::k_is_server_ready = true; |
771 | | |
772 | | // 7. load the deployed Java plugins, once the BE is otherwise serving. |
773 | | // |
774 | | // On its own thread and non-fatal on purpose: the point is that a plugin broken by a bad |
775 | | // deployment shows up in the log now instead of inside the first user query that needs |
776 | | // it, and a plugin that cannot load must not hold up or take down everything else. When |
777 | | // no plugin is deployed this starts no JVM and returns immediately. |
778 | | // |
779 | | // Joined on the way out rather than detached, so that a stop arriving while plugins are |
780 | | // still loading waits for them instead of running the global destructors underneath a |
781 | | // thread that is inside the JVM. Warming up is bounded - one JVM start plus one pass over |
782 | | // the plugin directory - and the Java side of it is a single call, so there is nothing to |
783 | | // interrupt halfway. |
784 | 7 | std::shared_ptr<doris::Thread> plugin_warmup_thread; |
785 | 7 | if (doris::config::enable_java_support && doris::config::java_plugin_warmup) { |
786 | 0 | EXIT_IF_ERROR(doris::Thread::create( |
787 | 0 | "Jni", "java_plugin_warmup", |
788 | 0 | []() { |
789 | | // Named background thread with a thread context of its own: everything it |
790 | | // allocates would otherwise be orphan memory, and the try/catch is what |
791 | | // keeps a directory that becomes unreadable mid-iteration from reaching |
792 | | // std::terminate (directory_iterator::operator++ throws). |
793 | 0 | SCOPED_INIT_THREAD_CONTEXT(); |
794 | 0 | try { |
795 | 0 | if (Status status = doris::Jni::PluginRegistry::warmup(); !status.ok()) { |
796 | 0 | LOG(WARNING) << "failed to warm up Java plugins: " << status; |
797 | 0 | } |
798 | 0 | } catch (const std::exception& e) { |
799 | 0 | LOG(WARNING) << "failed to warm up Java plugins: " << e.what(); |
800 | 0 | } catch (...) { |
801 | 0 | LOG(WARNING) << "failed to warm up Java plugins: unknown exception"; |
802 | 0 | } |
803 | 0 | }, |
804 | 0 | &plugin_warmup_thread)); |
805 | 0 | } |
806 | | |
807 | 1.58k | while (!doris::k_doris_exit) { |
808 | | #if defined(LEAK_SANITIZER) |
809 | | __lsan_do_leak_check(); |
810 | | #endif |
811 | 1.57k | sleep(3); |
812 | 1.57k | } |
813 | 7 | doris::k_is_server_ready = false; |
814 | 7 | LOG(INFO) << "Doris main exiting."; |
815 | 7 | #if defined(LLVM_PROFILE) |
816 | 7 | __llvm_profile_write_file(); |
817 | 7 | LOG(INFO) << "Flush profile file."; |
818 | 7 | #endif |
819 | | // For graceful shutdown, need to wait for all running queries to stop |
820 | 7 | exec_env->wait_for_all_tasks_done(); |
821 | | |
822 | 7 | if (!doris::config::enable_graceful_exit_check) { |
823 | | // If not in memleak check mode, no need to wait all objects de-constructed normally, just exit. |
824 | | // It will make sure that graceful shutdown can be done definitely. |
825 | 0 | LOG(INFO) << "Doris main exited."; |
826 | 0 | google::FlushLogFiles(google::GLOG_INFO); |
827 | 0 | _exit(0); // Do not call exit(0), it will wait for all objects de-constructed normally |
828 | 0 | return 0; |
829 | 0 | } |
830 | | // Before anything is torn down: the warmup thread may still be inside the JVM, and it |
831 | | // reaches BE state that the destructors below free. The fast path above does not need |
832 | | // this - _exit() runs no destructor at all. |
833 | 7 | if (plugin_warmup_thread != nullptr) { |
834 | 0 | plugin_warmup_thread->join(); |
835 | 0 | LOG(INFO) << "Java plugin warmup stopped"; |
836 | 0 | } |
837 | 7 | daemon.stop(); |
838 | 7 | flight_starter->stop(); |
839 | 7 | flight_starter->join(); |
840 | 7 | LOG(INFO) << "Flight server stopped."; |
841 | 7 | heartbeat_thrift_starter->stop(); |
842 | 7 | heartbeat_thrift_starter->join(); |
843 | 7 | LOG(INFO) << "Heartbeat server stopped"; |
844 | | // The stream load recorder manager writes its audit records through this BE's own http |
845 | | // service, so it has to be stopped while that service is still up. Otherwise an audit |
846 | | // load that is in flight here can never be answered, and it blocks the join() done by |
847 | | // SAFE_STOP(_stream_load_recorder_manager) in ExecEnv::destroy() for up to the stream |
848 | | // load timeout, which is longer than the grace period of stop_be.sh --grace. |
849 | 7 | if (auto* recorder_manager = exec_env->stream_load_recorder_manager(); |
850 | 7 | recorder_manager != nullptr) { |
851 | 3 | recorder_manager->stop(); |
852 | 3 | } |
853 | | // TODO(zhiqiang): http_service |
854 | 7 | http_starter->stop(); |
855 | 7 | http_starter->join(); |
856 | 7 | LOG(INFO) << "Http service stopped"; |
857 | 7 | backend_thrift_starter->stop(); |
858 | 7 | backend_thrift_starter->join(); |
859 | 7 | LOG(INFO) << "Be server stopped"; |
860 | 7 | brpc_starter->stop(); |
861 | 7 | brpc_starter->join(); |
862 | 7 | LOG(INFO) << "Brpc service stopped"; |
863 | 7 | service.reset(); |
864 | 7 | LOG(INFO) << "Backend Service stopped"; |
865 | 7 | exec_env->destroy(); |
866 | 7 | LOG(INFO) << "All service stopped, doris main exited."; |
867 | 7 | return 0; |
868 | 7 | } |
869 | | |
870 | 0 | static void help(const char* progname) { |
871 | 0 | printf("%s is the Doris backend server.\n\n", progname); |
872 | 0 | printf("Usage:\n %s [OPTION]...\n\n", progname); |
873 | 0 | printf("Options:\n"); |
874 | 0 | printf(" -v, --version output version information, then exit\n"); |
875 | 0 | printf(" -?, --help show this help, then exit\n"); |
876 | 0 | } |