Coverage Report

Created: 2026-09-03 01:30

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/load/memtable/memtable_flush_executor.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,
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
#pragma once
19
20
#include <atomic>
21
#include <condition_variable>
22
#include <cstdint>
23
#include <iosfwd>
24
#include <memory>
25
#include <utility>
26
#include <vector>
27
28
#include "common/status.h"
29
#include "load/delta_writer/delta_writer_context.h"
30
#include "load/memtable/memtable.h"
31
#include "util/threadpool.h"
32
33
namespace doris {
34
35
class DataDir;
36
class MemTable;
37
class MemTableMemoryLimiter;
38
class Block;
39
class GroupRowsetWriter;
40
class OlapTableSchemaParam;
41
struct RowsetWriterContext;
42
class RowsetWriter;
43
class SystemMetrics;
44
class WorkloadGroup;
45
46
// the statistic of a certain flush handler.
47
// use atomic because it may be updated by multi threads
48
struct FlushStatistic {
49
    std::atomic_uint64_t flush_time_ns = 0;
50
    std::atomic_uint64_t flush_submit_count = 0;
51
    std::atomic_int64_t flush_running_count = 0;
52
    std::atomic_uint64_t flush_finish_count = 0;
53
    std::atomic_uint64_t flush_size_bytes = 0;
54
    std::atomic_uint64_t flush_disk_size_bytes = 0;
55
    std::atomic_uint64_t flush_wait_time_ns = 0;
56
};
57
58
struct SharedMemtable {
59
    std::shared_ptr<MemTable> memtable;
60
    int32_t segment_id = 0;
61
62
    ~SharedMemtable();
63
64
    std::once_flag block_once;
65
    Status block_status;
66
    std::shared_ptr<Block> block;
67
    // The group rowset writer whose context owns the allocated-lsn map. Held weakly: the
68
    // flush tasks only hold the owning FlushToken by a weak_ptr, so the last reference may
69
    // die (destroying the writer and the map) while this memtable is still alive. The
70
    // destructor must therefore lock before touching the context and skip the cleanup when
71
    // the writer is already gone (the map died with it and has nothing left to remove).
72
    std::weak_ptr<RowsetWriter> rowset_writer;
73
    bool has_allocated_lsns = false;
74
75
    std::atomic<int> finished_sub_task_count {0};
76
    // data + binlog
77
    std::atomic<int> total_sub_task_count {2};
78
79
4
    int add_finished_sub_task() { return finished_sub_task_count.fetch_add(1); }
80
81
0
    std::string debug_string() const {
82
0
        return "PartOfGroupMemtableFlushTask{segment_id=" + std::to_string(segment_id) +
83
0
               ", finished_sub_task_count=" + std::to_string(finished_sub_task_count.load()) +
84
0
               ", total_sub_task_count=" + std::to_string(total_sub_task_count.load()) + "}";
85
0
    }
86
};
87
88
std::ostream& operator<<(std::ostream& os, const FlushStatistic& stat);
89
90
// A thin wrapper of ThreadPoolToken to submit task.
91
// For a tablet, there may be multiple memtables, which will be flushed to disk
92
// one by one in the order of generation.
93
// If a memtable flush fails, then:
94
// 1. Immediately disallow submission of any subsequent memtable
95
// 2. For the memtables that have already been submitted, there is no need to flush,
96
//    because the entire job will definitely fail;
97
class FlushToken : public std::enable_shared_from_this<FlushToken> {
98
    ENABLE_FACTORY_CREATOR(FlushToken);
99
100
public:
101
    FlushToken(ThreadPool* thread_pool, std::shared_ptr<WorkloadGroup> wg_sptr)
102
19
            : _flush_status(Status::OK()), _thread_pool(thread_pool), _wg_wptr(wg_sptr) {}
103
104
    Status submit(std::shared_ptr<MemTable> mem_table);
105
106
    // error has happens, so we cancel this token
107
    // And remove all tasks in the queue.
108
    void cancel();
109
110
    // wait all tasks in token to be completed.
111
    Status wait();
112
113
    // get flush operations' statistics
114
42
    const FlushStatistic& get_stats() const { return _stats; }
115
116
19
    void set_rowset_writer(std::shared_ptr<RowsetWriter> rowset_writer) {
117
19
        _rowset_writer = rowset_writer;
118
19
    }
119
120
19
    void set_table_schema_param(std::shared_ptr<OlapTableSchemaParam> table_schema_param) {
121
19
        _table_schema_param = std::move(table_schema_param);
122
19
    }
123
124
30
    const MemTableStat& memtable_stat() { return _memtable_stat; }
125
126
private:
127
34
    void _shutdown_flush_token() { _shutdown.store(true); }
128
36
    bool _is_shutdown() { return _shutdown.load(); }
129
    void _wait_submit_task_finish();
130
    void _wait_running_task_finish();
131
132
private:
133
    friend class MemtableFlushTask;
134
    friend class PartOfGroupMemtableFlushTask;
135
136
    Status _submit_sub_tasks(ThreadPool* pool, std::vector<std::shared_ptr<Runnable>> sub_tasks);
137
138
    void _flush_memtable_impl(RowsetWriter* flush_writer, MemTable* memtable, int32_t segment_id,
139
                              int64_t submit_task_time, SharedMemtable* shared_memtable = nullptr);
140
141
    void _flush_memtable(std::shared_ptr<MemTable> memtable_ptr, int32_t segment_id,
142
                         int64_t submit_task_time);
143
144
    void _flush_group_memtable(std::shared_ptr<SharedMemtable> shared_memtable,
145
                               WriteRequestType write_req_type, int64_t submit_task_time);
146
147
    Status _memtable2block(MemTable* memtable, SharedMemtable* shared_memtable,
148
                           std::shared_ptr<Block>& flush_block);
149
150
    Status _try_reserve_memory(const std::shared_ptr<ResourceContext>& resource_context,
151
                               int64_t size);
152
153
    // Records the current flush status of the tablet.
154
    // Note: Once its value is set to Failed, it cannot return to SUCCESS.
155
    std::shared_mutex _flush_status_lock;
156
    Status _flush_status;
157
158
    FlushStatistic _stats;
159
160
    std::shared_ptr<RowsetWriter> _rowset_writer = nullptr;
161
162
    std::shared_ptr<OlapTableSchemaParam> _table_schema_param = nullptr;
163
164
    MemTableStat _memtable_stat;
165
166
    std::atomic<bool> _shutdown = false;
167
    ThreadPool* _thread_pool = nullptr;
168
169
    std::mutex _mutex;
170
    std::condition_variable _submit_task_finish_cond;
171
    std::condition_variable _running_task_finish_cond;
172
173
    std::weak_ptr<WorkloadGroup> _wg_wptr;
174
};
175
176
// MemTableFlushExecutor is responsible for flushing memtables to disk.
177
// It encapsulate a ThreadPool to handle all tasks.
178
// Usage Example:
179
//      ...
180
//      std::shared_ptr<FlushHandler> flush_handler;
181
//      memTableFlushExecutor.create_flush_token(&flush_handler);
182
//      ...
183
//      flush_token->submit(memtable)
184
//      ...
185
class MemTableFlushExecutor {
186
public:
187
165
    MemTableFlushExecutor() = default;
188
165
    ~MemTableFlushExecutor() {
189
165
        _flush_pool->shutdown();
190
165
        _high_prio_flush_pool->shutdown();
191
165
    }
192
193
    // init should be called after storage engine is opened,
194
    // because it needs path hash of each data dir.
195
    void init(int num_disk);
196
197
    Status create_flush_token(std::shared_ptr<FlushToken>& flush_token,
198
                              std::shared_ptr<RowsetWriter> rowset_writer, bool is_high_priority,
199
                              std::shared_ptr<WorkloadGroup> wg_sptr,
200
                              std::shared_ptr<OlapTableSchemaParam> table_schema_param = nullptr);
201
202
    // return true if it already has any flushing task
203
0
    bool check_and_inc_has_any_flushing_task() {
204
        // need to use CAS instead of only `if (0 == _flushing_task_count)` statement,
205
        // to avoid concurrent entries both pass the if statement
206
0
        int expected_count = 0;
207
0
        if (!_flushing_task_count.compare_exchange_strong(expected_count, 1)) {
208
0
            return true;
209
0
        }
210
0
        DCHECK(expected_count == 0 && _flushing_task_count == 1);
211
0
        return false;
212
0
    }
213
214
0
    void inc_flushing_task() { _flushing_task_count++; }
215
216
0
    void dec_flushing_task() { _flushing_task_count--; }
217
218
10
    ThreadPool* flush_pool() { return _flush_pool.get(); }
219
220
0
    ThreadPool* high_prio_flush_pool() { return _high_prio_flush_pool.get(); }
221
222
    void update_memtable_flush_threads();
223
224
    // Returns {min_threads, max_threads} for a flush thread pool.
225
    // thread_num_per_store is used as the baseline when adaptive mode is off.
226
    static std::pair<int, int> calc_flush_thread_count(int num_cpus, int num_disk,
227
                                                       int thread_num_per_store);
228
229
private:
230
    std::unique_ptr<ThreadPool> _flush_pool;
231
    std::unique_ptr<ThreadPool> _high_prio_flush_pool;
232
    std::atomic<int> _flushing_task_count = 0;
233
    int _num_disk = 0;
234
};
235
236
} // namespace doris