Coverage Report

Created: 2026-10-09 18:00

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/util/blocking_queue.hpp
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
// This file is copied from
18
// https://github.com/apache/impala/blob/branch-2.9.0/be/src/util/blocking-queue.hpp
19
// and modified by Doris
20
21
#pragma once
22
23
#include <unistd.h>
24
25
#include <atomic>
26
#include <condition_variable>
27
#include <list>
28
#include <mutex>
29
30
#include "common/logging.h"
31
#include "util/stopwatch.hpp"
32
33
#ifdef BE_TEST
34
#include "cpp/sync_point.h"
35
#endif
36
37
namespace doris {
38
// Fixed capacity FIFO queue, where both BlockingGet and BlockingPut operations block
39
// if the queue is empty or full, respectively.
40
template <typename T>
41
class BlockingQueue {
42
public:
43
    BlockingQueue(uint32_t max_elements)
44
138
            : _shutdown(false),
45
138
              _max_elements(max_elements),
46
138
              _total_get_wait_time(0),
47
138
              _total_put_wait_time(0),
48
138
              _get_waiting(0),
49
138
              _put_waiting(0) {}
_ZN5doris13BlockingQueueINS_14WorkThreadPoolILb0EE4TaskEEC2Ej
Line
Count
Source
44
114
            : _shutdown(false),
45
114
              _max_elements(max_elements),
46
114
              _total_get_wait_time(0),
47
114
              _total_put_wait_time(0),
48
114
              _get_waiting(0),
49
114
              _put_waiting(0) {}
_ZN5doris13BlockingQueueIiEC2Ej
Line
Count
Source
44
8
            : _shutdown(false),
45
8
              _max_elements(max_elements),
46
8
              _total_get_wait_time(0),
47
8
              _total_put_wait_time(0),
48
8
              _get_waiting(0),
49
8
              _put_waiting(0) {}
_ZN5doris13BlockingQueueIlEC2Ej
Line
Count
Source
44
2
            : _shutdown(false),
45
2
              _max_elements(max_elements),
46
2
              _total_get_wait_time(0),
47
2
              _total_put_wait_time(0),
48
2
              _get_waiting(0),
49
2
              _put_waiting(0) {}
_ZN5doris13BlockingQueueISt10shared_ptrIN5arrow11RecordBatchEEEC2Ej
Line
Count
Source
44
12
            : _shutdown(false),
45
12
              _max_elements(max_elements),
46
12
              _total_get_wait_time(0),
47
12
              _total_put_wait_time(0),
48
12
              _get_waiting(0),
49
12
              _put_waiting(0) {}
_ZN5doris13BlockingQueueIPN7RdKafka7MessageEEC2Ej
Line
Count
Source
44
2
            : _shutdown(false),
45
2
              _max_elements(max_elements),
46
2
              _total_get_wait_time(0),
47
2
              _total_put_wait_time(0),
48
2
              _get_waiting(0),
49
2
              _put_waiting(0) {}
Unexecuted instantiation: _ZN5doris13BlockingQueueISt10shared_ptrIN3Aws7Kinesis5Model6RecordEEEC2Ej
50
51
    // Get an element from the queue, waiting indefinitely for one to become available.
52
    // Returns false if we were shut down prior to getting the element, and there
53
    // are no more elements available.
54
120k
    bool blocking_get(T* out) { return controlled_blocking_get(out, MAX_CV_WAIT_TIMEOUT_MS); }
_ZN5doris13BlockingQueueINS_14WorkThreadPoolILb0EE4TaskEE12blocking_getEPS3_
Line
Count
Source
54
379
    bool blocking_get(T* out) { return controlled_blocking_get(out, MAX_CV_WAIT_TIMEOUT_MS); }
_ZN5doris13BlockingQueueIiE12blocking_getEPi
Line
Count
Source
54
119k
    bool blocking_get(T* out) { return controlled_blocking_get(out, MAX_CV_WAIT_TIMEOUT_MS); }
_ZN5doris13BlockingQueueIlE12blocking_getEPl
Line
Count
Source
54
4
    bool blocking_get(T* out) { return controlled_blocking_get(out, MAX_CV_WAIT_TIMEOUT_MS); }
_ZN5doris13BlockingQueueISt10shared_ptrIN5arrow11RecordBatchEEE12blocking_getEPS4_
Line
Count
Source
54
4
    bool blocking_get(T* out) { return controlled_blocking_get(out, MAX_CV_WAIT_TIMEOUT_MS); }
_ZN5doris13BlockingQueueIPN7RdKafka7MessageEE12blocking_getEPS3_
Line
Count
Source
54
2
    bool blocking_get(T* out) { return controlled_blocking_get(out, MAX_CV_WAIT_TIMEOUT_MS); }
Unexecuted instantiation: _ZN5doris13BlockingQueueISt10shared_ptrIN3Aws7Kinesis5Model6RecordEEE12blocking_getEPS6_
55
56
    // Blocking_get and blocking_put may cause deadlock,
57
    // but we still don't find root cause,
58
    // introduce condition variable wait timeout to avoid blocking queue deadlock temporarily.
59
120k
    bool controlled_blocking_get(T* out, int64_t cv_wait_timeout_ms) {
60
120k
        MonotonicStopWatch timer;
61
120k
        timer.start();
62
120k
        std::unique_lock<std::mutex> unique_lock(_lock);
63
120k
#ifdef BE_TEST
64
120k
        TEST_SYNC_POINT("BlockingQueue::controlled_blocking_get::after_lock");
65
120k
#endif
66
125k
        while (!(_shutdown || !_list.empty())) {
67
5.28k
            ++_get_waiting;
68
5.28k
#ifdef BE_TEST
69
5.28k
            TEST_SYNC_POINT("BlockingQueue::controlled_blocking_get::before_wait");
70
5.28k
#endif
71
5.28k
            _get_cv.wait_for(unique_lock, std::chrono::milliseconds(cv_wait_timeout_ms));
72
5.28k
            DCHECK_GT(_get_waiting, 0);
73
5.28k
            --_get_waiting;
74
5.28k
        }
75
120k
        _total_get_wait_time += timer.elapsed_time();
76
77
120k
        if (!_list.empty()) {
78
100k
            *out = _list.front();
79
100k
            _list.pop_front();
80
100k
            const bool has_put_waiter = _put_waiting > 0;
81
100k
            unique_lock.unlock();
82
100k
            if (has_put_waiter) {
83
4
                _put_cv.notify_one();
84
4
            }
85
100k
            return true;
86
100k
        } else {
87
20.1k
            assert(_shutdown);
88
20.2k
            return false;
89
20.1k
        }
90
120k
    }
_ZN5doris13BlockingQueueINS_14WorkThreadPoolILb0EE4TaskEE23controlled_blocking_getEPS3_l
Line
Count
Source
59
380
    bool controlled_blocking_get(T* out, int64_t cv_wait_timeout_ms) {
60
380
        MonotonicStopWatch timer;
61
380
        timer.start();
62
380
        std::unique_lock<std::mutex> unique_lock(_lock);
63
380
#ifdef BE_TEST
64
380
        TEST_SYNC_POINT("BlockingQueue::controlled_blocking_get::after_lock");
65
380
#endif
66
757
        while (!(_shutdown || !_list.empty())) {
67
377
            ++_get_waiting;
68
377
#ifdef BE_TEST
69
377
            TEST_SYNC_POINT("BlockingQueue::controlled_blocking_get::before_wait");
70
377
#endif
71
377
            _get_cv.wait_for(unique_lock, std::chrono::milliseconds(cv_wait_timeout_ms));
72
377
            DCHECK_GT(_get_waiting, 0);
73
377
            --_get_waiting;
74
377
        }
75
380
        _total_get_wait_time += timer.elapsed_time();
76
77
380
        if (!_list.empty()) {
78
118
            *out = _list.front();
79
118
            _list.pop_front();
80
118
            const bool has_put_waiter = _put_waiting > 0;
81
118
            unique_lock.unlock();
82
118
            if (has_put_waiter) {
83
0
                _put_cv.notify_one();
84
0
            }
85
118
            return true;
86
262
        } else {
87
262
            assert(_shutdown);
88
262
            return false;
89
262
        }
90
380
    }
_ZN5doris13BlockingQueueIiE23controlled_blocking_getEPil
Line
Count
Source
59
119k
    bool controlled_blocking_get(T* out, int64_t cv_wait_timeout_ms) {
60
119k
        MonotonicStopWatch timer;
61
119k
        timer.start();
62
119k
        std::unique_lock<std::mutex> unique_lock(_lock);
63
119k
#ifdef BE_TEST
64
119k
        TEST_SYNC_POINT("BlockingQueue::controlled_blocking_get::after_lock");
65
119k
#endif
66
124k
        while (!(_shutdown || !_list.empty())) {
67
4.89k
            ++_get_waiting;
68
4.89k
#ifdef BE_TEST
69
4.89k
            TEST_SYNC_POINT("BlockingQueue::controlled_blocking_get::before_wait");
70
4.89k
#endif
71
4.89k
            _get_cv.wait_for(unique_lock, std::chrono::milliseconds(cv_wait_timeout_ms));
72
4.89k
            DCHECK_GT(_get_waiting, 0);
73
4.89k
            --_get_waiting;
74
4.89k
        }
75
119k
        _total_get_wait_time += timer.elapsed_time();
76
77
119k
        if (!_list.empty()) {
78
100k
            *out = _list.front();
79
100k
            _list.pop_front();
80
100k
            const bool has_put_waiter = _put_waiting > 0;
81
100k
            unique_lock.unlock();
82
100k
            if (has_put_waiter) {
83
4
                _put_cv.notify_one();
84
4
            }
85
100k
            return true;
86
100k
        } else {
87
19.9k
            assert(_shutdown);
88
20.0k
            return false;
89
19.9k
        }
90
119k
    }
_ZN5doris13BlockingQueueIlE23controlled_blocking_getEPll
Line
Count
Source
59
4
    bool controlled_blocking_get(T* out, int64_t cv_wait_timeout_ms) {
60
4
        MonotonicStopWatch timer;
61
4
        timer.start();
62
4
        std::unique_lock<std::mutex> unique_lock(_lock);
63
4
#ifdef BE_TEST
64
4
        TEST_SYNC_POINT("BlockingQueue::controlled_blocking_get::after_lock");
65
4
#endif
66
4
        while (!(_shutdown || !_list.empty())) {
67
0
            ++_get_waiting;
68
0
#ifdef BE_TEST
69
0
            TEST_SYNC_POINT("BlockingQueue::controlled_blocking_get::before_wait");
70
0
#endif
71
0
            _get_cv.wait_for(unique_lock, std::chrono::milliseconds(cv_wait_timeout_ms));
72
0
            DCHECK_GT(_get_waiting, 0);
73
0
            --_get_waiting;
74
0
        }
75
4
        _total_get_wait_time += timer.elapsed_time();
76
77
4
        if (!_list.empty()) {
78
2
            *out = _list.front();
79
2
            _list.pop_front();
80
2
            const bool has_put_waiter = _put_waiting > 0;
81
2
            unique_lock.unlock();
82
2
            if (has_put_waiter) {
83
0
                _put_cv.notify_one();
84
0
            }
85
2
            return true;
86
2
        } else {
87
2
            assert(_shutdown);
88
2
            return false;
89
2
        }
90
4
    }
_ZN5doris13BlockingQueueISt10shared_ptrIN5arrow11RecordBatchEEE23controlled_blocking_getEPS4_l
Line
Count
Source
59
4
    bool controlled_blocking_get(T* out, int64_t cv_wait_timeout_ms) {
60
4
        MonotonicStopWatch timer;
61
4
        timer.start();
62
4
        std::unique_lock<std::mutex> unique_lock(_lock);
63
4
#ifdef BE_TEST
64
4
        TEST_SYNC_POINT("BlockingQueue::controlled_blocking_get::after_lock");
65
4
#endif
66
4
        while (!(_shutdown || !_list.empty())) {
67
0
            ++_get_waiting;
68
0
#ifdef BE_TEST
69
0
            TEST_SYNC_POINT("BlockingQueue::controlled_blocking_get::before_wait");
70
0
#endif
71
0
            _get_cv.wait_for(unique_lock, std::chrono::milliseconds(cv_wait_timeout_ms));
72
0
            DCHECK_GT(_get_waiting, 0);
73
0
            --_get_waiting;
74
0
        }
75
4
        _total_get_wait_time += timer.elapsed_time();
76
77
4
        if (!_list.empty()) {
78
4
            *out = _list.front();
79
4
            _list.pop_front();
80
4
            const bool has_put_waiter = _put_waiting > 0;
81
4
            unique_lock.unlock();
82
4
            if (has_put_waiter) {
83
0
                _put_cv.notify_one();
84
0
            }
85
4
            return true;
86
4
        } else {
87
0
            assert(_shutdown);
88
0
            return false;
89
0
        }
90
4
    }
_ZN5doris13BlockingQueueIPN7RdKafka7MessageEE23controlled_blocking_getEPS3_l
Line
Count
Source
59
4
    bool controlled_blocking_get(T* out, int64_t cv_wait_timeout_ms) {
60
4
        MonotonicStopWatch timer;
61
4
        timer.start();
62
4
        std::unique_lock<std::mutex> unique_lock(_lock);
63
4
#ifdef BE_TEST
64
4
        TEST_SYNC_POINT("BlockingQueue::controlled_blocking_get::after_lock");
65
4
#endif
66
16
        while (!(_shutdown || !_list.empty())) {
67
12
            ++_get_waiting;
68
12
#ifdef BE_TEST
69
12
            TEST_SYNC_POINT("BlockingQueue::controlled_blocking_get::before_wait");
70
12
#endif
71
12
            _get_cv.wait_for(unique_lock, std::chrono::milliseconds(cv_wait_timeout_ms));
72
12
            DCHECK_GT(_get_waiting, 0);
73
12
            --_get_waiting;
74
12
        }
75
4
        _total_get_wait_time += timer.elapsed_time();
76
77
4
        if (!_list.empty()) {
78
0
            *out = _list.front();
79
0
            _list.pop_front();
80
0
            const bool has_put_waiter = _put_waiting > 0;
81
0
            unique_lock.unlock();
82
0
            if (has_put_waiter) {
83
0
                _put_cv.notify_one();
84
0
            }
85
0
            return true;
86
4
        } else {
87
4
            assert(_shutdown);
88
4
            return false;
89
4
        }
90
4
    }
Unexecuted instantiation: _ZN5doris13BlockingQueueISt10shared_ptrIN3Aws7Kinesis5Model6RecordEEE23controlled_blocking_getEPS6_l
91
92
    // Puts an element into the queue, waiting indefinitely until there is space.
93
    // If the queue is shut down, returns false.
94
99.8k
    bool blocking_put(const T& val) { return controlled_blocking_put(val, MAX_CV_WAIT_TIMEOUT_MS); }
Unexecuted instantiation: _ZN5doris13BlockingQueueINS_14WorkThreadPoolILb0EE4TaskEE12blocking_putERKS3_
_ZN5doris13BlockingQueueISt10shared_ptrIN5arrow11RecordBatchEEE12blocking_putERKS4_
Line
Count
Source
94
8
    bool blocking_put(const T& val) { return controlled_blocking_put(val, MAX_CV_WAIT_TIMEOUT_MS); }
_ZN5doris13BlockingQueueIiE12blocking_putERKi
Line
Count
Source
94
99.8k
    bool blocking_put(const T& val) { return controlled_blocking_put(val, MAX_CV_WAIT_TIMEOUT_MS); }
_ZN5doris13BlockingQueueIlE12blocking_putERKl
Line
Count
Source
94
4
    bool blocking_put(const T& val) { return controlled_blocking_put(val, MAX_CV_WAIT_TIMEOUT_MS); }
95
96
    // Blocking_get and blocking_put may cause deadlock,
97
    // but we still don't find root cause,
98
    // introduce condition variable wait timeout to avoid blocking queue deadlock temporarily.
99
99.8k
    bool controlled_blocking_put(const T& val, int64_t cv_wait_timeout_ms) {
100
99.8k
        MonotonicStopWatch timer;
101
99.8k
        timer.start();
102
99.8k
        std::unique_lock<std::mutex> unique_lock(_lock);
103
100k
        while (!(_shutdown || _list.size() < _max_elements)) {
104
4
            ++_put_waiting;
105
4
#ifdef BE_TEST
106
4
            TEST_SYNC_POINT("BlockingQueue::controlled_blocking_put::before_wait");
107
4
#endif
108
4
            _put_cv.wait_for(unique_lock, std::chrono::milliseconds(cv_wait_timeout_ms));
109
4
            DCHECK_GT(_put_waiting, 0);
110
4
            --_put_waiting;
111
4
        }
112
99.8k
        _total_put_wait_time += timer.elapsed_time();
113
114
99.8k
        if (_shutdown) {
115
2
            return false;
116
2
        }
117
118
99.8k
        _list.push_back(val);
119
99.8k
        const bool has_get_waiter = _get_waiting > 0;
120
99.8k
        unique_lock.unlock();
121
99.8k
        if (has_get_waiter) {
122
18.3k
            _get_cv.notify_one();
123
18.3k
        }
124
99.8k
        return true;
125
99.8k
    }
Unexecuted instantiation: _ZN5doris13BlockingQueueINS_14WorkThreadPoolILb0EE4TaskEE23controlled_blocking_putERKS3_l
_ZN5doris13BlockingQueueISt10shared_ptrIN5arrow11RecordBatchEEE23controlled_blocking_putERKS4_l
Line
Count
Source
99
8
    bool controlled_blocking_put(const T& val, int64_t cv_wait_timeout_ms) {
100
8
        MonotonicStopWatch timer;
101
8
        timer.start();
102
8
        std::unique_lock<std::mutex> unique_lock(_lock);
103
8
        while (!(_shutdown || _list.size() < _max_elements)) {
104
0
            ++_put_waiting;
105
0
#ifdef BE_TEST
106
0
            TEST_SYNC_POINT("BlockingQueue::controlled_blocking_put::before_wait");
107
0
#endif
108
0
            _put_cv.wait_for(unique_lock, std::chrono::milliseconds(cv_wait_timeout_ms));
109
0
            DCHECK_GT(_put_waiting, 0);
110
0
            --_put_waiting;
111
0
        }
112
8
        _total_put_wait_time += timer.elapsed_time();
113
114
8
        if (_shutdown) {
115
0
            return false;
116
0
        }
117
118
8
        _list.push_back(val);
119
8
        const bool has_get_waiter = _get_waiting > 0;
120
8
        unique_lock.unlock();
121
8
        if (has_get_waiter) {
122
0
            _get_cv.notify_one();
123
0
        }
124
8
        return true;
125
8
    }
_ZN5doris13BlockingQueueIiE23controlled_blocking_putERKil
Line
Count
Source
99
99.7k
    bool controlled_blocking_put(const T& val, int64_t cv_wait_timeout_ms) {
100
99.7k
        MonotonicStopWatch timer;
101
99.7k
        timer.start();
102
99.7k
        std::unique_lock<std::mutex> unique_lock(_lock);
103
100k
        while (!(_shutdown || _list.size() < _max_elements)) {
104
4
            ++_put_waiting;
105
4
#ifdef BE_TEST
106
4
            TEST_SYNC_POINT("BlockingQueue::controlled_blocking_put::before_wait");
107
4
#endif
108
4
            _put_cv.wait_for(unique_lock, std::chrono::milliseconds(cv_wait_timeout_ms));
109
4
            DCHECK_GT(_put_waiting, 0);
110
4
            --_put_waiting;
111
4
        }
112
99.7k
        _total_put_wait_time += timer.elapsed_time();
113
114
99.7k
        if (_shutdown) {
115
0
            return false;
116
0
        }
117
118
99.7k
        _list.push_back(val);
119
99.7k
        const bool has_get_waiter = _get_waiting > 0;
120
99.7k
        unique_lock.unlock();
121
99.7k
        if (has_get_waiter) {
122
18.3k
            _get_cv.notify_one();
123
18.3k
        }
124
99.7k
        return true;
125
99.7k
    }
_ZN5doris13BlockingQueueIlE23controlled_blocking_putERKll
Line
Count
Source
99
4
    bool controlled_blocking_put(const T& val, int64_t cv_wait_timeout_ms) {
100
4
        MonotonicStopWatch timer;
101
4
        timer.start();
102
4
        std::unique_lock<std::mutex> unique_lock(_lock);
103
4
        while (!(_shutdown || _list.size() < _max_elements)) {
104
0
            ++_put_waiting;
105
0
#ifdef BE_TEST
106
0
            TEST_SYNC_POINT("BlockingQueue::controlled_blocking_put::before_wait");
107
0
#endif
108
0
            _put_cv.wait_for(unique_lock, std::chrono::milliseconds(cv_wait_timeout_ms));
109
0
            DCHECK_GT(_put_waiting, 0);
110
0
            --_put_waiting;
111
0
        }
112
4
        _total_put_wait_time += timer.elapsed_time();
113
114
4
        if (_shutdown) {
115
2
            return false;
116
2
        }
117
118
2
        _list.push_back(val);
119
2
        const bool has_get_waiter = _get_waiting > 0;
120
2
        unique_lock.unlock();
121
2
        if (has_get_waiter) {
122
0
            _get_cv.notify_one();
123
0
        }
124
2
        return true;
125
4
    }
Unexecuted instantiation: _ZN5doris13BlockingQueueIPN7RdKafka7MessageEE23controlled_blocking_putERKS3_l
Unexecuted instantiation: _ZN5doris13BlockingQueueISt10shared_ptrIN3Aws7Kinesis5Model6RecordEEE23controlled_blocking_putERKS6_l
126
127
    // Return false if queue full or has been shutdown.
128
150
    bool try_put(const T& val) {
129
150
        MonotonicStopWatch timer;
130
150
        timer.start();
131
150
        std::unique_lock<std::mutex> unique_lock(_lock);
132
150
#ifdef BE_TEST
133
150
        TEST_SYNC_POINT("BlockingQueue::try_put::after_lock");
134
150
#endif
135
150
        _total_put_wait_time += timer.elapsed_time();
136
137
150
        if (_shutdown || _list.size() >= _max_elements) {
138
10
            return false;
139
10
        }
140
141
140
        _list.push_back(val);
142
140
        const bool has_get_waiter = _get_waiting > 0;
143
140
        unique_lock.unlock();
144
140
        if (has_get_waiter) {
145
119
            _get_cv.notify_one();
146
119
        }
147
140
        return true;
148
150
    }
_ZN5doris13BlockingQueueINS_14WorkThreadPoolILb0EE4TaskEE7try_putERKS3_
Line
Count
Source
128
144
    bool try_put(const T& val) {
129
144
        MonotonicStopWatch timer;
130
144
        timer.start();
131
144
        std::unique_lock<std::mutex> unique_lock(_lock);
132
144
#ifdef BE_TEST
133
144
        TEST_SYNC_POINT("BlockingQueue::try_put::after_lock");
134
144
#endif
135
144
        _total_put_wait_time += timer.elapsed_time();
136
137
144
        if (_shutdown || _list.size() >= _max_elements) {
138
10
            return false;
139
10
        }
140
141
134
        _list.push_back(val);
142
134
        const bool has_get_waiter = _get_waiting > 0;
143
134
        unique_lock.unlock();
144
134
        if (has_get_waiter) {
145
115
            _get_cv.notify_one();
146
115
        }
147
134
        return true;
148
144
    }
_ZN5doris13BlockingQueueIiE7try_putERKi
Line
Count
Source
128
6
    bool try_put(const T& val) {
129
6
        MonotonicStopWatch timer;
130
6
        timer.start();
131
6
        std::unique_lock<std::mutex> unique_lock(_lock);
132
6
#ifdef BE_TEST
133
6
        TEST_SYNC_POINT("BlockingQueue::try_put::after_lock");
134
6
#endif
135
6
        _total_put_wait_time += timer.elapsed_time();
136
137
6
        if (_shutdown || _list.size() >= _max_elements) {
138
0
            return false;
139
0
        }
140
141
6
        _list.push_back(val);
142
6
        const bool has_get_waiter = _get_waiting > 0;
143
6
        unique_lock.unlock();
144
6
        if (has_get_waiter) {
145
4
            _get_cv.notify_one();
146
4
        }
147
6
        return true;
148
6
    }
149
150
    // Shut down the queue. Wakes up all threads waiting on BlockingGet or BlockingPut.
151
162
    void shutdown() {
152
162
        {
153
162
            std::lock_guard<std::mutex> guard(_lock);
154
162
            _shutdown = true;
155
162
        }
156
157
162
        _get_cv.notify_all();
158
162
        _put_cv.notify_all();
159
162
    }
_ZN5doris13BlockingQueueINS_14WorkThreadPoolILb0EE4TaskEE8shutdownEv
Line
Count
Source
151
144
    void shutdown() {
152
144
        {
153
144
            std::lock_guard<std::mutex> guard(_lock);
154
144
            _shutdown = true;
155
144
        }
156
157
144
        _get_cv.notify_all();
158
144
        _put_cv.notify_all();
159
144
    }
_ZN5doris13BlockingQueueIlE8shutdownEv
Line
Count
Source
151
2
    void shutdown() {
152
2
        {
153
2
            std::lock_guard<std::mutex> guard(_lock);
154
2
            _shutdown = true;
155
2
        }
156
157
2
        _get_cv.notify_all();
158
2
        _put_cv.notify_all();
159
2
    }
_ZN5doris13BlockingQueueIiE8shutdownEv
Line
Count
Source
151
6
    void shutdown() {
152
6
        {
153
6
            std::lock_guard<std::mutex> guard(_lock);
154
6
            _shutdown = true;
155
6
        }
156
157
6
        _get_cv.notify_all();
158
6
        _put_cv.notify_all();
159
6
    }
_ZN5doris13BlockingQueueISt10shared_ptrIN5arrow11RecordBatchEEE8shutdownEv
Line
Count
Source
151
4
    void shutdown() {
152
4
        {
153
4
            std::lock_guard<std::mutex> guard(_lock);
154
4
            _shutdown = true;
155
4
        }
156
157
4
        _get_cv.notify_all();
158
4
        _put_cv.notify_all();
159
4
    }
_ZN5doris13BlockingQueueIPN7RdKafka7MessageEE8shutdownEv
Line
Count
Source
151
6
    void shutdown() {
152
6
        {
153
6
            std::lock_guard<std::mutex> guard(_lock);
154
6
            _shutdown = true;
155
6
        }
156
157
6
        _get_cv.notify_all();
158
6
        _put_cv.notify_all();
159
6
    }
Unexecuted instantiation: _ZN5doris13BlockingQueueISt10shared_ptrIN3Aws7Kinesis5Model6RecordEEE8shutdownEv
160
161
1.03k
    uint32_t get_size() const {
162
1.03k
        std::lock_guard<std::mutex> l(_lock);
163
1.03k
        return static_cast<uint32_t>(_list.size());
164
1.03k
    }
_ZNK5doris13BlockingQueueINS_14WorkThreadPoolILb0EE4TaskEE8get_sizeEv
Line
Count
Source
161
1.03k
    uint32_t get_size() const {
162
1.03k
        std::lock_guard<std::mutex> l(_lock);
163
1.03k
        return static_cast<uint32_t>(_list.size());
164
1.03k
    }
Unexecuted instantiation: _ZNK5doris13BlockingQueueISt10shared_ptrIN5arrow11RecordBatchEEE8get_sizeEv
_ZNK5doris13BlockingQueueIPN7RdKafka7MessageEE8get_sizeEv
Line
Count
Source
161
2
    uint32_t get_size() const {
162
2
        std::lock_guard<std::mutex> l(_lock);
163
2
        return static_cast<uint32_t>(_list.size());
164
2
    }
Unexecuted instantiation: _ZNK5doris13BlockingQueueISt10shared_ptrIN3Aws7Kinesis5Model6RecordEEE8get_sizeEv
165
166
641
    uint32_t get_capacity() const { return _max_elements; }
167
168
#ifdef BE_TEST
169
16
    size_t get_waiting_count_for_test() const {
170
16
        std::lock_guard<std::mutex> guard(_lock);
171
16
        return _get_waiting;
172
16
    }
173
174
16
    size_t put_waiting_count_for_test() const {
175
16
        std::lock_guard<std::mutex> guard(_lock);
176
16
        return _put_waiting;
177
16
    }
178
#endif
179
180
    // Returns the total amount of time threads have blocked in BlockingGet.
181
640
    uint64_t total_get_wait_time() const { return _total_get_wait_time; }
182
183
    // Returns the total amount of time threads have blocked in BlockingPut.
184
642
    uint64_t total_put_wait_time() const { return _total_put_wait_time; }
185
186
private:
187
    static constexpr int64_t MAX_CV_WAIT_TIMEOUT_MS = 60 * 60 * 1000; // 1 hour
188
    bool _shutdown;
189
    const int _max_elements;
190
    std::condition_variable _get_cv; // 'get' callers wait on this
191
    std::condition_variable _put_cv; // 'put' callers wait on this
192
    // _lock guards access to _shutdown, _get_waiting, _put_waiting, _list, total_get_wait_time, and total_put_wait_time
193
    mutable std::mutex _lock;
194
    std::list<T> _list;
195
    std::atomic<uint64_t> _total_get_wait_time;
196
    std::atomic<uint64_t> _total_put_wait_time;
197
    // Number of threads currently inside the corresponding wait_for() call. The waiter that
198
    // increments a counter is also responsible for decrementing it after every wakeup reason.
199
    // Notifiers only inspect these counters and never consume a waiter's registration.
200
    size_t _get_waiting;
201
    size_t _put_waiting;
202
};
203
} // namespace doris