Coverage Report

Created: 2026-09-29 14:33

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/format/transformer/merge_partitioner.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 "format/transformer/merge_partitioner.h"
19
20
#include <algorithm>
21
#include <cstdint>
22
23
#include "common/cast_set.h"
24
#include "common/config.h"
25
#include "common/logging.h"
26
#include "common/status.h"
27
#include "core/block/block.h"
28
#include "core/column/column_const.h"
29
#include "core/column/column_nullable.h"
30
#include "core/column/column_vector.h"
31
#include "exec/sink/sink_common.h"
32
#include "format/transformer/iceberg_partition_function.h"
33
34
namespace doris {
35
36
MergePartitioner::MergePartitioner(size_t partition_count, const TMergePartitionInfo& merge_info,
37
                                   bool use_new_shuffle_hash_method)
38
6
        : PartitionerBase(static_cast<HashValType>(partition_count)),
39
6
          _merge_info(merge_info),
40
6
          _use_new_shuffle_hash_method(use_new_shuffle_hash_method),
41
6
          _insert_random(merge_info.insert_random) {}
42
43
6
Status MergePartitioner::init(const std::vector<TExpr>& /*texprs*/) {
44
6
    VExprContextSPtr op_ctx;
45
6
    RETURN_IF_ERROR(VExpr::create_expr_tree(_merge_info.operation_expr, op_ctx));
46
6
    _operation_expr_ctxs.emplace_back(std::move(op_ctx));
47
48
6
    std::vector<TExpr> insert_exprs;
49
6
    std::vector<TIcebergPartitionField> insert_fields;
50
6
    if (_merge_info.__isset.insert_partition_exprs) {
51
1
        insert_exprs = _merge_info.insert_partition_exprs;
52
1
    }
53
6
    if (_merge_info.__isset.insert_partition_fields) {
54
3
        insert_fields = _merge_info.insert_partition_fields;
55
3
    }
56
6
    if (!insert_exprs.empty() || !insert_fields.empty()) {
57
4
        _insert_partition_function = std::make_unique<IcebergInsertPartitionFunction>(
58
4
                _partition_count, _hash_method(), std::move(insert_exprs),
59
4
                std::move(insert_fields));
60
4
        RETURN_IF_ERROR(_insert_partition_function->init({}));
61
4
    }
62
63
6
    if (_merge_info.__isset.delete_partition_exprs && !_merge_info.delete_partition_exprs.empty()) {
64
1
        _delete_partition_function = std::make_unique<IcebergDeletePartitionFunction>(
65
1
                _partition_count, _hash_method(), _merge_info.delete_partition_exprs);
66
1
        RETURN_IF_ERROR(_delete_partition_function->init({}));
67
1
    }
68
6
    return Status::OK();
69
6
}
70
71
6
Status MergePartitioner::prepare(RuntimeState* state, const RowDescriptor& row_desc) {
72
6
    RETURN_IF_ERROR(VExpr::prepare(_operation_expr_ctxs, state, row_desc));
73
6
    if (_insert_partition_function != nullptr) {
74
4
        RETURN_IF_ERROR(_insert_partition_function->prepare(state, row_desc));
75
4
    }
76
6
    if (_delete_partition_function != nullptr) {
77
1
        RETURN_IF_ERROR(_delete_partition_function->prepare(state, row_desc));
78
1
    }
79
6
    return Status::OK();
80
6
}
81
82
6
Status MergePartitioner::open(RuntimeState* state) {
83
6
    RETURN_IF_ERROR(VExpr::open(_operation_expr_ctxs, state));
84
6
    if (_insert_partition_function != nullptr) {
85
4
        RETURN_IF_ERROR(_insert_partition_function->open(state));
86
4
        if (auto* insert_function =
87
4
                    dynamic_cast<IcebergInsertPartitionFunction*>(_insert_partition_function.get());
88
4
            insert_function != nullptr && insert_function->fallback_to_random()) {
89
1
            _insert_random = true;
90
1
        }
91
4
    }
92
6
    if (_delete_partition_function != nullptr) {
93
1
        RETURN_IF_ERROR(_delete_partition_function->open(state));
94
1
    }
95
6
    _init_insert_scaling(state);
96
6
    return Status::OK();
97
6
}
98
99
6
Status MergePartitioner::close(RuntimeState* /*state*/) {
100
6
    return Status::OK();
101
6
}
102
103
6
Status MergePartitioner::do_partitioning(RuntimeState* state, Block* block) const {
104
6
    const size_t rows = block->rows();
105
6
    if (rows == 0) {
106
0
        _channel_ids.clear();
107
0
        return Status::OK();
108
0
    }
109
110
6
    const size_t column_to_keep = block->columns();
111
6
    if (_operation_expr_ctxs.empty()) {
112
0
        return Status::InternalError("Merge partitioning missing operation expression");
113
0
    }
114
115
6
    int op_idx = -1;
116
6
    RETURN_IF_ERROR(_operation_expr_ctxs[0]->execute(block, &op_idx));
117
6
    if (op_idx < 0 || op_idx >= block->columns()) {
118
0
        return Status::InternalError("Merge partitioning missing operation column");
119
0
    }
120
6
    if (op_idx >= cast_set<int>(column_to_keep)) {
121
0
        return Status::InternalError("Merge partitioning requires operation column in input block");
122
0
    }
123
124
6
    const auto& op_column = block->get_by_position(op_idx).column;
125
6
    const auto* op_data = remove_nullable(op_column).get();
126
6
    std::vector<int8_t> ops(rows);
127
6
    bool has_insert = false;
128
6
    bool has_delete = false;
129
6
    bool has_update = false;
130
23
    for (size_t i = 0; i < rows; ++i) {
131
17
        int8_t op = static_cast<int8_t>(op_data->get_int(i));
132
17
        ops[i] = op;
133
17
        if (is_insert_op(op)) {
134
14
            has_insert = true;
135
14
        }
136
17
        if (is_delete_op(op)) {
137
4
            has_delete = true;
138
4
        }
139
17
        if (op == kUpdateOperation) {
140
1
            has_update = true;
141
1
        }
142
17
    }
143
144
6
    if (has_insert && !_insert_random && _insert_partition_function == nullptr) {
145
1
        return Status::InternalError("Merge partitioning insert exprs are empty");
146
1
    }
147
5
    if (has_delete && _delete_partition_function == nullptr) {
148
1
        return Status::InternalError("Merge partitioning delete exprs are empty");
149
1
    }
150
151
4
    std::vector<uint32_t> insert_hashes;
152
4
    std::vector<uint32_t> delete_hashes;
153
4
    const size_t insert_partition_count =
154
4
            _enable_insert_rebalance ? _insert_partition_count : _partition_count;
155
4
    if (has_insert && !_insert_random) {
156
3
        RETURN_IF_ERROR(_insert_partition_function->get_partitions(
157
3
                state, block, insert_partition_count, insert_hashes));
158
3
    }
159
4
    if (has_delete) {
160
1
        RETURN_IF_ERROR(_delete_partition_function->get_partitions(state, block, _partition_count,
161
1
                                                                   delete_hashes));
162
1
    }
163
4
    if (has_insert) {
164
4
        if (_insert_random) {
165
1
            if (_non_partition_scaling_threshold > 0) {
166
0
                _insert_data_processed += static_cast<int64_t>(block->bytes());
167
0
                if (_insert_writer_count < static_cast<int>(_partition_count) &&
168
0
                    _insert_data_processed >=
169
0
                            _insert_writer_count * _non_partition_scaling_threshold) {
170
0
                    _insert_writer_count++;
171
0
                }
172
1
            } else {
173
1
                _insert_writer_count = static_cast<int>(_partition_count);
174
1
            }
175
3
        } else if (_enable_insert_rebalance) {
176
1
            RETURN_IF_ERROR(_apply_insert_rebalance(ops, insert_hashes, block->bytes()));
177
1
        }
178
4
    }
179
180
4
    Block::erase_useless_column(block, column_to_keep);
181
182
4
    _channel_ids.resize(rows);
183
19
    for (size_t i = 0; i < rows; ++i) {
184
15
        const int8_t op = ops[i];
185
15
        if (op == kUpdateOperation) {
186
1
            _channel_ids[i] = delete_hashes[i];
187
1
            continue;
188
1
        }
189
14
        if (is_insert_op(op)) {
190
12
            _channel_ids[i] = _insert_random ? _next_rr_channel() : insert_hashes[i];
191
12
        } else if (is_delete_op(op)) {
192
2
            _channel_ids[i] = delete_hashes[i];
193
2
        } else {
194
0
            return Status::InternalError("Unknown Iceberg merge operation {}", op);
195
0
        }
196
14
    }
197
198
4
    if (has_update) {
199
5
        for (size_t col_idx = 0; col_idx < block->columns(); ++col_idx) {
200
4
            block->replace_by_position_if_const(col_idx);
201
4
        }
202
203
1
        auto mutable_columns_guard = block->mutate_columns_scoped();
204
1
        MutableColumns& mutable_columns = mutable_columns_guard.mutable_columns();
205
1
        MutableColumnPtr& op_mut = mutable_columns[op_idx];
206
1
        ColumnInt8* op_values_col = nullptr;
207
1
        if (auto* nullable_col = check_and_get_column<ColumnNullable>(op_mut.get())) {
208
0
            op_values_col =
209
0
                    check_and_get_column<ColumnInt8>(nullable_col->get_nested_column_ptr().get());
210
1
        } else {
211
1
            op_values_col = check_and_get_column<ColumnInt8>(op_mut.get());
212
1
        }
213
1
        if (op_values_col == nullptr) {
214
0
            return Status::InternalError("Merge operation column must be tinyint");
215
0
        }
216
1
        auto& op_values = op_values_col->get_data();
217
        // First pass: collect update row indices and mark original rows as DELETE.
218
1
        std::vector<size_t> update_rows;
219
6
        for (size_t row = 0; row < rows; ++row) {
220
5
            if (ops[row] != kUpdateOperation) {
221
4
                continue;
222
4
            }
223
1
            op_values[row] = kUpdateDeleteOperation;
224
1
            update_rows.push_back(row);
225
1
        }
226
        // Second pass: extract only the update rows into a temporary column,
227
        // then batch-append from it. This avoids cloning the entire column.
228
5
        for (size_t col_idx = 0; col_idx < mutable_columns.size(); ++col_idx) {
229
4
            auto tmp = mutable_columns[col_idx]->clone_empty();
230
4
            for (size_t row : update_rows) {
231
4
                tmp->insert_from(*mutable_columns[col_idx], row);
232
4
            }
233
4
            mutable_columns[col_idx]->insert_range_from(*tmp, 0, tmp->size());
234
4
        }
235
        // Mark the newly appended rows as INSERT and assign their channels.
236
1
        DCHECK(_insert_random || !insert_hashes.empty());
237
1
        const size_t appended_update_begin = rows;
238
2
        for (size_t idx = 0; idx < update_rows.size(); ++idx) {
239
1
            const size_t row = update_rows[idx];
240
1
            op_values[appended_update_begin + idx] = kUpdateInsertOperation;
241
1
            const uint32_t insert_channel =
242
1
                    _insert_random ? _next_rr_channel() : insert_hashes[row];
243
1
            _channel_ids.push_back(insert_channel);
244
1
        }
245
1
    }
246
247
4
    return Status::OK();
248
4
}
249
250
0
Status MergePartitioner::clone(RuntimeState* state, std::unique_ptr<PartitionerBase>& partitioner) {
251
0
    auto* new_partitioner =
252
0
            new MergePartitioner(_partition_count, _merge_info, _use_new_shuffle_hash_method);
253
0
    partitioner.reset(new_partitioner);
254
0
    RETURN_IF_ERROR(
255
0
            _clone_expr_ctxs(state, _operation_expr_ctxs, new_partitioner->_operation_expr_ctxs));
256
0
    if (_insert_partition_function != nullptr) {
257
0
        RETURN_IF_ERROR(_insert_partition_function->clone(
258
0
                state, new_partitioner->_insert_partition_function));
259
0
    }
260
0
    if (_delete_partition_function != nullptr) {
261
0
        RETURN_IF_ERROR(_delete_partition_function->clone(
262
0
                state, new_partitioner->_delete_partition_function));
263
0
    }
264
0
    new_partitioner->_insert_random = _insert_random;
265
0
    new_partitioner->_rr_offset = _rr_offset;
266
0
    return Status::OK();
267
0
}
268
269
Status MergePartitioner::_apply_insert_rebalance(const std::vector<int8_t>& ops,
270
                                                 std::vector<uint32_t>& insert_hashes,
271
1
                                                 size_t block_bytes) const {
272
1
    if (!_enable_insert_rebalance || _insert_writer_assigner == nullptr) {
273
0
        return Status::OK();
274
0
    }
275
1
    if (insert_hashes.empty() || _insert_partition_count == 0) {
276
0
        return Status::OK();
277
0
    }
278
1
    std::vector<uint8_t> mask(ops.size(), 0);
279
6
    for (size_t i = 0; i < ops.size(); ++i) {
280
5
        if (is_insert_op(ops[i])) {
281
3
            mask[i] = 1;
282
3
        }
283
5
    }
284
1
    return _insert_writer_assigner->assign(insert_hashes, &mask, ops.size(), block_bytes,
285
1
                                           insert_hashes);
286
1
}
287
288
6
void MergePartitioner::_init_insert_scaling(RuntimeState* state) {
289
6
    _enable_insert_rebalance = false;
290
6
    _insert_partition_count = 0;
291
6
    _insert_data_processed = 0;
292
6
    _insert_writer_count = 1;
293
6
    _insert_writer_assigner.reset();
294
6
    _non_partition_scaling_threshold =
295
6
            config::table_sink_non_partition_write_scaling_data_processed_threshold;
296
297
6
    if (_partition_count == 0) {
298
0
        return;
299
0
    }
300
6
    if (_insert_random) {
301
2
        return;
302
2
    }
303
4
    if (_insert_partition_function == nullptr) {
304
1
        return;
305
1
    }
306
307
3
    int max_partitions_per_writer =
308
3
            config::table_sink_partition_write_max_partition_nums_per_writer;
309
3
    if (max_partitions_per_writer <= 0) {
310
2
        return;
311
2
    }
312
1
    _insert_partition_count = _partition_count * max_partitions_per_writer;
313
1
    if (_insert_partition_count == 0) {
314
0
        return;
315
0
    }
316
317
1
    int task_num = state == nullptr ? 0 : state->task_num();
318
1
    int64_t min_partition_threshold = scale_writer_threshold_by_task(
319
1
            config::table_sink_partition_write_min_partition_data_processed_rebalance_threshold,
320
1
            task_num);
321
1
    int64_t min_data_threshold = scale_writer_threshold_by_task(
322
1
            config::table_sink_partition_write_min_data_processed_rebalance_threshold, task_num);
323
324
1
    _insert_writer_assigner = std::make_unique<SkewedWriterAssigner>(
325
1
            static_cast<int>(_insert_partition_count), static_cast<int>(_partition_count), 1,
326
1
            min_partition_threshold, min_data_threshold);
327
1
    _enable_insert_rebalance = true;
328
1
}
329
330
4
uint32_t MergePartitioner::_next_rr_channel() const {
331
4
    uint32_t writer_count = static_cast<uint32_t>(_partition_count);
332
4
    if (_insert_random && _insert_writer_count > 0) {
333
4
        writer_count = std::min<uint32_t>(static_cast<uint32_t>(_partition_count),
334
4
                                          static_cast<uint32_t>(_insert_writer_count));
335
4
    }
336
4
    if (writer_count == 0) {
337
0
        return 0;
338
0
    }
339
4
    const uint32_t channel = _rr_offset % writer_count;
340
4
    _rr_offset = (_rr_offset + 1) % writer_count;
341
4
    return channel;
342
4
}
343
344
Status MergePartitioner::_clone_expr_ctxs(RuntimeState* state, const VExprContextSPtrs& src,
345
0
                                          VExprContextSPtrs& dst) const {
346
0
    dst.resize(src.size());
347
0
    for (size_t i = 0; i < src.size(); ++i) {
348
0
        RETURN_IF_ERROR(src[i]->clone(state, dst[i]));
349
0
    }
350
0
    return Status::OK();
351
0
}
352
353
} // namespace doris