Coverage Report

Created: 2026-10-04 12:55

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
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