Coverage Report

Created: 2026-07-20 10:10

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/exec/operator/data_queue.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 "exec/operator/data_queue.h"
19
20
#include <glog/logging.h>
21
22
#include <algorithm>
23
#include <utility>
24
25
#include "common/thread_safety_annotations.h"
26
#include "core/block/block.h"
27
#include "exec/pipeline/dependency.h"
28
#include "fmt/format.h"
29
30
namespace doris {
31
32
8.24k
void SubQueue::try_pop(std::unique_ptr<Block>* output_block) {
33
8.24k
    LockGuard l(queue_lock);
34
8.24k
    if (!blocks.empty()) {
35
8.23k
        *output_block = std::move(blocks.front());
36
8.23k
        blocks.pop_front();
37
8.23k
        bytes_in_queue -= (*output_block)->allocated_bytes();
38
8.23k
        blocks_in_queue -= 1;
39
8.23k
        if (blocks.empty()) {
40
7.17k
            sink_dependency->set_ready();
41
7.17k
        }
42
8.23k
    }
43
8.24k
}
44
45
8.34k
bool SubQueue::try_push(std::unique_ptr<Block> block, std::atomic_uint32_t& total_counter) {
46
8.34k
    LockGuard l(queue_lock);
47
8.34k
    if (is_finished) {
48
72
        return false;
49
72
    }
50
8.27k
    total_counter++;
51
8.27k
    bytes_in_queue += block->allocated_bytes();
52
8.27k
    blocks.emplace_back(std::move(block));
53
8.27k
    blocks_in_queue += 1;
54
8.27k
    if (static_cast<int64_t>(blocks.size()) > max_blocks_in_queue.load()) {
55
1.08k
        sink_dependency->block();
56
1.08k
    }
57
8.27k
    return true;
58
8.34k
}
59
60
bool SubQueue::mark_finished(std::atomic_uint32_t& unfinished_counter,
61
10.8k
                             std::atomic_bool& all_finished) {
62
10.8k
    LockGuard l(queue_lock);
63
10.8k
    if (is_finished) {
64
5.37k
        return false;
65
5.37k
    }
66
5.49k
    is_finished = true;
67
5.49k
    if (unfinished_counter.fetch_sub(1) == 1) {
68
2.60k
        all_finished = true;
69
2.60k
    }
70
5.49k
    return true;
71
10.8k
}
72
73
5.46k
void SubQueue::clear_blocks() {
74
5.46k
    bool need_set_always_ready = false;
75
5.46k
    {
76
5.46k
        LockGuard l(queue_lock);
77
5.46k
        if (!blocks.empty()) {
78
31
            blocks.clear();
79
31
            bytes_in_queue = 0;
80
31
            blocks_in_queue = 0;
81
31
            need_set_always_ready = true;
82
31
        }
83
5.46k
    }
84
    // Notify outside of queue_lock to keep lock ordering simple.
85
5.46k
    if (need_set_always_ready) {
86
31
        sink_dependency->set_always_ready();
87
31
    }
88
5.46k
}
89
90
2.61k
DataQueue::DataQueue(int child_count) : _sub_queues(child_count), _child_count(child_count) {
91
5.52k
    for (auto& sub : _sub_queues) {
92
5.52k
        sub = std::make_unique<SubQueue>();
93
5.52k
    }
94
2.61k
    _un_finished_counter = child_count;
95
2.61k
}
96
97
3.67k
bool DataQueue::has_more_data() const {
98
3.67k
    return _cur_blocks_total_nums.load() > 0;
99
3.67k
}
100
101
5
std::string DataQueue::debug_string() const {
102
5
    return fmt::format("(is_all_finish = {}, has_data = {})", is_all_finish(), has_more_data());
103
5
}
104
105
void DataQueue::set_source_dependency(std::shared_ptr<Dependency> source_dependency)
106
2.58k
        NO_THREAD_SAFETY_ANALYSIS {
107
2.58k
    _source_dependency = std::move(source_dependency);
108
2.58k
}
109
110
5.48k
void DataQueue::set_sink_dependency(Dependency* sink_dependency, int child_idx) {
111
5.48k
    _sub_queues[child_idx]->sink_dependency = sink_dependency;
112
5.48k
}
113
114
5.49k
void DataQueue::set_max_blocks_in_sub_queue(int64_t max_blocks) {
115
13.2k
    for (auto& sub : _sub_queues) {
116
13.2k
        sub->max_blocks_in_queue = max_blocks;
117
13.2k
    }
118
5.49k
}
119
120
1
void DataQueue::set_low_memory_mode() {
121
1
    _is_low_memory_mode = true;
122
3
    for (auto& sub : _sub_queues) {
123
3
        sub->max_blocks_in_queue = 1;
124
3
    }
125
1
    clear_free_blocks();
126
1
}
127
128
8.16k
std::unique_ptr<Block> DataQueue::get_free_block(int child_idx) {
129
8.16k
    auto& sub = *_sub_queues[child_idx];
130
8.16k
    {
131
8.16k
        LockGuard l(sub.free_lock);
132
8.16k
        if (!sub.free_blocks.empty()) {
133
1.72k
            auto block = std::move(sub.free_blocks.front());
134
1.72k
            sub.free_blocks.pop_front();
135
1.72k
            return block;
136
1.72k
        }
137
8.16k
    }
138
139
6.44k
    return Block::create_unique();
140
8.16k
}
141
142
8.06k
void DataQueue::push_free_block(DataQueueBlock&& queue_block) {
143
8.06k
    if (!queue_block.block) {
144
0
        return;
145
0
    }
146
8.06k
    DCHECK(queue_block.block->rows() == 0);
147
148
8.06k
    if (!_is_low_memory_mode) {
149
8.06k
        auto& sub = *_sub_queues[queue_block.child_idx];
150
8.06k
        LockGuard l(sub.free_lock);
151
8.06k
        sub.free_blocks.emplace_back(std::move(queue_block.block));
152
8.06k
    }
153
8.06k
}
154
155
2.57k
void DataQueue::clear_free_blocks() {
156
5.46k
    for (auto& sub : _sub_queues) {
157
5.46k
        LockGuard l(sub->free_lock);
158
5.46k
        std::deque<std::unique_ptr<Block>> tmp_queue;
159
5.46k
        sub->free_blocks.swap(tmp_queue);
160
5.46k
    }
161
2.57k
}
162
163
2.57k
void DataQueue::terminate() {
164
8.03k
    for (int i = 0; i < _child_count; ++i) {
165
5.46k
        mark_finish(i);
166
5.46k
        _sub_queues[i]->clear_blocks();
167
5.46k
    }
168
2.57k
    _cur_blocks_total_nums = 0;
169
2.57k
    clear_free_blocks();
170
2.57k
    set_source_always_ready();
171
2.57k
}
172
173
8.23k
Result<DataQueueBlock> DataQueue::get_block_from_queue() {
174
8.23k
    DataQueueBlock result;
175
8.23k
    const int start_idx = (_flag_queue_idx + 1) % _child_count;
176
11.6k
    for (int offset = 0; offset < _child_count; ++offset) {
177
11.5k
        const int idx = (start_idx + offset) % _child_count;
178
11.5k
        if (_sub_queues[idx]->blocks_in_queue.load() == 0) {
179
3.36k
            continue;
180
3.36k
        }
181
182
8.23k
        auto& sub = *_sub_queues[idx];
183
8.23k
        sub.try_pop(&result.block);
184
8.23k
        if (!result.block) {
185
0
            continue;
186
0
        }
187
8.23k
        result.child_idx = idx;
188
8.23k
        _flag_queue_idx = idx;
189
8.23k
        auto old_total = _cur_blocks_total_nums.fetch_sub(1);
190
8.23k
        if (old_total == 1) {
191
6.29k
            set_source_block();
192
6.29k
        }
193
8.23k
        break;
194
8.23k
    }
195
196
    // A producer enqueues its final block before marking the child finished. Observe completion
197
    // first, then recheck queued data to avoid reporting EOS while the final block is still queued.
198
8.23k
    result.eos = is_all_finish() && !has_more_data();
199
8.23k
    return result;
200
8.23k
}
201
202
8.34k
Status DataQueue::push_block(std::unique_ptr<Block> block, int child_idx, bool eos) {
203
8.34k
    DCHECK(block || eos);
204
8.34k
    if (!block && !eos) {
205
0
        return Status::OK();
206
0
    }
207
208
8.34k
    if (block) {
209
8.32k
        auto& sub = *_sub_queues[child_idx];
210
        // total_counter is incremented inside try_push under queue_lock, only when the
211
        // block is actually enqueued. This ensures get_block_from_queue() always observes
212
        // _cur_blocks_total_nums >= 1 when it successfully pops a block, with no risk of
213
        // underflow or the need for a rollback on failure.
214
8.32k
        if (!sub.try_push(std::move(block), _cur_blocks_total_nums)) {
215
71
            return Status::EndOfFile("SubQueue already finished");
216
71
        }
217
8.32k
    }
218
219
8.27k
    if (eos) {
220
5.41k
        mark_finish(child_idx);
221
5.41k
    }
222
8.27k
    set_source_ready();
223
8.27k
    return Status::OK();
224
8.34k
}
225
226
10.8k
void DataQueue::mark_finish(int child_idx) {
227
10.8k
    auto& sub = *_sub_queues[child_idx];
228
10.8k
    if (!sub.mark_finished(_un_finished_counter, _is_all_finished)) {
229
5.37k
        return;
230
5.37k
    }
231
10.8k
}
232
233
14.5k
bool DataQueue::is_all_finish() const {
234
14.5k
    return _is_all_finished;
235
14.5k
}
236
237
8.29k
void DataQueue::set_source_ready() {
238
8.29k
    LockGuard lc(_source_lock);
239
8.29k
    if (_source_dependency) {
240
8.29k
        _source_dependency->set_ready();
241
8.29k
    }
242
8.29k
}
243
244
2.57k
void DataQueue::set_source_always_ready() {
245
2.57k
    LockGuard lc(_source_lock);
246
2.57k
    if (_source_dependency) {
247
2.57k
        _source_dependency->set_always_ready();
248
2.57k
    }
249
2.57k
}
250
251
6.29k
void DataQueue::set_source_block() {
252
    // Re-check under _source_lock to avoid blocking the source when a concurrent push
253
    // has already added new blocks (or all children have finished) since we observed
254
    // the counter drop to zero.
255
6.29k
    LockGuard lc(_source_lock);
256
6.29k
    if (_source_dependency && _cur_blocks_total_nums == 0 && !is_all_finish()) {
257
3.77k
        _source_dependency->block();
258
3.77k
    }
259
6.29k
}
260
261
} // namespace doris