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            durability_barrier: self
247                .wal
248                .durability_barrier(region.region_id, &region.provider),
249        }
250    }
251}
252
253pub(crate) fn region_memtable_usage(version: &VersionRef) -> usize {
254    version.memtables.mutable_usage() + version.memtables.immutables_usage()
255}
256
257pub(crate) fn region_write_buffer_size(
258    version: &VersionRef,
259    default_region_write_buffer_size: ReadableSize,
260) -> Option<usize> {
261    let size = version
262        .options
263        .write_buffer_size
264        .unwrap_or(default_region_write_buffer_size);
265    (size.as_bytes() > 0).then_some(size.as_bytes() as usize)
266}
267
268/// Returns the pressure state of a region's write buffer.
269fn region_write_buffer_status(
270    version: &VersionRef,
271    default_region_write_buffer_size: ReadableSize,
272    stalled_request_size: usize,
273) -> RegionWriteBufferStatus {
274    let Some(write_buffer_size) =
275        region_write_buffer_size(version, default_region_write_buffer_size)
276    else {
277        return RegionWriteBufferStatus::default();
278    };
279
280    let mutable_usage = version.memtables.mutable_usage();
281    let memory_usage = region_memtable_usage(version);
282    region_write_buffer_status_from_usage(
283        write_buffer_size,
284        mutable_usage,
285        memory_usage,
286        stalled_request_size,
287    )
288}
289
290fn region_write_buffer_status_from_usage(
291    write_buffer_size: usize,
292    mutable_usage: usize,
293    memory_usage: usize,
294    stalled_request_size: usize,
295) -> RegionWriteBufferStatus {
296    let mutable_limit = std::cmp::max(1, write_buffer_size / 2);
297    let should_stall = memory_usage >= write_buffer_size;
298    let reject_limit = write_buffer_size.saturating_mul(2);
299    let should_reject = memory_usage.saturating_add(stalled_request_size) >= reject_limit;
300
301    RegionWriteBufferStatus {
302        should_flush: mutable_usage >= mutable_limit || should_stall,
303        should_stall,
304        should_reject,
305    }
306}
307
308impl<S: LogStore> RegionWorkerLoop<S> {
309    /// Handles manual flush request.
310    pub(crate) fn handle_flush_request(
311        &mut self,
312        region_id: RegionId,
313        request: RegionFlushRequest,
314        sender: OptionOutputTx,
315    ) {
316        let region = match self.regions.flushable_region(region_id) {
317            Ok(region) => region,
318            Err(e) => {
319                sender.send(Err(e));
320                return;
321            }
322        };
323
324        // `update_topic_latest_entry_id` updates `topic_latest_entry_id` when memtables are empty.
325        // But the flush is skipped if memtables are empty. Thus should update the `topic_latest_entry_id`
326        // when handling flush request instead of in `schedule_flush` or `flush_finished`.
327        self.update_topic_latest_entry_id(&region);
328
329        let reason = resolve_flush_reason(request.reason, region.is_downgrading());
330        let mut task =
331            self.new_flush_task(&region, reason, request.row_group_size, self.config.clone());
332        task.push_sender(sender);
333        if let Err(e) =
334            self.flush_scheduler
335                .schedule_flush(region.region_id, &region.version_control, task)
336        {
337            error!(e; "Failed to schedule flush task for region {}", region.region_id);
338        }
339    }
340
341    /// Flushes regions periodically.
342    pub(crate) fn flush_periodically(&mut self) -> Result<()> {
343        let regions = self.regions.list_regions();
344        let now = self.time_provider.current_time_millis();
345
346        for region in &regions {
347            if self.flush_scheduler.is_flush_requested(region.region_id) || !region.is_writable() {
348                // Already flushing or not writable.
349                continue;
350            }
351            self.update_topic_latest_entry_id(region);
352
353            let auto_flush_interval = region
354                .version()
355                .options
356                .auto_flush_interval_or(self.config.auto_flush_interval);
357            let min_last_flush_time = now.saturating_sub(
358                auto_flush_interval
359                    .as_millis()
360                    .try_into()
361                    .unwrap_or(i64::MAX),
362            );
363
364            if region.last_flush_millis() < min_last_flush_time {
365                // If flush time of this region is earlier than `min_last_flush_time`, we can flush this region.
366                let task = self.new_flush_task(
367                    region,
368                    FlushReason::Periodically,
369                    None,
370                    self.config.clone(),
371                );
372                self.flush_scheduler.schedule_flush(
373                    region.region_id,
374                    &region.version_control,
375                    task,
376                )?;
377            }
378        }
379
380        Ok(())
381    }
382
383    /// On region flush job finished.
384    pub(crate) async fn handle_flush_finished(
385        &mut self,
386        region_id: RegionId,
387        mut request: FlushFinished,
388    ) {
389        // Notifies other workers. Even the remaining steps of this method fail we still
390        // wake up other workers as we have released some memory by flush.
391        self.notify_group();
392
393        let region = match self.regions.get_region(region_id) {
394            Some(region) => region,
395            None => {
396                request.on_failure(RegionNotFoundSnafu { region_id }.build());
397                return;
398            }
399        };
400
401        if request.is_staging {
402            // Skip the region metadata update.
403            info!(
404                "Skipping region metadata update for region {} in staging mode",
405                region_id
406            );
407            region.version_control.apply_edit(
408                None,
409                &request.memtables_to_remove,
410                region.file_purger.clone(),
411            );
412        } else {
413            region.version_control.apply_edit(
414                Some(request.edit.clone()),
415                &request.memtables_to_remove,
416                region.file_purger.clone(),
417            );
418        }
419
420        region.update_flush_millis();
421        // Update topic latest entry id as soon as possible after flush to make sure the prunable entry id is updated timely,
422        // which is important for remote WAL pruning.
423        self.update_topic_latest_entry_id(&region);
424
425        // Delete wal.
426        info!(
427            "Region {} flush finished, tries to bump wal to {}",
428            region_id, request.flushed_entry_id
429        );
430        if let Err(e) = self
431            .wal
432            .obsolete(region_id, request.flushed_entry_id, &region.provider)
433            .await
434        {
435            error!(e; "Failed to write wal, region: {}", region_id);
436            request.on_failure(e);
437            return;
438        }
439
440        let flush_on_close = request.flush_reason == FlushReason::Closing;
441
442        let index_build_file_metas = std::mem::take(&mut request.edit.files_to_add);
443
444        // In async mode, create indexes after flush.
445        if self.config.index.build_mode == IndexBuildMode::Async {
446            self.handle_rebuild_index(
447                BuildIndexRequest {
448                    region_id,
449                    build_type: IndexBuildType::Flush,
450                    file_metas: index_build_file_metas,
451                },
452                OptionOutputTx::new(None),
453            )
454            .await;
455        }
456
457        if flush_on_close {
458            self.remove_region(region_id).await;
459            info!("Region {} closed after flush", region_id);
460            request.on_success();
461            self.listener.on_flush_success(region_id);
462            return;
463        }
464
465        // Notifies waiters and observes the flush timer.
466        request.on_success();
467        // Handle pending requests for the region.
468        if let Some((mut ddl_requests, mut write_requests, mut bulk_writes)) =
469            self.flush_scheduler.on_flush_success(region_id)
470        {
471            // Perform DDLs first because they require empty memtables.
472            self.handle_ddl_requests(&mut ddl_requests).await;
473            if self.flush_scheduler.is_flush_requested(region_id) {
474                // The DDL may schedule another flush, e.g. a close-time flush after writes
475                // arrived in the mutable memtable during the previous flush. Keep pending
476                // writes fenced until that flush reaches its terminal state instead of
477                // accepting them while the DDL is still in progress.
478                for write_request in write_requests {
479                    self.flush_scheduler
480                        .add_write_request_to_pending(write_request);
481                }
482                for bulk_write in bulk_writes {
483                    self.flush_scheduler.add_bulk_request_to_pending(bulk_write);
484                }
485                self.listener.on_flush_success(region_id);
486                return;
487            }
488            // A pending close DDL may have removed the region. Reject queued writes as
489            // not found, then stop instead of scheduling compaction for a closed region.
490            if !self.regions.is_region_exists(region_id) {
491                self.handle_write_requests(&mut write_requests, &mut bulk_writes, false)
492                    .await;
493                self.listener.on_flush_success(region_id);
494                return;
495            }
496            // Handle pending write requests, we don't stall these requests.
497            self.handle_write_requests(&mut write_requests, &mut bulk_writes, false)
498                .await;
499        }
500        // Maybe flush worker again.
501        self.maybe_flush_worker();
502        // Handle stalled requests.
503        self.handle_stalled_requests().await;
504        // Schedules compaction.
505        self.schedule_compaction(&region).await;
506
507        self.listener.on_flush_success(region_id);
508    }
509
510    /// Updates the latest entry id since flush of the region.
511    /// **This is only used for remote WAL pruning.**
512    pub(crate) fn update_topic_latest_entry_id(&mut self, region: &MitoRegionRef) {
513        if region.provider.is_remote_wal() && region.version().memtables.is_empty() {
514            let latest_offset = self
515                .wal
516                .store()
517                .latest_entry_id(&region.provider)
518                .unwrap_or(0);
519            let topic_last_entry_id = region.topic_latest_entry_id.load(Ordering::Relaxed);
520
521            if latest_offset > topic_last_entry_id {
522                region
523                    .topic_latest_entry_id
524                    .store(latest_offset, Ordering::Relaxed);
525                debug!(
526                    "Region {} latest entry id updated to {}",
527                    region.region_id, latest_offset
528                );
529            }
530        }
531    }
532}
533
534#[cfg(test)]
535mod tests {
536    use super::*;
537
538    #[test]
539    fn test_resolve_flush_reason_uses_request_reason() {
540        assert_eq!(
541            resolve_flush_reason(Some(RegionFlushReason::RegionMigration), true),
542            FlushReason::RegionMigration
543        );
544        assert_eq!(
545            resolve_flush_reason(Some(RegionFlushReason::Repartition), false),
546            FlushReason::Repartition
547        );
548        assert_eq!(
549            resolve_flush_reason(Some(RegionFlushReason::RemoteWalPrune), false),
550            FlushReason::RemoteWalPrune
551        );
552        assert_eq!(
553            resolve_flush_reason(Some(RegionFlushReason::Closing), false),
554            FlushReason::Closing
555        );
556        assert_eq!(
557            resolve_flush_reason(Some(RegionFlushReason::Downgrading), false),
558            FlushReason::Downgrading
559        );
560    }
561
562    #[test]
563    fn test_resolve_flush_reason_fallback_unchanged() {
564        assert_eq!(resolve_flush_reason(None, true), FlushReason::Downgrading);
565        assert_eq!(resolve_flush_reason(None, false), FlushReason::Manual);
566    }
567
568    #[test]
569    fn test_region_write_buffer_status_boundaries() {
570        assert_eq!(
571            RegionWriteBufferStatus::default(),
572            region_write_buffer_status_from_usage(100, 49, 99, 0)
573        );
574        assert_eq!(
575            RegionWriteBufferStatus {
576                should_flush: true,
577                should_stall: false,
578                should_reject: false,
579            },
580            region_write_buffer_status_from_usage(100, 50, 99, 0)
581        );
582        assert_eq!(
583            RegionWriteBufferStatus {
584                should_flush: true,
585                should_stall: true,
586                should_reject: false,
587            },
588            region_write_buffer_status_from_usage(100, 50, 100, 99)
589        );
590        assert_eq!(
591            RegionWriteBufferStatus {
592                should_flush: true,
593                should_stall: true,
594                should_reject: true,
595            },
596            region_write_buffer_status_from_usage(100, 50, 100, 100)
597        );
598    }
599
600    #[test]
601    fn test_region_write_buffer_status_saturates_reject_limit() {
602        assert_eq!(
603            RegionWriteBufferStatus {
604                should_flush: true,
605                should_stall: true,
606                should_reject: true,
607            },
608            region_write_buffer_status_from_usage(usize::MAX, usize::MAX, usize::MAX, usize::MAX,)
609        );
610    }
611}