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