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 |