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 durability_barrier: self
247 .wal
248 .durability_barrier(region.region_id, ®ion.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
268fn 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 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 self.update_topic_latest_entry_id(®ion);
328
329 let reason = resolve_flush_reason(request.reason, region.is_downgrading());
330 let mut task =
331 self.new_flush_task(®ion, 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, ®ion.version_control, task)
336 {
337 error!(e; "Failed to schedule flush task for region {}", region.region_id);
338 }
339 }
340
341 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 ®ions {
347 if self.flush_scheduler.is_flush_requested(region.region_id) || !region.is_writable() {
348 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 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 ®ion.version_control,
375 task,
376 )?;
377 }
378 }
379
380 Ok(())
381 }
382
383 pub(crate) async fn handle_flush_finished(
385 &mut self,
386 region_id: RegionId,
387 mut request: FlushFinished,
388 ) {
389 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 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 self.update_topic_latest_entry_id(®ion);
424
425 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, ®ion.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 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 request.on_success();
467 if let Some((mut ddl_requests, mut write_requests, mut bulk_writes)) =
469 self.flush_scheduler.on_flush_success(region_id)
470 {
471 self.handle_ddl_requests(&mut ddl_requests).await;
473 if self.flush_scheduler.is_flush_requested(region_id) {
474 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 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 self.handle_write_requests(&mut write_requests, &mut bulk_writes, false)
498 .await;
499 }
500 self.maybe_flush_worker();
502 self.handle_stalled_requests().await;
504 self.schedule_compaction(®ion).await;
506
507 self.listener.on_flush_success(region_id);
508 }
509
510 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(®ion.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}