Skip to main content

mito2/worker/
handle_flush.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Handling flush related requests.
16
17use std::collections::HashSet;
18use std::sync::Arc;
19use std::sync::atomic::Ordering;
20
21use common_base::readable_size::ReadableSize;
22use common_telemetry::{debug, error, info};
23use store_api::logstore::LogStore;
24use store_api::region_request::{RegionFlushReason, RegionFlushRequest};
25use store_api::storage::RegionId;
26
27use crate::config::{IndexBuildMode, MitoConfig};
28use crate::error::{RegionNotFoundSnafu, Result};
29use crate::flush::{FlushReason, RegionFlushTask};
30use crate::region::MitoRegionRef;
31use crate::region::version::VersionRef;
32use crate::request::{BuildIndexRequest, FlushFailed, FlushFinished, OnFailure, OptionOutputTx};
33use crate::sst::index::IndexBuildType;
34use crate::worker::RegionWorkerLoop;
35
36#[derive(Debug, Default)]
37pub(crate) struct RegionWriteBufferPressure {
38    pub(crate) stalled_region_ids: HashSet<RegionId>,
39    pub(crate) rejected_region_ids: HashSet<RegionId>,
40}
41
42#[derive(Debug, Default, PartialEq, Eq)]
43struct RegionWriteBufferStatus {
44    should_flush: bool,
45    should_stall: bool,
46    should_reject: bool,
47}
48
49fn resolve_flush_reason(
50    request_reason: Option<RegionFlushReason>,
51    is_downgrading: bool,
52) -> FlushReason {
53    match request_reason {
54        Some(reason) => FlushReason::from(reason),
55        None if is_downgrading => FlushReason::Downgrading,
56        None => FlushReason::Manual,
57    }
58}
59
60impl<S: LogStore> RegionWorkerLoop<S> {
61    /// On region flush job failed.
62    pub(crate) async fn handle_flush_failed(&mut self, region_id: RegionId, request: FlushFailed) {
63        let mut pending_ddls = self.flush_scheduler.on_flush_failed(region_id, request.err);
64        if !pending_ddls.is_empty() {
65            info!(
66                "Flush terminated for region {}, handling {} pending lifecycle DDL requests",
67                region_id,
68                pending_ddls.len()
69            );
70            self.handle_ddl_requests(&mut pending_ddls).await;
71        }
72        // Maybe flush worker again.
73        self.maybe_flush_worker();
74
75        // Handle stalled requests.
76        self.handle_stalled_requests().await;
77    }
78
79    /// Checks whether the engine reaches flush threshold. If so, finds regions in this
80    /// worker to flush.
81    pub(crate) fn maybe_flush_worker(&mut self) {
82        if !self.write_buffer_manager.should_flush_engine() {
83            debug!("No need to flush worker");
84            // No need to flush worker.
85            return;
86        }
87
88        // If the engine needs flush, each worker will find some regions to flush. We might
89        // flush more memory than expect but it should be acceptable.
90        if let Err(e) = self.flush_regions_on_engine_full() {
91            error!(e; "Failed to flush worker");
92        }
93    }
94
95    /// Finds some regions to flush to reduce write buffer usage.
96    fn flush_regions_on_engine_full(&mut self) -> Result<()> {
97        let regions = self.regions.list_regions();
98        let now = self.time_provider.current_time_millis();
99        let mut pending_regions = vec![];
100
101        for region in &regions {
102            if self.flush_scheduler.is_flush_requested(region.region_id) || !region.is_writable() {
103                // Already flushing or not writable.
104                continue;
105            }
106
107            let version = region.version();
108            let region_memtable_size = region_memtable_usage(&version);
109
110            let auto_flush_interval = version
111                .options
112                .auto_flush_interval_or(self.config.auto_flush_interval);
113            let min_last_flush_time = now.saturating_sub(
114                auto_flush_interval
115                    .as_millis()
116                    .try_into()
117                    .unwrap_or(i64::MAX),
118            );
119
120            if region.last_flush_millis() < min_last_flush_time {
121                // If flush time of this region is earlier than `min_last_flush_time`, we can flush this region.
122                let task =
123                    self.new_flush_task(region, FlushReason::EngineFull, None, self.config.clone());
124                self.flush_scheduler.schedule_flush(
125                    region.region_id,
126                    &region.version_control,
127                    task,
128                )?;
129            } else if region_memtable_size > 0 {
130                // We should only consider regions with memtable size > 0 to flush.
131                pending_regions.push((region, region_memtable_size));
132            }
133        }
134        pending_regions.sort_unstable_by_key(|(_, size)| std::cmp::Reverse(*size));
135        // The flush target is the mutable memtable limit (half of the global buffer).
136        // When memory is full, we aggressively flush regions until usage drops below this target,
137        // not just below the full limit.
138        let target_memory_usage = self.write_buffer_manager.flush_limit();
139        let mut memory_usage = self.write_buffer_manager.memory_usage();
140
141        #[cfg(test)]
142        {
143            debug!(
144                "Flushing regions on engine full, target memory usage: {}, memory usage: {}, pending regions: {:?}",
145                target_memory_usage,
146                memory_usage,
147                pending_regions
148                    .iter()
149                    .map(|(region, mem_size)| (region.region_id, mem_size))
150                    .collect::<Vec<_>>()
151            );
152        }
153        // Iterate over pending regions in descending order of their memory size and schedule flush tasks
154        // for each region until the overall memory usage drops below the flush limit.
155        for (region, region_mem_size) in pending_regions.into_iter() {
156            // Make sure the first region is always flushed.
157            if memory_usage < target_memory_usage {
158                // Stop flushing regions if memory usage is already below the flush limit
159                break;
160            }
161            let task =
162                self.new_flush_task(region, FlushReason::EngineFull, None, self.config.clone());
163            debug!("Scheduling flush task for region {}", region.region_id);
164            // Schedule a flush task for the current region
165            self.flush_scheduler
166                .schedule_flush(region.region_id, &region.version_control, task)?;
167            // Reduce memory usage by the region's size, ensuring it doesn't go negative
168            memory_usage = memory_usage.saturating_sub(region_mem_size);
169        }
170
171        Ok(())
172    }
173
174    /// Flushes write regions that exceed their flush threshold and returns region pressure.
175    pub(crate) fn maybe_flush_write_regions(
176        &mut self,
177        region_ids: HashSet<RegionId>,
178    ) -> RegionWriteBufferPressure {
179        let mut pressure = RegionWriteBufferPressure::default();
180
181        for region_id in region_ids {
182            let Some(region) = self.regions.get_region(region_id) else {
183                continue;
184            };
185            if !region.is_writable() {
186                continue;
187            }
188
189            let status = region_write_buffer_status(
190                &region.version(),
191                self.config.default_region_write_buffer_size,
192                self.stalled_requests.estimated_size(&region_id),
193            );
194
195            if status.should_flush && !self.flush_scheduler.is_flush_requested(region.region_id) {
196                let task = self.new_flush_task(
197                    &region,
198                    FlushReason::RegionFull,
199                    None,
200                    self.config.clone(),
201                );
202                if let Err(e) = self.flush_scheduler.schedule_flush(
203                    region.region_id,
204                    &region.version_control,
205                    task,
206                ) {
207                    error!(e; "Failed to schedule flush task for region {}", region.region_id);
208                }
209            }
210
211            if status.should_reject {
212                pressure.rejected_region_ids.insert(region_id);
213            } else if status.should_stall {
214                pressure.stalled_region_ids.insert(region_id);
215            }
216        }
217
218        pressure
219    }
220
221    /// Creates a flush task with specific `reason` for the `region`.
222    pub(crate) fn new_flush_task(
223        &self,
224        region: &MitoRegionRef,
225        reason: FlushReason,
226        row_group_size: Option<usize>,
227        engine_config: Arc<MitoConfig>,
228    ) -> RegionFlushTask {
229        RegionFlushTask {
230            region_id: region.region_id,
231            reason,
232            senders: Vec::new(),
233            request_sender: self.sender.clone(),
234            access_layer: region.access_layer.clone(),
235            listener: self.listener.clone(),
236            engine_config,
237            // The request-provided size (e.g. manual flush) takes precedence over the
238            // region option.
239            row_group_size: row_group_size.or(region.version().options.max_row_group_row_count),
240            cache_manager: self.cache_manager.clone(),
241            manifest_ctx: region.manifest_ctx.clone(),
242            index_options: region.version().options.index_options.clone(),
243            flush_semaphore: self.flush_semaphore.clone(),
244            is_staging: region.is_staging(),
245            partition_expr: region.maybe_staging_partition_expr_str(),
246        }
247    }
248}
249
250pub(crate) fn region_memtable_usage(version: &VersionRef) -> usize {
251    version.memtables.mutable_usage() + version.memtables.immutables_usage()
252}
253
254pub(crate) fn region_write_buffer_size(
255    version: &VersionRef,
256    default_region_write_buffer_size: ReadableSize,
257) -> Option<usize> {
258    let size = version
259        .options
260        .write_buffer_size
261        .unwrap_or(default_region_write_buffer_size);
262    (size.as_bytes() > 0).then_some(size.as_bytes() as usize)
263}
264
265/// Returns the pressure state of a region's write buffer.
266fn region_write_buffer_status(
267    version: &VersionRef,
268    default_region_write_buffer_size: ReadableSize,
269    stalled_request_size: usize,
270) -> RegionWriteBufferStatus {
271    let Some(write_buffer_size) =
272        region_write_buffer_size(version, default_region_write_buffer_size)
273    else {
274        return RegionWriteBufferStatus::default();
275    };
276
277    let mutable_usage = version.memtables.mutable_usage();
278    let memory_usage = region_memtable_usage(version);
279    region_write_buffer_status_from_usage(
280        write_buffer_size,
281        mutable_usage,
282        memory_usage,
283        stalled_request_size,
284    )
285}
286
287fn region_write_buffer_status_from_usage(
288    write_buffer_size: usize,
289    mutable_usage: usize,
290    memory_usage: usize,
291    stalled_request_size: usize,
292) -> RegionWriteBufferStatus {
293    let mutable_limit = std::cmp::max(1, write_buffer_size / 2);
294    let should_stall = memory_usage >= write_buffer_size;
295    let reject_limit = write_buffer_size.saturating_mul(2);
296    let should_reject = memory_usage.saturating_add(stalled_request_size) >= reject_limit;
297
298    RegionWriteBufferStatus {
299        should_flush: mutable_usage >= mutable_limit || should_stall,
300        should_stall,
301        should_reject,
302    }
303}
304
305impl<S: LogStore> RegionWorkerLoop<S> {
306    /// Handles manual flush request.
307    pub(crate) fn handle_flush_request(
308        &mut self,
309        region_id: RegionId,
310        request: RegionFlushRequest,
311        sender: OptionOutputTx,
312    ) {
313        let region = match self.regions.flushable_region(region_id) {
314            Ok(region) => region,
315            Err(e) => {
316                sender.send(Err(e));
317                return;
318            }
319        };
320
321        // `update_topic_latest_entry_id` updates `topic_latest_entry_id` when memtables are empty.
322        // But the flush is skipped if memtables are empty. Thus should update the `topic_latest_entry_id`
323        // when handling flush request instead of in `schedule_flush` or `flush_finished`.
324        self.update_topic_latest_entry_id(&region);
325
326        let reason = resolve_flush_reason(request.reason, region.is_downgrading());
327        let mut task =
328            self.new_flush_task(&region, reason, request.row_group_size, self.config.clone());
329        task.push_sender(sender);
330        if let Err(e) =
331            self.flush_scheduler
332                .schedule_flush(region.region_id, &region.version_control, task)
333        {
334            error!(e; "Failed to schedule flush task for region {}", region.region_id);
335        }
336    }
337
338    /// Flushes regions periodically.
339    pub(crate) fn flush_periodically(&mut self) -> Result<()> {
340        let regions = self.regions.list_regions();
341        let now = self.time_provider.current_time_millis();
342
343        for region in &regions {
344            if self.flush_scheduler.is_flush_requested(region.region_id) || !region.is_writable() {
345                // Already flushing or not writable.
346                continue;
347            }
348            self.update_topic_latest_entry_id(region);
349
350            let auto_flush_interval = region
351                .version()
352                .options
353                .auto_flush_interval_or(self.config.auto_flush_interval);
354            let min_last_flush_time = now.saturating_sub(
355                auto_flush_interval
356                    .as_millis()
357                    .try_into()
358                    .unwrap_or(i64::MAX),
359            );
360
361            if region.last_flush_millis() < min_last_flush_time {
362                // If flush time of this region is earlier than `min_last_flush_time`, we can flush this region.
363                let task = self.new_flush_task(
364                    region,
365                    FlushReason::Periodically,
366                    None,
367                    self.config.clone(),
368                );
369                self.flush_scheduler.schedule_flush(
370                    region.region_id,
371                    &region.version_control,
372                    task,
373                )?;
374            }
375        }
376
377        Ok(())
378    }
379
380    /// On region flush job finished.
381    pub(crate) async fn handle_flush_finished(
382        &mut self,
383        region_id: RegionId,
384        mut request: FlushFinished,
385    ) {
386        // Notifies other workers. Even the remaining steps of this method fail we still
387        // wake up other workers as we have released some memory by flush.
388        self.notify_group();
389
390        let region = match self.regions.get_region(region_id) {
391            Some(region) => region,
392            None => {
393                request.on_failure(RegionNotFoundSnafu { region_id }.build());
394                return;
395            }
396        };
397
398        if request.is_staging {
399            // Skip the region metadata update.
400            info!(
401                "Skipping region metadata update for region {} in staging mode",
402                region_id
403            );
404            region.version_control.apply_edit(
405                None,
406                &request.memtables_to_remove,
407                region.file_purger.clone(),
408            );
409        } else {
410            region.version_control.apply_edit(
411                Some(request.edit.clone()),
412                &request.memtables_to_remove,
413                region.file_purger.clone(),
414            );
415        }
416
417        region.update_flush_millis();
418        // Update topic latest entry id as soon as possible after flush to make sure the prunable entry id is updated timely,
419        // which is important for remote WAL pruning.
420        self.update_topic_latest_entry_id(&region);
421
422        // Delete wal.
423        info!(
424            "Region {} flush finished, tries to bump wal to {}",
425            region_id, request.flushed_entry_id
426        );
427        if let Err(e) = self
428            .wal
429            .obsolete(region_id, request.flushed_entry_id, &region.provider)
430            .await
431        {
432            error!(e; "Failed to write wal, region: {}", region_id);
433            request.on_failure(e);
434            return;
435        }
436
437        let flush_on_close = request.flush_reason == FlushReason::Closing;
438
439        let index_build_file_metas = std::mem::take(&mut request.edit.files_to_add);
440
441        // In async mode, create indexes after flush.
442        if self.config.index.build_mode == IndexBuildMode::Async {
443            self.handle_rebuild_index(
444                BuildIndexRequest {
445                    region_id,
446                    build_type: IndexBuildType::Flush,
447                    file_metas: index_build_file_metas,
448                },
449                OptionOutputTx::new(None),
450            )
451            .await;
452        }
453
454        if flush_on_close {
455            self.remove_region(region_id).await;
456            info!("Region {} closed after flush", region_id);
457            request.on_success();
458            self.listener.on_flush_success(region_id);
459            return;
460        }
461
462        // Notifies waiters and observes the flush timer.
463        request.on_success();
464        // Handle pending requests for the region.
465        if let Some((mut ddl_requests, mut write_requests, mut bulk_writes)) =
466            self.flush_scheduler.on_flush_success(region_id)
467        {
468            // Perform DDLs first because they require empty memtables.
469            self.handle_ddl_requests(&mut ddl_requests).await;
470            if self.flush_scheduler.is_flush_requested(region_id) {
471                // The DDL may schedule another flush, e.g. a close-time flush after writes
472                // arrived in the mutable memtable during the previous flush. Keep pending
473                // writes fenced until that flush reaches its terminal state instead of
474                // accepting them while the DDL is still in progress.
475                for write_request in write_requests {
476                    self.flush_scheduler
477                        .add_write_request_to_pending(write_request);
478                }
479                for bulk_write in bulk_writes {
480                    self.flush_scheduler.add_bulk_request_to_pending(bulk_write);
481                }
482                self.listener.on_flush_success(region_id);
483                return;
484            }
485            // A pending close DDL may have removed the region. Reject queued writes as
486            // not found, then stop instead of scheduling compaction for a closed region.
487            if !self.regions.is_region_exists(region_id) {
488                self.handle_write_requests(&mut write_requests, &mut bulk_writes, false)
489                    .await;
490                self.listener.on_flush_success(region_id);
491                return;
492            }
493            // Handle pending write requests, we don't stall these requests.
494            self.handle_write_requests(&mut write_requests, &mut bulk_writes, false)
495                .await;
496        }
497        // Maybe flush worker again.
498        self.maybe_flush_worker();
499        // Handle stalled requests.
500        self.handle_stalled_requests().await;
501        // Schedules compaction.
502        self.schedule_compaction(&region).await;
503
504        self.listener.on_flush_success(region_id);
505    }
506
507    /// Updates the latest entry id since flush of the region.
508    /// **This is only used for remote WAL pruning.**
509    pub(crate) fn update_topic_latest_entry_id(&mut self, region: &MitoRegionRef) {
510        if region.provider.is_remote_wal() && region.version().memtables.is_empty() {
511            let latest_offset = self
512                .wal
513                .store()
514                .latest_entry_id(&region.provider)
515                .unwrap_or(0);
516            let topic_last_entry_id = region.topic_latest_entry_id.load(Ordering::Relaxed);
517
518            if latest_offset > topic_last_entry_id {
519                region
520                    .topic_latest_entry_id
521                    .store(latest_offset, Ordering::Relaxed);
522                debug!(
523                    "Region {} latest entry id updated to {}",
524                    region.region_id, latest_offset
525                );
526            }
527        }
528    }
529}
530
531#[cfg(test)]
532mod tests {
533    use super::*;
534
535    #[test]
536    fn test_resolve_flush_reason_uses_request_reason() {
537        assert_eq!(
538            resolve_flush_reason(Some(RegionFlushReason::RegionMigration), true),
539            FlushReason::RegionMigration
540        );
541        assert_eq!(
542            resolve_flush_reason(Some(RegionFlushReason::Repartition), false),
543            FlushReason::Repartition
544        );
545        assert_eq!(
546            resolve_flush_reason(Some(RegionFlushReason::RemoteWalPrune), false),
547            FlushReason::RemoteWalPrune
548        );
549        assert_eq!(
550            resolve_flush_reason(Some(RegionFlushReason::Closing), false),
551            FlushReason::Closing
552        );
553        assert_eq!(
554            resolve_flush_reason(Some(RegionFlushReason::Downgrading), false),
555            FlushReason::Downgrading
556        );
557    }
558
559    #[test]
560    fn test_resolve_flush_reason_fallback_unchanged() {
561        assert_eq!(resolve_flush_reason(None, true), FlushReason::Downgrading);
562        assert_eq!(resolve_flush_reason(None, false), FlushReason::Manual);
563    }
564
565    #[test]
566    fn test_region_write_buffer_status_boundaries() {
567        assert_eq!(
568            RegionWriteBufferStatus::default(),
569            region_write_buffer_status_from_usage(100, 49, 99, 0)
570        );
571        assert_eq!(
572            RegionWriteBufferStatus {
573                should_flush: true,
574                should_stall: false,
575                should_reject: false,
576            },
577            region_write_buffer_status_from_usage(100, 50, 99, 0)
578        );
579        assert_eq!(
580            RegionWriteBufferStatus {
581                should_flush: true,
582                should_stall: true,
583                should_reject: false,
584            },
585            region_write_buffer_status_from_usage(100, 50, 100, 99)
586        );
587        assert_eq!(
588            RegionWriteBufferStatus {
589                should_flush: true,
590                should_stall: true,
591                should_reject: true,
592            },
593            region_write_buffer_status_from_usage(100, 50, 100, 100)
594        );
595    }
596
597    #[test]
598    fn test_region_write_buffer_status_saturates_reject_limit() {
599        assert_eq!(
600            RegionWriteBufferStatus {
601                should_flush: true,
602                should_stall: true,
603                should_reject: true,
604            },
605            region_write_buffer_status_from_usage(usize::MAX, usize::MAX, usize::MAX, usize::MAX,)
606        );
607    }
608}