be/src/runtime/thread_context.h
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 | | #pragma once |
19 | | |
20 | | #include <bthread/bthread.h> |
21 | | #include <bthread/types.h> |
22 | | |
23 | | #include <memory> |
24 | | #include <string> |
25 | | #include <thread> |
26 | | |
27 | | #include "common/exception.h" |
28 | | #include "common/logging.h" |
29 | | #include "common/macros.h" |
30 | | #include "common/query_log_context.h" |
31 | | #include "runtime/memory/mem_tracker_limiter.h" |
32 | | #include "runtime/memory/thread_mem_tracker_mgr.h" |
33 | | #include "util/defer_op.h" // IWYU pragma: keep |
34 | | |
35 | | // Used to tracking query/load/compaction/e.g. execution thread memory usage. |
36 | | // This series of methods saves some information to the thread local context of the current worker thread, |
37 | | // including MemTracker, QueryID, etc. Use CONSUME_THREAD_MEM_TRACKER/RELEASE_THREAD_MEM_TRACKER in the code segment where |
38 | | // the macro is located to record the memory into MemTracker. |
39 | | // Not use it in rpc done.run(), because bthread_setspecific may have errors when UBSAN compiles. |
40 | | |
41 | | // Attach to query/load/compaction/e.g. when thread starts. |
42 | | // This will save some info about a working thread in the thread context. |
43 | | // Looking forward to tracking memory during thread execution into MemTrackerLimiter. |
44 | 320k | #define SCOPED_ATTACH_TASK(arg1) auto VARNAME_LINENUM(attach_task) = AttachTask(arg1) |
45 | | |
46 | | // If the current thread is not executing a Task, such as a StorageEngine thread, |
47 | | // use SCOPED_INIT_THREAD_CONTEXT to initialize ThreadContext. |
48 | | #define SCOPED_INIT_THREAD_CONTEXT() \ |
49 | 448 | auto VARNAME_LINENUM(scoped_tls_itc) = doris::ScopedInitThreadContext() |
50 | | |
51 | | // Switch resource context in thread context, used after SCOPED_ATTACH_TASK. |
52 | | // If just want to switch mem tracker, use SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER first. |
53 | | #define SCOPED_SWITCH_RESOURCE_CONTEXT(arg1) \ |
54 | | auto VARNAME_LINENUM(switch_resource_context) = doris::SwitchResourceContext(arg1) |
55 | | |
56 | | // Switch MemTrackerLimiter for count memory during thread execution. |
57 | | // Used after SCOPED_ATTACH_TASK, in order to count the memory into another |
58 | | // MemTrackerLimiter instead of the MemTrackerLimiter added by the attach task. |
59 | | #define SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(arg1) \ |
60 | 996k | auto VARNAME_LINENUM(switch_mem_tracker) = doris::SwitchThreadMemTrackerLimiter(arg1) |
61 | | |
62 | | // Looking forward to tracking memory during thread execution into MemTracker. |
63 | | // Usually used to record query more detailed memory, including ExecNode operators. |
64 | | #define SCOPED_CONSUME_MEM_TRACKER(mem_tracker) \ |
65 | 3.00k | auto VARNAME_LINENUM(add_mem_consumer) = doris::AddThreadMemTrackerConsumer(mem_tracker) |
66 | | |
67 | | #define DEFER_RELEASE_RESERVED() \ |
68 | 1.26M | Defer VARNAME_LINENUM(defer) { \ |
69 | 1.26M | [&]() { doris::thread_context()->thread_mem_tracker_mgr->shrink_reserved(); }};unity_8_cxx.cxx:_ZZN5doris12PipelineTask7executeEPbENK3$_2clEv Line | Count | Source | 69 | 1.24M | [&]() { doris::thread_context()->thread_mem_tracker_mgr->shrink_reserved(); }}; |
unity_8_cxx.cxx:_ZZN5doris12PipelineTask7executeEPbENK3$_3clEv Line | Count | Source | 69 | 16.8k | [&]() { doris::thread_context()->thread_mem_tracker_mgr->shrink_reserved(); }}; |
|
70 | | |
71 | | // Count a code segment memory |
72 | | // Usage example: |
73 | | // int64_t peak_mem = 0; |
74 | | // { |
75 | | // SCOPED_PEAK_MEM(&peak_mem); |
76 | | // xxxx |
77 | | // } |
78 | | // LOG(INFO) << *peak_mem; |
79 | | #define SCOPED_PEAK_MEM(peak_mem) \ |
80 | 7.78k | auto VARNAME_LINENUM(scope_peak_mem) = doris::ScopedPeakMem(peak_mem) |
81 | | |
82 | | #define SCOPED_SKIP_MEMORY_CHECK() \ |
83 | 1.50M | auto VARNAME_LINENUM(scope_skip_memory_check) = doris::ScopeSkipMemoryCheck() |
84 | | |
85 | | #define SKIP_LARGE_MEMORY_CHECK(...) \ |
86 | | do { \ |
87 | | doris::ThreadLocalHandle::create_thread_local_if_not_exits(); \ |
88 | | doris::thread_context()->thread_mem_tracker_mgr->skip_large_memory_check++; \ |
89 | | DEFER({ \ |
90 | | doris::thread_context()->thread_mem_tracker_mgr->skip_large_memory_check--; \ |
91 | | doris::ThreadLocalHandle::del_thread_local_if_count_is_zero(); \ |
92 | | }); \ |
93 | | __VA_ARGS__; \ |
94 | | } while (0) |
95 | | |
96 | | #define LIMIT_LOCAL_SCAN_IO(data_dir, bytes_read) \ |
97 | 72.5k | std::shared_ptr<IOThrottle> iot = nullptr; \ |
98 | 72.5k | auto* t_ctx = doris::thread_context(); \ |
99 | 72.5k | if (t_ctx->is_attach_task() && t_ctx->resource_ctx()->workload_group() != nullptr) { \ |
100 | 66.1k | iot = t_ctx->resource_ctx()->workload_group()->get_local_scan_io_throttle(data_dir); \ |
101 | 66.1k | } \ |
102 | 72.5k | if (iot) { \ |
103 | 66.1k | iot->acquire(-1); \ |
104 | 66.1k | } \ |
105 | 72.5k | Defer defer { \ |
106 | 426k | [&]() { \ |
107 | 426k | if (iot) { \ |
108 | 66.1k | iot->update_next_io_time(*bytes_read); \ |
109 | 66.1k | t_ctx->resource_ctx()->workload_group()->update_local_scan_io(data_dir, \ |
110 | 66.1k | *bytes_read); \ |
111 | 66.1k | } \ |
112 | 426k | } \ unity_2_cxx.cxx:_ZZN5doris2io15LocalFileReader12read_at_implEmNS_5SliceEPmPKNS0_9IOContextEENK3$_0clEv Line | Count | Source | 106 | 426k | [&]() { \ | 107 | 426k | if (iot) { \ | 108 | 66.1k | iot->update_next_io_time(*bytes_read); \ | 109 | 66.1k | t_ctx->resource_ctx()->workload_group()->update_local_scan_io(data_dir, \ | 110 | 66.1k | *bytes_read); \ | 111 | 66.1k | } \ | 112 | 426k | } \ |
unity_2_cxx.cxx:_ZZN5doris2io15LocalFileReader18read_at_iobuf_implEmmPN5butil5IOBufEPmPKNS0_9IOContextEENK3$_0clEv Line | Count | Source | 106 | 7 | [&]() { \ | 107 | 7 | if (iot) { \ | 108 | 0 | iot->update_next_io_time(*bytes_read); \ | 109 | 0 | t_ctx->resource_ctx()->workload_group()->update_local_scan_io(data_dir, \ | 110 | 0 | *bytes_read); \ | 111 | 0 | } \ | 112 | 7 | } \ |
|
113 | 72.5k | } |
114 | | |
115 | | #define LIMIT_REMOTE_SCAN_IO(bytes_read) \ |
116 | 53 | std::shared_ptr<IOThrottle> iot = nullptr; \ |
117 | 53 | auto* t_ctx = doris::thread_context(); \ |
118 | 53 | if (t_ctx->is_attach_task() && t_ctx->resource_ctx()->workload_group() != nullptr) { \ |
119 | 3 | iot = t_ctx->resource_ctx()->workload_group()->get_remote_scan_io_throttle(); \ |
120 | 3 | } \ |
121 | 53 | if (iot) { \ |
122 | 3 | iot->acquire(-1); \ |
123 | 3 | } \ |
124 | 53 | Defer defer { \ |
125 | 53 | [&]() { \ |
126 | 53 | if (iot) { \ |
127 | 3 | iot->update_next_io_time(*bytes_read); \ |
128 | 3 | t_ctx->resource_ctx()->workload_group()->update_remote_scan_io(*bytes_read); \ |
129 | 3 | } \ |
130 | 53 | } \ Unexecuted instantiation: unity_2_cxx.cxx:_ZZN5doris2io14HdfsFileReader15do_read_at_implEmNS_5SliceEPmPKNS0_9IOContextEENK3$_0clEv Unexecuted instantiation: unity_2_cxx.cxx:_ZZN5doris2io12S3FileReader12read_at_implEmNS_5SliceEPmPKNS0_9IOContextEENK3$_0clEv unity_1_cxx.cxx:_ZZN5doris2io19PeerFileCacheReader12fetch_blocksERKSt6vectorISt10shared_ptrINS0_9FileBlockEESaIS5_EEPNS0_15PeerFetchResultEmPKNS0_9IOContextEblNSt7__cxx1112basic_stringIcSt11char_traitsIcESaIcEEEENK3$_1clEv Line | Count | Source | 125 | 42 | [&]() { \ | 126 | 42 | if (iot) { \ | 127 | 3 | iot->update_next_io_time(*bytes_read); \ | 128 | 3 | t_ctx->resource_ctx()->workload_group()->update_remote_scan_io(*bytes_read); \ | 129 | 3 | } \ | 130 | 42 | } \ |
unity_1_cxx.cxx:_ZZN5doris2io14PrefetchBuffer11read_bufferEmPKcmPmENK3$_1clEv Line | Count | Source | 125 | 11 | [&]() { \ | 126 | 11 | if (iot) { \ | 127 | 0 | iot->update_next_io_time(*bytes_read); \ | 128 | 0 | t_ctx->resource_ctx()->workload_group()->update_remote_scan_io(*bytes_read); \ | 129 | 0 | } \ | 130 | 11 | } \ |
|
131 | 53 | } |
132 | | |
133 | | namespace doris { |
134 | | |
135 | | class ThreadContext; |
136 | | class MemTracker; |
137 | | class QueryContext; |
138 | | class ResourceContext; |
139 | | class RuntimeState; |
140 | | class SwitchResourceContext; |
141 | | |
142 | | extern bthread_key_t btls_key; |
143 | | |
144 | | // Initialize btls_key exactly once before it is used by any thread context. The key has |
145 | | // process lifetime so an existing bthread context never becomes invalid while the BE is running. |
146 | | void init_thread_context_btls_key(); |
147 | | |
148 | | static std::string NO_THREAD_CONTEXT_MSG = |
149 | | "Current thread not exist ThreadContext, usually after the thread is started, using " |
150 | | "SCOPED_ATTACH_TASK macro to create a ThreadContext and bind a Task."; |
151 | | |
152 | | // Is true after ThreadContext construction. |
153 | | inline thread_local bool pthread_context_ptr_init = false; |
154 | | inline thread_local constinit ThreadContext* thread_context_ptr = nullptr; |
155 | | |
156 | | // The thread context saves some info about a working thread. |
157 | | // 2 required info: |
158 | | // 1. thread_id: Current thread id, Auto generated. |
159 | | // 2. type(abolished): The type is a enum value indicating which type of task current thread is running. |
160 | | // For example: QUERY, LOAD, COMPACTION, ... |
161 | | // 3. task id: A unique id to identify this task. maybe query id, load job id, etc. |
162 | | // 4. ThreadMemTrackerMgr |
163 | | // |
164 | | // There may be other optional info to be added later. |
165 | | // |
166 | | // Note: Keep the class simple and only add properties. |
167 | | class ThreadContext { |
168 | | public: |
169 | 398k | ThreadContext() { thread_mem_tracker_mgr = std::make_unique<ThreadMemTrackerMgr>(); } |
170 | | |
171 | 394k | ~ThreadContext() = default; |
172 | | |
173 | | void attach_task(const std::shared_ptr<ResourceContext>& rc); |
174 | | |
175 | 320k | void detach_task() { |
176 | 320k | resource_ctx_.reset(); |
177 | 320k | thread_mem_tracker_mgr->detach_limiter_tracker(); |
178 | 320k | thread_mem_tracker_mgr->disable_wait_gc(); |
179 | 320k | } |
180 | | |
181 | 875k | bool is_attach_task() { return resource_ctx_ != nullptr; } |
182 | | |
183 | 224k | std::shared_ptr<ResourceContext> resource_ctx() { |
184 | 224k | #ifndef BE_TEST |
185 | 224k | DCHECK(is_attach_task()); |
186 | 224k | #endif |
187 | 224k | if (is_attach_task()) { |
188 | 224k | return resource_ctx_; |
189 | 224k | } |
190 | 10 | return _make_orphan_resource_ctx(); |
191 | 224k | } |
192 | | |
193 | 16 | static std::string get_thread_id() { |
194 | 16 | std::stringstream ss; |
195 | 16 | ss << std::this_thread::get_id(); |
196 | 16 | return ss.str(); |
197 | 16 | } |
198 | | // Note that if set global Memory Hook, After thread_mem_tracker_mgr is initialized, |
199 | | // the current thread Hook starts to consume/release mem_tracker. |
200 | | // the use of shared_ptr will cause a crash. The guess is that there is an |
201 | | // intermediate state during the copy construction of shared_ptr. Shared_ptr is not equal |
202 | | // to nullptr, but the object it points to is not initialized. At this time, when the memory |
203 | | // is released somewhere, the hook is triggered to cause the crash. |
204 | | std::unique_ptr<ThreadMemTrackerMgr> thread_mem_tracker_mgr; |
205 | | |
206 | | int thread_local_handle_count = 0; |
207 | | |
208 | | private: |
209 | | friend class SwitchResourceContext; |
210 | | |
211 | | // Cold fallback for threads without an attached task; defined in the .cpp |
212 | | // so this header does not need the full ResourceContext / ExecEnv types. |
213 | | static std::shared_ptr<ResourceContext> _make_orphan_resource_ctx(); |
214 | | |
215 | | std::shared_ptr<ResourceContext> resource_ctx_; |
216 | | }; |
217 | | |
218 | | class ThreadLocalHandle { |
219 | | public: |
220 | 2.84M | static void create_thread_local_if_not_exits() { |
221 | 2.84M | init_thread_context_btls_key(); |
222 | 2.84M | if (bthread_self() == 0) { |
223 | 2.84M | if (!pthread_context_ptr_init) { |
224 | 396k | thread_context_ptr = new ThreadContext(); |
225 | 396k | pthread_context_ptr_init = true; |
226 | 396k | } |
227 | 2.84M | DCHECK(thread_context_ptr != nullptr); |
228 | 2.84M | thread_context_ptr->thread_local_handle_count++; |
229 | 2.84M | } else { |
230 | | // Avoid calling bthread_getspecific frequently to get bthread local. |
231 | | // Very frequent bthread_getspecific will slow, but create_thread_local_if_not_exits is not expected to be much. |
232 | | // Cache the pointer of bthread local in pthead local. |
233 | 6.27k | auto* bthread_context = static_cast<ThreadContext*>(bthread_getspecific(btls_key)); |
234 | 6.27k | if (bthread_context == nullptr) { |
235 | | // If bthread_context == nullptr: |
236 | | // 1. First call to bthread_getspecific (and before any bthread_setspecific) returns NULL |
237 | | // 2. There are not enough reusable btls in btls pool. |
238 | | // else if bthread_context != nullptr: |
239 | | // 1. A new bthread starts, but get a reuses btls. |
240 | 1.58k | bthread_context = new ThreadContext; |
241 | | // The brpc server should respond as quickly as possible. |
242 | 1.58k | bthread_context->thread_mem_tracker_mgr->disable_wait_gc(); |
243 | | // set the data so that next time bthread_getspecific in the thread returns the data. |
244 | 1.58k | CHECK(0 == bthread_setspecific(btls_key, bthread_context) || doris::k_doris_exit); |
245 | 1.58k | } |
246 | 6.27k | DCHECK(bthread_context != nullptr); |
247 | 6.27k | bthread_context->thread_local_handle_count++; |
248 | 6.27k | } |
249 | 2.84M | } |
250 | | |
251 | | // `create_thread_local_if_not_exits` and `del_thread_local_if_count_is_zero` should be used in pairs, |
252 | | // `del_thread_local_if_count_is_zero` should only be called if `create_thread_local_if_not_exits` returns true |
253 | 2.84M | static void del_thread_local_if_count_is_zero() { |
254 | 2.84M | if (pthread_context_ptr_init) { |
255 | | // in pthread |
256 | 2.83M | thread_context_ptr->thread_local_handle_count--; |
257 | 2.83M | if (thread_context_ptr->thread_local_handle_count == 0) { |
258 | 392k | pthread_context_ptr_init = false; |
259 | 392k | delete doris::thread_context_ptr; |
260 | 392k | thread_context_ptr = nullptr; |
261 | 392k | } |
262 | 2.83M | } else if (bthread_self() != 0) { |
263 | | // in bthread |
264 | 8.01k | auto* bthread_context = static_cast<ThreadContext*>(bthread_getspecific(btls_key)); |
265 | 8.01k | DCHECK(bthread_context != nullptr); |
266 | 8.01k | bthread_context->thread_local_handle_count--; |
267 | 18.4E | } else { |
268 | 18.4E | throw Exception(Status::FatalError("__builtin_unreachable")); |
269 | 18.4E | } |
270 | 2.84M | } |
271 | | }; |
272 | | |
273 | | // must call create_thread_local_if_not_exits() before use thread_context(). |
274 | 38.3M | static ThreadContext* thread_context() { |
275 | 38.3M | if (pthread_context_ptr_init) { |
276 | | // in pthread |
277 | 38.3M | DCHECK(bthread_self() == 0); |
278 | 38.3M | DCHECK(thread_context_ptr != nullptr); |
279 | 38.3M | return thread_context_ptr; |
280 | 38.3M | } |
281 | 16.3k | if (bthread_self() != 0) { |
282 | | // in bthread |
283 | | // bthread switching pthread may be very frequent, remember not to use lock or other time-consuming operations. |
284 | 16.3k | auto* bthread_context = static_cast<ThreadContext*>(bthread_getspecific(btls_key)); |
285 | 16.3k | DCHECK(bthread_context != nullptr && bthread_context->thread_local_handle_count > 0); |
286 | 16.3k | return bthread_context; |
287 | 16.3k | } |
288 | | // It means that use thread_context() but this thread not attached a query/load using SCOPED_ATTACH_TASK macro. |
289 | 18.4E | throw Exception(Status::FatalError("{}", doris::NO_THREAD_CONTEXT_MSG)); |
290 | 16.1k | } Unexecuted instantiation: doris_main.cpp:_ZN5dorisL14thread_contextEv unity_0_cxx.cxx:_ZN5dorisL14thread_contextEv Line | Count | Source | 274 | 26.8M | static ThreadContext* thread_context() { | 275 | 26.8M | if (pthread_context_ptr_init) { | 276 | | // in pthread | 277 | 26.8M | DCHECK(bthread_self() == 0); | 278 | 26.8M | DCHECK(thread_context_ptr != nullptr); | 279 | 26.8M | return thread_context_ptr; | 280 | 26.8M | } | 281 | 197 | if (bthread_self() != 0) { | 282 | | // in bthread | 283 | | // bthread switching pthread may be very frequent, remember not to use lock or other time-consuming operations. | 284 | 197 | auto* bthread_context = static_cast<ThreadContext*>(bthread_getspecific(btls_key)); | 285 | 197 | DCHECK(bthread_context != nullptr && bthread_context->thread_local_handle_count > 0); | 286 | 197 | return bthread_context; | 287 | 197 | } | 288 | | // It means that use thread_context() but this thread not attached a query/load using SCOPED_ATTACH_TASK macro. | 289 | 18.4E | throw Exception(Status::FatalError("{}", doris::NO_THREAD_CONTEXT_MSG)); | 290 | 51 | } |
Unexecuted instantiation: config.cpp:_ZN5dorisL14thread_contextEv unity_2_cxx.cxx:_ZN5dorisL14thread_contextEv Line | Count | Source | 274 | 439k | static ThreadContext* thread_context() { | 275 | 439k | if (pthread_context_ptr_init) { | 276 | | // in pthread | 277 | 439k | DCHECK(bthread_self() == 0); | 278 | 439k | DCHECK(thread_context_ptr != nullptr); | 279 | 439k | return thread_context_ptr; | 280 | 439k | } | 281 | 18.4E | if (bthread_self() != 0) { | 282 | | // in bthread | 283 | | // bthread switching pthread may be very frequent, remember not to use lock or other time-consuming operations. | 284 | 0 | auto* bthread_context = static_cast<ThreadContext*>(bthread_getspecific(btls_key)); | 285 | 0 | DCHECK(bthread_context != nullptr && bthread_context->thread_local_handle_count > 0); | 286 | 0 | return bthread_context; | 287 | 0 | } | 288 | | // It means that use thread_context() but this thread not attached a query/load using SCOPED_ATTACH_TASK macro. | 289 | 18.4E | throw Exception(Status::FatalError("{}", doris::NO_THREAD_CONTEXT_MSG)); | 290 | 18.4E | } |
unity_3_cxx.cxx:_ZN5dorisL14thread_contextEv Line | Count | Source | 274 | 3.24M | static ThreadContext* thread_context() { | 275 | 3.24M | if (pthread_context_ptr_init) { | 276 | | // in pthread | 277 | 3.22M | DCHECK(bthread_self() == 0); | 278 | 3.22M | DCHECK(thread_context_ptr != nullptr); | 279 | 3.22M | return thread_context_ptr; | 280 | 3.22M | } | 281 | 16.1k | if (bthread_self() != 0) { | 282 | | // in bthread | 283 | | // bthread switching pthread may be very frequent, remember not to use lock or other time-consuming operations. | 284 | 16.1k | auto* bthread_context = static_cast<ThreadContext*>(bthread_getspecific(btls_key)); | 285 | 16.1k | DCHECK(bthread_context != nullptr && bthread_context->thread_local_handle_count > 0); | 286 | 16.1k | return bthread_context; | 287 | 16.1k | } | 288 | | // It means that use thread_context() but this thread not attached a query/load using SCOPED_ATTACH_TASK macro. | 289 | 18.4E | throw Exception(Status::FatalError("{}", doris::NO_THREAD_CONTEXT_MSG)); | 290 | 16.0k | } |
Unexecuted instantiation: unity_5_cxx.cxx:_ZN5dorisL14thread_contextEv Unexecuted instantiation: unity_13_cxx.cxx:_ZN5dorisL14thread_contextEv unity_11_cxx.cxx:_ZN5dorisL14thread_contextEv Line | Count | Source | 274 | 307 | static ThreadContext* thread_context() { | 275 | 307 | if (pthread_context_ptr_init) { | 276 | | // in pthread | 277 | 307 | DCHECK(bthread_self() == 0); | 278 | 307 | DCHECK(thread_context_ptr != nullptr); | 279 | 307 | return thread_context_ptr; | 280 | 307 | } | 281 | 0 | if (bthread_self() != 0) { | 282 | | // in bthread | 283 | | // bthread switching pthread may be very frequent, remember not to use lock or other time-consuming operations. | 284 | 0 | auto* bthread_context = static_cast<ThreadContext*>(bthread_getspecific(btls_key)); | 285 | 0 | DCHECK(bthread_context != nullptr && bthread_context->thread_local_handle_count > 0); | 286 | 0 | return bthread_context; | 287 | 0 | } | 288 | | // It means that use thread_context() but this thread not attached a query/load using SCOPED_ATTACH_TASK macro. | 289 | 0 | throw Exception(Status::FatalError("{}", doris::NO_THREAD_CONTEXT_MSG)); | 290 | 0 | } |
Unexecuted instantiation: unity_7_cxx.cxx:_ZN5dorisL14thread_contextEv Unexecuted instantiation: unity_14_cxx.cxx:_ZN5dorisL14thread_contextEv Unexecuted instantiation: unity_6_cxx.cxx:_ZN5dorisL14thread_contextEv Unexecuted instantiation: unity_12_cxx.cxx:_ZN5dorisL14thread_contextEv unity_10_cxx.cxx:_ZN5dorisL14thread_contextEv Line | Count | Source | 274 | 39.3k | static ThreadContext* thread_context() { | 275 | 39.3k | if (pthread_context_ptr_init) { | 276 | | // in pthread | 277 | 39.3k | DCHECK(bthread_self() == 0); | 278 | 39.3k | DCHECK(thread_context_ptr != nullptr); | 279 | 39.3k | return thread_context_ptr; | 280 | 39.3k | } | 281 | 18.4E | if (bthread_self() != 0) { | 282 | | // in bthread | 283 | | // bthread switching pthread may be very frequent, remember not to use lock or other time-consuming operations. | 284 | 0 | auto* bthread_context = static_cast<ThreadContext*>(bthread_getspecific(btls_key)); | 285 | 0 | DCHECK(bthread_context != nullptr && bthread_context->thread_local_handle_count > 0); | 286 | 0 | return bthread_context; | 287 | 0 | } | 288 | | // It means that use thread_context() but this thread not attached a query/load using SCOPED_ATTACH_TASK macro. | 289 | 18.4E | throw Exception(Status::FatalError("{}", doris::NO_THREAD_CONTEXT_MSG)); | 290 | 18.4E | } |
Unexecuted instantiation: unity_9_cxx.cxx:_ZN5dorisL14thread_contextEv unity_8_cxx.cxx:_ZN5dorisL14thread_contextEv Line | Count | Source | 274 | 1.29M | static ThreadContext* thread_context() { | 275 | 1.29M | if (pthread_context_ptr_init) { | 276 | | // in pthread | 277 | 1.29M | DCHECK(bthread_self() == 0); | 278 | 1.29M | DCHECK(thread_context_ptr != nullptr); | 279 | 1.29M | return thread_context_ptr; | 280 | 1.29M | } | 281 | 0 | if (bthread_self() != 0) { | 282 | | // in bthread | 283 | | // bthread switching pthread may be very frequent, remember not to use lock or other time-consuming operations. | 284 | 0 | auto* bthread_context = static_cast<ThreadContext*>(bthread_getspecific(btls_key)); | 285 | 0 | DCHECK(bthread_context != nullptr && bthread_context->thread_local_handle_count > 0); | 286 | 0 | return bthread_context; | 287 | 0 | } | 288 | | // It means that use thread_context() but this thread not attached a query/load using SCOPED_ATTACH_TASK macro. | 289 | 0 | throw Exception(Status::FatalError("{}", doris::NO_THREAD_CONTEXT_MSG)); | 290 | 0 | } |
unity_1_cxx.cxx:_ZN5dorisL14thread_contextEv Line | Count | Source | 274 | 6.48M | static ThreadContext* thread_context() { | 275 | 6.48M | if (pthread_context_ptr_init) { | 276 | | // in pthread | 277 | 6.48M | DCHECK(bthread_self() == 0); | 278 | 6.48M | DCHECK(thread_context_ptr != nullptr); | 279 | 6.48M | return thread_context_ptr; | 280 | 6.48M | } | 281 | 20 | if (bthread_self() != 0) { | 282 | | // in bthread | 283 | | // bthread switching pthread may be very frequent, remember not to use lock or other time-consuming operations. | 284 | 16 | auto* bthread_context = static_cast<ThreadContext*>(bthread_getspecific(btls_key)); | 285 | 16 | DCHECK(bthread_context != nullptr && bthread_context->thread_local_handle_count > 0); | 286 | 16 | return bthread_context; | 287 | 16 | } | 288 | | // It means that use thread_context() but this thread not attached a query/load using SCOPED_ATTACH_TASK macro. | 289 | 4 | throw Exception(Status::FatalError("{}", doris::NO_THREAD_CONTEXT_MSG)); | 290 | 20 | } |
Unexecuted instantiation: unity_4_cxx.cxx:_ZN5dorisL14thread_contextEv Unexecuted instantiation: hashjoin_build_sink.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: join_build_sink_operator.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: operator.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: partitioned_aggregation_sink_operator.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: scan_operator.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: vfile_result_writer.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: unity_31_cxx.cxx:_ZN5dorisL14thread_contextEv Unexecuted instantiation: unity_30_cxx.cxx:_ZN5dorisL14thread_contextEv Unexecuted instantiation: unity_29_cxx.cxx:_ZN5dorisL14thread_contextEv Unexecuted instantiation: unity_28_cxx.cxx:_ZN5dorisL14thread_contextEv Unexecuted instantiation: unity_27_cxx.cxx:_ZN5dorisL14thread_contextEv Unexecuted instantiation: unity_26_cxx.cxx:_ZN5dorisL14thread_contextEv Unexecuted instantiation: unity_25_cxx.cxx:_ZN5dorisL14thread_contextEv Unexecuted instantiation: unity_24_cxx.cxx:_ZN5dorisL14thread_contextEv Unexecuted instantiation: unity_22_cxx.cxx:_ZN5dorisL14thread_contextEv Unexecuted instantiation: unity_21_cxx.cxx:_ZN5dorisL14thread_contextEv unity_20_cxx.cxx:_ZN5dorisL14thread_contextEv Line | Count | Source | 274 | 121 | static ThreadContext* thread_context() { | 275 | 121 | if (pthread_context_ptr_init) { | 276 | | // in pthread | 277 | 121 | DCHECK(bthread_self() == 0); | 278 | 121 | DCHECK(thread_context_ptr != nullptr); | 279 | 121 | return thread_context_ptr; | 280 | 121 | } | 281 | 0 | if (bthread_self() != 0) { | 282 | | // in bthread | 283 | | // bthread switching pthread may be very frequent, remember not to use lock or other time-consuming operations. | 284 | 0 | auto* bthread_context = static_cast<ThreadContext*>(bthread_getspecific(btls_key)); | 285 | 0 | DCHECK(bthread_context != nullptr && bthread_context->thread_local_handle_count > 0); | 286 | 0 | return bthread_context; | 287 | 0 | } | 288 | | // It means that use thread_context() but this thread not attached a query/load using SCOPED_ATTACH_TASK macro. | 289 | 0 | throw Exception(Status::FatalError("{}", doris::NO_THREAD_CONTEXT_MSG)); | 290 | 0 | } |
Unexecuted instantiation: unity_19_cxx.cxx:_ZN5dorisL14thread_contextEv Unexecuted instantiation: unity_18_cxx.cxx:_ZN5dorisL14thread_contextEv unity_17_cxx.cxx:_ZN5dorisL14thread_contextEv Line | Count | Source | 274 | 2 | static ThreadContext* thread_context() { | 275 | 2 | if (pthread_context_ptr_init) { | 276 | | // in pthread | 277 | 2 | DCHECK(bthread_self() == 0); | 278 | 2 | DCHECK(thread_context_ptr != nullptr); | 279 | 2 | return thread_context_ptr; | 280 | 2 | } | 281 | 0 | if (bthread_self() != 0) { | 282 | | // in bthread | 283 | | // bthread switching pthread may be very frequent, remember not to use lock or other time-consuming operations. | 284 | 0 | auto* bthread_context = static_cast<ThreadContext*>(bthread_getspecific(btls_key)); | 285 | 0 | DCHECK(bthread_context != nullptr && bthread_context->thread_local_handle_count > 0); | 286 | 0 | return bthread_context; | 287 | 0 | } | 288 | | // It means that use thread_context() but this thread not attached a query/load using SCOPED_ATTACH_TASK macro. | 289 | 0 | throw Exception(Status::FatalError("{}", doris::NO_THREAD_CONTEXT_MSG)); | 290 | 0 | } |
Unexecuted instantiation: unity_16_cxx.cxx:_ZN5dorisL14thread_contextEv Unexecuted instantiation: unity_15_cxx.cxx:_ZN5dorisL14thread_contextEv Unexecuted instantiation: ai_functions.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: function_array_aggregation.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: function_date_or_datetime_computation.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: function_datetime_floor_ceil.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: in.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: uuid.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: arrow_row_batch.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: arrow_stream_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: csv_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: new_plain_binary_line_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: jni_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: new_json_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: vorc_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: column_type_convert.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: vparquet_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: schema_desc.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: vparquet_group_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: vparquet_column_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: vparquet_column_chunk_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: vparquet_page_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: parquet_column_convert.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: vparquet_page_index.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: parquet_block_split_bloom_filter.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: deletion_vector_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: es_http_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: hive_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: hive_orc_nested_column_utils.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: hudi_jni_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: hudi_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: iceberg_delete_file_reader_helper.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: iceberg_position_delete_sys_table_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: iceberg_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: iceberg_orc_nested_column_utils.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: equality_delete.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: iceberg_sys_table_jni_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: jdbc_jni_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: max_compute_jni_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: paimon_jni_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: paimon_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: parquet_metadata_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: parquet_utils.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: remote_doris_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: table_format_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: table_schema_change_helper.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: transactional_hive_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: trino_connector_jni_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: text_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: merge_partitioner.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: iceberg_partition_function.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: vcsv_transformer.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: vfile_format_transformer_factory.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: viceberg_parquet_writer.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: vjni_format_transformer.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: vorc_transformer.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: vparquet_writer.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: adbc_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: hdfs_file_system.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: s3_file_system.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: unity_23_cxx.cxx:_ZN5dorisL14thread_contextEv Unexecuted instantiation: inverted_index_compound_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: inverted_index_fs_directory.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: zone_map_index.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: vcollect_iterator.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: predicate_creator_comparison.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: predicate_creator_in_list_in.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: predicate_creator_in_list_not_in.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: tablet.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: engine_clone_task.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: descriptors.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: backend_service.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: point_query_executor.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: be_server_starter_factory.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: http_service.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: internal_service.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: arrow_flight_batch_reader.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: thrift_util.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: load_stream.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: routine_load_task_executor.cpp:_ZN5dorisL14thread_contextEv Unexecuted instantiation: cloud_compaction_action.cpp:_ZN5dorisL14thread_contextEv |
291 | | |
292 | | class ScopedPeakMem { |
293 | | public: |
294 | | explicit ScopedPeakMem(int64_t* peak_mem) |
295 | 7.79k | : _peak_mem(peak_mem), |
296 | 7.79k | _mem_tracker("ScopedPeakMem:" + UniqueId::gen_uid().to_string()) { |
297 | 7.79k | ThreadLocalHandle::create_thread_local_if_not_exits(); |
298 | 7.79k | thread_context()->thread_mem_tracker_mgr->push_consumer_tracker(&_mem_tracker); |
299 | 7.79k | } |
300 | | |
301 | 7.79k | ~ScopedPeakMem() { |
302 | 7.79k | thread_context()->thread_mem_tracker_mgr->pop_consumer_tracker(); |
303 | 7.79k | *_peak_mem += _mem_tracker.peak_consumption(); |
304 | 7.79k | ThreadLocalHandle::del_thread_local_if_count_is_zero(); |
305 | 7.79k | } |
306 | | |
307 | | private: |
308 | | int64_t* _peak_mem; |
309 | | MemTracker _mem_tracker; |
310 | | }; |
311 | | |
312 | | // only hold thread context in scope. |
313 | | class ScopedInitThreadContext { |
314 | | public: |
315 | 596 | explicit ScopedInitThreadContext() { ThreadLocalHandle::create_thread_local_if_not_exits(); } |
316 | | |
317 | 587 | ~ScopedInitThreadContext() { ThreadLocalHandle::del_thread_local_if_count_is_zero(); } |
318 | | }; |
319 | | |
320 | | class AttachTask { |
321 | | public: |
322 | | // you must use ResourceCtx or MemTracker initialization. |
323 | | explicit AttachTask() = delete; |
324 | | |
325 | | explicit AttachTask(const std::shared_ptr<ResourceContext>& rc); |
326 | | |
327 | | // Shortcut attach task, initialize an empty resource context, and set the memory tracker. |
328 | | explicit AttachTask(const std::shared_ptr<MemTrackerLimiter>& mem_tracker); |
329 | | |
330 | | // is query or load, initialize with memory tracker, query id and workload group wptr. |
331 | | explicit AttachTask(RuntimeState* runtime_state); |
332 | | |
333 | | explicit AttachTask(QueryContext* query_ctx); |
334 | | |
335 | | void init(const std::shared_ptr<ResourceContext>& rc); |
336 | | |
337 | | ~AttachTask(); |
338 | | |
339 | | private: |
340 | | ScopedQueryLogContext _query_log_scope; |
341 | | }; |
342 | | |
343 | | class SwitchResourceContext { |
344 | | public: |
345 | | explicit SwitchResourceContext(const std::shared_ptr<ResourceContext>& rc); |
346 | | |
347 | | ~SwitchResourceContext(); |
348 | | |
349 | | private: |
350 | | std::shared_ptr<ResourceContext> old_resource_ctx_ {nullptr}; |
351 | | ScopedQueryLogContext _query_log_scope; |
352 | | }; |
353 | | |
354 | | class SwitchThreadMemTrackerLimiter { |
355 | | public: |
356 | | explicit SwitchThreadMemTrackerLimiter( |
357 | | const std::shared_ptr<doris::MemTrackerLimiter>& mem_tracker); |
358 | | |
359 | | ~SwitchThreadMemTrackerLimiter(); |
360 | | |
361 | | private: |
362 | | bool is_switched_ {false}; |
363 | | }; |
364 | | |
365 | | class AddThreadMemTrackerConsumer { |
366 | | public: |
367 | | // The owner and user of MemTracker are in the same thread, and the raw pointer is faster. |
368 | | // If mem_tracker is nullptr, do nothing. |
369 | | explicit AddThreadMemTrackerConsumer(MemTracker* mem_tracker); |
370 | | |
371 | | // The owner and user of MemTracker are in different threads. If mem_tracker is nullptr, do nothing. |
372 | | explicit AddThreadMemTrackerConsumer(const std::shared_ptr<MemTracker>& mem_tracker); |
373 | | |
374 | | ~AddThreadMemTrackerConsumer(); |
375 | | |
376 | | private: |
377 | | std::shared_ptr<MemTracker> _mem_tracker; // Avoid mem_tracker being released midway. |
378 | | bool _need_pop = false; |
379 | | }; |
380 | | |
381 | | class ScopeSkipMemoryCheck { |
382 | | public: |
383 | 1.50M | explicit ScopeSkipMemoryCheck() { |
384 | 1.50M | ThreadLocalHandle::create_thread_local_if_not_exits(); |
385 | 1.50M | doris::thread_context()->thread_mem_tracker_mgr->skip_memory_check++; |
386 | 1.50M | } |
387 | | |
388 | 1.50M | ~ScopeSkipMemoryCheck() { |
389 | 1.50M | doris::thread_context()->thread_mem_tracker_mgr->skip_memory_check--; |
390 | 1.50M | ThreadLocalHandle::del_thread_local_if_count_is_zero(); |
391 | 1.50M | } |
392 | | }; |
393 | | |
394 | | // Basic macros for mem tracker, usually do not need to be modified and used. |
395 | | // must call create_thread_local_if_not_exits() before use thread_context(). |
396 | | #define CONSUME_THREAD_MEM_TRACKER(size) \ |
397 | 13.7M | do { \ |
398 | 13.7M | if (size == 0) { \ |
399 | 115k | break; \ |
400 | 115k | } \ |
401 | 13.7M | if (doris::pthread_context_ptr_init) { \ |
402 | 13.1M | DCHECK(bthread_self() == 0); \ |
403 | 13.1M | doris::thread_context_ptr->thread_mem_tracker_mgr->consume(size); \ |
404 | 13.1M | } else if (bthread_self() != 0) { \ |
405 | 0 | auto* bthread_context = \ |
406 | 0 | static_cast<doris::ThreadContext*>(bthread_getspecific(doris::btls_key)); \ |
407 | 0 | DCHECK(bthread_context != nullptr); \ |
408 | 0 | if (bthread_context != nullptr) { \ |
409 | 0 | bthread_context->thread_mem_tracker_mgr->consume(size); \ |
410 | 0 | } else { \ |
411 | 0 | doris::ExecEnv::GetInstance()->orphan_mem_tracker()->consume_no_update_peak(size); \ |
412 | 0 | } \ |
413 | 525k | } else if (doris::ExecEnv::ready()) { \ |
414 | 0 | DCHECK(doris::k_doris_exit || !doris::config::enable_memory_orphan_check) \ |
415 | 0 | << doris::NO_THREAD_CONTEXT_MSG; \ |
416 | 0 | doris::ExecEnv::GetInstance()->orphan_mem_tracker()->consume_no_update_peak(size); \ |
417 | 0 | } \ |
418 | 13.6M | } while (0) |
419 | 6.93M | #define RELEASE_THREAD_MEM_TRACKER(size) CONSUME_THREAD_MEM_TRACKER(-size) |
420 | | |
421 | | } // namespace doris |