Coverage Report

Created: 2026-08-17 22:43

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/root/doris/be/src/exec/partitioner/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 "exec/partitioner/partitioner.h"
19
20
#include "common/cast_set.h"
21
#include "common/status.h"
22
#include "core/column/column_const.h"
23
#include "exec/exchange/local_exchange_sink_operator.h"
24
#include "exec/exchange/vdata_stream_sender.h"
25
#include "runtime/thread_context.h"
26
27
namespace doris {
28
#include "common/compile_check_begin.h"
29
30
template <typename ChannelIds>
31
38
Status Crc32HashPartitioner<ChannelIds>::do_partitioning(RuntimeState* state, Block* block) const {
32
38
    size_t rows = block->rows();
33
34
38
    if (rows > 0) {
35
38
        auto column_to_keep = block->columns();
36
37
38
        int result_size = cast_set<int>(_partition_expr_ctxs.size());
38
38
        std::vector<int> result(result_size);
39
40
38
        _initialize_hash_vals(rows);
41
38
        auto* __restrict hashes = _hash_vals.data();
42
38
        RETURN_IF_ERROR(_get_partition_column_result(block, result));
43
75
        for (int j = 0; j < result_size; ++j) {
44
37
            const auto& [col, is_const] = unpack_if_const(block->get_by_position(result[j]).column);
45
37
            if (is_const) {
46
0
                continue;
47
0
            }
48
37
            _do_hash(col, hashes, j);
49
37
        }
50
51
3.14M
        for (size_t i = 0; i < rows; i++) {
52
3.14M
            hashes[i] = ChannelIds()(hashes[i], _partition_count);
53
3.14M
        }
54
55
38
        Block::erase_useless_column(block, column_to_keep);
56
38
    }
57
38
    return Status::OK();
58
38
}
_ZNK5doris20Crc32HashPartitionerINS_17ShuffleChannelIdsEE15do_partitioningEPNS_12RuntimeStateEPNS_5BlockE
Line
Count
Source
31
23
Status Crc32HashPartitioner<ChannelIds>::do_partitioning(RuntimeState* state, Block* block) const {
32
23
    size_t rows = block->rows();
33
34
23
    if (rows > 0) {
35
23
        auto column_to_keep = block->columns();
36
37
23
        int result_size = cast_set<int>(_partition_expr_ctxs.size());
38
23
        std::vector<int> result(result_size);
39
40
23
        _initialize_hash_vals(rows);
41
23
        auto* __restrict hashes = _hash_vals.data();
42
23
        RETURN_IF_ERROR(_get_partition_column_result(block, result));
43
45
        for (int j = 0; j < result_size; ++j) {
44
22
            const auto& [col, is_const] = unpack_if_const(block->get_by_position(result[j]).column);
45
22
            if (is_const) {
46
0
                continue;
47
0
            }
48
22
            _do_hash(col, hashes, j);
49
22
        }
50
51
240
        for (size_t i = 0; i < rows; i++) {
52
217
            hashes[i] = ChannelIds()(hashes[i], _partition_count);
53
217
        }
54
55
23
        Block::erase_useless_column(block, column_to_keep);
56
23
    }
57
23
    return Status::OK();
58
23
}
_ZNK5doris20Crc32HashPartitionerINS_24SpillPartitionChannelIdsEE15do_partitioningEPNS_12RuntimeStateEPNS_5BlockE
Line
Count
Source
31
6
Status Crc32HashPartitioner<ChannelIds>::do_partitioning(RuntimeState* state, Block* block) const {
32
6
    size_t rows = block->rows();
33
34
6
    if (rows > 0) {
35
6
        auto column_to_keep = block->columns();
36
37
6
        int result_size = cast_set<int>(_partition_expr_ctxs.size());
38
6
        std::vector<int> result(result_size);
39
40
6
        _initialize_hash_vals(rows);
41
6
        auto* __restrict hashes = _hash_vals.data();
42
6
        RETURN_IF_ERROR(_get_partition_column_result(block, result));
43
12
        for (int j = 0; j < result_size; ++j) {
44
6
            const auto& [col, is_const] = unpack_if_const(block->get_by_position(result[j]).column);
45
6
            if (is_const) {
46
0
                continue;
47
0
            }
48
6
            _do_hash(col, hashes, j);
49
6
        }
50
51
3.14M
        for (size_t i = 0; i < rows; i++) {
52
3.14M
            hashes[i] = ChannelIds()(hashes[i], _partition_count);
53
3.14M
        }
54
55
6
        Block::erase_useless_column(block, column_to_keep);
56
6
    }
57
6
    return Status::OK();
58
6
}
_ZNK5doris20Crc32HashPartitionerINS_26SpillRePartitionChannelIdsEE15do_partitioningEPNS_12RuntimeStateEPNS_5BlockE
Line
Count
Source
31
8
Status Crc32HashPartitioner<ChannelIds>::do_partitioning(RuntimeState* state, Block* block) const {
32
8
    size_t rows = block->rows();
33
34
8
    if (rows > 0) {
35
8
        auto column_to_keep = block->columns();
36
37
8
        int result_size = cast_set<int>(_partition_expr_ctxs.size());
38
8
        std::vector<int> result(result_size);
39
40
8
        _initialize_hash_vals(rows);
41
8
        auto* __restrict hashes = _hash_vals.data();
42
8
        RETURN_IF_ERROR(_get_partition_column_result(block, result));
43
16
        for (int j = 0; j < result_size; ++j) {
44
8
            const auto& [col, is_const] = unpack_if_const(block->get_by_position(result[j]).column);
45
8
            if (is_const) {
46
0
                continue;
47
0
            }
48
8
            _do_hash(col, hashes, j);
49
8
        }
50
51
28
        for (size_t i = 0; i < rows; i++) {
52
20
            hashes[i] = ChannelIds()(hashes[i], _partition_count);
53
20
        }
54
55
8
        Block::erase_useless_column(block, column_to_keep);
56
8
    }
57
8
    return Status::OK();
58
8
}
_ZNK5doris20Crc32HashPartitionerINS_15ShiftChannelIdsEE15do_partitioningEPNS_12RuntimeStateEPNS_5BlockE
Line
Count
Source
31
1
Status Crc32HashPartitioner<ChannelIds>::do_partitioning(RuntimeState* state, Block* block) const {
32
1
    size_t rows = block->rows();
33
34
1
    if (rows > 0) {
35
1
        auto column_to_keep = block->columns();
36
37
1
        int result_size = cast_set<int>(_partition_expr_ctxs.size());
38
1
        std::vector<int> result(result_size);
39
40
1
        _initialize_hash_vals(rows);
41
1
        auto* __restrict hashes = _hash_vals.data();
42
1
        RETURN_IF_ERROR(_get_partition_column_result(block, result));
43
2
        for (int j = 0; j < result_size; ++j) {
44
1
            const auto& [col, is_const] = unpack_if_const(block->get_by_position(result[j]).column);
45
1
            if (is_const) {
46
0
                continue;
47
0
            }
48
1
            _do_hash(col, hashes, j);
49
1
        }
50
51
6
        for (size_t i = 0; i < rows; i++) {
52
5
            hashes[i] = ChannelIds()(hashes[i], _partition_count);
53
5
        }
54
55
1
        Block::erase_useless_column(block, column_to_keep);
56
1
    }
57
1
    return Status::OK();
58
1
}
59
60
template <typename ChannelIds>
61
void Crc32HashPartitioner<ChannelIds>::_do_hash(const ColumnPtr& column,
62
36
                                                HashValType* __restrict result, int idx) const {
63
36
    column->update_crcs_with_value(
64
36
            result, _partition_expr_ctxs[idx]->root()->data_type()->get_primitive_type(),
65
36
            cast_set<HashValType>(column->size()));
66
36
}
_ZNK5doris20Crc32HashPartitionerINS_17ShuffleChannelIdsEE8_do_hashERKNS_3COWINS_7IColumnEE13immutable_ptrIS4_EEPji
Line
Count
Source
62
22
                                                HashValType* __restrict result, int idx) const {
63
22
    column->update_crcs_with_value(
64
22
            result, _partition_expr_ctxs[idx]->root()->data_type()->get_primitive_type(),
65
22
            cast_set<HashValType>(column->size()));
66
22
}
_ZNK5doris20Crc32HashPartitionerINS_24SpillPartitionChannelIdsEE8_do_hashERKNS_3COWINS_7IColumnEE13immutable_ptrIS4_EEPji
Line
Count
Source
62
6
                                                HashValType* __restrict result, int idx) const {
63
6
    column->update_crcs_with_value(
64
6
            result, _partition_expr_ctxs[idx]->root()->data_type()->get_primitive_type(),
65
6
            cast_set<HashValType>(column->size()));
66
6
}
_ZNK5doris20Crc32HashPartitionerINS_26SpillRePartitionChannelIdsEE8_do_hashERKNS_3COWINS_7IColumnEE13immutable_ptrIS4_EEPji
Line
Count
Source
62
8
                                                HashValType* __restrict result, int idx) const {
63
8
    column->update_crcs_with_value(
64
8
            result, _partition_expr_ctxs[idx]->root()->data_type()->get_primitive_type(),
65
8
            cast_set<HashValType>(column->size()));
66
8
}
Unexecuted instantiation: _ZNK5doris20Crc32HashPartitionerINS_15ShiftChannelIdsEE8_do_hashERKNS_3COWINS_7IColumnEE13immutable_ptrIS4_EEPji
67
68
template <typename ChannelIds>
69
Status Crc32HashPartitioner<ChannelIds>::clone(RuntimeState* state,
70
9
                                               std::unique_ptr<PartitionerBase>& partitioner) {
71
9
    auto* new_partitioner = new Crc32HashPartitioner<ChannelIds>(_partition_count);
72
9
    partitioner.reset(new_partitioner);
73
9
    return _clone_expr_ctxs(state, new_partitioner->_partition_expr_ctxs);
74
9
}
Unexecuted instantiation: _ZN5doris20Crc32HashPartitionerINS_17ShuffleChannelIdsEE5cloneEPNS_12RuntimeStateERSt10unique_ptrINS_15PartitionerBaseESt14default_deleteIS6_EE
_ZN5doris20Crc32HashPartitionerINS_24SpillPartitionChannelIdsEE5cloneEPNS_12RuntimeStateERSt10unique_ptrINS_15PartitionerBaseESt14default_deleteIS6_EE
Line
Count
Source
70
5
                                               std::unique_ptr<PartitionerBase>& partitioner) {
71
5
    auto* new_partitioner = new Crc32HashPartitioner<ChannelIds>(_partition_count);
72
5
    partitioner.reset(new_partitioner);
73
5
    return _clone_expr_ctxs(state, new_partitioner->_partition_expr_ctxs);
74
5
}
_ZN5doris20Crc32HashPartitionerINS_26SpillRePartitionChannelIdsEE5cloneEPNS_12RuntimeStateERSt10unique_ptrINS_15PartitionerBaseESt14default_deleteIS6_EE
Line
Count
Source
70
4
                                               std::unique_ptr<PartitionerBase>& partitioner) {
71
4
    auto* new_partitioner = new Crc32HashPartitioner<ChannelIds>(_partition_count);
72
4
    partitioner.reset(new_partitioner);
73
4
    return _clone_expr_ctxs(state, new_partitioner->_partition_expr_ctxs);
74
4
}
Unexecuted instantiation: _ZN5doris20Crc32HashPartitionerINS_15ShiftChannelIdsEE5cloneEPNS_12RuntimeStateERSt10unique_ptrINS_15PartitionerBaseESt14default_deleteIS6_EE
75
76
void Crc32CHashPartitioner::_do_hash(const ColumnPtr& column, HashValType* __restrict result,
77
1
                                     int idx) const {
78
1
    column->update_crc32c_batch(result, nullptr);
79
1
}
80
81
Status Crc32CHashPartitioner::clone(RuntimeState* state,
82
0
                                    std::unique_ptr<PartitionerBase>& partitioner) {
83
0
    auto* new_partitioner = new Crc32CHashPartitioner(_partition_count);
84
0
    partitioner.reset(new_partitioner);
85
0
    return _clone_expr_ctxs(state, new_partitioner->_partition_expr_ctxs);
86
0
}
87
88
HashPartitionFunction::HashPartitionFunction(HashValType partition_count,
89
                                             ShuffleHashMethod hash_method)
