Coverage Report

Created: 2026-08-06 11:46

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/runtime/result_block_buffer.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 "runtime/result_block_buffer.h"
19
20
#include <gen_cpp/Data_types.h>
21
#include <gen_cpp/PaloInternalService_types.h>
22
#include <gen_cpp/internal_service.pb.h>
23
#include <glog/logging.h>
24
#include <google/protobuf/stubs/callback.h>
25
// IWYU pragma: no_include <bits/chrono.h>
26
#include <chrono> // IWYU pragma: keep
27
#include <limits>
28
#include <ostream>
29
#include <string>
30
#include <utility>
31
#include <vector>
32
33
#include "arrow/type_fwd.h"
34
#include "common/config.h"
35
#include "core/block/block.h"
36
#include "exec/pipeline/dependency.h"
37
#include "exec/sink/writer/varrow_flight_result_writer.h"
38
#include "exec/sink/writer/vmysql_result_writer.h"
39
#include "runtime/runtime_profile.h"
40
#include "runtime/thread_context.h"
41
#include "util/thrift_util.h"
42
43
namespace doris {
44
45
template <typename ResultCtxType>
46
ResultBlockBuffer<ResultCtxType>::ResultBlockBuffer(TUniqueId id, RuntimeState* state,
47
                                                    int buffer_size)
48
273k
        : _fragment_id(std::move(id)),
49
273k
          _is_close(false),
50
273k
          _batch_size(state->batch_size()),
51
273k
          _timezone(state->timezone()),
52
273k
          _be_exec_version(state->be_exec_version()),
53
273k
          _fragment_transmission_compression_type(state->fragement_transmission_compression_type()),
54
273k
          _buffer_limit(buffer_size) {
55
273k
    _mem_tracker = MemTrackerLimiter::create_shared(
56
273k
            MemTrackerLimiter::Type::QUERY,
57
273k
            fmt::format("ResultBlockBuffer#FragmentInstanceId={}", print_id(_fragment_id)));
58
273k
}
_ZN5doris17ResultBlockBufferINS_22GetArrowResultBatchCtxEEC2ENS_9TUniqueIdEPNS_12RuntimeStateEi
Line
Count
Source
48
113
        : _fragment_id(std::move(id)),
49
113
          _is_close(false),
50
113
          _batch_size(state->batch_size()),
51
113
          _timezone(state->timezone()),
52
113
          _be_exec_version(state->be_exec_version()),
53
113
          _fragment_transmission_compression_type(state->fragement_transmission_compression_type()),
54
113
          _buffer_limit(buffer_size) {
55
113
    _mem_tracker = MemTrackerLimiter::create_shared(
56
113
            MemTrackerLimiter::Type::QUERY,
57
113
            fmt::format("ResultBlockBuffer#FragmentInstanceId={}", print_id(_fragment_id)));
58
113
}
_ZN5doris17ResultBlockBufferINS_17GetResultBatchCtxEEC2ENS_9TUniqueIdEPNS_12RuntimeStateEi
Line
Count
Source
48
273k
        : _fragment_id(std::move(id)),
49
273k
          _is_close(false),
50
273k
          _batch_size(state->batch_size()),
51
273k
          _timezone(state->timezone()),
52
273k
          _be_exec_version(state->be_exec_version()),
53
273k
          _fragment_transmission_compression_type(state->fragement_transmission_compression_type()),
54
273k
          _buffer_limit(buffer_size) {
55
273k
    _mem_tracker = MemTrackerLimiter::create_shared(
56
273k
            MemTrackerLimiter::Type::QUERY,
57
273k
            fmt::format("ResultBlockBuffer#FragmentInstanceId={}", print_id(_fragment_id)));
58
273k
}
59
60
template <typename ResultCtxType>
61
Status ResultBlockBuffer<ResultCtxType>::close(const TUniqueId& id, Status exec_status,
62
425k
                                               int64_t num_rows, bool& is_fully_closed) {
63
425k
    std::unique_lock<std::mutex> l(_lock);
64
425k
    _returned_rows.fetch_add(num_rows);
65
    // close will be called multiple times and error status needs to be collected.
66
425k
    if (!exec_status.ok()) {
67
1.33k
        _status = exec_status;
68
1.33k
    }
69
70
425k
    auto it = _result_sink_dependencies.find(id);
71
426k
    if (it != _result_sink_dependencies.end()) {
72
426k
        it->second->set_always_ready();
73
426k
        _result_sink_dependencies.erase(it);
74
18.4E
    } else {
75
18.4E
        _status = Status::InternalError("Instance {} is not found in ResultBlockBuffer",
76
18.4E
                                        print_id(id));
77
18.4E
    }
78
425k
    if (!_result_sink_dependencies.empty()) {
79
        // Still waiting for other instances to finish; this is not the final close.
80
151k
        is_fully_closed = false;
81
151k
        return _status;
82
151k
    }
83
84
    // All instances have closed: the buffer is now fully closed.
85
274k
    is_fully_closed = true;
86
274k
    _is_close = true;
87
274k
    _arrow_data_arrival.notify_all();
88
89
274k
    if (!_waiting_rpc.empty()) {
90
161k
        if (_status.ok()) {
91
159k
            for (auto& ctx : _waiting_rpc) {
92
159k
                ctx->on_close(_packet_num, _returned_rows);
93
159k
            }
94
159k
        } else {
95
1.37k
            for (auto& ctx : _waiting_rpc) {
96
1.30k
                ctx->on_failure(_status);
97
1.30k
            }
98
1.37k
        }
99
161k
        _waiting_rpc.clear();
100
161k
    }
101
102
274k
    return _status;
103
425k
}
_ZN5doris17ResultBlockBufferINS_22GetArrowResultBatchCtxEE5closeERKNS_9TUniqueIdENS_6StatusElRb
Line
Count
Source
62
111
                                               int64_t num_rows, bool& is_fully_closed) {
63
111
    std::unique_lock<std::mutex> l(_lock);
64
111
    _returned_rows.fetch_add(num_rows);
65
    // close will be called multiple times and error status needs to be collected.
66
111
    if (!exec_status.ok()) {
67
2
        _status = exec_status;
68
2
    }
69
70
111
    auto it = _result_sink_dependencies.find(id);
71
111
    if (it != _result_sink_dependencies.end()) {
72
110
        it->second->set_always_ready();
73
110
        _result_sink_dependencies.erase(it);
74
110
    } else {
75
1
        _status = Status::InternalError("Instance {} is not found in ResultBlockBuffer",
76
1
                                        print_id(id));
77
1
    }
78
111
    if (!_result_sink_dependencies.empty()) {
79
        // Still waiting for other instances to finish; this is not the final close.
80
1
        is_fully_closed = false;
81
1
        return _status;
82
1
    }
83
84
    // All instances have closed: the buffer is now fully closed.
85
110
    is_fully_closed = true;
86
110
    _is_close = true;
87
110
    _arrow_data_arrival.notify_all();
88
89
110
    if (!_waiting_rpc.empty()) {
90
2
        if (_status.ok()) {
91
1
            for (auto& ctx : _waiting_rpc) {
92
1
                ctx->on_close(_packet_num, _returned_rows);
93
1
            }
94
1
        } else {
95
1
            for (auto& ctx : _waiting_rpc) {
96
1
                ctx->on_failure(_status);
97
1
            }
98
1
        }
99
2
        _waiting_rpc.clear();
100
2
    }
101
102
110
    return _status;
103
111
}
_ZN5doris17ResultBlockBufferINS_17GetResultBatchCtxEE5closeERKNS_9TUniqueIdENS_6StatusElRb
Line
Count
Source
62
425k
                                               int64_t num_rows, bool& is_fully_closed) {
63
425k
    std::unique_lock<std::mutex> l(_lock);
64
425k
    _returned_rows.fetch_add(num_rows);
65
    // close will be called multiple times and error status needs to be collected.
66
425k
    if (!exec_status.ok()) {
67
1.33k
        _status = exec_status;
68
1.33k
    }
69
70
425k
    auto it = _result_sink_dependencies.find(id);
71
426k
    if (it != _result_sink_dependencies.end()) {
72
426k
        it->second->set_always_ready();
73
426k
        _result_sink_dependencies.erase(it);
74
18.4E
    } else {
75
18.4E
        _status = Status::InternalError("Instance {} is not found in ResultBlockBuffer",
76
18.4E
                                        print_id(id));
77
18.4E
    }
78
425k
    if (!_result_sink_dependencies.empty()) {
79
        // Still waiting for other instances to finish; this is not the final close.
80
151k
        is_fully_closed = false;
81
151k
        return _status;
82
151k
    }
83
84
    // All instances have closed: the buffer is now fully closed.
85
273k
    is_fully_closed = true;
86
273k
    _is_close = true;
87
273k
    _arrow_data_arrival.notify_all();
88
89
273k
    if (!_waiting_rpc.empty()) {
90
161k
        if (_status.ok()) {
91
159k
            for (auto& ctx : _waiting_rpc) {
92
159k
                ctx->on_close(_packet_num, _returned_rows);
93
159k
            }
94
159k
        } else {
95
1.36k
            for (auto& ctx : _waiting_rpc) {
96
1.30k
                ctx->on_failure(_status);
97
1.30k
            }
98
1.36k
        }
99
161k
        _waiting_rpc.clear();
100
161k
    }
101
102
273k
    return _status;
103
425k
}
104
105
template <typename ResultCtxType>
106
224k
void ResultBlockBuffer<ResultCtxType>::cancel(const Status& reason) {
107
224k
    std::unique_lock<std::mutex> l(_lock);
108
224k
    SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(_mem_tracker);
109
224k
    if (_status.ok()) {
110
223k
        _status = reason;
111
223k
    }
112
224k
    _arrow_data_arrival.notify_all();
113
224k
    for (auto& ctx : _waiting_rpc) {
114
2
        ctx->on_failure(reason);
115
2
    }
116
224k
    _waiting_rpc.clear();
117
224k
    _update_dependency();
118
224k
    _result_batch_queue.clear();
119
224k
}
_ZN5doris17ResultBlockBufferINS_22GetArrowResultBatchCtxEE6cancelERKNS_6StatusE
Line
Count
Source
106
109
void ResultBlockBuffer<ResultCtxType>::cancel(const Status& reason) {
107
109
    std::unique_lock<std::mutex> l(_lock);
108
109
    SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(_mem_tracker);
109
109
    if (_status.ok()) {
110
109
        _status = reason;
111
109
    }
112
109
    _arrow_data_arrival.notify_all();
113
109
    for (auto& ctx : _waiting_rpc) {
114
1
        ctx->on_failure(reason);
115
1
    }
116
109
    _waiting_rpc.clear();
117
109
    _update_dependency();
118
109
    _result_batch_queue.clear();
119
109
}
_ZN5doris17ResultBlockBufferINS_17GetResultBatchCtxEE6cancelERKNS_6StatusE
Line
Count
Source
106
224k
void ResultBlockBuffer<ResultCtxType>::cancel(const Status& reason) {
107
224k
    std::unique_lock<std::mutex> l(_lock);
108
224k
    SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(_mem_tracker);
109
224k
    if (_status.ok()) {
110
223k
        _status = reason;
111
223k
    }
112
224k
    _arrow_data_arrival.notify_all();
113
224k
    for (auto& ctx : _waiting_rpc) {
114
1
        ctx->on_failure(reason);
115
1
    }
116
224k
    _waiting_rpc.clear();
117
224k
    _update_dependency();
118
224k
    _result_batch_queue.clear();
119
224k
}
120
121
template <typename ResultCtxType>
122
void ResultBlockBuffer<ResultCtxType>::set_dependency(
123
423k
        const TUniqueId& id, std::shared_ptr<Dependency> result_sink_dependency) {
124
423k
    std::unique_lock<std::mutex> l(_lock);
125
423k
    _result_sink_dependencies[id] = result_sink_dependency;
126
423k
    _update_dependency();
127
423k
}
_ZN5doris17ResultBlockBufferINS_22GetArrowResultBatchCtxEE14set_dependencyERKNS_9TUniqueIdESt10shared_ptrINS_10DependencyEE
Line
Count
Source
123
113
        const TUniqueId& id, std::shared_ptr<Dependency> result_sink_dependency) {
124
113
    std::unique_lock<std::mutex> l(_lock);
125
113
    _result_sink_dependencies[id] = result_sink_dependency;
126
113
    _update_dependency();
127
113
}
_ZN5doris17ResultBlockBufferINS_17GetResultBatchCtxEE14set_dependencyERKNS_9TUniqueIdESt10shared_ptrINS_10DependencyEE
Line
Count
Source
123
423k
        const TUniqueId& id, std::shared_ptr<Dependency> result_sink_dependency) {
124
423k
    std::unique_lock<std::mutex> l(_lock);
125
423k
    _result_sink_dependencies[id] = result_sink_dependency;
126
423k
    _update_dependency();
127
423k
}
128
129
template <typename ResultCtxType>
130
1.27M
void ResultBlockBuffer<ResultCtxType>::_update_dependency() {
131
1.27M
    if (!_status.ok()) {
132
224k
        for (auto it : _result_sink_dependencies) {
133
6
            it.second->set_ready();
134
6
        }
135
224k
        return;
136
224k
    }
137
138
1.84M
    for (auto it : _result_sink_dependencies) {
139
1.84M
        if (_instance_rows[it.first] > _batch_size) {
140
11
            it.second->block();
141
1.84M
        } else {
142
1.84M
            it.second->set_ready();
143
1.84M
        }
144
1.84M
    }
145
1.04M
}
_ZN5doris17ResultBlockBufferINS_22GetArrowResultBatchCtxEE18_update_dependencyEv
Line
Count
Source
130
721
void ResultBlockBuffer<ResultCtxType>::_update_dependency() {
131
721
    if (!_status.ok()) {
132
111
        for (auto it : _result_sink_dependencies) {
133
3
            it.second->set_ready();
134
3
        }
135
111
        return;
136
111
    }
137
138
610
    for (auto it : _result_sink_dependencies) {
139
409
        if (_instance_rows[it.first] > _batch_size) {
140
4
            it.second->block();
141
405
        } else {
142
405
            it.second->set_ready();
143
405
        }
144
409
    }
145
610
}
_ZN5doris17ResultBlockBufferINS_17GetResultBatchCtxEE18_update_dependencyEv
Line
Count
Source
130
1.27M
void ResultBlockBuffer<ResultCtxType>::_update_dependency() {
131
1.27M
    if (!_status.ok()) {
132
224k
        for (auto it : _result_sink_dependencies) {
133
3
            it.second->set_ready();
134
3
        }
135
224k
        return;
136
224k
    }
137
138
1.84M
    for (auto it : _result_sink_dependencies) {
139
1.84M
        if (_instance_rows[it.first] > _batch_size) {
140
7
            it.second->block();
141
1.84M
        } else {
142
1.84M
            it.second->set_ready();
143
1.84M
        }
144
1.84M
    }
145
1.04M
}
146
147
template <typename ResultCtxType>
148
445k
Status ResultBlockBuffer<ResultCtxType>::get_batch(std::shared_ptr<ResultCtxType> ctx) {
149
445k
    std::lock_guard<std::mutex> l(_lock);
150
445k
    SCOPED_ATTACH_TASK(_mem_tracker);
151
446k
    Defer defer {[&]() { _update_dependency(); }};
_ZZN5doris17ResultBlockBufferINS_22GetArrowResultBatchCtxEE9get_batchESt10shared_ptrIS1_EENKUlvE_clEv
Line
Count
Source
151
8
    Defer defer {[&]() { _update_dependency(); }};
_ZZN5doris17ResultBlockBufferINS_17GetResultBatchCtxEE9get_batchESt10shared_ptrIS1_EENKUlvE_clEv
Line
Count
Source
151
446k
    Defer defer {[&]() { _update_dependency(); }};
152
445k
    if (!_status.ok()) {
153
9
        ctx->on_failure(_status);
154
9
        return _status;
155
9
    }
156
445k
    if (!_result_batch_queue.empty()) {
157
9.50k
        auto result = _result_batch_queue.front();
158
9.50k
        _result_batch_queue.pop_front();
159
9.69k
        for (auto it : _instance_rows_in_queue.front()) {
160
9.69k
            _instance_rows[it.first] -= it.second;
161
9.69k
        }
162
9.50k
        _instance_rows_in_queue.pop_front();
163
9.50k
        RETURN_IF_ERROR(ctx->on_data(result, _packet_num, this));
164
9.50k
        _packet_num++;
165
9.50k
        return Status::OK();
166
9.50k
    }
167
436k
    if (_is_close) {
168
110k
        if (!_status.ok()) {
169
0
            ctx->on_failure(_status);
170
0
            return Status::OK();
171
0
        }
172
110k
        ctx->on_close(_packet_num, _returned_rows);
173
110k
        LOG(INFO) << fmt::format(
174
110k
                "ResultBlockBuffer finished, fragment_id={}, is_close={}, is_cancelled={}, "
175
110k
                "packet_num={}, peak_memory_usage={}",
176
110k
                print_id(_fragment_id), _is_close, !_status.ok(), _packet_num,
177
110k
                _mem_tracker->peak_consumption());
178
110k
        return Status::OK();
179
110k
    }
180
    // no ready data, push ctx to waiting list
181
325k
    _waiting_rpc.push_back(ctx);
182
325k
    return Status::OK();
183
436k
}
_ZN5doris17ResultBlockBufferINS_22GetArrowResultBatchCtxEE9get_batchESt10shared_ptrIS1_E
Line
Count
Source
148
8
Status ResultBlockBuffer<ResultCtxType>::get_batch(std::shared_ptr<ResultCtxType> ctx) {
149
8
    std::lock_guard<std::mutex> l(_lock);
150
8
    SCOPED_ATTACH_TASK(_mem_tracker);
151
8
    Defer defer {[&]() { _update_dependency(); }};
152
8
    if (!_status.ok()) {
153
1
        ctx->on_failure(_status);
154
1
        return _status;
155
1
    }
156
7
    if (!_result_batch_queue.empty()) {
157
2
        auto result = _result_batch_queue.front();
158
2
        _result_batch_queue.pop_front();
159
2
        for (auto it : _instance_rows_in_queue.front()) {
160
2
            _instance_rows[it.first] -= it.second;
161
2
        }
162
2
        _instance_rows_in_queue.pop_front();
163
2
        RETURN_IF_ERROR(ctx->on_data(result, _packet_num, this));
164
2
        _packet_num++;
165
2
        return Status::OK();
166
2
    }
167
5
    if (_is_close) {
168
1
        if (!_status.ok()) {
169
0
            ctx->on_failure(_status);
170
0
            return Status::OK();
171
0
        }
172
1
        ctx->on_close(_packet_num, _returned_rows);
173
1
        LOG(INFO) << fmt::format(
174
1
                "ResultBlockBuffer finished, fragment_id={}, is_close={}, is_cancelled={}, "
175
1
                "packet_num={}, peak_memory_usage={}",
176
1
                print_id(_fragment_id), _is_close, !_status.ok(), _packet_num,
177
1
                _mem_tracker->peak_consumption());
178
1
        return Status::OK();
179
1
    }
180
    // no ready data, push ctx to waiting list
181
4
    _waiting_rpc.push_back(ctx);
182
4
    return Status::OK();
183
5
}
_ZN5doris17ResultBlockBufferINS_17GetResultBatchCtxEE9get_batchESt10shared_ptrIS1_E
Line
Count
Source
148
445k
Status ResultBlockBuffer<ResultCtxType>::get_batch(std::shared_ptr<ResultCtxType> ctx) {
149
445k
    std::lock_guard<std::mutex> l(_lock);
150
445k
    SCOPED_ATTACH_TASK(_mem_tracker);
151
445k
    Defer defer {[&]() { _update_dependency(); }};
152
445k
    if (!_status.ok()) {
153
8
        ctx->on_failure(_status);
154
8
        return _status;
155
8
    }
156
445k
    if (!_result_batch_queue.empty()) {
157
9.50k
        auto result = _result_batch_queue.front();
158
9.50k
        _result_batch_queue.pop_front();
159
9.69k
        for (auto it : _instance_rows_in_queue.front()) {
160
9.69k
            _instance_rows[it.first] -= it.second;
161
9.69k
        }
162
9.50k
        _instance_rows_in_queue.pop_front();
163
9.50k
        RETURN_IF_ERROR(ctx->on_data(result, _packet_num, this));
164
9.50k
        _packet_num++;
165
9.50k
        return Status::OK();
166
9.50k
    }
167
436k
    if (_is_close) {
168
110k
        if (!_status.ok()) {
169
0
            ctx->on_failure(_status);
170
0
            return Status::OK();
171
0
        }
172
110k
        ctx->on_close(_packet_num, _returned_rows);
173
110k
        LOG(INFO) << fmt::format(
174
110k
                "ResultBlockBuffer finished, fragment_id={}, is_close={}, is_cancelled={}, "
175
110k
                "packet_num={}, peak_memory_usage={}",
176
110k
                print_id(_fragment_id), _is_close, !_status.ok(), _packet_num,
177
110k
                _mem_tracker->peak_consumption());
178
110k
        return Status::OK();
179
110k
    }
180
    // no ready data, push ctx to waiting list
181
325k
    _waiting_rpc.push_back(ctx);
182
325k
    return Status::OK();
183
436k
}
184
185
template <typename ResultCtxType>
186
Status ResultBlockBuffer<ResultCtxType>::add_batch(RuntimeState* state,
187
177k
                                                   std::shared_ptr<InBlockType>& result) {
188
177k
    std::unique_lock<std::mutex> l(_lock);
189
190
177k
    if (!_status.ok()) {
191
2
        return _status;
192
2
    }
193
194
177k
    if (_waiting_rpc.empty()) {
195
11.8k
        auto sz = 0;
196
11.8k
        auto num_rows = 0;
197
11.8k
        size_t batch_size = 0;
198
11.8k
        if constexpr (std::is_same_v<InBlockType, Block>) {
199
200
            num_rows = cast_set<int>(result->rows());
200
200
            batch_size = result->bytes();
201
11.6k
        } else if constexpr (std::is_same_v<InBlockType, TFetchDataResult>) {
202
11.6k
            num_rows = cast_set<int>(result->result_batch.rows.size());
203
19.5k
            for (const auto& row : result->result_batch.rows) {
204
19.5k
                batch_size += row.size();
205
19.5k
            }
206
11.6k
        }
207
11.8k
        if (!_result_batch_queue.empty()) {
208
1.97k
            if constexpr (std::is_same_v<InBlockType, Block>) {
209
16
                sz = cast_set<int>(_result_batch_queue.back()->rows());
210
1.96k
            } else if constexpr (std::is_same_v<InBlockType, TFetchDataResult>) {
211
1.96k
                sz = cast_set<int>(_result_batch_queue.back()->result_batch.rows.size());
212
1.96k
            }
213
1.97k
            if (sz + num_rows < _buffer_limit &&
214
1.97k
                (batch_size + _last_batch_bytes) <= config::thrift_max_message_size) {
215
1.97k
                if constexpr (std::is_same_v<InBlockType, Block>) {
216
16
                    auto last_block = _result_batch_queue.back();
217
16
                    auto mutable_columns_guard = last_block->mutate_columns_scoped();
218
16
                    auto& mutable_columns = mutable_columns_guard.mutable_columns();
219
72
                    for (size_t i = 0; i < last_block->columns(); i++) {
220
56
                        mutable_columns[i]->insert_range_from(*result->get_by_position(i).column, 0,
221
56
                                                              num_rows);
222
56
                    }
223
1.96k
                } else {
224
1.96k
                    std::vector<std::string>& back_rows =
225
1.96k
                            _result_batch_queue.back()->result_batch.rows;
226
1.96k
                    std::vector<std::string>& result_rows = result->result_batch.rows;
227
1.96k
                    back_rows.insert(back_rows.end(), std::make_move_iterator(result_rows.begin()),
228
1.96k
                                     std::make_move_iterator(result_rows.end()));
229
1.96k
                }
230
1.97k
                _last_batch_bytes += batch_size;
231
1.97k
            } else {
232
0
                _instance_rows_in_queue.emplace_back();
233
0
                _result_batch_queue.push_back(std::move(result));
234
0
                _last_batch_bytes = batch_size;
235
0
                _arrow_data_arrival
236
0
                        .notify_one(); // Only valid for get_arrow_batch(std::shared_ptr<Block>,)
237
0
            }
238
9.90k
        } else {
239
9.90k
            _instance_rows_in_queue.emplace_back();
240
9.90k
            _result_batch_queue.push_back(std::move(result));
241
9.90k
            _last_batch_bytes = batch_size;
242
9.90k
            _arrow_data_arrival
243
9.90k
                    .notify_one(); // Only valid for get_arrow_batch(std::shared_ptr<Block>,)
244
9.90k
        }
245
11.8k
        _instance_rows[state->fragment_instance_id()] += num_rows;
246
11.8k
        _instance_rows_in_queue.back()[state->fragment_instance_id()] += num_rows;
247
165k
    } else {
248
165k
        auto ctx = _waiting_rpc.front();
249
165k
        _waiting_rpc.pop_front();
250
165k
        RETURN_IF_ERROR(ctx->on_data(result, _packet_num, this));
251
165k
        _packet_num++;
252
165k
    }
253
254
177k
    _update_dependency();
255
177k
    return Status::OK();
256
177k
}
_ZN5doris17ResultBlockBufferINS_22GetArrowResultBatchCtxEE9add_batchEPNS_12RuntimeStateERSt10shared_ptrINS_5BlockEE
Line
Count
Source
187
202
                                                   std::shared_ptr<InBlockType>& result) {
188
202
    std::unique_lock<std::mutex> l(_lock);
189
190
202
    if (!_status.ok()) {
191
1
        return _status;
192
1
    }
193
194
201
    if (_waiting_rpc.empty()) {
195
200
        auto sz = 0;
196
200
        auto num_rows = 0;
197
200
        size_t batch_size = 0;
198
200
        if constexpr (std::is_same_v<InBlockType, Block>) {
199
200
            num_rows = cast_set<int>(result->rows());
200
200
            batch_size = result->bytes();
201
        } else if constexpr (std::is_same_v<InBlockType, TFetchDataResult>) {
202
            num_rows = cast_set<int>(result->result_batch.rows.size());
203
            for (const auto& row : result->result_batch.rows) {
204
                batch_size += row.size();
205
            }
206
        }
207
200
        if (!_result_batch_queue.empty()) {
208
16
            if constexpr (std::is_same_v<InBlockType, Block>) {
209
16
                sz = cast_set<int>(_result_batch_queue.back()->rows());
210
            } else if constexpr (std::is_same_v<InBlockType, TFetchDataResult>) {
211
                sz = cast_set<int>(_result_batch_queue.back()->result_batch.rows.size());
212
            }
213
16
            if (sz + num_rows < _buffer_limit &&
214
16
                (batch_size + _last_batch_bytes) <= config::thrift_max_message_size) {
215
16
                if constexpr (std::is_same_v<InBlockType, Block>) {
216
16
                    auto last_block = _result_batch_queue.back();
217
16
                    auto mutable_columns_guard = last_block->mutate_columns_scoped();
218
16
                    auto& mutable_columns = mutable_columns_guard.mutable_columns();
219
72
                    for (size_t i = 0; i < last_block->columns(); i++) {
220
56
                        mutable_columns[i]->insert_range_from(*result->get_by_position(i).column, 0,
221
56
                                                              num_rows);
222
56
                    }
223
                } else {
224
                    std::vector<std::string>& back_rows =
225
                            _result_batch_queue.back()->result_batch.rows;
226
                    std::vector<std::string>& result_rows = result->result_batch.rows;
227
                    back_rows.insert(back_rows.end(), std::make_move_iterator(result_rows.begin()),
228
                                     std::make_move_iterator(result_rows.end()));
229
                }
230
16
                _last_batch_bytes += batch_size;
231
16
            } else {
232
0
                _instance_rows_in_queue.emplace_back();
233
0
                _result_batch_queue.push_back(std::move(result));
234
0
                _last_batch_bytes = batch_size;
235
0
                _arrow_data_arrival
236
0
                        .notify_one(); // Only valid for get_arrow_batch(std::shared_ptr<Block>,)
237
0
            }
238
184
        } else {
239
184
            _instance_rows_in_queue.emplace_back();
240
184
            _result_batch_queue.push_back(std::move(result));
241
184
            _last_batch_bytes = batch_size;
242
184
            _arrow_data_arrival
243
184
                    .notify_one(); // Only valid for get_arrow_batch(std::shared_ptr<Block>,)
244
184
        }
245
200
        _instance_rows[state->fragment_instance_id()] += num_rows;
246
200
        _instance_rows_in_queue.back()[state->fragment_instance_id()] += num_rows;
247
200
    } else {
248
1
        auto ctx = _waiting_rpc.front();
249
1
        _waiting_rpc.pop_front();
250
1
        RETURN_IF_ERROR(ctx->on_data(result, _packet_num, this));
251
1
        _packet_num++;
252
1
    }
253
254
201
    _update_dependency();
255
201
    return Status::OK();
256
201
}
_ZN5doris17ResultBlockBufferINS_17GetResultBatchCtxEE9add_batchEPNS_12RuntimeStateERSt10shared_ptrINS_16TFetchDataResultEE
Line
Count
Source
187
176k
                                                   std::shared_ptr<InBlockType>& result) {
188
176k
    std::unique_lock<std::mutex> l(_lock);
189
190
176k
    if (!_status.ok()) {
191
1
        return _status;
192
1
    }
193
194
176k
    if (_waiting_rpc.empty()) {
195
11.6k
        auto sz = 0;
196
11.6k
        auto num_rows = 0;
197
11.6k
        size_t batch_size = 0;
198
        if constexpr (std::is_same_v<InBlockType, Block>) {
199
            num_rows = cast_set<int>(result->rows());
200
            batch_size = result->bytes();
201
11.6k
        } else if constexpr (std::is_same_v<InBlockType, TFetchDataResult>) {
202
11.6k
            num_rows = cast_set<int>(result->result_batch.rows.size());
203
19.5k
            for (const auto& row : result->result_batch.rows) {
204
19.5k
                batch_size += row.size();
205
19.5k
            }
206
11.6k
        }
207
11.6k
        if (!_result_batch_queue.empty()) {
208
            if constexpr (std::is_same_v<InBlockType, Block>) {
209
                sz = cast_set<int>(_result_batch_queue.back()->rows());
210
1.96k
            } else if constexpr (std::is_same_v<InBlockType, TFetchDataResult>) {
211
1.96k
                sz = cast_set<int>(_result_batch_queue.back()->result_batch.rows.size());
212
1.96k
            }
213
1.96k
            if (sz + num_rows < _buffer_limit &&
214
1.96k
                (batch_size + _last_batch_bytes) <= config::thrift_max_message_size) {
215
                if constexpr (std::is_same_v<InBlockType, Block>) {
216
                    auto last_block = _result_batch_queue.back();
217
                    auto mutable_columns_guard = last_block->mutate_columns_scoped();
218
                    auto& mutable_columns = mutable_columns_guard.mutable_columns();
219
                    for (size_t i = 0; i < last_block->columns(); i++) {
220
                        mutable_columns[i]->insert_range_from(*result->get_by_position(i).column, 0,
221
                                                              num_rows);
222
                    }
223
1.96k
                } else {
224
1.96k
                    std::vector<std::string>& back_rows =
225
1.96k
                            _result_batch_queue.back()->result_batch.rows;
226
1.96k
                    std::vector<std::string>& result_rows = result->result_batch.rows;
227
1.96k
                    back_rows.insert(back_rows.end(), std::make_move_iterator(result_rows.begin()),
228
1.96k
                                     std::make_move_iterator(result_rows.end()));
229
1.96k
                }
230
1.96k
                _last_batch_bytes += batch_size;
231
1.96k
            } else {
232
0
                _instance_rows_in_queue.emplace_back();
233
0
                _result_batch_queue.push_back(std::move(result));
234
0
                _last_batch_bytes = batch_size;
235
0
                _arrow_data_arrival
236
0
                        .notify_one(); // Only valid for get_arrow_batch(std::shared_ptr<Block>,)
237
0
            }
238
9.71k
        } else {
239
9.71k
            _instance_rows_in_queue.emplace_back();
240
9.71k
            _result_batch_queue.push_back(std::move(result));
241
9.71k
            _last_batch_bytes = batch_size;
242
9.71k
            _arrow_data_arrival
243
9.71k
                    .notify_one(); // Only valid for get_arrow_batch(std::shared_ptr<Block>,)
244
9.71k
        }
245
11.6k
        _instance_rows[state->fragment_instance_id()] += num_rows;
246
11.6k
        _instance_rows_in_queue.back()[state->fragment_instance_id()] += num_rows;
247
165k
    } else {
248
165k
        auto ctx = _waiting_rpc.front();
249
165k
        _waiting_rpc.pop_front();
250
165k
        RETURN_IF_ERROR(ctx->on_data(result, _packet_num, this));
251
165k
        _packet_num++;
252
165k
    }
253
254
176k
    _update_dependency();
255
176k
    return Status::OK();
256
176k
}
257
258
template class ResultBlockBuffer<GetArrowResultBatchCtx>;
259
template class ResultBlockBuffer<GetResultBatchCtx>;
260
261
} // namespace doris