Coverage Report

Created: 2026-07-25 15:41

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/exec/operator/repeat_operator.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/repeat_operator.h"
19
20
#include <memory>
21
22
#include "common/logging.h"
23
#include "core/assert_cast.h"
24
#include "core/block/block.h"
25
#include "core/block/column_with_type_and_name.h"
26
#include "exec/operator/operator.h"
27
28
namespace doris {
29
class RuntimeState;
30
} // namespace doris
31
32
namespace doris {
33
34
RepeatLocalState::RepeatLocalState(RuntimeState* state, OperatorXBase* parent)
35
6
        : Base(state, parent), _child_block(Block::create_unique()), _repeat_id_idx(0) {}
36
37
6
Status RepeatLocalState::open(RuntimeState* state) {
38
6
    SCOPED_TIMER(exec_time_counter());
39
6
    SCOPED_TIMER(_open_timer);
40
6
    RETURN_IF_ERROR(Base::open(state));
41
6
    auto& p = _parent->cast<Parent>();
42
6
    _expr_ctxs.resize(p._expr_ctxs.size());
43
16
    for (size_t i = 0; i < _expr_ctxs.size(); i++) {
44
10
        RETURN_IF_ERROR(p._expr_ctxs[i]->clone(state, _expr_ctxs[i]));
45
10
    }
46
6
    return Status::OK();
47
6
}
48
49
6
Status RepeatLocalState::init(RuntimeState* state, LocalStateInfo& info) {
50
6
    RETURN_IF_ERROR(Base::init(state, info));
51
6
    SCOPED_TIMER(exec_time_counter());
52
6
    SCOPED_TIMER(_init_timer);
53
6
    _evaluate_input_timer = ADD_TIMER(custom_profile(), "EvaluateInputDataTime");
54
6
    _get_repeat_data_timer = ADD_TIMER(custom_profile(), "GetRepeatDataTime");
55
6
    _filter_timer = ADD_TIMER(custom_profile(), "FilterTime");
56
6
    return Status::OK();
57
6
}
58
59
0
Status RepeatOperatorX::init(const TPlanNode& tnode, RuntimeState* state) {
60
0
    RETURN_IF_ERROR(OperatorXBase::init(tnode, state));
61
0
    RETURN_IF_ERROR(VExpr::create_expr_trees(tnode.repeat_node.exprs, _expr_ctxs));
62
0
    for (const auto& slot_idx : _grouping_list) {
63
0
        if (slot_idx.size() < _repeat_id_list_size) {
64
0
            return Status::InternalError(
65
0
                    "grouping_list size {} is less than repeat_id_list size {}", slot_idx.size(),
66
0
                    _repeat_id_list_size);
67
0
        }
68
0
    }
69
0
    return Status::OK();
70
0
}
71
72
0
Status RepeatOperatorX::prepare(RuntimeState* state) {
73
0
    VLOG_CRITICAL << "VRepeatNode::open";
74
0
    RETURN_IF_ERROR(OperatorXBase::prepare(state));
75
0
    const auto* output_tuple_desc = state->desc_tbl().get_tuple_descriptor(_output_tuple_id);
76
0
    if (output_tuple_desc == nullptr) {
77
0
        return Status::InternalError("Failed to get tuple descriptor.");
78
0
    }
79
0
    for (const auto& slot_desc : output_tuple_desc->slots()) {
80
0
        _output_slots.push_back(slot_desc);
81
0
    }
82
0
    RETURN_IF_ERROR(VExpr::prepare(_expr_ctxs, state, _child->row_desc()));
83
0
    RETURN_IF_ERROR(VExpr::open(_expr_ctxs, state));
84
0
    return Status::OK();
85
0
}
86
87
RepeatOperatorX::RepeatOperatorX(ObjectPool* pool, const TPlanNode& tnode, int operator_id,
88
                                 const DescriptorTbl& descs)
89
0
        : Base(pool, tnode, operator_id, descs),
90
0
          _slot_id_set_list(tnode.repeat_node.slot_id_set_list),
91
0
          _all_slot_ids(tnode.repeat_node.all_slot_ids),
92
0
          _repeat_id_list_size(tnode.repeat_node.repeat_id_list.size()),
93
0
          _grouping_list(tnode.repeat_node.grouping_list),
94
0
          _output_tuple_id(tnode.repeat_node.output_tuple_id) {};
95
96
// The control logic of RepeatOperator is
97
// push a block, output _repeat_id_list_size blocks
98
// In the output block, the first part of the columns comes from the input block's columns, and the latter part of the columns is the grouping_id
99
// If there is no expr, there is only grouping_id
100
// If there is an expr, the first part of the columns in the output block uses _all_slot_ids and _slot_id_set_list to control whether it is null
101
24
bool RepeatOperatorX::need_more_input_data(RuntimeState* state) const {
102
24
    auto& local_state = state->get_local_state(operator_id())->cast<RepeatLocalState>();
103
24
    return !local_state._child_block->rows() && !local_state._child_eos;
104
24
}
105
106
Status RepeatLocalState::get_repeated_block(Block* input_block, int repeat_id_idx,
107
8
                                            Block* output_block) {
108
8
    auto& p = _parent->cast<RepeatOperatorX>();
109
8
    DCHECK(input_block != nullptr);
110
8
    DCHECK_EQ(output_block->rows(), 0);
111
112
8
    size_t input_column_size = input_block->columns();
113
8
    size_t output_column_size = p._output_slots.size();
114
8
    DCHECK_LT(input_column_size, output_column_size);
115
8
    auto scoped_mutable_block =
116
8
            VectorizedUtils::build_scoped_mutable_mem_reuse_block(output_block, p._output_slots);
117
8
    auto& m_block = scoped_mutable_block.mutable_block();
118
8
    auto& output_columns = m_block.mutable_columns();
119
    /* Fill all slots according to child, for example:select tc1,tc2,sum(tc3) from t1 group by grouping sets((tc1),(tc2));
120
     * insert into t1 values(1,2,1),(1,3,1),(2,1,1),(3,1,1);
121
     * slot_id_set_list=[[0],[1]],repeat_id_idx=0,
122
     * child_block 1,2,1 | 1,3,1 | 2,1,1 | 3,1,1
123
     * output_block 1,null,1,1 | 1,null,1,1 | 2,nul,1,1 | 3,null,1,1
124
     */
125
8
    size_t cur_col = 0;
126
28
    for (size_t i = 0; i < input_column_size; i++) {
127
20
        const ColumnWithTypeAndName& src_column = input_block->get_by_position(i);
128
20
        const auto slot_id = p._output_slots[cur_col]->id();
129
20
        const bool is_repeat_slot = p._all_slot_ids.contains(slot_id);
130
20
        const bool is_set_null_slot = !p._slot_id_set_list[repeat_id_idx].contains(slot_id);
131
20
        const auto row_size = src_column.column->size();
132
20
        ColumnPtr src = src_column.column;
133
20
        if (is_repeat_slot) {
134
12
            DCHECK(p._output_slots[cur_col]->is_nullable());
135
12
            auto* nullable_column = assert_cast<ColumnNullable*>(output_columns[cur_col].get());
136
12
            if (is_set_null_slot) {
137
                // is_set_null_slot = true, output all null
138
2
                nullable_column->insert_many_defaults(row_size);
139
10
            } else {
140
10
                if (!src_column.type->is_nullable()) {
141
6
                    nullable_column->get_nested_column().insert_range_from(*src_column.column, 0,
142
6
                                                                           row_size);
143
6
                    nullable_column->push_false_to_nullmap(row_size);
144
6
                } else {
145
4
                    nullable_column->insert_range_from(*src_column.column, 0, row_size);
146
4
                }
147
10
            }
148
12
        } else {
149
8
            output_columns[cur_col]->insert_range_from(*src_column.column, 0, row_size);
150
8
        }
151
20
        cur_col++;
152
20
    }
153
154
8
    const auto rows = input_block->rows();
155
    // Fill grouping ID to block
156
8
    RETURN_IF_ERROR(add_grouping_id_column(rows, cur_col, output_columns, repeat_id_idx));
157
158
8
    DCHECK_EQ(cur_col, output_column_size);
159
8
    return Status::OK();
160
8
}
161
162
Status RepeatLocalState::add_grouping_id_column(std::size_t rows, std::size_t& cur_col,
163
12
                                                MutableColumns& columns, int repeat_id_idx) {
164
12
    auto& p = _parent->cast<RepeatOperatorX>();
165
48
    for (auto slot_idx = 0; slot_idx < p._grouping_list.size(); slot_idx++) {
166
36
        DCHECK_LT(slot_idx, p._output_slots.size());
167
36
        int64_t val = p._grouping_list[slot_idx][repeat_id_idx];
168
36
        auto* column_ptr = columns[cur_col].get();
169
36
        DCHECK(!p._output_slots[cur_col]->is_nullable());
170
36
        auto* col = assert_cast<ColumnInt64*>(column_ptr);
171
36
        col->insert_many_vals(val, rows);
172
36
        cur_col++;
173
36
    }
174
12
    return Status::OK();
175
12
}
176
177
6
Status RepeatOperatorX::push(RuntimeState* state, Block* input_block, bool eos) const {
178
6
    auto& local_state = get_local_state(state);
179
6
    SCOPED_TIMER(local_state._evaluate_input_timer);
180
6
    local_state._child_eos = eos;
181
6
    auto& intermediate_block = local_state._intermediate_block;
182
6
    auto& expr_ctxs = local_state._expr_ctxs;
183
6
    DCHECK(!intermediate_block || intermediate_block->rows() == 0);
184
6
    if (input_block->rows() > 0) {
185
6
        SCOPED_PEAK_MEM(&local_state._estimate_memory_usage);
186
6
        intermediate_block = Block::create_unique();
187
188
10
        for (auto& expr : expr_ctxs) {
189
10
            ColumnWithTypeAndName result_data;
190
10
            RETURN_IF_ERROR(expr->execute(input_block, result_data));
191
10
            result_data.column = result_data.column->convert_to_full_column_if_const();
192
10
            intermediate_block->insert(result_data);
193
10
        }
194
6
        DCHECK_EQ(expr_ctxs.size(), intermediate_block->columns());
195
6
    }
196
197
6
    return Status::OK();
198
6
}
199
200
12
Status RepeatOperatorX::pull(doris::RuntimeState* state, Block* output_block, bool* eos) const {
201
12
    auto& local_state = get_local_state(state);
202
12
    auto& _repeat_id_idx = local_state._repeat_id_idx;
203
12
    auto& _child_block = *local_state._child_block;
204
12
    auto& _child_eos = local_state._child_eos;
205
12
    auto& _intermediate_block = local_state._intermediate_block;
206
12
    RETURN_IF_CANCELLED(state);
207
208
12
    SCOPED_PEAK_MEM(&local_state._estimate_memory_usage);
209
210
12
    DCHECK(_repeat_id_idx >= 0);
211
36
    for (const std::vector<int64_t>& v : _grouping_list) {
212
36
        DCHECK(_repeat_id_idx <= (int)v.size());
213
36
    }
214
12
    DCHECK(output_block->rows() == 0);
215
216
12
    {
217
12
        SCOPED_TIMER(local_state._get_repeat_data_timer);
218
        // Each pull increases _repeat_id_idx by one until _repeat_id_idx equals _repeat_id_list_size
219
        // Then clear the data of _intermediate_block and _child_block, and set _repeat_id_idx to 0
220
        // need_more_input_data will check if _child_block is empty
221
12
        if (_intermediate_block && _intermediate_block->rows() > 0) {
222
8
            RETURN_IF_ERROR(local_state.get_repeated_block(_intermediate_block.get(),
223
8
                                                           _repeat_id_idx, output_block));
224
225
8
            _repeat_id_idx++;
226
227
8
            if (_repeat_id_idx >= _repeat_id_list_size) {
228
4
                _intermediate_block->clear();
229
4
                _child_block.clear_column_data(_child->row_desc().num_materialized_slots());
230
4
                _repeat_id_idx = 0;
231
4
            }
232
8
        } else if (local_state._expr_ctxs.empty()) {
233
4
            auto scoped_mutable_block = VectorizedUtils::build_scoped_mutable_mem_reuse_block(
234
4
                    output_block, _output_slots);
235
4
            auto& m_block = scoped_mutable_block.mutable_block();
236
4
            auto rows = _child_block.rows();
237
4
            auto& columns = m_block.mutable_columns();
238
239
4
            std::size_t cur_col = 0;
240
4
            RETURN_IF_ERROR(
241
4
                    local_state.add_grouping_id_column(rows, cur_col, columns, _repeat_id_idx));
242
4
            _repeat_id_idx++;
243
244
4
            if (_repeat_id_idx >= _repeat_id_list_size) {
245
2
                _intermediate_block->clear();
246
2
                _child_block.clear_column_data(_child->row_desc().num_materialized_slots());
247
2
                _repeat_id_idx = 0;
248
2
            }
249
4
        }
250
12
    }
251
252
12
    {
253
12
        SCOPED_TIMER(local_state._filter_timer);
254
12
        RETURN_IF_ERROR(VExprContext::filter_block(local_state._conjuncts, output_block,
255
12
                                                   output_block->columns()));
256
12
    }
257
258
12
    *eos = _child_eos && _child_block.rows() == 0;
259
12
    local_state.reached_limit(output_block, eos);
260
12
    return Status::OK();
261
12
}
262
263
} // namespace doris