Coverage Report

Created: 2026-09-28 19:52

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/exec/partitioner/writer_assigner.h
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, software
12
// distributed under the License is distributed on an "AS IS" BASIS,
13
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14
// See the License for the specific language governing permissions and
15
// limitations under the License.
16
17
#pragma once
18
19
#include <cstddef>
20
#include <cstdint>
21
#include <memory>
22
#include <vector>
23
24
#include "common/status.h"
25
26
namespace doris {
27
class SkewedPartitionRebalancer;
28
}
29
30
namespace doris {
31
32
// Maps logical partitions computed by a PartitionFunction to Doris exchange channels.
33
class WriterAssigner {
34
public:
35
2
    virtual ~WriterAssigner() = default;
36
37
    virtual Status assign(const std::vector<uint32_t>& partition_ids,
38
                          const std::vector<uint8_t>* mask, size_t rows, size_t block_bytes,
39
                          std::vector<uint32_t>& writer_ids) = 0;
40
};
41
42
// Preserves stable ownership: one logical partition always maps to one writer id.
43
class IdentityWriterAssigner final : public WriterAssigner {
44
public:
45
2
    explicit IdentityWriterAssigner(uint32_t writer_count) : _writer_count(writer_count) {}
46
47
    Status assign(const std::vector<uint32_t>& partition_ids, const std::vector<uint8_t>* mask,
48
                  size_t rows, size_t block_bytes, std::vector<uint32_t>& writer_ids) override;
49
50
private:
51
    uint32_t _writer_count;
52
};
53
54
// Allows a hot logical partition to use multiple writers while retaining the existing
55
// ScaleWriter affinity and rebalance behavior.
56
class SkewedWriterAssigner final : public WriterAssigner {
57
public:
58
    SkewedWriterAssigner(int partition_count, int task_count, int task_bucket_count,
59
                         long min_partition_data_processed_rebalance_threshold,
60
                         long min_data_processed_rebalance_threshold);
61
62
    ~SkewedWriterAssigner() override;
63
64
    Status assign(const std::vector<uint32_t>& partition_ids, const std::vector<uint8_t>* mask,
65
                  size_t rows, size_t block_bytes, std::vector<uint32_t>& writer_ids) override;
66
67
private:
68
    int _get_next_writer_id(uint32_t partition_id);
69
70
    std::unique_ptr<SkewedPartitionRebalancer> _rebalancer;
71
    int _writer_count;
72
    std::vector<int> _partition_row_counts;
73
    std::vector<int> _partition_writer_ids;
74
    std::vector<int> _partition_writer_indexes;
75
};
76
77
// Scale table-sink thresholds by local pipeline task count while preserving the historical
78
// behavior for very small values.
79
int64_t scale_writer_threshold_by_task(int64_t value, int task_num);
80
81
} // namespace doris