1use 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 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 self.maybe_flush_worker();
74
75 self.handle_stalled_requests().await;
77 }
78
79 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 return;
86 }
87
88 if let Err(e) = self.flush_regions_on_engine_full() {
91 error!(e; "Failed to flush worker");
92 }
93 }
94
95 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 ®ions {
102 if self.flush_scheduler.is_flush_requested(region.region_id) || !region.is_writable() {
103 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 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 ®ion.version_control,
127 task,
128 )?;
129 } else if region_memtable_size > 0 {
130 pending_regions.push((region, region_memtable_size));
132 }
133 }
134 pending_regions.sort_unstable_by_key(|(_, size)| std::cmp::Reverse(*size));
135 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 for (region, region_mem_size) in pending_regions.into_iter() {
156 if memory_usage < target_memory_usage {
158 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 self.flush_scheduler
166 .schedule_flush(region.region_id, ®ion.version_control, task)?;
167 memory_usage = memory_usage.saturating_sub(region_mem_size);
169 }
170
171 Ok(())
172 }
173
174 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 ®ion.version(),
191 self.config.default_region_write_buffer_size,
192 self.stalled_requests.estimated_size(®ion_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 ®ion,
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 ®ion.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 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 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
265fn 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 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 self.update_topic_latest_entry_id(®ion);
325
326 let reason = resolve_flush_reason(request.reason, region.is_downgrading());
327 let mut task =
328 self.new_flush_task(®ion, 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, ®ion.version_control, task)
333 {
334 error!(e; "Failed to schedule flush task for region {}", region.region_id);
335 }
336 }
337
338 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 ®ions {
344 if self.flush_scheduler.is_flush_requested(region.region_id) || !region.is_writable() {
345 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 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 ®ion.version_control,
372 task,
373 )?;
374 }
375 }
376
377 Ok(())
378 }
379
380 pub(crate) async fn handle_flush_finished(
382 &mut self,
383 region_id: RegionId,
384 mut request: FlushFinished,
385 ) {
386 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 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 self.update_topic_latest_entry_id(®ion);
421
422 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, ®ion.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 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 request.on_success();
464 if let Some((mut ddl_requests, mut write_requests, mut bulk_writes)) =
466 self.flush_scheduler.on_flush_success(region_id)
467 {
468 self.handle_ddl_requests(&mut ddl_requests).await;
470 if self.flush_scheduler.is_flush_requested(region_id) {
471 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 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 self.handle_write_requests(&mut write_requests, &mut bulk_writes, false)
495 .await;
496 }
497 self.maybe_flush_worker();
499 self.handle_stalled_requests().await;
501 self.schedule_compaction(®ion).await;
503
504 self.listener.on_flush_success(region_id);
505 }
506
507 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(®ion.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}