be/src/util/jni_scan_heap_gate.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 <condition_variable> |
21 | | #include <cstdint> |
22 | | #include <deque> |
23 | | #include <functional> |
24 | | #include <mutex> |
25 | | |
26 | | namespace doris { |
27 | | |
28 | | // Admits Java scanners by the JVM heap they declare they will hold. |
29 | | // |
30 | | // Every JNI reader in BE runs its Java scanner in the one JVM BE starts. A few kinds hold a large part |
31 | | // of its heap from before their first batch until they close: a paimon merge read keeps a row group of |
32 | | // every file it merges, a fluss primary-key bucket read replays its change log into a map, a fluss |
33 | | // union read keeps its whole log tail. Each fits; opened together - the sixteen scanners of one scan, |
34 | | // or a few queries at once - they can run the heap out. |
35 | | // |
36 | | // The connector that plans such a range can say how much heap its reader will hold |
37 | | // (TFileRangeDesc.jni_heap_bytes), and does so only for a statement that asks with the session variable |
38 | | // enable_jni_heap_admission. It is off by default: then no reader declares anything, the gate admits |
39 | | // nothing, and a scan that needs more heap than the JVM has fails with OutOfMemoryError and says how to |
40 | | // give the JVM more. |
41 | | // |
42 | | // A reader that declared its heap asks the gate before it opens its Java scanner. The gate keeps an |
43 | | // account of what the readers it admitted declared, admits the next one while the account plus this |
44 | | // reader stays within jni_scanner_heap_budget_ratio of the JVM's maximum heap, and takes a reader's |
45 | | // share back when its scanner closes. Nothing here looks at the heap itself: what other users of the |
46 | | // JVM hold, and garbage nobody has collected yet, are not the gate's business. |
47 | | // |
48 | | // Readers that do not fit wait in the order they came, so a large one is not passed over for ever by |
49 | | // small ones arriving behind it. Three things end a wait whatever the account says: no reader holding |
50 | | // a share (nobody would give one back for this reader, so one that declares more than the whole budget |
51 | | // runs alone), the caller asking to stop waiting (a cancelled query), and the wait outlasting |
52 | | // jni_scanner_heap_max_wait_ms. |
53 | | class JniScanHeapGate { |
54 | | public: |
55 | | // Held by a reader from before its Java scanner opens until that scanner is closed. |
56 | | class Permit { |
57 | | public: |
58 | 52 | Permit() = default; |
59 | 52 | ~Permit() { release(); } |
60 | | Permit(const Permit&) = delete; |
61 | | Permit& operator=(const Permit&) = delete; |
62 | | |
63 | | // The reader's Java scanner is closed. Idempotent; a permit never held is a no-op. |
64 | | void release(); |
65 | 19 | bool held() const { return _gate != nullptr; } |
66 | | |
67 | | private: |
68 | | friend class JniScanHeapGate; |
69 | | JniScanHeapGate* _gate = nullptr; |
70 | | int64_t _bytes = 0; |
71 | | }; |
72 | | |
73 | | // `budget` answers how many bytes the admitted readers may declare together. It is asked every |
74 | | // time the gate looks, so a change of jni_scanner_heap_budget_ratio applies at once. |
75 | | explicit JniScanHeapGate(std::function<int64_t()> budget); |
76 | | |
77 | | // The gate every JNI reader of this process shares: its budget is jni_scanner_heap_budget_ratio |
78 | | // of the -Xmx the JVM was started with. |
79 | | static JniScanHeapGate* instance(); |
80 | | |
81 | | // Waits until `bytes` fit (see the class comment) and hands the reader `permit` for them. |
82 | | // `stop_waiting` is asked whenever the gate looks again; once it answers true the reader is |
83 | | // admitted without further waiting, and is expected to notice the stop itself. `wait_ns` gets the |
84 | | // time spent here. |
85 | | void acquire(int64_t bytes, const std::function<bool()>& stop_waiting, Permit* permit, |
86 | | int64_t* wait_ns); |
87 | | |
88 | | int64_t admitted_bytes() const; |
89 | | int64_t holders() const; |
90 | | int64_t waiters() const; |
91 | | |
92 | | private: |
93 | | // Caller holds _lock. |
94 | | bool _fits(uint64_t ticket, int64_t bytes) const; |
95 | | void _release(int64_t bytes); |
96 | | |
97 | | const std::function<int64_t()> _budget; |
98 | | mutable std::mutex _lock; |
99 | | std::condition_variable _cv; |
100 | | // What the admitted readers declared, and how many of them there are. |
101 | | int64_t _admitted_bytes = 0; |
102 | | int64_t _holders = 0; |
103 | | // The readers waiting, in the order they came: only the first may be admitted by the account. |
104 | | std::deque<uint64_t> _waiting; |
105 | | uint64_t _next_ticket = 0; |
106 | | }; |
107 | | |
108 | | } // namespace doris |