1use std::collections::HashMap;
18use std::num::NonZeroU64;
19use std::sync::atomic::{AtomicUsize, Ordering};
20use std::sync::{Arc, Mutex};
21use std::time::Instant;
22
23use bytes::Bytes;
24use common_base::cancellation::CancellableFuture;
25use common_telemetry::{debug, error, info};
26use datatypes::arrow::datatypes::SchemaRef;
27use datatypes::extension::json::is_json2_extension_type;
28use partition::expr::PartitionExpr;
29use smallvec::{SmallVec, smallvec};
30use snafu::ResultExt;
31use store_api::region_request::RegionFlushReason;
32use store_api::storage::{RegionId, SequenceNumber};
33use strum::IntoStaticStr;
34use tokio::sync::{Semaphore, mpsc, watch};
35
36use crate::access_layer::{
37 AccessLayerRef, Metrics, OperationType, SstInfoArray, SstWriteRequest, WriteType,
38};
39use crate::cache::CacheManagerRef;
40use crate::config::MitoConfig;
41use crate::engine::region_hook::SstFileInfo;
42use crate::error::{
43 Error, FlushCancelledSnafu, FlushRegionSnafu, JoinSnafu, RegionBusySnafu, RegionClosedSnafu,
44 RegionDroppedSnafu, RegionTruncatedSnafu, Result,
45};
46use crate::manifest::action::{RegionEdit, RegionMetaAction, RegionMetaActionList};
47use crate::memtable::bulk::ENCODE_ROW_THRESHOLD;
48use crate::memtable::bulk::json_align::Json2Aligner;
49use crate::memtable::{BoxedRecordBatchIterator, EncodedRange, MemtableRanges, RangesOptions};
50use crate::metrics::{
51 FLUSH_BYTES_TOTAL, FLUSH_ELAPSED, FLUSH_FAILURE_TOTAL, FLUSH_FILE_TOTAL, FLUSH_REQUESTS_TOTAL,
52 INFLIGHT_FLUSH_COUNT,
53};
54use crate::read::FlatSource;
55use crate::read::flat_dedup::{FlatDedupIterator, FlatLastNonNull, FlatLastRow};
56use crate::read::flat_merge::FlatMergeIterator;
57use crate::region::options::{IndexOptions, MergeMode, RegionOptions};
58use crate::region::version::{VersionControlData, VersionControlRef, VersionRef};
59use crate::region::{ManifestContextRef, RegionLeaderState, RegionRoleState, parse_partition_expr};
60use crate::request::{
61 BackgroundNotify, DdlRequest, FlushFailed, FlushFinished, OnFailure, OptionOutputTx, OutputTx,
62 SenderBulkRequest, SenderDdlRequest, SenderWriteRequest, WorkerRequest, WorkerRequestWithTime,
63};
64use crate::schedule::CancellableTaskState;
65use crate::schedule::scheduler::{Job, SchedulerRef};
66use crate::sst::file::{FileMeta, UncommittedSsts};
67use crate::sst::parquet::metadata::extract_primary_key_range;
68use crate::sst::parquet::{
69 DEFAULT_READ_BATCH_SIZE, DEFAULT_ROW_GROUP_SIZE, SstInfo, WriteOptions, flat_format,
70};
71use crate::sst::{FlatSchemaOptions, FormatType, to_flat_sst_arrow_schema};
72use crate::worker::WorkerListener;
73
74pub trait WriteBufferManager: Send + Sync + std::fmt::Debug {
78 fn should_flush_engine(&self) -> bool;
80
81 fn should_stall(&self) -> bool;
83
84 fn reserve_mem(&self, mem: usize);
86
87 fn schedule_free_mem(&self, mem: usize);
92
93 fn free_mem(&self, mem: usize);
95
96 fn memory_usage(&self) -> usize;
98
99 fn flush_limit(&self) -> usize;
104}
105
106pub type WriteBufferManagerRef = Arc<dyn WriteBufferManager>;
107
108#[derive(Debug)]
113pub struct WriteBufferManagerImpl {
114 global_write_buffer_size: usize,
116 mutable_limit: usize,
118 memory_used: AtomicUsize,
120 memory_active: AtomicUsize,
122 notifier: Option<watch::Sender<()>>,
125}
126
127impl WriteBufferManagerImpl {
128 pub fn new(global_write_buffer_size: usize) -> Self {
130 Self {
131 global_write_buffer_size,
132 mutable_limit: Self::get_mutable_limit(global_write_buffer_size),
133 memory_used: AtomicUsize::new(0),
134 memory_active: AtomicUsize::new(0),
135 notifier: None,
136 }
137 }
138
139 pub fn with_notifier(mut self, notifier: watch::Sender<()>) -> Self {
141 self.notifier = Some(notifier);
142 self
143 }
144
145 pub fn mutable_usage(&self) -> usize {
147 self.memory_active.load(Ordering::Relaxed)
148 }
149
150 fn get_mutable_limit(global_write_buffer_size: usize) -> usize {
152 global_write_buffer_size / 2
154 }
155}
156
157impl WriteBufferManager for WriteBufferManagerImpl {
158 fn should_flush_engine(&self) -> bool {
159 let mutable_memtable_memory_usage = self.memory_active.load(Ordering::Relaxed);
160 if mutable_memtable_memory_usage >= self.mutable_limit {
161 debug!(
162 "Engine should flush (over mutable limit), mutable_usage: {}, memory_usage: {}, mutable_limit: {}, global_limit: {}",
163 mutable_memtable_memory_usage,
164 self.memory_usage(),
165 self.mutable_limit,
166 self.global_write_buffer_size,
167 );
168 return true;
169 }
170
171 let memory_usage = self.memory_used.load(Ordering::Relaxed);
172 if memory_usage >= self.global_write_buffer_size {
173 return true;
174 }
175
176 false
177 }
178
179 fn should_stall(&self) -> bool {
180 self.memory_usage() >= self.global_write_buffer_size
181 }
182
183 fn reserve_mem(&self, mem: usize) {
184 self.memory_used.fetch_add(mem, Ordering::Relaxed);
185 self.memory_active.fetch_add(mem, Ordering::Relaxed);
186 }
187
188 fn schedule_free_mem(&self, mem: usize) {
189 self.memory_active.fetch_sub(mem, Ordering::Relaxed);
190 }
191
192 fn free_mem(&self, mem: usize) {
193 self.memory_used.fetch_sub(mem, Ordering::Relaxed);
194 if let Some(notifier) = &self.notifier {
195 let _ = notifier.send(());
199 }
200 }
201
202 fn memory_usage(&self) -> usize {
203 self.memory_used.load(Ordering::Relaxed)
204 }
205
206 fn flush_limit(&self) -> usize {
207 self.mutable_limit
208 }
209}
210
211#[derive(Debug, IntoStaticStr, Clone, Copy, PartialEq, Eq)]
213pub enum FlushReason {
214 EngineFull,
216 RegionFull,
218 Manual,
220 Alter,
222 Periodically,
224 Downgrading,
226 EnterStaging,
228 RegionMigration,
230 Repartition,
232 RemoteWalPrune,
234 Closing,
236}
237
238impl FlushReason {
239 fn as_str(&self) -> &'static str {
241 self.into()
242 }
243}
244
245impl From<RegionFlushReason> for FlushReason {
246 fn from(reason: RegionFlushReason) -> Self {
247 match reason {
248 RegionFlushReason::RegionMigration => FlushReason::RegionMigration,
249 RegionFlushReason::Repartition => FlushReason::Repartition,
250 RegionFlushReason::RemoteWalPrune => FlushReason::RemoteWalPrune,
251 RegionFlushReason::Closing => FlushReason::Closing,
252 RegionFlushReason::Downgrading => FlushReason::Downgrading,
253 }
254 }
255}
256
257pub(crate) struct RegionFlushTask {
259 pub(crate) region_id: RegionId,
261 pub(crate) reason: FlushReason,
263 pub(crate) senders: Vec<OutputTx>,
265 pub(crate) request_sender: mpsc::Sender<WorkerRequestWithTime>,
267
268 pub(crate) access_layer: AccessLayerRef,
269 pub(crate) listener: WorkerListener,
270 pub(crate) engine_config: Arc<MitoConfig>,
271 pub(crate) row_group_size: Option<usize>,
272 pub(crate) cache_manager: CacheManagerRef,
273 pub(crate) manifest_ctx: ManifestContextRef,
274
275 pub(crate) index_options: IndexOptions,
277 pub(crate) flush_semaphore: Arc<Semaphore>,
279 pub(crate) is_staging: bool,
281 pub(crate) partition_expr: Option<String>,
285}
286
287struct FlushTaskWaiters {
288 region_id: RegionId,
289 senders: Mutex<Vec<OutputTx>>,
290}
291
292impl FlushTaskWaiters {
293 fn new(region_id: RegionId, senders: Vec<OutputTx>) -> Self {
294 Self {
295 region_id,
296 senders: Mutex::new(senders),
297 }
298 }
299
300 fn take(&self) -> Vec<OutputTx> {
301 std::mem::take(&mut *self.senders.lock().unwrap())
302 }
303
304 fn on_failure(&self, err: Arc<Error>) {
305 for sender in self.take() {
306 sender.send(Err(err.clone()).context(FlushRegionSnafu {
307 region_id: self.region_id,
308 }));
309 }
310 }
311}
312
313impl Drop for FlushTaskWaiters {
314 fn drop(&mut self) {
315 self.on_failure(Arc::new(
316 RegionBusySnafu {
317 region_id: self.region_id,
318 }
319 .build(),
320 ));
321 }
322}
323
324impl RegionFlushTask {
325 pub(crate) fn push_sender(&mut self, mut sender: OptionOutputTx) {
327 if let Some(sender) = sender.take_inner() {
328 self.senders.push(sender);
329 }
330 }
331
332 fn on_success(self) {
334 for sender in self.senders {
335 sender.send(Ok(0));
336 }
337 }
338
339 fn on_failure(&mut self, err: Arc<Error>) {
341 for sender in self.senders.drain(..) {
342 sender.send(Err(err.clone()).context(FlushRegionSnafu {
343 region_id: self.region_id,
344 }));
345 }
346 }
347
348 fn into_flush_job(
352 mut self,
353 version_control: &VersionControlRef,
354 state: CancellableTaskState,
355 ) -> (Job, Arc<FlushTaskWaiters>) {
356 let version_data = version_control.current();
359 let waiters = Arc::new(FlushTaskWaiters::new(
360 self.region_id,
361 std::mem::take(&mut self.senders),
362 ));
363 let job_waiters = waiters.clone();
364
365 let job = Box::pin(async move {
366 self.senders = job_waiters.take();
367 INFLIGHT_FLUSH_COUNT.inc();
368 self.do_flush(version_data, state).await;
369 INFLIGHT_FLUSH_COUNT.dec();
370 });
371 (job, waiters)
372 }
373
374 async fn do_flush(&mut self, version_data: VersionControlData, state: CancellableTaskState) {
376 let timer = FLUSH_ELAPSED.with_label_values(&["total"]).start_timer();
377 let uncommitted = UncommittedSsts::new(
378 self.region_id,
379 self.access_layer.clone(),
380 Some(self.cache_manager.clone()),
381 );
382 self.listener.on_flush_begin(self.region_id).await;
383 let flush_result = if state.is_cancelled() {
384 FlushCancelledSnafu.fail()
385 } else {
386 self.flush_memtables(&version_data, &state, &uncommitted)
387 .await
388 };
389
390 let worker_request = match flush_result {
391 Ok(edit) => {
392 let memtables_to_remove = version_data
393 .version
394 .memtables
395 .immutables()
396 .iter()
397 .map(|m| m.id())
398 .collect();
399 let flush_finished = FlushFinished {
400 region_id: self.region_id,
401 flush_reason: self.reason,
402 flushed_entry_id: version_data.last_entry_id,
404 senders: std::mem::take(&mut self.senders),
405 _timer: timer,
406 edit,
407 memtables_to_remove,
408 is_staging: self.is_staging,
409 };
410 WorkerRequest::Background {
411 region_id: self.region_id,
412 notify: BackgroundNotify::FlushFinished(flush_finished),
413 }
414 }
415 Err(e) => {
416 let err = Arc::new(e);
417 let failed = FlushFailed { err: err.clone() };
418 if failed.is_cancelled() {
419 info!("Flush cancelled for region {}", self.region_id);
420 } else {
421 error!(err; "Failed to flush region {}", self.region_id);
422 }
423 timer.stop_and_discard();
425 uncommitted.cleanup().await;
426
427 self.on_failure(err.clone());
428 WorkerRequest::Background {
429 region_id: self.region_id,
430 notify: BackgroundNotify::FlushFailed(failed),
431 }
432 }
433 };
434 self.send_worker_request(worker_request).await;
435 }
436
437 async fn flush_memtables(
440 &self,
441 version_data: &VersionControlData,
442 state: &CancellableTaskState,
443 uncommitted: &UncommittedSsts,
444 ) -> Result<RegionEdit> {
445 let version = &version_data.version;
448 let timer = FLUSH_ELAPSED
449 .with_label_values(&["flush_memtables"])
450 .start_timer();
451
452 let mut write_opts = WriteOptions {
453 write_buffer_size: self.engine_config.sst_write_buffer_size,
454 ..Default::default()
455 };
456 if let Some(row_group_size) = self.row_group_size {
457 write_opts.row_group_size = row_group_size;
458 }
459
460 let DoFlushMemtablesResult {
461 file_metas,
462 flushed_bytes,
463 series_count,
464 encoded_part_count,
465 flush_metrics,
466 sst_infos,
467 } = self
468 .do_flush_memtables(version, write_opts, state, uncommitted)
469 .await?;
470
471 if !file_metas.is_empty() {
472 FLUSH_BYTES_TOTAL.inc_by(flushed_bytes);
473 }
474
475 let mut file_ids = Vec::with_capacity(file_metas.len());
476 let mut total_rows = 0;
477 let mut total_bytes = 0;
478 for meta in &file_metas {
479 file_ids.push(meta.file_id);
480 total_rows += meta.num_rows;
481 total_bytes += meta.file_size;
482 }
483 info!(
484 "Successfully flush memtables, region: {}, reason: {}, files: {:?}, series count: {}, total_rows: {}, total_bytes: {}, cost: {:?}, encoded_part_count: {}, metrics: {:?}",
485 self.region_id,
486 self.reason.as_str(),
487 file_ids,
488 series_count,
489 total_rows,
490 total_bytes,
491 timer.stop_and_record(),
492 encoded_part_count,
493 flush_metrics,
494 );
495 flush_metrics.observe();
496
497 let hook = self.manifest_ctx.hook();
498 if let Some(hook) = &hook {
499 let files: Vec<SstFileInfo<'_>> = sst_infos
500 .iter()
501 .zip(file_metas.iter())
502 .map(|(sst_info, file_meta)| SstFileInfo {
503 sst_info_ref: sst_info,
504 file_meta,
505 })
506 .collect();
507 hook.on_sst_files_written(self.region_id, &version.metadata, &files)
508 .await;
509 }
510
511 let edit = RegionEdit {
512 files_to_add: file_metas,
513 files_to_remove: Vec::new(),
514 timestamp_ms: Some(chrono::Utc::now().timestamp_millis()),
515 compaction_time_window: None,
516 flushed_entry_id: Some(version_data.last_entry_id),
518 flushed_sequence: Some(version_data.committed_sequence),
519 committed_sequence: None,
520 };
521 info!(
522 "Applying {edit:?} to region {}, is_staging: {}",
523 self.region_id, self.is_staging
524 );
525
526 let action_list = RegionMetaActionList::with_action(RegionMetaAction::Edit(edit.clone()));
527
528 if !state.mark_commit_started() {
530 return FlushCancelledSnafu.fail();
531 }
532 self.listener.on_flush_commit_begin(self.region_id).await;
533
534 let expected_state = if matches!(self.reason, FlushReason::Downgrading) {
535 RegionLeaderState::Downgrading
536 } else {
537 let current_state = self.manifest_ctx.current_state();
539 if current_state == RegionRoleState::Leader(RegionLeaderState::Staging) {
540 RegionLeaderState::Staging
541 } else {
542 RegionLeaderState::Writable
543 }
544 };
545 let manifest_version = match self
546 .manifest_ctx
547 .update_manifest(expected_state, action_list, self.is_staging)
548 .await
549 {
550 Ok(manifest_version) => {
551 uncommitted.disarm_cleanup();
552 manifest_version
553 }
554 Err(e) => {
555 if e.may_have_persisted_manifest_update() {
556 uncommitted.disarm_cleanup();
557 } else {
558 info!(
559 "Cleaning uncommitted SSTs because the manifest update was not persisted, region: {}, job: flush, error: {:?}",
560 self.region_id, e
561 );
562 uncommitted.cleanup().await;
563 }
564 return Err(e);
565 }
566 };
567 info!(
568 "Successfully update manifest version to {manifest_version}, region: {}, is_staging: {}, reason: {}",
569 self.region_id,
570 self.is_staging,
571 self.reason.as_str()
572 );
573
574 Ok(edit)
575 }
576
577 async fn do_flush_memtables(
578 &self,
579 version: &VersionRef,
580 write_opts: WriteOptions,
581 state: &CancellableTaskState,
582 uncommitted: &UncommittedSsts,
583 ) -> Result<DoFlushMemtablesResult> {
584 let memtables = version.memtables.immutables();
585 let mut file_metas = Vec::with_capacity(memtables.len());
586 let mut flushed_bytes = 0;
587 let mut series_count = 0;
588 let mut encoded_part_count = 0;
589 let mut flush_metrics = Metrics::new(WriteType::Flush);
590 let partition_expr = parse_partition_expr(self.partition_expr.as_deref())?;
591 let hook = self.manifest_ctx.hook();
592 let mut all_sst_infos = Vec::new();
593 for mem in memtables {
594 if mem.is_empty() {
595 continue;
597 }
598
599 let compact_start = std::time::Instant::now();
601 if let Err(e) = mem.compact(true) {
602 common_telemetry::error!(e; "Failed to compact memtable before flush");
603 }
604 let compact_cost = compact_start.elapsed();
605 flush_metrics.compact_memtable += compact_cost;
606
607 let mem_stats = mem.stats();
608 let batch_size = crate::batch_size::estimate_batch_size([(
609 mem_stats.num_rows() as u64,
610 mem_stats.bytes_allocated() as u64,
611 )]);
612 let mem_ranges =
614 mem.ranges(None, RangesOptions::for_flush().with_batch_size(batch_size))?;
615 let num_mem_ranges = mem_ranges.ranges.len();
616
617 let num_mem_rows = mem_ranges.num_rows();
619 let memtable_series_count = mem_ranges.series_count();
620 let memtable_id = mem.id();
621 series_count += memtable_series_count;
624
625 let flush_start = Instant::now();
626 let FlushFlatMemResult {
627 num_encoded,
628 num_sources,
629 results,
630 } = self
631 .flush_flat_mem_ranges(version, &write_opts, mem_ranges, state, uncommitted)
632 .await?;
633 encoded_part_count += num_encoded;
634 for (source_idx, result) in results.into_iter().enumerate() {
635 let (max_sequence, ssts_written, metrics) = result?;
636 if ssts_written.is_empty() {
637 continue;
639 }
640
641 common_telemetry::debug!(
642 "Region {} flush one memtable {} {}/{}, metrics: {:?}",
643 self.region_id,
644 memtable_id,
645 source_idx,
646 num_sources,
647 metrics
648 );
649
650 flush_metrics = flush_metrics.merge(metrics);
651
652 for sst_info in &ssts_written {
653 flushed_bytes += sst_info.file_size;
654 let pk_range = sst_info
655 .file_metadata
656 .as_ref()
657 .and_then(|meta| extract_primary_key_range(meta, &version.metadata));
658 file_metas.push(Self::new_file_meta(
659 self.region_id,
660 max_sequence,
661 sst_info,
662 partition_expr.clone(),
663 pk_range,
664 ));
665 }
666 if hook.is_some() {
667 all_sst_infos.extend(ssts_written);
668 }
669 }
670
671 common_telemetry::debug!(
672 "Region {} flush {} memtables for {}, num_mem_ranges: {}, num_encoded: {}, num_rows: {}, flush_cost: {:?}, compact_cost: {:?}",
673 self.region_id,
674 num_sources,
675 memtable_id,
676 num_mem_ranges,
677 num_encoded,
678 num_mem_rows,
679 flush_start.elapsed(),
680 compact_cost,
681 );
682 }
683
684 Ok(DoFlushMemtablesResult {
685 file_metas,
686 flushed_bytes,
687 series_count,
688 encoded_part_count,
689 flush_metrics,
690 sst_infos: all_sst_infos,
691 })
692 }
693
694 async fn flush_flat_mem_ranges(
695 &self,
696 version: &VersionRef,
697 write_opts: &WriteOptions,
698 mem_ranges: MemtableRanges,
699 state: &CancellableTaskState,
700 uncommitted: &UncommittedSsts,
701 ) -> Result<FlushFlatMemResult> {
702 let batch_schema = to_flat_sst_arrow_schema(
703 &version.metadata,
704 &FlatSchemaOptions::from_encoding(version.metadata.primary_key_encoding),
705 );
706 let field_column_start =
707 flat_format::field_column_start(&version.metadata, batch_schema.fields().len());
708 let flat_sources = memtable_flat_sources(
709 batch_schema,
710 mem_ranges,
711 &version.options,
712 field_column_start,
713 )?;
714 let mut tasks = Vec::with_capacity(flat_sources.encoded.len() + flat_sources.sources.len());
715 let num_encoded = flat_sources.encoded.len();
716 for (source, max_sequence) in flat_sources.sources {
717 let write_request = self.new_write_request(version, max_sequence, source);
718 let access_layer = self.access_layer.clone();
719 let write_opts = write_opts.clone();
720 let semaphore = self.flush_semaphore.clone();
721 let uncommitted = uncommitted.clone();
722 let task = common_runtime::spawn_global(async move {
723 let _permit = semaphore.acquire().await.unwrap();
724 let mut metrics = Metrics::new(WriteType::Flush);
725 let ssts = access_layer
726 .write_sst(write_request, &write_opts, &mut metrics)
727 .await?;
728 uncommitted.track(&ssts);
729 FLUSH_FILE_TOTAL.inc_by(ssts.len() as u64);
730 Ok((max_sequence, ssts, metrics))
731 });
732 tasks.push(task);
733 }
734 for (encoded, max_sequence) in flat_sources.encoded {
735 let access_layer = self.access_layer.clone();
736 let cache_manager = self.cache_manager.clone();
737 let region_id = version.metadata.region_id;
738 let semaphore = self.flush_semaphore.clone();
739 let uncommitted = uncommitted.clone();
740 let task = common_runtime::spawn_global(async move {
741 let _permit = semaphore.acquire().await.unwrap();
742 let metrics = access_layer
743 .put_sst(&encoded.data, region_id, &encoded.sst_info, &cache_manager)
744 .await?;
745 uncommitted.track(std::slice::from_ref(&encoded.sst_info));
746 FLUSH_FILE_TOTAL.inc();
747 Ok((max_sequence, smallvec![encoded.sst_info], metrics))
748 });
749 tasks.push(task);
750 }
751 let num_sources = tasks.len();
752 let abort_handles = tasks
753 .iter()
754 .map(|task| task.abort_handle())
755 .collect::<Vec<_>>();
756 let join_all = futures::future::join_all(tasks);
757 tokio::pin!(join_all);
758 let results = match CancellableFuture::new(join_all.as_mut(), state.cancel_handle()).await {
759 Ok(results) => results
760 .into_iter()
761 .map(|result| result.context(JoinSnafu))
762 .collect::<Result<Vec<_>>>()?,
763 Err(_) => {
764 for handle in abort_handles {
765 handle.abort();
766 }
767 let _ = join_all.await;
770 return FlushCancelledSnafu.fail();
771 }
772 };
773 Ok(FlushFlatMemResult {
774 num_encoded,
775 num_sources,
776 results,
777 })
778 }
779
780 fn new_file_meta(
781 region_id: RegionId,
782 max_sequence: u64,
783 sst_info: &SstInfo,
784 partition_expr: Option<PartitionExpr>,
785 primary_key_range: Option<(Bytes, Bytes)>,
786 ) -> FileMeta {
787 let (primary_key_min, primary_key_max) = match primary_key_range {
788 Some((min, max)) => (Some(min), Some(max)),
789 None => (None, None),
790 };
791 FileMeta {
792 region_id,
793 file_id: sst_info.file_id,
794 time_range: sst_info.time_range,
795 level: 0,
796 file_size: sst_info.file_size,
797 max_row_group_uncompressed_size: sst_info.max_row_group_uncompressed_size,
798 available_indexes: sst_info.index_metadata.build_available_indexes(),
799 indexes: sst_info.index_metadata.build_indexes(),
800 index_file_size: sst_info.index_metadata.file_size,
801 index_version: 0,
802 num_rows: sst_info.num_rows as u64,
803 num_row_groups: sst_info.num_row_groups,
804 sequence: NonZeroU64::new(max_sequence),
805 partition_expr,
806 num_series: sst_info.num_series,
807 primary_key_min,
808 primary_key_max,
809 }
810 }
811
812 fn new_write_request(
813 &self,
814 version: &VersionRef,
815 max_sequence: u64,
816 source: FlatSource,
817 ) -> SstWriteRequest {
818 let flat_format = version
819 .options
820 .sst_format
821 .map(|f| f == FormatType::Flat)
822 .unwrap_or(self.engine_config.default_flat_format);
823 SstWriteRequest {
824 op_type: OperationType::Flush,
825 metadata: version.metadata.clone(),
826 source,
827 cache_manager: self.cache_manager.clone(),
828 storage: version.options.storage.clone(),
829 max_sequence: Some(max_sequence),
830 sst_write_format: if flat_format {
831 FormatType::Flat
832 } else {
833 FormatType::PrimaryKey
834 },
835 index_options: self.index_options.clone(),
836 index_config: self.engine_config.index.clone(),
837 inverted_index_config: self.engine_config.inverted_index.clone(),
838 fulltext_index_config: self.engine_config.fulltext_index.clone(),
839 bloom_filter_index_config: self.engine_config.bloom_filter_index.clone(),
840 #[cfg(feature = "vector_index")]
841 vector_index_config: self.engine_config.vector_index.clone(),
842 }
843 }
844
845 pub(crate) async fn send_worker_request(&self, request: WorkerRequest) {
847 if let Err(e) = self
848 .request_sender
849 .send(WorkerRequestWithTime::new(request))
850 .await
851 {
852 let request = e.0.request;
853 error!(
854 "Failed to notify flush job status for region {}, request: {:?}",
855 self.region_id, request
856 );
857 if let WorkerRequest::Background {
858 notify: BackgroundNotify::FlushFinished(mut finished),
859 ..
860 } = request
861 {
862 finished.on_failure(
863 RegionClosedSnafu {
864 region_id: self.region_id,
865 }
866 .build(),
867 );
868 }
869 }
870 }
871
872 fn merge(&mut self, mut other: RegionFlushTask) {
874 assert_eq!(self.region_id, other.region_id);
875 self.senders.append(&mut other.senders);
877 }
878}
879
880struct FlushFlatMemResult {
881 num_encoded: usize,
882 num_sources: usize,
883 results: Vec<Result<(SequenceNumber, SstInfoArray, Metrics)>>,
884}
885
886struct DoFlushMemtablesResult {
887 file_metas: Vec<FileMeta>,
888 flushed_bytes: u64,
889 series_count: usize,
890 encoded_part_count: usize,
891 flush_metrics: Metrics,
892 sst_infos: Vec<SstInfo>,
893}
894
895struct FlatSources {
896 sources: SmallVec<[(FlatSource, SequenceNumber); 4]>,
897 encoded: SmallVec<[(EncodedRange, SequenceNumber); 4]>,
898}
899
900fn memtable_flat_sources(
902 schema: SchemaRef,
903 mem_ranges: MemtableRanges,
904 options: &RegionOptions,
905 field_column_start: usize,
906) -> Result<FlatSources> {
907 let MemtableRanges { ranges } = mem_ranges;
908 let mut flat_sources = FlatSources {
909 sources: SmallVec::new(),
910 encoded: SmallVec::new(),
911 };
912
913 if ranges.len() == 1 {
914 debug!("Flushing single flat range");
915
916 let only_range = ranges.into_values().next().unwrap();
917 let max_sequence = only_range.stats().max_sequence();
918 if let Some(encoded) = only_range.encoded() {
919 flat_sources.encoded.push((encoded, max_sequence));
920 } else {
921 let iter = only_range.build_record_batch_iter(None, None)?;
922 let iter = maybe_dedup_one(
925 options.append_mode,
926 options.merge_mode(),
927 field_column_start,
928 iter,
929 );
930 flat_sources
931 .sources
932 .push((FlatSource::new_iter(schema, iter), max_sequence));
933 };
934 } else {
935 let min_flush_rows = *ENCODE_ROW_THRESHOLD;
936 let total_rows: usize = ranges
938 .values()
939 .filter(|r| r.encoded().is_none())
940 .map(|r| r.num_rows())
941 .sum();
942 debug!(
943 "Flushing multiple flat ranges, total_rows: {}, min_flush_rows: {}, num_ranges: {}",
944 total_rows,
945 min_flush_rows,
946 ranges.len()
947 );
948 let mut rows_remaining = total_rows;
949 let mut last_iter_rows = 0;
950 let num_ranges = ranges.len();
951 let mut input_iters = Vec::with_capacity(num_ranges);
952 let mut current_ranges = Vec::new();
953
954 let has_json2 = schema.fields().iter().any(is_json2_extension_type);
955 let mut json_align_schemas = if has_json2 {
956 Some(Vec::with_capacity(num_ranges))
957 } else {
958 None
959 };
960
961 for (_range_id, range) in ranges {
962 if let Some(encoded) = range.encoded() {
963 let max_sequence = range.stats().max_sequence();
964 flat_sources.encoded.push((encoded, max_sequence));
965 continue;
966 }
967
968 if let Some(schemas) = json_align_schemas.as_mut() {
970 let schema = range
971 .record_batch_schema_hint()
972 .unwrap_or_else(|| schema.clone());
973 schemas.push(schema);
974 }
975
976 let iter = range.build_record_batch_iter(None, None)?;
977 input_iters.push(iter);
978 let range_rows = range.num_rows();
979 last_iter_rows += range_rows;
980 rows_remaining -= range_rows;
981 current_ranges.push(range);
982
983 if last_iter_rows >= min_flush_rows
986 && (rows_remaining == 0 || rows_remaining >= DEFAULT_ROW_GROUP_SIZE)
987 {
988 debug!(
989 "Flush batch ready, rows: {}, min_rows: {}, num_iters: {}, remaining: {}",
990 last_iter_rows,
991 min_flush_rows,
992 input_iters.len(),
993 rows_remaining
994 );
995
996 let max_sequence = current_ranges
998 .iter()
999 .map(|r| r.stats().max_sequence())
1000 .max()
1001 .unwrap_or(0);
1002 let batch_size =
1003 crate::batch_size::estimate_batch_size(current_ranges.iter().map(|range| {
1004 let stats = range.stats();
1005 (stats.num_rows() as u64, stats.bytes_allocated() as u64)
1006 }));
1007
1008 let input_iters =
1009 std::mem::replace(&mut input_iters, Vec::with_capacity(num_ranges));
1010 let (schema, input_iters) = maybe_align_json2_iters(
1011 schema.clone(),
1012 json_align_schemas.take(),
1013 input_iters,
1014 )?;
1015
1016 let maybe_dedup = merge_and_dedup_with_batch_size(
1017 &schema,
1018 options.append_mode,
1019 options.merge_mode(),
1020 field_column_start,
1021 input_iters,
1022 batch_size,
1023 )?;
1024
1025 flat_sources
1026 .sources
1027 .push((FlatSource::new_iter(schema, maybe_dedup), max_sequence));
1028 last_iter_rows = 0;
1029 current_ranges.clear();
1030
1031 json_align_schemas = if has_json2 {
1032 Some(Vec::with_capacity(num_ranges))
1033 } else {
1034 None
1035 };
1036 }
1037 }
1038
1039 if !input_iters.is_empty() {
1041 debug!(
1042 "Flush remaining batch, rows: {}, min_rows: {}, num_iters: {}, remaining: {}",
1043 last_iter_rows,
1044 min_flush_rows,
1045 input_iters.len(),
1046 rows_remaining
1047 );
1048
1049 let (schema, input_iters) =
1050 maybe_align_json2_iters(schema, json_align_schemas, input_iters)?;
1051
1052 let max_sequence = current_ranges
1053 .iter()
1054 .map(|r| r.stats().max_sequence())
1055 .max()
1056 .unwrap_or(0);
1057 let batch_size =
1058 crate::batch_size::estimate_batch_size(current_ranges.iter().map(|range| {
1059 let stats = range.stats();
1060 (stats.num_rows() as u64, stats.bytes_allocated() as u64)
1061 }));
1062
1063 let maybe_dedup = merge_and_dedup_with_batch_size(
1064 &schema,
1065 options.append_mode,
1066 options.merge_mode(),
1067 field_column_start,
1068 input_iters,
1069 batch_size,
1070 )?;
1071
1072 flat_sources
1073 .sources
1074 .push((FlatSource::new_iter(schema, maybe_dedup), max_sequence));
1075 }
1076 }
1077
1078 Ok(flat_sources)
1079}
1080
1081fn maybe_align_json2_iters(
1082 schema: SchemaRef,
1083 schemas: Option<Vec<SchemaRef>>,
1084 input_iters: Vec<BoxedRecordBatchIterator>,
1085) -> Result<(SchemaRef, Vec<BoxedRecordBatchIterator>)> {
1086 let Some(schemas) = schemas else {
1087 return Ok((schema, input_iters));
1088 };
1089
1090 let aligner = Json2Aligner::try_new(schemas)?;
1091 let input_iters = input_iters
1092 .into_iter()
1093 .map(|input_iter| aligner.wrap_iter(input_iter))
1094 .collect();
1095
1096 Ok((aligner.schema().clone(), input_iters))
1097}
1098
1099pub fn merge_and_dedup(
1144 schema: &SchemaRef,
1145 append_mode: bool,
1146 merge_mode: MergeMode,
1147 field_column_start: usize,
1148 input_iters: Vec<BoxedRecordBatchIterator>,
1149) -> Result<BoxedRecordBatchIterator> {
1150 merge_and_dedup_with_batch_size(
1151 schema,
1152 append_mode,
1153 merge_mode,
1154 field_column_start,
1155 input_iters,
1156 DEFAULT_READ_BATCH_SIZE,
1157 )
1158}
1159
1160pub fn merge_and_dedup_with_batch_size(
1166 schema: &SchemaRef,
1167 append_mode: bool,
1168 merge_mode: MergeMode,
1169 field_column_start: usize,
1170 input_iters: Vec<BoxedRecordBatchIterator>,
1171 batch_size: usize,
1172) -> Result<BoxedRecordBatchIterator> {
1173 let batch_size = batch_size.max(1);
1174 let merge_iter = FlatMergeIterator::new(schema.clone(), input_iters, batch_size)?;
1175 let maybe_dedup = if append_mode {
1176 Box::new(merge_iter) as _
1178 } else {
1179 match merge_mode {
1181 MergeMode::LastRow => {
1182 Box::new(FlatDedupIterator::new(merge_iter, FlatLastRow::new(false))) as _
1183 }
1184 MergeMode::LastNonNull => Box::new(FlatDedupIterator::new(
1185 merge_iter,
1186 FlatLastNonNull::new(field_column_start, false),
1187 )) as _,
1188 }
1189 };
1190 Ok(maybe_dedup)
1191}
1192
1193pub fn maybe_dedup_one(
1194 append_mode: bool,
1195 merge_mode: MergeMode,
1196 field_column_start: usize,
1197 input_iter: BoxedRecordBatchIterator,
1198) -> BoxedRecordBatchIterator {
1199 if append_mode {
1200 input_iter
1202 } else {
1203 match merge_mode {
1205 MergeMode::LastRow => {
1206 Box::new(FlatDedupIterator::new(input_iter, FlatLastRow::new(false)))
1207 }
1208 MergeMode::LastNonNull => Box::new(FlatDedupIterator::new(
1209 input_iter,
1210 FlatLastNonNull::new(field_column_start, false),
1211 )),
1212 }
1213 }
1214}
1215
1216pub(crate) struct FlushScheduler {
1218 region_status: HashMap<RegionId, FlushStatus>,
1220 scheduler: SchedulerRef,
1222}
1223
1224impl FlushScheduler {
1225 pub(crate) fn new(scheduler: SchedulerRef) -> FlushScheduler {
1227 FlushScheduler {
1228 region_status: HashMap::new(),
1229 scheduler,
1230 }
1231 }
1232
1233 pub(crate) fn is_flush_requested(&self, region_id: RegionId) -> bool {
1235 self.region_status.contains_key(®ion_id)
1236 }
1237
1238 fn schedule_flush_task(
1239 &mut self,
1240 version_control: &VersionControlRef,
1241 mut task: RegionFlushTask,
1242 ) -> Result<CancellableTaskState> {
1243 let region_id = task.region_id;
1244
1245 if let Err(e) = version_control.freeze_mutable() {
1247 error!(e; "Failed to freeze the mutable memtable for region {}", region_id);
1248 task.on_failure(Arc::new(RegionBusySnafu { region_id }.build()));
1249
1250 return Err(e);
1251 }
1252 let state = CancellableTaskState::new();
1254 let (job, waiters) = task.into_flush_job(version_control, state.clone());
1255 if let Err(e) = self.scheduler.schedule(job) {
1256 error!(e; "Failed to schedule flush job for region {}", region_id);
1257 waiters.on_failure(Arc::new(RegionBusySnafu { region_id }.build()));
1258
1259 return Err(e);
1260 }
1261 Ok(state)
1262 }
1263
1264 pub(crate) fn schedule_flush(
1266 &mut self,
1267 region_id: RegionId,
1268 version_control: &VersionControlRef,
1269 mut task: RegionFlushTask,
1270 ) -> Result<()> {
1271 debug_assert_eq!(region_id, task.region_id);
1272
1273 let version = version_control.current().version;
1274 if version.memtables.is_empty() {
1275 debug_assert!(!self.region_status.contains_key(®ion_id));
1276 task.on_success();
1278 return Ok(());
1279 }
1280
1281 FLUSH_REQUESTS_TOTAL
1283 .with_label_values(&[task.reason.as_str()])
1284 .inc();
1285
1286 if let Some(flush_status) = self.region_status.get_mut(®ion_id) {
1288 if flush_status.has_pending_lifecycle_ddl() {
1289 task.on_failure(Arc::new(FlushCancelledSnafu.build()));
1290 return Ok(());
1291 }
1292 debug!("Merging flush task for region {}", region_id);
1294 flush_status.merge_task(task);
1295 return Ok(());
1296 }
1297
1298 let closing = task.reason == FlushReason::Closing;
1299 let state = self.schedule_flush_task(version_control, task)?;
1300
1301 let _ = self.region_status.insert(
1303 region_id,
1304 FlushStatus::new(region_id, version_control.clone(), state, closing),
1305 );
1306
1307 Ok(())
1308 }
1309
1310 pub(crate) fn on_flush_success(
1314 &mut self,
1315 region_id: RegionId,
1316 ) -> Option<(
1317 Vec<SenderDdlRequest>,
1318 Vec<SenderWriteRequest>,
1319 Vec<SenderBulkRequest>,
1320 )> {
1321 let flush_status = self.region_status.get_mut(®ion_id)?;
1322 if flush_status.pending_task.is_none() {
1324 debug!(
1327 "Region {} doesn't have any pending flush task, removing it from the status",
1328 region_id
1329 );
1330 let flush_status = self.region_status.remove(®ion_id).unwrap();
1331 return Some((
1332 flush_status.pending_ddls,
1333 flush_status.pending_writes,
1334 flush_status.pending_bulk_writes,
1335 ));
1336 }
1337
1338 let version_data = flush_status.version_control.current();
1340 if version_data.version.memtables.is_empty() {
1341 let task = flush_status.pending_task.take().unwrap();
1344 task.on_success();
1346 debug!(
1347 "Region {} has nothing to flush, removing it from the status",
1348 region_id
1349 );
1350 let flush_status = self.region_status.remove(®ion_id).unwrap();
1352 return Some((
1353 flush_status.pending_ddls,
1354 flush_status.pending_writes,
1355 flush_status.pending_bulk_writes,
1356 ));
1357 }
1358
1359 debug!("Scheduling pending flush task for region {}", region_id);
1361 let task = flush_status.pending_task.take().unwrap();
1363 let version_control = flush_status.version_control.clone();
1364 match self.schedule_flush_task(&version_control, task) {
1365 Ok(state) => {
1366 self.region_status.get_mut(®ion_id).unwrap().state = state;
1367 }
1368 Err(err) => {
1369 error!(
1370 err;
1371 "Flush succeeded for region {region_id}, but failed to schedule next flush for it."
1372 );
1373 let flush_status = self.region_status.remove(®ion_id).unwrap();
1374 flush_status.fail_all(Arc::new(RegionBusySnafu { region_id }.build()));
1375 return None;
1376 }
1377 }
1378 None
1380 }
1381
1382 pub(crate) fn on_flush_failed(
1387 &mut self,
1388 region_id: RegionId,
1389 err: Arc<Error>,
1390 ) -> Vec<SenderDdlRequest> {
1391 if matches!(err.as_ref(), Error::FlushCancelled { .. }) {
1392 info!("Region {} flush was cancelled", region_id);
1393 } else {
1394 error!(err; "Region {} failed to flush, cancel all pending tasks", region_id);
1395 FLUSH_FAILURE_TOTAL.inc();
1396 }
1397
1398 let Some(flush_status) = self.region_status.remove(®ion_id) else {
1400 return Vec::new();
1401 };
1402
1403 flush_status.on_failure(err)
1404 }
1405
1406 pub(crate) fn try_cancel_and_add_ddl<T>(
1408 &mut self,
1409 region_id: RegionId,
1410 sender: OptionOutputTx,
1411 request: T,
1412 into_ddl_request: impl FnOnce(T) -> DdlRequest,
1413 ) -> std::result::Result<(), (OptionOutputTx, T)> {
1414 let Some(status) = self.region_status.get_mut(®ion_id) else {
1415 return Err((sender, request));
1416 };
1417
1418 let cancel_result = status.state.request_cancel();
1419 debug!(
1420 "Requested flush cancellation for region {}, result: {:?}",
1421 region_id, cancel_result
1422 );
1423 status.pending_ddls.push(SenderDdlRequest {
1424 region_id,
1425 sender,
1426 request: into_ddl_request(request),
1427 });
1428 Ok(())
1429 }
1430
1431 pub(crate) fn on_region_dropped(&mut self, region_id: RegionId) {
1433 self.remove_region_on_failure(
1434 region_id,
1435 Arc::new(RegionDroppedSnafu { region_id }.build()),
1436 );
1437 }
1438
1439 pub(crate) fn on_region_closed(&mut self, region_id: RegionId) {
1441 let Some(flush_status) = self.region_status.remove(®ion_id) else {
1442 return;
1443 };
1444
1445 flush_status.on_region_closed(Arc::new(RegionClosedSnafu { region_id }.build()));
1446 }
1447
1448 pub(crate) fn on_region_truncated(&mut self, region_id: RegionId) {
1450 self.remove_region_on_failure(
1451 region_id,
1452 Arc::new(RegionTruncatedSnafu { region_id }.build()),
1453 );
1454 }
1455
1456 fn remove_region_on_failure(&mut self, region_id: RegionId, err: Arc<Error>) {
1457 let Some(flush_status) = self.region_status.remove(®ion_id) else {
1459 return;
1460 };
1461
1462 flush_status.fail_all(err);
1464 }
1465
1466 pub(crate) fn add_ddl_request_to_pending(&mut self, request: SenderDdlRequest) {
1471 let status = self.region_status.get_mut(&request.region_id).unwrap();
1472 status.pending_ddls.push(request);
1473 }
1474
1475 pub(crate) fn add_write_request_to_pending(&mut self, request: SenderWriteRequest) {
1480 let status = self
1481 .region_status
1482 .get_mut(&request.request.region_id)
1483 .unwrap();
1484 status.pending_writes.push(request);
1485 }
1486
1487 pub(crate) fn add_bulk_request_to_pending(&mut self, request: SenderBulkRequest) {
1492 let status = self.region_status.get_mut(&request.region_id).unwrap();
1493 status.pending_bulk_writes.push(request);
1494 }
1495
1496 pub(crate) fn has_pending_ddls(&self, region_id: RegionId) -> bool {
1498 self.region_status
1499 .get(®ion_id)
1500 .map(|status| !status.pending_ddls.is_empty() || status.closing)
1501 .unwrap_or(false)
1502 }
1503}
1504
1505impl Drop for FlushScheduler {
1506 fn drop(&mut self) {
1507 for (region_id, flush_status) in self.region_status.drain() {
1508 flush_status.fail_all(Arc::new(RegionClosedSnafu { region_id }.build()));
1510 }
1511 }
1512}
1513
1514struct FlushStatus {
1518 region_id: RegionId,
1520 version_control: VersionControlRef,
1522 state: CancellableTaskState,
1524 pending_task: Option<RegionFlushTask>,
1526 closing: bool,
1528 pending_ddls: Vec<SenderDdlRequest>,
1530 pending_writes: Vec<SenderWriteRequest>,
1532 pending_bulk_writes: Vec<SenderBulkRequest>,
1534}
1535
1536impl FlushStatus {
1537 fn new(
1538 region_id: RegionId,
1539 version_control: VersionControlRef,
1540 state: CancellableTaskState,
1541 closing: bool,
1542 ) -> FlushStatus {
1543 FlushStatus {
1544 region_id,
1545 version_control,
1546 state,
1547 pending_task: None,
1548 closing,
1549 pending_ddls: Vec::new(),
1550 pending_writes: Vec::new(),
1551 pending_bulk_writes: Vec::new(),
1552 }
1553 }
1554
1555 fn merge_task(&mut self, task: RegionFlushTask) {
1557 self.closing |= task.reason == FlushReason::Closing;
1558 if let Some(pending) = &mut self.pending_task {
1559 pending.merge(task);
1560 } else {
1561 self.pending_task = Some(task);
1562 }
1563 }
1564
1565 fn has_pending_lifecycle_ddl(&self) -> bool {
1566 self.pending_ddls
1567 .iter()
1568 .any(|ddl| matches!(ddl.request, DdlRequest::Drop(_) | DdlRequest::Truncate(_)))
1569 }
1570
1571 fn on_failure(self, err: Arc<Error>) -> Vec<SenderDdlRequest> {
1573 if let Some(mut task) = self.pending_task {
1574 task.on_failure(err.clone());
1575 }
1576 let mut lifecycle_ddls = Vec::new();
1577 for ddl in self.pending_ddls {
1578 if matches!(ddl.request, DdlRequest::Drop(_) | DdlRequest::Truncate(_)) {
1579 lifecycle_ddls.push(ddl);
1580 } else {
1581 ddl.sender.send(Err(err.clone()).context(FlushRegionSnafu {
1582 region_id: self.region_id,
1583 }));
1584 }
1585 }
1586 for write_req in self.pending_writes {
1587 write_req
1588 .sender
1589 .send(Err(err.clone()).context(FlushRegionSnafu {
1590 region_id: self.region_id,
1591 }));
1592 }
1593 for bulk_req in self.pending_bulk_writes {
1594 bulk_req
1595 .sender
1596 .send(Err(err.clone()).context(FlushRegionSnafu {
1597 region_id: self.region_id,
1598 }));
1599 }
1600 lifecycle_ddls
1601 }
1602
1603 fn fail_all(self, err: Arc<Error>) {
1604 let region_id = self.region_id;
1605 let ddls = self.on_failure(err.clone());
1606 for ddl in ddls {
1607 ddl.sender
1608 .send(Err(err.clone()).context(FlushRegionSnafu { region_id }));
1609 }
1610 }
1611
1612 fn on_region_closed(self, err: Arc<Error>) {
1613 if let Some(mut task) = self.pending_task {
1614 task.on_failure(err.clone());
1615 }
1616
1617 for ddl in self.pending_ddls {
1618 if matches!(ddl.request, DdlRequest::Close(_)) {
1619 ddl.sender.send(Ok(0));
1620 } else {
1621 ddl.sender.send(Err(err.clone()).context(FlushRegionSnafu {
1622 region_id: self.region_id,
1623 }));
1624 }
1625 }
1626
1627 for write_req in self.pending_writes {
1628 write_req
1629 .sender
1630 .send(Err(err.clone()).context(FlushRegionSnafu {
1631 region_id: self.region_id,
1632 }));
1633 }
1634 for bulk_req in self.pending_bulk_writes {
1635 bulk_req
1636 .sender
1637 .send(Err(err.clone()).context(FlushRegionSnafu {
1638 region_id: self.region_id,
1639 }));
1640 }
1641 }
1642}
1643
1644#[cfg(test)]
1645mod tests {
1646 use api::v1::{OpType, Rows};
1647 use common_error::ext::ErrorExt;
1648 use common_error::status_code::StatusCode;
1649 use mito_codec::row_converter::build_primary_key_codec;
1650 use tokio::sync::oneshot;
1651
1652 use super::*;
1653 use crate::cache::CacheManager;
1654 use crate::error::InvalidSchedulerStateSnafu;
1655 use crate::memtable::bulk::part::BulkPartConverter;
1656 use crate::memtable::time_series::TimeSeriesMemtableBuilder;
1657 use crate::memtable::{Memtable, RangesOptions};
1658 use crate::request::WriteRequest;
1659 use crate::schedule::scheduler::Scheduler;
1660 use crate::sst::{FlatSchemaOptions, to_flat_sst_arrow_schema};
1661 use crate::test_util::memtable_util::{build_key_values_with_ts_seq_values, metadata_for_test};
1662 use crate::test_util::scheduler_util::{SchedulerEnv, VecScheduler};
1663 use crate::test_util::version_util::{VersionControlBuilder, write_rows_to_version};
1664
1665 struct FailingScheduler;
1666
1667 #[async_trait::async_trait]
1668 impl Scheduler for FailingScheduler {
1669 fn schedule(&self, _job: Job) -> Result<()> {
1670 InvalidSchedulerStateSnafu.fail()
1671 }
1672
1673 async fn stop(&self, _await_termination: bool) -> Result<()> {
1674 Ok(())
1675 }
1676 }
1677
1678 fn new_test_flush_task(
1679 env: &SchedulerEnv,
1680 region_id: RegionId,
1681 reason: FlushReason,
1682 request_sender: mpsc::Sender<WorkerRequestWithTime>,
1683 manifest_ctx: ManifestContextRef,
1684 ) -> RegionFlushTask {
1685 RegionFlushTask {
1686 region_id,
1687 reason,
1688 senders: Vec::new(),
1689 request_sender,
1690 access_layer: env.access_layer.clone(),
1691 listener: WorkerListener::default(),
1692 engine_config: Arc::new(MitoConfig::default()),
1693 row_group_size: None,
1694 cache_manager: Arc::new(CacheManager::default()),
1695 manifest_ctx,
1696 index_options: IndexOptions::default(),
1697 flush_semaphore: Arc::new(Semaphore::new(2)),
1698 is_staging: false,
1699 partition_expr: None,
1700 }
1701 }
1702
1703 fn new_test_bulk_request(
1704 region_id: RegionId,
1705 ) -> (
1706 SenderBulkRequest,
1707 oneshot::Receiver<Result<store_api::region_request::AffectedRows>>,
1708 ) {
1709 let metadata = metadata_for_test();
1710 let schema = to_flat_sst_arrow_schema(
1711 &metadata,
1712 &FlatSchemaOptions::from_encoding(metadata.primary_key_encoding),
1713 );
1714 let pk_codec = build_primary_key_codec(&metadata);
1715 let mut converter = BulkPartConverter::new(&metadata, schema, 16, pk_codec, true);
1716 let kvs = build_key_values_with_ts_seq_values(
1717 &metadata,
1718 "bulk_key".to_string(),
1719 1,
1720 std::iter::once(1000i64),
1721 std::iter::once(Some(1.0f64)),
1722 1,
1723 );
1724 converter.append_key_values(&kvs).unwrap();
1725 let (sender, receiver) = oneshot::channel();
1726
1727 (
1728 SenderBulkRequest {
1729 sender: OptionOutputTx::from(sender),
1730 region_id,
1731 request: converter.convert().unwrap(),
1732 region_metadata: Some(metadata),
1733 partition_expr_version: None,
1734 },
1735 receiver,
1736 )
1737 }
1738
1739 fn new_test_write_request(
1740 region_id: RegionId,
1741 ) -> (
1742 SenderWriteRequest,
1743 oneshot::Receiver<Result<store_api::region_request::AffectedRows>>,
1744 ) {
1745 let (sender, receiver) = oneshot::channel();
1746 let request = WriteRequest::new(region_id, OpType::Put, Rows::default(), None).unwrap();
1747 (
1748 SenderWriteRequest {
1749 sender: OptionOutputTx::from(sender),
1750 request,
1751 },
1752 receiver,
1753 )
1754 }
1755
1756 #[test]
1757 fn test_get_mutable_limit() {
1758 assert_eq!(4, WriteBufferManagerImpl::get_mutable_limit(8));
1759 assert_eq!(5, WriteBufferManagerImpl::get_mutable_limit(10));
1760 assert_eq!(32, WriteBufferManagerImpl::get_mutable_limit(64));
1761 assert_eq!(0, WriteBufferManagerImpl::get_mutable_limit(0));
1762 }
1763
1764 #[test]
1765 fn test_over_mutable_limit() {
1766 let manager = WriteBufferManagerImpl::new(1000);
1768 manager.reserve_mem(400);
1769 assert!(!manager.should_flush_engine());
1770 assert!(!manager.should_stall());
1771
1772 manager.reserve_mem(400);
1774 assert!(manager.should_flush_engine());
1775
1776 manager.schedule_free_mem(400);
1778 assert!(!manager.should_flush_engine());
1779 assert_eq!(800, manager.memory_used.load(Ordering::Relaxed));
1780 assert_eq!(400, manager.memory_active.load(Ordering::Relaxed));
1781
1782 manager.free_mem(400);
1784 assert_eq!(400, manager.memory_used.load(Ordering::Relaxed));
1785 assert_eq!(400, manager.memory_active.load(Ordering::Relaxed));
1786 }
1787
1788 #[test]
1789 fn test_over_global() {
1790 let manager = WriteBufferManagerImpl::new(1000);
1792 manager.reserve_mem(1100);
1793 assert!(manager.should_stall());
1794 manager.schedule_free_mem(200);
1796 assert!(manager.should_flush_engine());
1797 assert!(manager.should_stall());
1798
1799 manager.schedule_free_mem(450);
1801 assert!(manager.should_flush_engine());
1802 assert!(manager.should_stall());
1803
1804 manager.reserve_mem(50);
1806 assert!(manager.should_flush_engine());
1807 manager.reserve_mem(100);
1808 assert!(manager.should_flush_engine());
1809 }
1810
1811 #[test]
1812 fn test_manager_notify() {
1813 let (sender, receiver) = watch::channel(());
1814 let manager = WriteBufferManagerImpl::new(1000).with_notifier(sender);
1815 manager.reserve_mem(500);
1816 assert!(!receiver.has_changed().unwrap());
1817 manager.schedule_free_mem(500);
1818 assert!(!receiver.has_changed().unwrap());
1819 manager.free_mem(500);
1820 assert!(receiver.has_changed().unwrap());
1821 }
1822
1823 #[tokio::test]
1824 async fn test_schedule_empty() {
1825 let env = SchedulerEnv::new().await;
1826 let (tx, _rx) = mpsc::channel(4);
1827 let mut scheduler = env.mock_flush_scheduler();
1828 let builder = VersionControlBuilder::new();
1829
1830 let version_control = Arc::new(builder.build());
1831 let (output_tx, output_rx) = oneshot::channel();
1832 let mut task = RegionFlushTask {
1833 region_id: builder.region_id(),
1834 reason: FlushReason::Manual,
1835 senders: Vec::new(),
1836 request_sender: tx,
1837 access_layer: env.access_layer.clone(),
1838 listener: WorkerListener::default(),
1839 engine_config: Arc::new(MitoConfig::default()),
1840 row_group_size: None,
1841 cache_manager: Arc::new(CacheManager::default()),
1842 manifest_ctx: env
1843 .mock_manifest_context(version_control.current().version.metadata.clone())
1844 .await,
1845 index_options: IndexOptions::default(),
1846 flush_semaphore: Arc::new(Semaphore::new(2)),
1847 is_staging: false,
1848 partition_expr: None,
1849 };
1850 task.push_sender(OptionOutputTx::from(output_tx));
1851 scheduler
1852 .schedule_flush(builder.region_id(), &version_control, task)
1853 .unwrap();
1854 assert!(scheduler.region_status.is_empty());
1855 let output = output_rx.await.unwrap().unwrap();
1856 assert_eq!(output, 0);
1857 assert!(scheduler.region_status.is_empty());
1858 }
1859
1860 #[tokio::test]
1861 async fn test_schedule_flush_failure_notifies_waiter() {
1862 let env = SchedulerEnv::new()
1863 .await
1864 .scheduler(Arc::new(FailingScheduler));
1865 let (tx, _rx) = mpsc::channel(4);
1866 let mut scheduler = env.mock_flush_scheduler();
1867 let mut builder = VersionControlBuilder::new();
1868 builder.set_memtable_builder(Arc::new(TimeSeriesMemtableBuilder::default()));
1869 let version_control = Arc::new(builder.build());
1870 let version_data = version_control.current();
1871 write_rows_to_version(&version_data.version, "host0", 0, 10);
1872 let manifest_ctx = env
1873 .mock_manifest_context(version_data.version.metadata.clone())
1874 .await;
1875 let (output_tx, output_rx) = oneshot::channel();
1876 let mut task = new_test_flush_task(
1877 &env,
1878 builder.region_id(),
1879 FlushReason::Manual,
1880 tx,
1881 manifest_ctx,
1882 );
1883 task.push_sender(OptionOutputTx::from(output_tx));
1884
1885 scheduler
1886 .schedule_flush(builder.region_id(), &version_control, task)
1887 .unwrap_err();
1888
1889 let err = output_rx
1890 .await
1891 .expect("waiter must receive explicit error")
1892 .unwrap_err();
1893 assert_eq!(err.status_code(), StatusCode::RegionBusy);
1894 }
1895
1896 #[tokio::test]
1897 async fn test_send_worker_request_failure_notifies_flush_finished_waiter() {
1898 let env = SchedulerEnv::new().await;
1899 let (tx, rx) = mpsc::channel(1);
1900 drop(rx);
1901 let builder = VersionControlBuilder::new();
1902 let version_control = Arc::new(builder.build());
1903 let manifest_ctx = env
1904 .mock_manifest_context(version_control.current().version.metadata.clone())
1905 .await;
1906 let task = new_test_flush_task(
1907 &env,
1908 builder.region_id(),
1909 FlushReason::Manual,
1910 tx,
1911 manifest_ctx,
1912 );
1913 let (output_tx, output_rx) = oneshot::channel();
1914 let request = WorkerRequest::Background {
1915 region_id: builder.region_id(),
1916 notify: BackgroundNotify::FlushFinished(FlushFinished {
1917 region_id: builder.region_id(),
1918 flush_reason: FlushReason::Manual,
1919 flushed_entry_id: 0,
1920 senders: vec![OutputTx::new(output_tx)],
1921 _timer: FLUSH_ELAPSED.with_label_values(&["total"]).start_timer(),
1922 edit: RegionEdit {
1923 files_to_add: Vec::new(),
1924 files_to_remove: Vec::new(),
1925 timestamp_ms: None,
1926 compaction_time_window: None,
1927 flushed_entry_id: None,
1928 flushed_sequence: None,
1929 committed_sequence: None,
1930 },
1931 memtables_to_remove: smallvec![],
1932 is_staging: false,
1933 }),
1934 };
1935
1936 task.send_worker_request(request).await;
1937
1938 let output = output_rx.await.expect("waiter must receive explicit error");
1939 assert!(output.is_err());
1940 }
1941
1942 #[tokio::test]
1943 async fn test_flush_waiters_drop_notifies_waiter() {
1944 let region_id = RegionId::new(1, 1);
1945 let (output_tx, output_rx) = oneshot::channel();
1946 let waiters = FlushTaskWaiters::new(region_id, vec![OutputTx::new(output_tx)]);
1947
1948 drop(waiters);
1949
1950 let err = output_rx
1951 .await
1952 .expect("waiter must receive explicit error")
1953 .unwrap_err();
1954 assert_eq!(err.status_code(), StatusCode::RegionBusy);
1955 }
1956
1957 #[tokio::test]
1958 async fn test_flush_failure_notifies_pending_bulk_writes() {
1959 let region_id = RegionId::new(1, 1);
1960 let version_control = Arc::new(VersionControlBuilder::new().build());
1961 let (bulk_req, output_rx) = new_test_bulk_request(region_id);
1962 let status = FlushStatus {
1963 region_id,
1964 version_control,
1965 state: CancellableTaskState::new(),
1966 pending_task: None,
1967 closing: false,
1968 pending_ddls: Vec::new(),
1969 pending_writes: Vec::new(),
1970 pending_bulk_writes: vec![bulk_req],
1971 };
1972
1973 let pending_ddls = status.on_failure(Arc::new(RegionClosedSnafu { region_id }.build()));
1974 assert!(pending_ddls.is_empty());
1975
1976 let err = output_rx
1977 .await
1978 .expect("pending bulk write must receive explicit error")
1979 .unwrap_err();
1980 assert_eq!(err.status_code(), StatusCode::Cancelled);
1981 }
1982
1983 #[tokio::test]
1984 async fn test_flush_failure_retains_lifecycle_ddls_and_fails_pending_writes() {
1985 let region_id = RegionId::new(1, 1);
1986 let version_control = Arc::new(VersionControlBuilder::new().build());
1987 let (write_req, write_rx) = new_test_write_request(region_id);
1988 let (bulk_req, bulk_rx) = new_test_bulk_request(region_id);
1989 let (truncate_tx, mut truncate_rx) = oneshot::channel();
1990 let (drop_tx, mut drop_rx) = oneshot::channel();
1991 let status = FlushStatus {
1992 region_id,
1993 version_control,
1994 state: CancellableTaskState::new(),
1995 pending_task: None,
1996 closing: false,
1997 pending_ddls: vec![
1998 SenderDdlRequest {
1999 region_id,
2000 sender: OptionOutputTx::from(truncate_tx),
2001 request: DdlRequest::Truncate(
2002 store_api::region_request::RegionTruncateRequest::All,
2003 ),
2004 },
2005 SenderDdlRequest {
2006 region_id,
2007 sender: OptionOutputTx::from(drop_tx),
2008 request: DdlRequest::Drop(store_api::region_request::RegionDropRequest {
2009 fast_path: false,
2010 force: false,
2011 partial_drop: false,
2012 }),
2013 },
2014 ],
2015 pending_writes: vec![write_req],
2016 pending_bulk_writes: vec![bulk_req],
2017 };
2018
2019 let ddls = status.on_failure(Arc::new(RegionBusySnafu { region_id }.build()));
2020 assert_eq!(2, ddls.len());
2021 assert!(matches!(ddls[0].request, DdlRequest::Truncate(_)));
2022 assert!(matches!(ddls[1].request, DdlRequest::Drop(_)));
2023
2024 let write_err = write_rx
2025 .await
2026 .expect("pending write must receive explicit error")
2027 .unwrap_err();
2028 assert_eq!(write_err.status_code(), StatusCode::RegionBusy);
2029 let bulk_err = bulk_rx
2030 .await
2031 .expect("pending bulk write must receive explicit error")
2032 .unwrap_err();
2033 assert_eq!(bulk_err.status_code(), StatusCode::RegionBusy);
2034
2035 assert!(truncate_rx.try_recv().is_err());
2036 assert!(drop_rx.try_recv().is_err());
2037 for ddl in ddls {
2038 ddl.sender.send(Ok(0));
2039 }
2040 assert_eq!(0, truncate_rx.await.unwrap().unwrap());
2041 assert_eq!(0, drop_rx.await.unwrap().unwrap());
2042 }
2043
2044 #[tokio::test]
2045 async fn test_uncommitted_ssts_cleanup_finalized_file() {
2046 let env = SchedulerEnv::new().await;
2047 let region_id = RegionId::new(1, 1);
2048 let file_id = store_api::storage::FileId::random();
2049 let path = crate::sst::location::sst_file_path(
2050 env.access_layer.table_dir(),
2051 crate::sst::file::RegionFileId::new(region_id, file_id),
2052 env.access_layer.path_type(),
2053 );
2054 env.access_layer
2055 .object_store()
2056 .write(&path, Bytes::from_static(b"sst"))
2057 .await
2058 .unwrap();
2059 assert!(env.access_layer.object_store().exists(&path).await.unwrap());
2060
2061 let uncommitted = UncommittedSsts::new(region_id, env.access_layer.clone(), None);
2062 uncommitted.track(&[SstInfo {
2063 file_id,
2064 ..Default::default()
2065 }]);
2066 assert_eq!(1, uncommitted.num_tracked_files());
2067 uncommitted.cleanup_for_test().await.unwrap();
2068 assert_eq!(0, uncommitted.num_tracked_files());
2069
2070 let entries = env
2073 .access_layer
2074 .object_store()
2075 .list(&env.access_layer.build_region_dir(region_id))
2076 .await
2077 .unwrap();
2078 assert!(entries.iter().all(|entry| entry.path() != path));
2079 }
2080
2081 #[tokio::test]
2082 async fn test_region_closed_notifies_pending_bulk_writes() {
2083 let region_id = RegionId::new(1, 1);
2084 let version_control = Arc::new(VersionControlBuilder::new().build());
2085 let (bulk_req, output_rx) = new_test_bulk_request(region_id);
2086 let status = FlushStatus {
2087 region_id,
2088 version_control,
2089 state: CancellableTaskState::new(),
2090 pending_task: None,
2091 closing: false,
2092 pending_ddls: Vec::new(),
2093 pending_writes: Vec::new(),
2094 pending_bulk_writes: vec![bulk_req],
2095 };
2096
2097 status.on_region_closed(Arc::new(RegionClosedSnafu { region_id }.build()));
2098
2099 let err = output_rx
2100 .await
2101 .expect("pending bulk write must receive explicit error")
2102 .unwrap_err();
2103 assert_eq!(err.status_code(), StatusCode::Cancelled);
2104 }
2105
2106 #[tokio::test]
2107 async fn test_schedule_pending_request() {
2108 let job_scheduler = Arc::new(VecScheduler::default());
2109 let env = SchedulerEnv::new().await.scheduler(job_scheduler.clone());
2110 let (tx, _rx) = mpsc::channel(4);
2111 let mut scheduler = env.mock_flush_scheduler();
2112 let mut builder = VersionControlBuilder::new();
2113 builder.set_memtable_builder(Arc::new(TimeSeriesMemtableBuilder::default()));
2115 let version_control = Arc::new(builder.build());
2116 let version_data = version_control.current();
2118 write_rows_to_version(&version_data.version, "host0", 0, 10);
2119 let manifest_ctx = env
2120 .mock_manifest_context(version_data.version.metadata.clone())
2121 .await;
2122 let mut tasks: Vec<_> = (0..3)
2124 .map(|_| RegionFlushTask {
2125 region_id: builder.region_id(),
2126 reason: FlushReason::Manual,
2127 senders: Vec::new(),
2128 request_sender: tx.clone(),
2129 access_layer: env.access_layer.clone(),
2130 listener: WorkerListener::default(),
2131 engine_config: Arc::new(MitoConfig::default()),
2132 row_group_size: None,
2133 cache_manager: Arc::new(CacheManager::default()),
2134 manifest_ctx: manifest_ctx.clone(),
2135 index_options: IndexOptions::default(),
2136 flush_semaphore: Arc::new(Semaphore::new(2)),
2137 is_staging: false,
2138 partition_expr: None,
2139 })
2140 .collect();
2141 let task = tasks.pop().unwrap();
2143 scheduler
2144 .schedule_flush(builder.region_id(), &version_control, task)
2145 .unwrap();
2146 assert_eq!(1, scheduler.region_status.len());
2148 assert_eq!(1, job_scheduler.num_jobs());
2149 let version_data = version_control.current();
2151 assert_eq!(0, version_data.version.memtables.immutables()[0].id());
2152 let output_rxs: Vec<_> = tasks
2154 .into_iter()
2155 .map(|mut task| {
2156 let (output_tx, output_rx) = oneshot::channel();
2157 task.push_sender(OptionOutputTx::from(output_tx));
2158 scheduler
2159 .schedule_flush(builder.region_id(), &version_control, task)
2160 .unwrap();
2161 output_rx
2162 })
2163 .collect();
2164 version_control.apply_edit(
2166 Some(RegionEdit {
2167 files_to_add: Vec::new(),
2168 files_to_remove: Vec::new(),
2169 timestamp_ms: None,
2170 compaction_time_window: None,
2171 flushed_entry_id: None,
2172 flushed_sequence: None,
2173 committed_sequence: None,
2174 }),
2175 &[0],
2176 builder.file_purger(),
2177 );
2178 scheduler.on_flush_success(builder.region_id());
2179 assert_eq!(1, job_scheduler.num_jobs());
2181 assert!(scheduler.region_status.is_empty());
2183 for output_rx in output_rxs {
2184 let output = output_rx.await.unwrap().unwrap();
2185 assert_eq!(output, 0);
2186 }
2187 }
2188
2189 #[test]
2191 fn test_memtable_flat_sources_single_range_append_mode_behavior() {
2192 let metadata = metadata_for_test();
2194 let schema = to_flat_sst_arrow_schema(
2195 &metadata,
2196 &FlatSchemaOptions::from_encoding(metadata.primary_key_encoding),
2197 );
2198
2199 let capacity = 16;
2202 let pk_codec = build_primary_key_codec(&metadata);
2203 let mut converter =
2204 BulkPartConverter::new(&metadata, schema.clone(), capacity, pk_codec, true);
2205 let kvs = build_key_values_with_ts_seq_values(
2206 &metadata,
2207 "dup_key".to_string(),
2208 1,
2209 vec![1000i64, 1000i64].into_iter(),
2210 vec![Some(1.0f64), Some(2.0f64)].into_iter(),
2211 1,
2212 );
2213 converter.append_key_values(&kvs).unwrap();
2214 let part = converter.convert().unwrap();
2215
2216 let build_ranges = |append_mode: bool| -> MemtableRanges {
2219 let memtable = crate::memtable::bulk::BulkMemtable::new(
2220 1,
2221 crate::memtable::bulk::BulkMemtableConfig::default(),
2222 metadata.clone(),
2223 None,
2224 None,
2225 append_mode,
2226 MergeMode::LastRow,
2227 );
2228 memtable.write_bulk(part.clone()).unwrap();
2229 memtable.ranges(None, RangesOptions::for_flush()).unwrap()
2230 };
2231
2232 {
2234 let mem_ranges = build_ranges(false);
2235 assert_eq!(1, mem_ranges.ranges.len());
2236
2237 let options = RegionOptions {
2238 append_mode: false,
2239 merge_mode: Some(MergeMode::LastRow),
2240 ..Default::default()
2241 };
2242
2243 let flat_sources = memtable_flat_sources(
2244 schema.clone(),
2245 mem_ranges,
2246 &options,
2247 metadata.primary_key.len(),
2248 )
2249 .unwrap();
2250 assert!(flat_sources.encoded.is_empty());
2251 assert_eq!(1, flat_sources.sources.len());
2252
2253 let mut total_rows = 0usize;
2255 for (source, _sequence) in flat_sources.sources {
2256 total_rows += source
2257 .take_iter()
2258 .map(|x| x.unwrap().num_rows())
2259 .sum::<usize>();
2260 }
2261 assert_eq!(1, total_rows, "dedup should keep a single row");
2262 }
2263
2264 {
2266 let mem_ranges = build_ranges(true);
2267 assert_eq!(1, mem_ranges.ranges.len());
2268
2269 let options = RegionOptions {
2270 append_mode: true,
2271 ..Default::default()
2272 };
2273
2274 let flat_sources =
2275 memtable_flat_sources(schema, mem_ranges, &options, metadata.primary_key.len())
2276 .unwrap();
2277 assert!(flat_sources.encoded.is_empty());
2278 assert_eq!(1, flat_sources.sources.len());
2279
2280 let mut total_rows = 0usize;
2281 for (source, _sequence) in flat_sources.sources {
2282 total_rows += source
2283 .take_iter()
2284 .map(|x| x.unwrap().num_rows())
2285 .sum::<usize>();
2286 }
2287 assert_eq!(2, total_rows, "append_mode should preserve duplicates");
2288 }
2289 }
2290
2291 #[tokio::test]
2292 async fn test_schedule_pending_request_on_flush_success() {
2293 common_telemetry::init_default_ut_logging();
2294 let job_scheduler = Arc::new(VecScheduler::default());
2295 let env = SchedulerEnv::new().await.scheduler(job_scheduler.clone());
2296 let (tx, _rx) = mpsc::channel(4);
2297 let mut scheduler = env.mock_flush_scheduler();
2298 let mut builder = VersionControlBuilder::new();
2299 builder.set_memtable_builder(Arc::new(TimeSeriesMemtableBuilder::default()));
2301 let version_control = Arc::new(builder.build());
2302 let version_data = version_control.current();
2304 write_rows_to_version(&version_data.version, "host0", 0, 10);
2305 let manifest_ctx = env
2306 .mock_manifest_context(version_data.version.metadata.clone())
2307 .await;
2308 let mut tasks: Vec<_> = (0..2)
2310 .map(|_| RegionFlushTask {
2311 region_id: builder.region_id(),
2312 reason: FlushReason::Manual,
2313 senders: Vec::new(),
2314 request_sender: tx.clone(),
2315 access_layer: env.access_layer.clone(),
2316 listener: WorkerListener::default(),
2317 engine_config: Arc::new(MitoConfig::default()),
2318 row_group_size: None,
2319 cache_manager: Arc::new(CacheManager::default()),
2320 manifest_ctx: manifest_ctx.clone(),
2321 index_options: IndexOptions::default(),
2322 flush_semaphore: Arc::new(Semaphore::new(2)),
2323 is_staging: false,
2324 partition_expr: None,
2325 })
2326 .collect();
2327 let task = tasks.pop().unwrap();
2329 scheduler
2330 .schedule_flush(builder.region_id(), &version_control, task)
2331 .unwrap();
2332 assert_eq!(1, scheduler.region_status.len());
2334 assert_eq!(1, job_scheduler.num_jobs());
2335 let task = tasks.pop().unwrap();
2337 scheduler
2338 .schedule_flush(builder.region_id(), &version_control, task)
2339 .unwrap();
2340 assert!(
2341 scheduler
2342 .region_status
2343 .get(&builder.region_id())
2344 .unwrap()
2345 .pending_task
2346 .is_some()
2347 );
2348
2349 let version_data = version_control.current();
2351 assert_eq!(0, version_data.version.memtables.immutables()[0].id());
2352 version_control.apply_edit(
2354 Some(RegionEdit {
2355 files_to_add: Vec::new(),
2356 files_to_remove: Vec::new(),
2357 timestamp_ms: None,
2358 compaction_time_window: None,
2359 flushed_entry_id: None,
2360 flushed_sequence: None,
2361 committed_sequence: None,
2362 }),
2363 &[0],
2364 builder.file_purger(),
2365 );
2366 write_rows_to_version(&version_data.version, "host1", 0, 10);
2367 scheduler.on_flush_success(builder.region_id());
2368 assert_eq!(2, job_scheduler.num_jobs());
2369 assert!(
2371 scheduler
2372 .region_status
2373 .get(&builder.region_id())
2374 .unwrap()
2375 .pending_task
2376 .is_none()
2377 );
2378 }
2379
2380 #[tokio::test]
2381 async fn test_schedule_pending_request_failure_drains_pending_ddls() {
2382 common_telemetry::init_default_ut_logging();
2383 let job_scheduler = Arc::new(VecScheduler::default());
2384 let env = SchedulerEnv::new().await.scheduler(job_scheduler.clone());
2385 let (tx, _rx) = mpsc::channel(4);
2386 let mut scheduler = env.mock_flush_scheduler();
2387 let mut builder = VersionControlBuilder::new();
2388 builder.set_memtable_builder(Arc::new(TimeSeriesMemtableBuilder::default()));
2389 let version_control = Arc::new(builder.build());
2390
2391 let version_data = version_control.current();
2392 write_rows_to_version(&version_data.version, "host0", 0, 10);
2393 let manifest_ctx = env
2394 .mock_manifest_context(version_data.version.metadata.clone())
2395 .await;
2396
2397 let task = new_test_flush_task(
2398 &env,
2399 builder.region_id(),
2400 FlushReason::Manual,
2401 tx.clone(),
2402 manifest_ctx.clone(),
2403 );
2404 scheduler
2405 .schedule_flush(builder.region_id(), &version_control, task)
2406 .unwrap();
2407
2408 let task = new_test_flush_task(
2409 &env,
2410 builder.region_id(),
2411 FlushReason::Closing,
2412 tx,
2413 manifest_ctx,
2414 );
2415 scheduler
2416 .schedule_flush(builder.region_id(), &version_control, task)
2417 .unwrap();
2418
2419 let (sender, receiver) = oneshot::channel();
2420 scheduler.add_ddl_request_to_pending(SenderDdlRequest {
2421 sender: OptionOutputTx::from(sender),
2422 region_id: builder.region_id(),
2423 request: DdlRequest::Close(store_api::region_request::RegionCloseRequest::default()),
2424 });
2425
2426 let version_data = version_control.current();
2427 version_control.apply_edit(
2428 Some(RegionEdit {
2429 files_to_add: Vec::new(),
2430 files_to_remove: Vec::new(),
2431 timestamp_ms: None,
2432 compaction_time_window: None,
2433 flushed_entry_id: None,
2434 flushed_sequence: None,
2435 committed_sequence: None,
2436 }),
2437 &[0],
2438 builder.file_purger(),
2439 );
2440 write_rows_to_version(&version_data.version, "host1", 0, 10);
2441
2442 scheduler.scheduler = Arc::new(FailingScheduler);
2443 scheduler.on_flush_success(builder.region_id());
2444
2445 assert!(scheduler.region_status.is_empty());
2446 let err = receiver
2447 .await
2448 .expect("pending DDL must be notified")
2449 .unwrap_err();
2450 assert_eq!(err.status_code(), StatusCode::RegionBusy);
2451 }
2452}