1use std::any::TypeId;
18use std::collections::HashMap;
19use std::sync::atomic::{AtomicI64, AtomicU64};
20use std::sync::{Arc, LazyLock};
21use std::time::Instant;
22
23use common_telemetry::{debug, error, info, warn};
24use common_wal::options::WalOptions;
25use futures::StreamExt;
26use futures::future::BoxFuture;
27use log_store::kafka::log_store::KafkaLogStore;
28use log_store::noop::log_store::NoopLogStore;
29use log_store::raft_engine::log_store::RaftEngineLogStore;
30use object_store::ObjectStore;
31use object_store::manager::ObjectStoreManagerRef;
32use object_store::util::{is_object_storage, normalize_dir};
33use parquet::file::metadata::PageIndexPolicy;
34use snafu::{OptionExt, ResultExt, ensure};
35use store_api::codec::PrimaryKeyEncoding;
36use store_api::logstore::LogStore;
37use store_api::logstore::provider::Provider;
38use store_api::metadata::{
39 ColumnMetadata, RegionMetadata, RegionMetadataBuilder, RegionMetadataRef,
40};
41use store_api::region_engine::RegionRole;
42use store_api::region_request::{PathType, RegionRequirements};
43use store_api::storage::{ColumnId, RegionId};
44use tokio::sync::Semaphore;
45
46use crate::access_layer::AccessLayer;
47use crate::cache::file_cache::{FileCache, FileType, IndexKey};
48use crate::cache::{CacheManagerRef, SstMetaPreparation, prepare_sst_meta};
49use crate::config::MitoConfig;
50use crate::engine::region_hook::RegionHookRef;
51use crate::error;
52use crate::error::{
53 EmptyRegionDirSnafu, InvalidMetadataSnafu, InvalidRegionOptionsSnafu, ObjectStoreNotFoundSnafu,
54 RegionCorruptedSnafu, Result, StaleLogEntrySnafu,
55};
56use crate::manifest::action::RegionManifest;
57use crate::manifest::manager::{RegionManifestManager, RegionManifestOptions};
58use crate::memtable::MemtableBuilderProvider;
59use crate::memtable::bulk::part::BulkPart;
60use crate::memtable::time_partition::{TimePartitions, TimePartitionsRef};
61use crate::metrics::{CACHE_FILL_DOWNLOADED_FILES, CACHE_FILL_PENDING_FILES};
62use crate::region::options::RegionOptions;
63use crate::region::version::{VersionBuilder, VersionControl, VersionControlRef};
64use crate::region::{
65 ManifestContext, ManifestStats, MitoRegion, MitoRegionRef, RegionLeaderState, RegionRoleState,
66 RegionStats,
67};
68use crate::region_write_ctx::RegionWriteCtx;
69use crate::request::OptionOutputTx;
70use crate::schedule::scheduler::SchedulerRef;
71use crate::sst::FormatType;
72use crate::sst::file::{FileHandle, RegionFileId, RegionIndexId};
73use crate::sst::file_purger::{FilePurgerRef, create_file_purger};
74use crate::sst::file_ref::FileReferenceManagerRef;
75use crate::sst::index::intermediate::IntermediateManager;
76use crate::sst::index::puffin_manager::PuffinManagerFactory;
77use crate::sst::location::{self, region_dir_from_table_dir};
78use crate::sst::parquet::metadata::{MetadataLoader, extract_primary_key_range};
79use crate::sst::parquet::reader::MetadataCacheMetrics;
80use crate::time_provider::TimeProviderRef;
81use crate::wal::entry_reader::WalEntryReader;
82use crate::wal::{EntryId, Wal};
83
84const PARQUET_META_PRELOAD_CONCURRENCY: usize = 8;
85
86static PARQUET_META_PRELOAD_SEMAPHORE: LazyLock<Semaphore> =
87 LazyLock::new(|| Semaphore::new(PARQUET_META_PRELOAD_CONCURRENCY));
88
89fn initial_pruned_entry_id(wal_options: &WalOptions) -> EntryId {
90 match wal_options {
91 WalOptions::Kafka(options) => options.initial_pruned_entry_id.unwrap_or(0),
92 WalOptions::RaftEngine | WalOptions::Noop => 0,
93 }
94}
95
96#[async_trait::async_trait]
102pub trait PartitionExprFetcher {
103 async fn fetch_expr(&self, region_id: RegionId) -> Option<String>;
104}
105
106pub type PartitionExprFetcherRef = Arc<dyn PartitionExprFetcher + Send + Sync>;
107
108pub(crate) struct RegionOpener {
110 region_id: RegionId,
111 metadata_builder: Option<RegionMetadataBuilder>,
112 memtable_builder_provider: MemtableBuilderProvider,
113 object_store_manager: ObjectStoreManagerRef,
114 table_dir: String,
115 path_type: PathType,
116 purge_scheduler: SchedulerRef,
117 options: Option<RegionOptions>,
118 cache_manager: Option<CacheManagerRef>,
119 skip_wal_replay: bool,
120 puffin_manager_factory: PuffinManagerFactory,
121 intermediate_manager: IntermediateManager,
122 time_provider: TimeProviderRef,
123 stats: ManifestStats,
124 wal_entry_reader: Option<Box<dyn WalEntryReader>>,
125 replay_checkpoint: Option<u64>,
126 file_ref_manager: FileReferenceManagerRef,
127 partition_expr_fetcher: PartitionExprFetcherRef,
128 hook: Option<RegionHookRef>,
129}
130
131impl RegionOpener {
132 #[allow(clippy::too_many_arguments)]
135 pub(crate) fn new(
136 region_id: RegionId,
137 table_dir: &str,
138 path_type: PathType,
139 memtable_builder_provider: MemtableBuilderProvider,
140 object_store_manager: ObjectStoreManagerRef,
141 purge_scheduler: SchedulerRef,
142 puffin_manager_factory: PuffinManagerFactory,
143 intermediate_manager: IntermediateManager,
144 time_provider: TimeProviderRef,
145 file_ref_manager: FileReferenceManagerRef,
146 partition_expr_fetcher: PartitionExprFetcherRef,
147 ) -> RegionOpener {
148 RegionOpener {
149 region_id,
150 metadata_builder: None,
151 memtable_builder_provider,
152 object_store_manager,
153 table_dir: normalize_dir(table_dir),
154 path_type,
155 purge_scheduler,
156 options: None,
157 cache_manager: None,
158 skip_wal_replay: false,
159 puffin_manager_factory,
160 intermediate_manager,
161 time_provider,
162 stats: Default::default(),
163 wal_entry_reader: None,
164 replay_checkpoint: None,
165 file_ref_manager,
166 partition_expr_fetcher,
167 hook: None,
168 }
169 }
170
171 pub(crate) fn hook(mut self, hook: Option<RegionHookRef>) -> Self {
173 self.hook = hook;
174 self
175 }
176
177 pub(crate) fn metadata_builder(mut self, builder: RegionMetadataBuilder) -> Self {
179 self.metadata_builder = Some(builder);
180 self
181 }
182
183 fn region_dir(&self) -> String {
185 region_dir_from_table_dir(&self.table_dir, self.region_id, self.path_type)
186 }
187
188 fn build_metadata(&mut self) -> Result<RegionMetadata> {
194 let options = self.options.as_ref().unwrap();
195 let mut metadata_builder = self.metadata_builder.take().unwrap();
196 metadata_builder.primary_key_encoding(options.primary_key_encoding());
197 metadata_builder.build().context(InvalidMetadataSnafu)
198 }
199
200 pub(crate) fn parse_options(self, options: HashMap<String, String>) -> Result<Self> {
202 let region_id = self.region_id;
203 self.options(RegionOptions::try_from_options(region_id, &options)?)
204 }
205
206 pub(crate) fn replay_checkpoint(mut self, replay_checkpoint: Option<u64>) -> Self {
208 self.replay_checkpoint = replay_checkpoint;
209 self
210 }
211
212 pub(crate) fn wal_entry_reader(
215 mut self,
216 wal_entry_reader: Option<Box<dyn WalEntryReader>>,
217 ) -> Self {
218 self.wal_entry_reader = wal_entry_reader;
219 self
220 }
221
222 pub(crate) fn options(mut self, options: RegionOptions) -> Result<Self> {
224 options.validate()?;
225 self.options = Some(options);
226 Ok(self)
227 }
228
229 pub(crate) fn ensure_region_requirements(
231 &self,
232 requirements: RegionRequirements,
233 ) -> Result<()> {
234 if !requirements.object_storage {
235 return Ok(());
236 }
237
238 let options = self.options.as_ref().context(InvalidRegionOptionsSnafu {
239 reason: "missing region options before requirement check".to_string(),
240 })?;
241 let object_store = get_object_store(&options.storage, &self.object_store_manager)?;
242
243 ensure!(
244 supports_open_region_object_storage_requirement(&object_store),
245 error::RegionRequirementSnafu {
246 region_id: self.region_id,
247 requirement: "object storage",
248 reason: "region data must be accessible from another datanode",
249 }
250 );
251
252 Ok(())
253 }
254
255 pub(crate) fn cache(mut self, cache_manager: Option<CacheManagerRef>) -> Self {
257 self.cache_manager = cache_manager;
258 self
259 }
260
261 pub(crate) fn skip_wal_replay(mut self, skip: bool) -> Self {
263 self.skip_wal_replay = skip;
264 self
265 }
266
267 pub(crate) async fn create_or_open<S: LogStore>(
274 mut self,
275 config: &MitoConfig,
276 wal: &Wal<S>,
277 ) -> Result<MitoRegionRef> {
278 let region_id = self.region_id;
279 let region_dir = self.region_dir();
280 let metadata = self.build_metadata()?;
281 match self.maybe_open(config, wal).await {
283 Ok(Some(region)) => {
284 let recovered = region.metadata();
285 let expect = &metadata;
287 check_recovered_region(
288 &recovered,
289 expect.region_id,
290 &expect.column_metadatas,
291 &expect.primary_key,
292 )?;
293 region.set_role(RegionRole::Leader);
295
296 return Ok(region);
297 }
298 Ok(None) => {
299 debug!(
300 "No data under directory {}, region_id: {}",
301 region_dir, self.region_id
302 );
303 }
304 Err(e) => {
305 warn!(e;
306 "Failed to open region {} before creating it, region_dir: {}",
307 self.region_id, region_dir
308 );
309 }
310 }
311 let mut options = self.options.take().unwrap();
313 let object_store = get_object_store(&options.storage, &self.object_store_manager)?;
314 let provider = self.provider::<S>(&options.wal_options)?;
315 let metadata = Arc::new(metadata);
316 let sst_format = if let Some(format) = options.sst_format {
318 format
319 } else if config.default_flat_format {
320 options.sst_format = Some(FormatType::Flat);
321 FormatType::Flat
322 } else {
323 options.sst_format = Some(FormatType::PrimaryKey);
325 FormatType::PrimaryKey
326 };
327 let mut region_manifest_options =
329 RegionManifestOptions::new(config, ®ion_dir, &object_store);
330 region_manifest_options.manifest_cache = self
332 .cache_manager
333 .as_ref()
334 .and_then(|cm| cm.write_cache())
335 .and_then(|wc| wc.manifest_cache());
336 let flushed_entry_id = provider
339 .initial_flushed_entry_id::<S>(wal.store())
340 .max(initial_pruned_entry_id(&options.wal_options));
341 let manifest_manager = RegionManifestManager::new(
342 metadata.clone(),
343 flushed_entry_id,
344 region_manifest_options,
345 sst_format,
346 &self.stats,
347 )
348 .await?;
349
350 let memtable_builder = self.memtable_builder_provider.builder_for_options(&options);
351 let part_duration = options.compaction.time_window();
352 let mutable = Arc::new(TimePartitions::new(
354 metadata.clone(),
355 memtable_builder.clone(),
356 0,
357 part_duration,
358 ));
359
360 debug!(
361 "Create region {} with options: {:?}, default_flat_format: {}",
362 region_id, options, config.default_flat_format
363 );
364
365 let version = VersionBuilder::new(metadata, mutable)
366 .options(options)
367 .flushed_entry_id(flushed_entry_id)
368 .build();
369 let version_control = Arc::new(VersionControl::new(version));
370 let access_layer = Arc::new(AccessLayer::new(
371 self.table_dir.clone(),
372 self.path_type,
373 object_store,
374 self.puffin_manager_factory,
375 self.intermediate_manager,
376 ));
377 let now = self.time_provider.current_time_millis();
378
379 Ok(Arc::new(MitoRegion {
380 region_id,
381 version_control,
382 access_layer: access_layer.clone(),
383 manifest_ctx: Arc::new(ManifestContext::new(
385 manifest_manager,
386 RegionRoleState::Leader(RegionLeaderState::Writable),
387 self.hook.clone(),
388 )),
389 file_purger: create_file_purger(
390 config.gc.enable,
391 self.path_type,
392 self.purge_scheduler,
393 access_layer,
394 self.cache_manager,
395 self.file_ref_manager.clone(),
396 ),
397 provider,
398 last_flush_millis: AtomicI64::new(now),
399 last_schedule_compaction_millis: AtomicI64::new(now),
400 time_provider: self.time_provider.clone(),
401 topic_latest_entry_id: AtomicU64::new(flushed_entry_id),
402 region_stats: RegionStats::new(),
403 stats: self.stats,
404 }))
405 }
406
407 pub(crate) async fn open<S: LogStore>(
411 mut self,
412 config: &MitoConfig,
413 wal: &Wal<S>,
414 ) -> Result<MitoRegionRef> {
415 let region_id = self.region_id;
416 let region_dir = self.region_dir();
417 let region = self
418 .maybe_open(config, wal)
419 .await?
420 .with_context(|| EmptyRegionDirSnafu {
421 region_id,
422 region_dir: ®ion_dir,
423 })?;
424
425 ensure!(
426 region.region_id == self.region_id,
427 RegionCorruptedSnafu {
428 region_id: self.region_id,
429 reason: format!(
430 "recovered region has different region id {}",
431 region.region_id
432 ),
433 }
434 );
435
436 Ok(region)
437 }
438
439 fn provider<S: LogStore>(&self, wal_options: &WalOptions) -> Result<Provider> {
440 provider_from_wal_options::<S>(self.region_id, wal_options)
441 }
442
443 async fn maybe_open<S: LogStore>(
445 &mut self,
446 config: &MitoConfig,
447 wal: &Wal<S>,
448 ) -> Result<Option<MitoRegionRef>> {
449 let now = Instant::now();
450 let mut region_options = self.options.as_ref().unwrap().clone();
451 let object_storage = get_object_store(®ion_options.storage, &self.object_store_manager)?;
452 let mut region_manifest_options =
453 RegionManifestOptions::new(config, &self.region_dir(), &object_storage);
454 region_manifest_options.manifest_cache = self
456 .cache_manager
457 .as_ref()
458 .and_then(|cm| cm.write_cache())
459 .and_then(|wc| wc.manifest_cache());
460 let Some(manifest_manager) =
461 RegionManifestManager::open(region_manifest_options, &self.stats).await?
462 else {
463 return Ok(None);
464 };
465
466 let manifest = manifest_manager.manifest();
468 let metadata = if manifest.metadata.partition_expr.is_none()
469 && let Some(expr_json) = self.partition_expr_fetcher.fetch_expr(self.region_id).await
470 {
471 let metadata = manifest.metadata.as_ref().clone();
472 let mut builder = RegionMetadataBuilder::from_existing(metadata);
473 builder.partition_expr_json(Some(expr_json));
474 Arc::new(builder.build().context(InvalidMetadataSnafu)?)
475 } else {
476 manifest.metadata.clone()
477 };
478 sanitize_region_options(&manifest, &mut region_options);
480
481 let region_id = self.region_id;
482 let provider = self.provider::<S>(®ion_options.wal_options)?;
483 let wal_entry_reader = self
484 .wal_entry_reader
485 .take()
486 .unwrap_or_else(|| wal.wal_entry_reader(&provider, region_id, None));
487 let on_region_opened = wal.on_region_opened();
488 let object_store = get_object_store(®ion_options.storage, &self.object_store_manager)?;
489
490 debug!(
491 "Open region {} at {} with options: {:?}",
492 region_id, self.table_dir, self.options
493 );
494
495 let access_layer = Arc::new(AccessLayer::new(
496 self.table_dir.clone(),
497 self.path_type,
498 object_store,
499 self.puffin_manager_factory.clone(),
500 self.intermediate_manager.clone(),
501 ));
502 let file_purger = create_file_purger(
503 config.gc.enable,
504 self.path_type,
505 self.purge_scheduler.clone(),
506 access_layer.clone(),
507 self.cache_manager.clone(),
508 self.file_ref_manager.clone(),
509 );
510 let memtable_builder = self
512 .memtable_builder_provider
513 .builder_for_options(®ion_options);
514 let part_duration = region_options
517 .compaction
518 .time_window()
519 .or(manifest.compaction_time_window);
520 let mutable = Arc::new(TimePartitions::new(
522 metadata.clone(),
523 memtable_builder.clone(),
524 0,
525 part_duration,
526 ));
527
528 let version_builder = version_builder_from_manifest(
530 &manifest,
531 metadata,
532 file_purger.clone(),
533 mutable,
534 region_options,
535 );
536 let version = version_builder.build();
537 let flushed_entry_id = version.flushed_entry_id;
538 let version_control = Arc::new(VersionControl::new(version));
539
540 let replay_from_entry_id = self
541 .replay_checkpoint
542 .unwrap_or_default()
543 .max(flushed_entry_id);
544 let topic_latest_entry_id = if !self.skip_wal_replay {
545 info!(
546 "Start replaying memtable at replay_from_entry_id: {} for region {}, manifest version: {}, flushed entry id: {}, elapsed: {:?}",
547 replay_from_entry_id,
548 region_id,
549 manifest.manifest_version,
550 flushed_entry_id,
551 now.elapsed()
552 );
553 replay_memtable(
554 &provider,
555 wal_entry_reader,
556 region_id,
557 replay_from_entry_id,
558 &version_control,
559 config.allow_stale_entries,
560 on_region_opened,
561 )
562 .await?;
563 if provider.is_remote_wal() && version_control.current().version.memtables.is_empty() {
567 wal.store()
568 .latest_entry_id(&provider)
569 .unwrap_or(replay_from_entry_id)
570 } else if provider.is_remote_wal() {
571 replay_from_entry_id
572 } else {
573 0
574 }
575 } else {
576 info!(
577 "Skip the WAL replay for region: {}, manifest version: {}, flushed_entry_id: {}, elapsed: {:?}",
578 region_id,
579 manifest.manifest_version,
580 flushed_entry_id,
581 now.elapsed()
582 );
583
584 if provider.is_remote_wal() {
585 replay_from_entry_id
586 } else {
587 0
588 }
589 };
590
591 if let Some(committed_in_manifest) = manifest.committed_sequence {
592 let committed_after_replay = version_control.committed_sequence();
593 if committed_in_manifest > committed_after_replay {
594 info!(
595 "Overriding committed sequence, region: {}, flushed_sequence: {}, committed_sequence: {} -> {}",
596 self.region_id,
597 version_control.current().version.flushed_sequence,
598 version_control.committed_sequence(),
599 committed_in_manifest
600 );
601 version_control.set_committed_sequence(committed_in_manifest);
602 }
603 }
604
605 let now = self.time_provider.current_time_millis();
606
607 let region = MitoRegion {
608 region_id: self.region_id,
609 version_control: version_control.clone(),
610 access_layer: access_layer.clone(),
611 manifest_ctx: Arc::new(ManifestContext::new(
613 manifest_manager,
614 RegionRoleState::Follower,
615 self.hook.clone(),
616 )),
617 file_purger,
618 provider: provider.clone(),
619 last_flush_millis: AtomicI64::new(now),
620 last_schedule_compaction_millis: AtomicI64::new(now),
621 time_provider: self.time_provider.clone(),
622 topic_latest_entry_id: AtomicU64::new(topic_latest_entry_id),
623 region_stats: RegionStats::new(),
624 stats: self.stats.clone(),
625 };
626
627 let region = Arc::new(region);
628
629 maybe_load_cache(®ion, config, &self.cache_manager);
630 maybe_preload_parquet_meta_cache(®ion, config, &self.cache_manager);
631
632 Ok(Some(region))
633 }
634}
635
636pub(crate) fn provider_from_wal_options<S: LogStore>(
637 region_id: RegionId,
638 wal_options: &WalOptions,
639) -> Result<Provider> {
640 match wal_options {
641 WalOptions::RaftEngine => {
642 ensure!(
643 TypeId::of::<RaftEngineLogStore>() == TypeId::of::<S>()
644 || TypeId::of::<NoopLogStore>() == TypeId::of::<S>(),
645 error::IncompatibleWalProviderChangeSnafu {
646 global: "`kafka`",
647 region: "`raft_engine`",
648 }
649 );
650 Ok(Provider::raft_engine_provider(region_id.as_u64()))
651 }
652 WalOptions::Kafka(options) => {
653 ensure!(
654 TypeId::of::<KafkaLogStore>() == TypeId::of::<S>()
655 || TypeId::of::<NoopLogStore>() == TypeId::of::<S>(),
656 error::IncompatibleWalProviderChangeSnafu {
657 global: "`raft_engine`",
658 region: "`kafka`",
659 }
660 );
661 Ok(Provider::kafka_provider(options.topic.clone()))
662 }
663 WalOptions::Noop => Ok(Provider::noop_provider()),
664 }
665}
666
667#[cfg(not(feature = "test-shared-fs-region-migration"))]
668fn supports_open_region_object_storage_requirement(object_store: &ObjectStore) -> bool {
669 is_object_storage(object_store)
670}
671
672#[cfg(feature = "test-shared-fs-region-migration")]
673fn supports_open_region_object_storage_requirement(object_store: &ObjectStore) -> bool {
674 is_object_storage(object_store)
679 || object_store.info().scheme() == object_store::services::FS_SCHEME
680}
681
682pub(crate) fn version_builder_from_manifest(
684 manifest: &RegionManifest,
685 metadata: RegionMetadataRef,
686 file_purger: FilePurgerRef,
687 mutable: TimePartitionsRef,
688 region_options: RegionOptions,
689) -> VersionBuilder {
690 VersionBuilder::new(metadata, mutable)
691 .add_files(file_purger, manifest.files.values().cloned())
692 .flushed_entry_id(manifest.flushed_entry_id)
693 .flushed_sequence(manifest.flushed_sequence)
694 .truncated_entry_id(manifest.truncated_entry_id)
695 .compaction_time_window(manifest.compaction_time_window)
696 .options(region_options)
697}
698
699pub(crate) fn sanitize_region_options(manifest: &RegionManifest, options: &mut RegionOptions) {
701 match options.sst_format {
706 Some(format) if format != manifest.sst_format => {
707 common_telemetry::warn!(
708 "Overriding SST format from {:?} (manifest) to {:?} (options) for region {}",
709 manifest.sst_format,
710 format,
711 manifest.metadata.region_id,
712 );
713 }
714 Some(_) => {}
715 None => {
716 options.sst_format = Some(manifest.sst_format);
717 }
718 }
719 if let Some(manifest_append_mode) = manifest.append_mode
720 && options.append_mode != manifest_append_mode
721 {
722 common_telemetry::warn!(
723 "Overriding append_mode from {} to {} for region {}",
724 options.append_mode,
725 manifest_append_mode,
726 manifest.metadata.region_id,
727 );
728 options.append_mode = manifest_append_mode;
729 }
730 if options.append_mode && options.merge_mode.take().is_some() {
731 common_telemetry::warn!(
732 "Ignoring merge_mode because append_mode is enabled for region {}",
733 manifest.metadata.region_id,
734 );
735 }
736}
737
738pub(crate) fn sanitize_open_request_options(options: &mut HashMap<String, String>) {
740 let append_mode_enabled = options
741 .get("append_mode")
742 .is_some_and(|v| matches!(v.trim().to_ascii_lowercase().as_str(), "true" | "1"));
743
744 if append_mode_enabled && options.remove("merge_mode").is_some() {
745 common_telemetry::warn!(
746 "Ignoring merge_mode in open request options because append_mode is enabled"
747 );
748 }
749}
750
751pub fn get_object_store(
753 name: &Option<String>,
754 object_store_manager: &ObjectStoreManagerRef,
755) -> Result<object_store::ObjectStore> {
756 if let Some(name) = name {
757 Ok(object_store_manager
758 .find(name)
759 .with_context(|| ObjectStoreNotFoundSnafu {
760 object_store: name.clone(),
761 })?
762 .clone())
763 } else {
764 Ok(object_store_manager.default_object_store().clone())
765 }
766}
767
768pub(crate) fn check_recovered_region(
770 recovered: &RegionMetadata,
771 region_id: RegionId,
772 column_metadatas: &[ColumnMetadata],
773 primary_key: &[ColumnId],
774) -> Result<()> {
775 if recovered.region_id != region_id {
776 error!(
777 "Recovered region {}, expect region {}",
778 recovered.region_id, region_id
779 );
780 return RegionCorruptedSnafu {
781 region_id,
782 reason: format!(
783 "recovered metadata has different region id {}",
784 recovered.region_id
785 ),
786 }
787 .fail();
788 }
789 if recovered.column_metadatas != column_metadatas {
790 error!(
791 "Unexpected schema in recovered region {}, recovered: {:?}, expect: {:?}",
792 recovered.region_id, recovered.column_metadatas, column_metadatas
793 );
794
795 return RegionCorruptedSnafu {
796 region_id,
797 reason: "recovered metadata has different schema",
798 }
799 .fail();
800 }
801 if recovered.primary_key != primary_key {
802 error!(
803 "Unexpected primary key in recovered region {}, recovered: {:?}, expect: {:?}",
804 recovered.region_id, recovered.primary_key, primary_key
805 );
806
807 return RegionCorruptedSnafu {
808 region_id,
809 reason: "recovered metadata has different primary key",
810 }
811 .fail();
812 }
813
814 Ok(())
815}
816
817pub(crate) async fn replay_memtable<F>(
819 provider: &Provider,
820 mut wal_entry_reader: Box<dyn WalEntryReader>,
821 region_id: RegionId,
822 flushed_entry_id: EntryId,
823 version_control: &VersionControlRef,
824 allow_stale_entries: bool,
825 on_region_opened: F,
826) -> Result<EntryId>
827where
828 F: FnOnce(RegionId, EntryId, &Provider) -> BoxFuture<Result<()>> + Send,
829{
830 let now = Instant::now();
831 let mut rows_replayed = 0;
832 let mut last_entry_id = flushed_entry_id;
835 let replay_from_entry_id = flushed_entry_id + 1;
836 let region_metadata = version_control.current().version.metadata.clone();
837
838 let mut wal_stream = wal_entry_reader.read(provider, replay_from_entry_id)?;
839 while let Some(res) = wal_stream.next().await {
840 let (entry_id, entry) = res?;
841 if entry_id <= flushed_entry_id {
842 warn!(
843 "Stale WAL entries read during replay, region id: {}, flushed entry id: {}, entry id read: {}",
844 region_id, flushed_entry_id, entry_id
845 );
846 ensure!(
847 allow_stale_entries,
848 StaleLogEntrySnafu {
849 region_id,
850 flushed_entry_id,
851 unexpected_entry_id: entry_id,
852 }
853 );
854 }
855 last_entry_id = last_entry_id.max(entry_id);
856
857 let mut region_write_ctx = RegionWriteCtx::new(
858 region_id,
859 version_control,
860 provider.clone(),
861 None,
863 );
864 for mutation in entry.mutations {
865 rows_replayed += mutation
866 .rows
867 .as_ref()
868 .map(|rows| rows.rows.len())
869 .unwrap_or(0);
870 region_write_ctx.push_mutation(
871 mutation.op_type,
872 mutation.rows,
873 mutation.write_hint,
874 OptionOutputTx::none(),
875 Some(mutation.sequence),
877 );
878 }
879
880 for bulk_entry in entry.bulk_entries {
881 let mut part = BulkPart::try_from(bulk_entry)?;
882 if region_metadata.primary_key_encoding != PrimaryKeyEncoding::Sparse {
889 part.fill_missing_columns(®ion_metadata)?;
890 }
891 rows_replayed += part.num_rows();
892 let bulk_sequence_from_wal = part.sequence;
894 ensure!(
895 region_write_ctx.push_bulk(
896 OptionOutputTx::none(),
897 part,
898 Some(bulk_sequence_from_wal)
899 ),
900 RegionCorruptedSnafu {
901 region_id,
902 reason: "unable to replay memtable with bulk entries",
903 }
904 );
905 }
906
907 region_write_ctx.set_next_entry_id(last_entry_id + 1);
909 region_write_ctx.write_memtable().await;
910 region_write_ctx.write_bulk().await;
911 }
912
913 (on_region_opened)(region_id, flushed_entry_id, provider).await?;
916
917 let series_count = version_control.current().series_count();
918 info!(
919 "Replay WAL for region: {}, provider: {:?}, rows recovered: {}, replay from entry id: {}, last entry id: {}, total timeseries replayed: {}, elapsed: {:?}",
920 region_id,
921 provider,
922 rows_replayed,
923 replay_from_entry_id,
924 last_entry_id,
925 series_count,
926 now.elapsed()
927 );
928 Ok(last_entry_id)
929}
930
931pub(crate) struct RegionLoadCacheTask {
933 region: MitoRegionRef,
934}
935
936impl RegionLoadCacheTask {
937 pub(crate) fn new(region: MitoRegionRef) -> Self {
938 Self { region }
939 }
940
941 pub(crate) async fn fill_cache(&self, file_cache: &FileCache) {
943 let region_id = self.region.region_id;
944 let table_dir = self.region.access_layer.table_dir();
945 let path_type = self.region.access_layer.path_type();
946 let object_store = self.region.access_layer.object_store();
947 let version_control = &self.region.version_control;
948
949 let mut files_to_download = Vec::new();
951 let mut files_already_cached = 0;
952
953 {
954 let version = version_control.current().version;
955 for level in version.ssts.levels() {
956 for file_handle in level.files.values() {
957 let file_meta = file_handle.meta_ref();
958 if file_meta.exists_index() {
959 let puffin_key = IndexKey::new(
960 file_meta.region_id,
961 file_meta.file_id,
962 FileType::Puffin(file_meta.index_version),
963 );
964
965 if !file_cache.contains_key(&puffin_key) {
966 files_to_download.push((
967 puffin_key,
968 file_meta.index_file_size,
969 file_meta.time_range.1, ));
971 } else {
972 files_already_cached += 1;
973 }
974 }
975 }
976 }
977 }
980
981 files_to_download.sort_by_key(|b| std::cmp::Reverse(b.2));
983
984 let total_files = files_to_download.len() as i64;
985
986 info!(
987 "Starting background index cache preload for region {}, total_files_to_download: {}, files_already_cached: {}",
988 region_id, total_files, files_already_cached
989 );
990
991 CACHE_FILL_PENDING_FILES.add(total_files);
992
993 let mut files_downloaded = 0;
994 let mut files_skipped = 0;
995
996 for (puffin_key, file_size, max_timestamp) in files_to_download {
997 let current_size = file_cache.puffin_cache_size();
998 let capacity = file_cache.puffin_cache_capacity();
999 let region_state = self.region.state();
1000 if !can_load_cache(region_state) {
1001 info!(
1002 "Stopping index cache by state: {:?}, region: {}, current_size: {}, capacity: {}",
1003 region_state, region_id, current_size, capacity
1004 );
1005 break;
1006 }
1007
1008 if current_size + file_size > capacity {
1010 info!(
1011 "Stopping index cache preload due to capacity limit, region: {}, file_id: {}, current_size: {}, file_size: {}, capacity: {}, file_timestamp: {:?}",
1012 region_id, puffin_key.file_id, current_size, file_size, capacity, max_timestamp
1013 );
1014 files_skipped = (total_files - files_downloaded) as usize;
1015 CACHE_FILL_PENDING_FILES.sub(total_files - files_downloaded);
1016 break;
1017 }
1018
1019 let index_version = if let FileType::Puffin(version) = puffin_key.file_type {
1020 version
1021 } else {
1022 unreachable!("`files_to_download` should only contains Puffin files");
1023 };
1024 let index_id = RegionIndexId::new(
1025 RegionFileId::new(puffin_key.region_id, puffin_key.file_id),
1026 index_version,
1027 );
1028
1029 let index_remote_path = location::index_file_path(table_dir, index_id, path_type);
1030
1031 match file_cache
1032 .download(puffin_key, &index_remote_path, object_store, file_size)
1033 .await
1034 {
1035 Ok(_) => {
1036 debug!(
1037 "Downloaded index file to write cache, region: {}, file_id: {}",
1038 region_id, puffin_key.file_id
1039 );
1040 files_downloaded += 1;
1041 CACHE_FILL_DOWNLOADED_FILES.inc_by(1);
1042 CACHE_FILL_PENDING_FILES.dec();
1043 }
1044 Err(e) => {
1045 warn!(
1046 e; "Failed to download index file to write cache, region: {}, file_id: {}",
1047 region_id, puffin_key.file_id
1048 );
1049 CACHE_FILL_PENDING_FILES.dec();
1050 }
1051 }
1052 }
1053
1054 info!(
1055 "Completed background cache fill task for region {}, total_files: {}, files_downloaded: {}, files_already_cached: {}, files_skipped: {}",
1056 region_id, total_files, files_downloaded, files_already_cached, files_skipped
1057 );
1058 }
1059}
1060
1061fn maybe_load_cache(
1063 region: &MitoRegionRef,
1064 config: &MitoConfig,
1065 cache_manager: &Option<CacheManagerRef>,
1066) {
1067 let Some(cache_manager) = cache_manager else {
1068 return;
1069 };
1070 let Some(write_cache) = cache_manager.write_cache() else {
1071 return;
1072 };
1073
1074 let preload_enabled = config.preload_index_cache;
1075 if !preload_enabled {
1076 return;
1077 }
1078
1079 let task = RegionLoadCacheTask::new(region.clone());
1080 write_cache.load_region_cache(task);
1081}
1082
1083#[allow(clippy::too_many_arguments)]
1094async fn preload_parquet_meta_cache_for_files(
1095 region_id: RegionId,
1096 cache_manager: CacheManagerRef,
1097 sst_meta_cache_capacity: u64,
1098 table_dir: String,
1099 path_type: PathType,
1100 object_store: object_store::ObjectStore,
1101 region_metadata: RegionMetadataRef,
1102 mut files: Vec<FileHandle>,
1103) -> usize {
1104 if !cache_manager.sst_meta_cache_enabled()
1105 || sst_meta_cache_capacity == 0
1106 || cache_manager.sst_meta_cache_is_full()
1107 {
1108 return 0;
1109 }
1110
1111 let allow_direct_load = object_store.info().scheme() == object_store::services::FS_SCHEME;
1112
1113 files.sort_by_key(|b| std::cmp::Reverse(b.meta_ref().time_range.1));
1115
1116 let mut loaded = 0usize;
1117 for file_handle in files {
1118 if cache_manager.sst_meta_cache_is_full() {
1120 break;
1121 }
1122
1123 let file_id = file_handle.file_id();
1124 let mut cache_metrics = MetadataCacheMetrics::default();
1125 if let Some(metadata) = cache_manager
1126 .get_compact_sst_meta_data(file_id, PageIndexPolicy::Optional)
1127 .await
1128 {
1129 if file_handle.primary_key_range().is_none()
1130 && let Some(primary_key_range) = extract_primary_key_range(
1131 metadata.parquet_metadata().as_ref(),
1132 ®ion_metadata,
1133 )
1134 {
1135 file_handle.set_primary_key_range(primary_key_range);
1136 }
1137 continue;
1138 }
1139
1140 if let Some(write_cache) = cache_manager.write_cache() {
1141 let key = IndexKey::new(file_id.region_id(), file_id.file_id(), FileType::Parquet);
1142 if let Some(metadata) = write_cache
1143 .file_cache()
1144 .get_sst_meta_data(key, &mut cache_metrics, PageIndexPolicy::Optional)
1145 .await
1146 {
1147 let decoded = metadata.decoded();
1148 if file_handle.primary_key_range().is_none()
1149 && let Some(primary_key_range) = extract_primary_key_range(
1150 decoded.parquet_metadata().as_ref(),
1151 ®ion_metadata,
1152 )
1153 {
1154 file_handle.set_primary_key_range(primary_key_range);
1155 }
1156 match metadata {
1157 SstMetaPreparation::Prepared(metadata) => {
1158 cache_manager.put_prepared_sst_meta(file_id, metadata, false);
1159 loaded += 1;
1160 }
1161 SstMetaPreparation::DecodedOnly { encoding_error, .. } => warn!(
1162 encoding_error;
1163 "Failed to encode file-cached SST metadata during preload, region: {}, file: {}",
1164 region_id,
1165 file_id.file_id()
1166 ),
1167 }
1168 continue;
1169 }
1170 }
1171
1172 if !allow_direct_load {
1173 continue;
1174 }
1175
1176 let file_size = file_handle.meta_ref().file_size;
1177 let file_path = file_handle.file_path(&table_dir, path_type);
1178 let mut loader = MetadataLoader::new(object_store.clone(), &file_path, file_size);
1179 loader.with_page_index_policy(PageIndexPolicy::Optional);
1180 match loader.load(&mut cache_metrics).await {
1181 Ok(metadata) => {
1182 if let Some(primary_key_range) =
1183 extract_primary_key_range(&metadata, ®ion_metadata)
1184 {
1185 file_handle.set_primary_key_range(primary_key_range);
1186 }
1187 match prepare_sst_meta(
1188 &file_path,
1189 metadata,
1190 None,
1193 PageIndexPolicy::Optional,
1194 )
1195 .await
1196 {
1197 Ok(SstMetaPreparation::Prepared(metadata)) => {
1198 cache_manager.put_prepared_sst_meta(file_id, metadata, false);
1201 loaded += 1;
1202 }
1203 Ok(SstMetaPreparation::DecodedOnly { encoding_error, .. }) => warn!(
1204 encoding_error;
1205 "Failed to encode preloaded parquet metadata, region: {}, file: {}",
1206 region_id,
1207 file_path
1208 ),
1209 Err(err) => {
1210 warn!(
1211 err; "Failed to prepare preloaded parquet metadata, region: {}, file: {}",
1212 region_id, file_path
1213 );
1214 }
1215 }
1216 }
1217 Err(err) => {
1218 warn!(
1220 err; "Failed to preload parquet metadata from local store, region: {}, file: {}",
1221 region_id, file_path
1222 );
1223 }
1224 }
1225 }
1226
1227 loaded
1228}
1229
1230fn maybe_preload_parquet_meta_cache(
1231 region: &MitoRegionRef,
1232 config: &MitoConfig,
1233 cache_manager: &Option<CacheManagerRef>,
1234) {
1235 let Some(cache_manager) = cache_manager else {
1236 return;
1237 };
1238 if !cache_manager.sst_meta_cache_enabled() {
1239 return;
1240 }
1241
1242 if config.sst_meta_cache_size.as_bytes() == 0 {
1244 return;
1245 }
1246 if !config.preload_index_cache {
1247 return;
1248 }
1249
1250 let region = region.clone();
1251 let cache_manager = cache_manager.clone();
1252 let sst_meta_cache_capacity = config.sst_meta_cache_size.as_bytes();
1253
1254 tokio::spawn(async move {
1255 let _permit = PARQUET_META_PRELOAD_SEMAPHORE.acquire().await.unwrap();
1257
1258 let region_id = region.region_id;
1259 let table_dir = region.access_layer.table_dir().to_string();
1260 let path_type = region.access_layer.path_type();
1261 let object_store = region.access_layer.object_store().clone();
1262 let region_metadata = region.version_control.current().version.metadata.clone();
1263
1264 let mut files = Vec::new();
1266 {
1267 let version = region.version_control.current().version;
1268 for level in version.ssts.levels() {
1269 for file_handle in level.files.values() {
1270 files.push(file_handle.clone());
1271 }
1272 }
1273 }
1274 let preloading_start = Instant::now();
1275 let loaded = preload_parquet_meta_cache_for_files(
1276 region_id,
1277 cache_manager,
1278 sst_meta_cache_capacity,
1279 table_dir,
1280 path_type,
1281 object_store,
1282 region_metadata,
1283 files,
1284 )
1285 .await;
1286 let preloading_cost = preloading_start.elapsed();
1287
1288 if loaded > 0 {
1289 info!(
1290 "Preloaded parquet metadata for region {}, loaded_files: {}, elapsed_ms: {}",
1291 region_id,
1292 loaded,
1293 preloading_cost.as_millis()
1294 );
1295 }
1296 });
1297}
1298
1299fn can_load_cache(state: RegionRoleState) -> bool {
1300 match state {
1301 RegionRoleState::Leader(RegionLeaderState::Writable)
1302 | RegionRoleState::Leader(RegionLeaderState::Staging)
1303 | RegionRoleState::Leader(RegionLeaderState::Altering)
1304 | RegionRoleState::Leader(RegionLeaderState::EnteringStaging)
1305 | RegionRoleState::Leader(RegionLeaderState::Editing)
1306 | RegionRoleState::Follower => true,
1307 RegionRoleState::Leader(RegionLeaderState::Downgrading)
1309 | RegionRoleState::Leader(RegionLeaderState::Dropping)
1310 | RegionRoleState::Leader(RegionLeaderState::Truncating) => false,
1311 }
1312}
1313
1314#[cfg(test)]
1315mod tests {
1316 use std::collections::HashMap;
1317 use std::sync::Arc;
1318
1319 use common_base::readable_size::ReadableSize;
1320 use common_test_util::temp_dir::create_temp_dir;
1321 use common_time::Timestamp;
1322 use common_wal::options::{KafkaWalOptions, WalOptions};
1323 use datatypes::arrow::array::{ArrayRef, BinaryArray, Int64Array};
1324 use datatypes::arrow::record_batch::RecordBatch;
1325 use object_store::ObjectStore;
1326 use object_store::services::{Fs, Memory, S3};
1327 use parquet::arrow::ArrowWriter;
1328 use parquet::file::metadata::{KeyValue, PageIndexPolicy};
1329 use parquet::file::properties::WriterProperties;
1330 use store_api::region_request::PathType;
1331 use store_api::storage::{FileId, RegionId};
1332
1333 use super::{
1334 initial_pruned_entry_id, preload_parquet_meta_cache_for_files, sanitize_region_options,
1335 supports_open_region_object_storage_requirement,
1336 };
1337 use crate::cache::CacheManager;
1338 use crate::cache::file_cache::{FileType, IndexKey};
1339 use crate::manifest::action::{RegionManifest, RemovedFilesRecord};
1340 use crate::region::options::RegionOptions;
1341 use crate::sst::FormatType;
1342 use crate::sst::file::{FileHandle, FileMeta};
1343 use crate::sst::file_purger::NoopFilePurger;
1344 use crate::sst::parquet::PARQUET_METADATA_KEY;
1345 use crate::test_util::TestEnv;
1346 use crate::test_util::sst_util::sst_region_metadata;
1347
1348 fn build_test_manifest(sst_format: FormatType) -> RegionManifest {
1349 RegionManifest {
1350 metadata: Arc::new(sst_region_metadata()),
1351 files: HashMap::new(),
1352 removed_files: RemovedFilesRecord::default(),
1353 flushed_entry_id: 0,
1354 flushed_sequence: 0,
1355 committed_sequence: None,
1356 manifest_version: 0,
1357 truncated_entry_id: None,
1358 compaction_time_window: None,
1359 sst_format,
1360 append_mode: None,
1361 }
1362 }
1363
1364 fn build_fs_object_store() -> ObjectStore {
1365 ObjectStore::new(Fs::default().root("/tmp"))
1366 .unwrap()
1367 .finish()
1368 }
1369
1370 #[test]
1371 fn test_initial_pruned_entry_id() {
1372 assert_eq!(0, initial_pruned_entry_id(&WalOptions::RaftEngine));
1373 assert_eq!(0, initial_pruned_entry_id(&WalOptions::Noop));
1374 assert_eq!(
1375 0,
1376 initial_pruned_entry_id(&WalOptions::Kafka(KafkaWalOptions::new(
1377 "test_topic".to_string()
1378 )))
1379 );
1380 assert_eq!(
1381 42,
1382 initial_pruned_entry_id(&WalOptions::Kafka(KafkaWalOptions {
1383 topic: "test_topic".to_string(),
1384 initial_pruned_entry_id: Some(42),
1385 }))
1386 );
1387 }
1388
1389 #[test]
1390 #[cfg(not(feature = "test-shared-fs-region-migration"))]
1391 fn test_open_requirement_rejects_fs_object_store() {
1392 let object_store = build_fs_object_store();
1393
1394 assert!(!supports_open_region_object_storage_requirement(
1395 &object_store
1396 ));
1397 }
1398
1399 #[test]
1400 #[cfg(feature = "test-shared-fs-region-migration")]
1401 fn test_open_requirement_accepts_shared_fs_object_store_for_tests() {
1402 let object_store = build_fs_object_store();
1403
1404 assert!(supports_open_region_object_storage_requirement(
1405 &object_store
1406 ));
1407 }
1408
1409 #[test]
1410 fn test_open_requirement_accepts_s3_object_store() {
1411 let object_store = ObjectStore::new(
1412 S3::default()
1413 .bucket("test-bucket")
1414 .region("us-east-1")
1415 .disable_ec2_metadata(),
1416 )
1417 .unwrap()
1418 .finish();
1419
1420 assert!(supports_open_region_object_storage_requirement(
1421 &object_store
1422 ));
1423 }
1424
1425 #[test]
1426 fn test_sanitize_region_options_options_format_wins() {
1427 let manifest = build_test_manifest(FormatType::PrimaryKey);
1430 let mut options = RegionOptions {
1431 sst_format: Some(FormatType::Flat),
1432 ..Default::default()
1433 };
1434 sanitize_region_options(&manifest, &mut options);
1435 assert_eq!(options.sst_format, Some(FormatType::Flat));
1436 }
1437
1438 #[test]
1439 fn test_sanitize_region_options_fills_from_manifest_when_unset() {
1440 let manifest = build_test_manifest(FormatType::Flat);
1443 let mut options = RegionOptions {
1444 sst_format: None,
1445 ..Default::default()
1446 };
1447 sanitize_region_options(&manifest, &mut options);
1448 assert_eq!(options.sst_format, Some(FormatType::Flat));
1449 }
1450
1451 fn sst_parquet_bytes(batch: &RecordBatch) -> Vec<u8> {
1452 let key_value_meta = KeyValue::new(
1453 PARQUET_METADATA_KEY.to_string(),
1454 sst_region_metadata().to_json().unwrap(),
1455 );
1456 let props = WriterProperties::builder()
1457 .set_key_value_metadata(Some(vec![key_value_meta]))
1458 .build();
1459
1460 let mut parquet_bytes = Vec::new();
1461 let mut writer =
1462 ArrowWriter::try_new(&mut parquet_bytes, batch.schema(), Some(props)).unwrap();
1463 writer.write(batch).unwrap();
1464 writer.close().unwrap();
1465
1466 parquet_bytes
1467 }
1468
1469 #[tokio::test]
1470 async fn test_preload_parquet_meta_cache_uses_file_cache() {
1471 let env = TestEnv::new().await;
1472
1473 let local_store = ObjectStore::new(Memory::default()).unwrap().finish();
1474 let write_cache = env
1475 .create_write_cache(local_store, ReadableSize::mb(1024))
1476 .await;
1477 let cache_manager = Arc::new(
1478 CacheManager::builder()
1479 .sst_meta_cache_size(1024 * 1024)
1480 .write_cache(Some(write_cache.clone()))
1481 .build(),
1482 );
1483
1484 let region_id = RegionId::new(1, 1);
1485 let file_id = FileId::random();
1486
1487 let col = Arc::new(Int64Array::from_iter_values([1, 2, 3])) as ArrayRef;
1488 let primary_key = Arc::new(BinaryArray::from_iter_values([b"a", b"b", b"c"])) as ArrayRef;
1489 let batch = RecordBatch::try_from_iter([
1490 ("col", col),
1491 (
1492 store_api::storage::consts::PRIMARY_KEY_COLUMN_NAME,
1493 primary_key,
1494 ),
1495 ])
1496 .unwrap();
1497 let parquet_bytes = sst_parquet_bytes(&batch);
1498 let file_size = parquet_bytes.len() as u64;
1499
1500 let file_meta = FileMeta {
1501 region_id,
1502 file_id,
1503 time_range: (Timestamp::new_millisecond(0), Timestamp::new_millisecond(1)),
1504 level: 0,
1505 file_size,
1506 max_row_group_uncompressed_size: 0,
1507 available_indexes: Default::default(),
1508 indexes: vec![],
1509 index_file_size: 0,
1510 index_version: 0,
1511 num_rows: 3,
1512 num_row_groups: 1,
1513 sequence: None,
1514 partition_expr: None,
1515 num_series: 0,
1516 ..Default::default()
1517 };
1518 let file_handle = FileHandle::new(file_meta, Arc::new(NoopFilePurger));
1519
1520 let table_dir = "test_table";
1521 let path_type = PathType::Bare;
1522 let remote_path = file_handle.file_path(table_dir, path_type);
1523
1524 let source_store = ObjectStore::new(Memory::default()).unwrap().finish();
1525 source_store
1526 .write(&remote_path, parquet_bytes)
1527 .await
1528 .unwrap();
1529
1530 let index_key = IndexKey::new(region_id, file_id, FileType::Parquet);
1532 write_cache
1533 .file_cache()
1534 .download(index_key, &remote_path, &source_store, file_size)
1535 .await
1536 .unwrap();
1537
1538 let region_file_id = file_handle.file_id();
1539 assert!(
1540 cache_manager
1541 .get_parquet_meta_data_from_mem_cache(region_file_id)
1542 .is_none()
1543 );
1544
1545 let loaded = preload_parquet_meta_cache_for_files(
1546 region_id,
1547 cache_manager.clone(),
1548 1024 * 1024,
1549 table_dir.to_string(),
1550 path_type,
1551 source_store.clone(),
1552 Arc::new(sst_region_metadata()),
1553 vec![file_handle.clone()],
1554 )
1555 .await;
1556
1557 assert_eq!(loaded, 1);
1559 assert!(
1560 cache_manager
1561 .get_compact_sst_meta_data(
1562 region_file_id,
1563 parquet::file::metadata::PageIndexPolicy::Optional,
1564 )
1565 .await
1566 .is_some()
1567 );
1568 assert!(file_handle.primary_key_range().is_some());
1569 }
1570
1571 #[tokio::test]
1572 async fn test_preload_parquet_meta_cache_skips_files_not_in_file_cache() {
1573 let cache_manager = Arc::new(
1574 CacheManager::builder()
1575 .sst_meta_cache_size(1024 * 1024)
1576 .build(),
1577 );
1578
1579 let region_id = RegionId::new(1, 1);
1580 let file_id = FileId::random();
1581
1582 let file_meta = FileMeta {
1584 region_id,
1585 file_id,
1586 time_range: (Timestamp::new_millisecond(0), Timestamp::new_millisecond(1)),
1587 level: 0,
1588 file_size: 0,
1589 max_row_group_uncompressed_size: 0,
1590 available_indexes: Default::default(),
1591 indexes: vec![],
1592 index_file_size: 0,
1593 index_version: 0,
1594 num_rows: 3,
1595 num_row_groups: 1,
1596 sequence: None,
1597 partition_expr: None,
1598 num_series: 0,
1599 ..Default::default()
1600 };
1601 let file_handle = FileHandle::new(file_meta, Arc::new(NoopFilePurger));
1602
1603 let table_dir = "test_table";
1604 let path_type = PathType::Bare;
1605 let remote_path = file_handle.file_path(table_dir, path_type);
1606
1607 let object_store = ObjectStore::new(Memory::default()).unwrap().finish();
1609 object_store
1610 .write(&remote_path, b"noop".as_slice())
1611 .await
1612 .unwrap();
1613
1614 let region_file_id = file_handle.file_id();
1615 assert!(
1616 cache_manager
1617 .get_parquet_meta_data_from_mem_cache(region_file_id)
1618 .is_none()
1619 );
1620
1621 let loaded = preload_parquet_meta_cache_for_files(
1622 region_id,
1623 cache_manager.clone(),
1624 1024 * 1024,
1625 table_dir.to_string(),
1626 path_type,
1627 object_store,
1628 Arc::new(sst_region_metadata()),
1629 vec![file_handle],
1630 )
1631 .await;
1632
1633 assert_eq!(loaded, 0);
1634 assert!(
1635 cache_manager
1636 .get_parquet_meta_data_from_mem_cache(region_file_id)
1637 .is_none()
1638 );
1639 }
1640
1641 #[tokio::test]
1642 async fn test_preload_parquet_meta_cache_loads_from_local_fs() {
1643 let cache_manager = Arc::new(
1644 CacheManager::builder()
1645 .sst_meta_cache_size(1024 * 1024)
1646 .build(),
1647 );
1648
1649 let region_id = RegionId::new(1, 1);
1650 let file_id = FileId::random();
1651
1652 let col = Arc::new(Int64Array::from_iter_values([1, 2, 3])) as ArrayRef;
1653 let batch = RecordBatch::try_from_iter([("col", col)]).unwrap();
1654 let parquet_bytes = sst_parquet_bytes(&batch);
1655
1656 let file_meta = FileMeta {
1659 region_id,
1660 file_id,
1661 time_range: (Timestamp::new_millisecond(0), Timestamp::new_millisecond(1)),
1662 level: 0,
1663 file_size: 0,
1664 max_row_group_uncompressed_size: 0,
1665 available_indexes: Default::default(),
1666 indexes: vec![],
1667 index_file_size: 0,
1668 index_version: 0,
1669 num_rows: 3,
1670 num_row_groups: 1,
1671 sequence: None,
1672 partition_expr: None,
1673 num_series: 0,
1674 ..Default::default()
1675 };
1676 let file_handle = FileHandle::new(file_meta, Arc::new(NoopFilePurger));
1677
1678 let table_dir = "test_table";
1679 let path_type = PathType::Bare;
1680 let file_path = file_handle.file_path(table_dir, path_type);
1681
1682 let root = create_temp_dir("parquet-meta-preload");
1683 let object_store = ObjectStore::new(Fs::default().root(root.path().to_str().unwrap()))
1684 .unwrap()
1685 .finish();
1686 object_store.write(&file_path, parquet_bytes).await.unwrap();
1687
1688 let region_file_id = file_handle.file_id();
1689 assert!(
1690 cache_manager
1691 .get_parquet_meta_data_from_mem_cache(region_file_id)
1692 .is_none()
1693 );
1694
1695 let loaded = preload_parquet_meta_cache_for_files(
1696 region_id,
1697 cache_manager.clone(),
1698 1024 * 1024,
1699 table_dir.to_string(),
1700 path_type,
1701 object_store,
1702 Arc::new(sst_region_metadata()),
1703 vec![file_handle],
1704 )
1705 .await;
1706
1707 assert_eq!(loaded, 1);
1708 assert!(
1709 cache_manager
1710 .get_compact_sst_meta_data(region_file_id, PageIndexPolicy::Optional)
1711 .await
1712 .is_some()
1713 );
1714 }
1715}