90
2
        : _partition_count(partition_count), _hash_method(hash_method) {}
91
92
2
Status HashPartitionFunction::init(const std::vector<TExpr>& texprs) {
93
2
    if (_hash_method == ShuffleHashMethod::CRC32C) {
94
1
        _partitioner = std::make_unique<Crc32CHashPartitioner>(_partition_count);
95
1
    } else {
96
1
        _partitioner = std::make_unique<Crc32HashPartitioner<ShuffleChannelIds>>(_partition_count);
97
1
    }
98
2
    return _partitioner->init(texprs);
99
2
}
100
101
2
Status HashPartitionFunction::prepare(RuntimeState* state, const RowDescriptor& row_desc) {
102
2
    return _partitioner->prepare(state, row_desc);
103
2
}
104
105
2
Status HashPartitionFunction::open(RuntimeState* state) {
106
2
    return _partitioner->open(state);
107
2
}
108
109
2
Status HashPartitionFunction::close(RuntimeState* state) {
110
2
    return _partitioner->close(state);
111
2
}
112
113
Status HashPartitionFunction::get_partitions(RuntimeState* state, Block* block,
114
                                             size_t partition_count,
115
2
                                             std::vector<HashValType>& partitions) const {
116
2
    if (partition_count != _partition_count) {
117
0
        return Status::InvalidArgument("Hash partition count {} does not match planned count {}",
118
0
                                       partition_count, _partition_count);
119
0
    }
120
2
    RETURN_IF_ERROR(_partitioner->do_partitioning(state, block));
121
2
    partitions = _partitioner->get_channel_ids();
122
2
    return Status::OK();
123
2
}
124
125
Status HashPartitionFunction::clone(RuntimeState* state,
126
0
                                    std::unique_ptr<PartitionFunction>& function) const {
127
0
    auto cloned = std::make_unique<HashPartitionFunction>(_partition_count, _hash_method);
128
0
    RETURN_IF_ERROR(_partitioner->clone(state, cloned->_partitioner));
129
0
    function = std::move(cloned);
130
0
    return Status::OK();
131
0
}
132
133
template class Crc32HashPartitioner<ShuffleChannelIds>;
134
template class Crc32HashPartitioner<SpillPartitionChannelIds>;
135
template class Crc32HashPartitioner<SpillRePartitionChannelIds>;
136
137
} // namespace doris