diff --git a/src/mito2/AGENTS.md b/src/mito2/AGENTS.md index 20ca9b53ca..b22ee758bb 100644 --- a/src/mito2/AGENTS.md +++ b/src/mito2/AGENTS.md @@ -26,7 +26,7 @@ snapshot isolation). It implements the `RegionEngine` trait from `store-api`. | `compaction` | `src/mito2/src/compaction/` | Compaction scheduler (`scheduler.rs` + `scheduler/`), TWCS picker, strict-window manual picker, compactor, memory control | | `access_layer` | `src/mito2/src/access_layer.rs` | SST read/write over the object store | | `sst` | `src/mito2/src/sst/` | Parquet format, file metadata, index layout | -| `series_index` | `src/mito2/src/series_index/` | Incremental writer and predicate searcher for aggregate series index files | +| `series_index` | `src/mito2/src/series_index/` | Series index writer/searcher, catalogs, immutable snapshots, file lifecycle, and background maintenance | | `read` | `src/mito2/src/read/` | `ScanRegion`, merge, dedup, projection, streaming | | `manifest` | `src/mito2/src/manifest/` | `RegionManifestManager`, manifest actions/edits | | `cache` | `src/mito2/src/cache.rs` | Write/file/page caches | diff --git a/src/mito2/src/config.rs b/src/mito2/src/config.rs index be74fa67c6..6ec790c52b 100644 --- a/src/mito2/src/config.rs +++ b/src/mito2/src/config.rs @@ -32,6 +32,7 @@ use crate::gc::GcConfig; use crate::sst::DEFAULT_WRITE_BUFFER_SIZE; const MULTIPART_UPLOAD_MINIMUM_SIZE: ReadableSize = ReadableSize::mb(5); +const DEFAULT_SERIES_INDEX_MAINTENANCE_INTERVAL: Duration = Duration::from_secs(5 * 60); /// Default maximum number of SST files to scan concurrently. pub(crate) const DEFAULT_MAX_CONCURRENT_SCAN_FILES: usize = 384; @@ -88,6 +89,14 @@ pub struct MitoConfig { // Background job configs: /// Max number of running background index build jobs (default: 1/8 of cpu cores). pub max_background_index_builds: usize, + // TODO: Document both series-index settings in the example configs and regenerate configuration + // docs before exposing the feature. + /// Under development; do not enable. Root directory for loading series indexes, currently stored + /// on the local filesystem. Empty disables the feature. Relative paths resolve under `data_home`. + pub experimental_series_index_root: String, + /// Interval between series-index maintenance runs (default 5 min). Zero uses the default. + #[serde(with = "humantime_serde")] + pub experimental_series_index_maintenance_interval: Duration, /// Max number of running background flush jobs (default: 1/2 of cpu cores). pub max_background_flushes: usize, /// Max number of running background compaction jobs (default: 1/4 of cpu cores). @@ -205,6 +214,9 @@ impl Default for MitoConfig { experimental_manifest_keep_removed_file_ttl: Duration::from_secs(60 * 60), compress_manifest: false, max_background_index_builds: divide_num_cpus(8), + experimental_series_index_root: String::new(), + experimental_series_index_maintenance_interval: + DEFAULT_SERIES_INDEX_MAINTENANCE_INTERVAL, max_background_flushes: divide_num_cpus(2), max_background_compactions: divide_num_cpus(4), max_background_purges: get_total_cpu_cores(), @@ -291,6 +303,24 @@ impl MitoConfig { self.max_background_purges = cpu_cores; } + if self + .experimental_series_index_maintenance_interval + .is_zero() + { + warn!("Sanitize series-index maintenance interval 0 to 5 minutes"); + self.experimental_series_index_maintenance_interval = + DEFAULT_SERIES_INDEX_MAINTENANCE_INTERVAL; + } + + if !self.experimental_series_index_root.trim().is_empty() + && Path::new(&self.experimental_series_index_root).is_relative() + { + self.experimental_series_index_root = Path::new(data_home) + .join(&self.experimental_series_index_root) + .display() + .to_string(); + } + if self.global_write_buffer_reject_size <= self.global_write_buffer_size { self.global_write_buffer_reject_size = self.global_write_buffer_size * 2; warn!( @@ -402,6 +432,36 @@ mod tests { assert_eq!(ReadableSize::mb(128), config.prefilter_result_cache_size); } + #[test] + fn test_series_index_config() { + assert!( + MitoConfig::default() + .experimental_series_index_root + .is_empty() + ); + let mut config: MitoConfig = toml::from_str( + "experimental_series_index_root = 'indexes' + experimental_series_index_maintenance_interval = '30s'", + ) + .unwrap(); + config.sanitize("/data").unwrap(); + assert_eq!(config.experimental_series_index_root, "/data/indexes"); + assert_eq!( + config.experimental_series_index_maintenance_interval, + Duration::from_secs(30) + ); + let restored: MitoConfig = toml::from_str(&toml::to_string(&config).unwrap()).unwrap(); + assert_eq!(config, restored); + + let mut config: MitoConfig = + toml::from_str("experimental_series_index_maintenance_interval = '0s'").unwrap(); + config.sanitize("/data").unwrap(); + assert_eq!( + config.experimental_series_index_maintenance_interval, + MitoConfig::default().experimental_series_index_maintenance_interval + ); + } + #[test] fn test_experimental_series_scan_v2_config() { assert!(MitoConfig::default().experimental_series_scan_v2); diff --git a/src/mito2/src/metrics.rs b/src/mito2/src/metrics.rs index 17e18c38f5..2d2decbde7 100644 --- a/src/mito2/src/metrics.rs +++ b/src/mito2/src/metrics.rs @@ -319,6 +319,12 @@ lazy_static! { // Index metrics. lazy_static! { // Index metrics. + /// Series and range index file operations. + pub static ref SERIES_INDEX_FILE_OPERATION_TOTAL: IntCounterVec = register_int_counter_vec!( + "greptime_mito_series_index_file_operation_total", + "series-index file operations", + &["index_type", "operation", "result"], + ).unwrap(); /// Number of stale index publications rejected at each publication stage. pub static ref INDEX_PUBLICATION_STALE_TOTAL: IntCounterVec = register_int_counter_vec!( diff --git a/src/mito2/src/region.rs b/src/mito2/src/region.rs index 072ec6ca62..10195f5680 100644 --- a/src/mito2/src/region.rs +++ b/src/mito2/src/region.rs @@ -58,6 +58,7 @@ use crate::manifest::action::{ use crate::manifest::manager::RegionManifestManager; use crate::region::version::{VersionControlRef, VersionRef}; use crate::request::{OnFailure, OptionOutputTx}; +use crate::series_index::{SeriesIndexVersion, SeriesIndexVersionControl}; use crate::sst::file::FileMeta; use crate::sst::file_purger::FilePurgerRef; use crate::sst::location::{index_file_path, sst_file_path}; @@ -150,6 +151,8 @@ pub struct MitoRegion { /// /// We MUST update the version control inside the write lock of the region manifest manager. pub(crate) version_control: VersionControlRef, + /// Snapshot controller for range and series indexes. + pub(crate) series_index_version_control: SeriesIndexVersionControl, /// SSTs accessor for this region. pub(crate) access_layer: AccessLayerRef, /// Context to maintain manifest for this region. @@ -238,6 +241,12 @@ impl StagingPartitionInfo { } impl MitoRegion { + /// Returns the current immutable series-index snapshot. + #[allow(dead_code)] // Used by the upcoming query integration. + pub(crate) fn series_index_version(&self) -> Arc { + self.series_index_version_control.current() + } + fn remove_region_metrics(&self) { let region_id = self.region_id.as_u64().to_string(); let labels = &[region_id.as_str()]; @@ -2040,6 +2049,7 @@ mod tests { MitoRegion { region_id: metadata.region_id, version_control, + series_index_version_control: Default::default(), access_layer: env.access_layer.clone(), manifest_ctx, file_purger: crate::test_util::new_noop_file_purger(), @@ -2543,6 +2553,7 @@ mod tests { let region = MitoRegion { region_id: metadata.region_id, version_control, + series_index_version_control: Default::default(), access_layer, manifest_ctx: manifest_ctx.clone(), file_purger: crate::test_util::new_noop_file_purger(), diff --git a/src/mito2/src/region/opener.rs b/src/mito2/src/region/opener.rs index ca9baeef34..ebabbbd295 100644 --- a/src/mito2/src/region/opener.rs +++ b/src/mito2/src/region/opener.rs @@ -61,6 +61,7 @@ use crate::memtable::bulk::part::BulkPart; use crate::memtable::time_partition::{TimePartitions, TimePartitionsRef}; use crate::memtable::{MemtableBuilderProvider, ensure_json2_not_use_time_series_memtable}; use crate::metrics::{CACHE_FILL_DOWNLOADED_FILES, CACHE_FILL_PENDING_FILES}; +use crate::read::series_candidate::is_sparse_metric_metadata; use crate::region::options::RegionOptions; use crate::region::version::{VersionBuilder, VersionControl, VersionControlRef}; use crate::region::{ @@ -70,6 +71,7 @@ use crate::region::{ use crate::region_write_ctx::RegionWriteCtx; use crate::request::OptionOutputTx; use crate::schedule::scheduler::SchedulerRef; +use crate::series_index::{IndexFilePurger, load_version_control}; use crate::sst::FormatType; use crate::sst::file::{FileHandle, RegionFileId, RegionIndexId}; use crate::sst::file_purger::{FilePurgerRef, create_file_purger}; @@ -79,6 +81,7 @@ use crate::sst::index::puffin_manager::PuffinManagerFactory; use crate::sst::location::{self, region_dir_from_table_dir}; use crate::sst::parquet::metadata::{MetadataLoader, extract_primary_key_range}; use crate::sst::parquet::reader::MetadataCacheMetrics; +use crate::sst::range_index::RangeIndexDeleter; use crate::time_provider::TimeProviderRef; use crate::wal::entry_reader::WalEntryReader; use crate::wal::{EntryId, Wal}; @@ -163,6 +166,8 @@ pub(crate) struct RegionOpener { file_ref_manager: FileReferenceManagerRef, partition_expr_fetcher: PartitionExprFetcherRef, hook: Option, + series_index_store: Option, + series_index_purger: Option, } impl RegionOpener { @@ -202,9 +207,23 @@ impl RegionOpener { file_ref_manager, partition_expr_fetcher, hook: None, + series_index_store: None, + series_index_purger: None, } } + /// Sets the store for companion range indexes. + pub(crate) fn series_index_store(mut self, store: Option) -> Self { + self.series_index_store = store; + self + } + + /// Sets the purger shared with the worker's series-index maintenance task. + pub(crate) fn series_index_purger(mut self, purger: Option) -> Self { + self.series_index_purger = purger; + self + } + /// Sets the region hook for observing manifest mutations. pub(crate) fn hook(mut self, hook: Option) -> Self { self.hook = hook; @@ -420,6 +439,7 @@ impl RegionOpener { Ok(Arc::new(MitoRegion { region_id, version_control, + series_index_version_control: Default::default(), access_layer: access_layer.clone(), // Region is writable after it is created. manifest_ctx: Arc::new(ManifestContext::new( @@ -434,6 +454,8 @@ impl RegionOpener { access_layer, self.cache_manager, self.file_ref_manager.clone(), + self.series_index_store + .map(|store| RangeIndexDeleter::new(store, region_id)), ), provider, last_flush_millis: AtomicI64::new(now), @@ -549,6 +571,9 @@ impl RegionOpener { access_layer.clone(), self.cache_manager.clone(), self.file_ref_manager.clone(), + self.series_index_store + .clone() + .map(|store| RangeIndexDeleter::new(store, region_id)), ); // We should sanitize the region options before creating a new memtable. let memtable_builder = self @@ -647,9 +672,20 @@ impl RegionOpener { let now = self.time_provider.current_time_millis(); + let series_index_version_control = + match (&self.series_index_store, &self.series_index_purger) { + (Some(store), Some(purger)) + if is_sparse_metric_metadata(&version_control.current().version.metadata) => + { + load_version_control(store, self.region_id, purger).await + } + _ => Default::default(), + }; + let region = MitoRegion { region_id: self.region_id, version_control: version_control.clone(), + series_index_version_control, access_layer: access_layer.clone(), // Region is always opened in read only mode. manifest_ctx: Arc::new(ManifestContext::new( diff --git a/src/mito2/src/series_index.rs b/src/mito2/src/series_index.rs index 605c6fec97..f08178456a 100644 --- a/src/mito2/src/series_index.rs +++ b/src/mito2/src/series_index.rs @@ -14,7 +14,15 @@ //! Series index writer and searcher. +// These components are consumed by the upcoming query and maintenance integration. +#[allow(dead_code)] +mod catalog; +#[allow(dead_code)] +mod purger; mod searcher; +mod task; +#[allow(dead_code)] +mod version; mod writer; use futures::stream::BoxStream; @@ -28,6 +36,10 @@ pub use writer::{ }; use crate::error::Result; +pub(crate) use crate::series_index::catalog::{delete_catalogs, load_version_control}; +pub(crate) use crate::series_index::purger::{IndexFilePurger, series_index_channel}; +pub(crate) use crate::series_index::task::{SeriesIndexTaskState, spawn_series_index_tasks}; +pub(crate) use crate::series_index::version::{SeriesIndexVersion, SeriesIndexVersionControl}; pub(crate) const MIN_TS_COLUMN: &str = "__series_min_ts"; pub(crate) const MAX_TS_COLUMN: &str = "__series_max_ts"; diff --git a/src/mito2/src/series_index/catalog.rs b/src/mito2/src/series_index/catalog.rs new file mode 100644 index 0000000000..a59c1926fc --- /dev/null +++ b/src/mito2/src/series_index/catalog.rs @@ -0,0 +1,257 @@ +// Copyright 2023 Greptime Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! Index catalog persistence, coverage metadata, and file paths. + +use common_telemetry::warn; +use common_time::Timestamp; +use object_store::{ErrorKind, ObjectStore}; +use parquet::file::metadata::KeyValue; +use serde::de::DeserializeOwned; +use serde::{Deserialize, Serialize}; +use snafu::ResultExt; +use store_api::storage::{FileId, RegionId}; + +use crate::error::{OpenDalSnafu, Result, SerdeJsonSnafu}; +use crate::series_index::purger::IndexFilePurger; +use crate::series_index::version::{ + SeriesIndexFileHandle, SeriesIndexVersion, SeriesIndexVersionControl, +}; +const SERIES_DIR: &str = "series"; +const RANGE_CATALOG: &str = "range-index.json"; +const SERIES_CATALOG: &str = "series-index.json"; +const SERIES_METADATA_KEY: &str = "greptime.series_index"; + +/// Self-describing coverage stored in a series-index Parquet footer. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub(crate) struct SeriesIndexEntry { + pub(crate) index_uuid: FileId, + /// Inclusive bucket start. + pub(crate) bucket_start: Timestamp, + /// Exclusive bucket end. + pub(crate) bucket_end: Timestamp, + pub(crate) source_file_ids: Vec, + pub(crate) min_file_sequence: u64, + pub(crate) max_file_sequence: u64, +} + +#[derive(Debug, Default, Serialize, Deserialize)] +pub(crate) struct SeriesIndexCatalog { + pub(crate) indexes: Vec, +} + +#[derive(Debug, Default, Serialize, Deserialize)] +pub(crate) struct RangeIndexCatalog { + pub(crate) indexes: Vec, +} + +pub(crate) fn range_catalog_path(region_id: RegionId) -> String { + format!("{}/{RANGE_CATALOG}", region_id.as_u64()) +} + +pub(crate) fn series_index_path(region_id: RegionId, index_uuid: FileId) -> String { + format!("{}/{SERIES_DIR}/{index_uuid}.parquet", region_id.as_u64()) +} + +pub(crate) fn series_catalog_path(region_id: RegionId) -> String { + format!("{}/{SERIES_CATALOG}", region_id.as_u64()) +} + +pub(crate) fn series_metadata(entry: &SeriesIndexEntry) -> Result> { + Ok(vec![KeyValue::new( + SERIES_METADATA_KEY.to_string(), + Some(serde_json::to_string(entry).context(SerdeJsonSnafu)?), + )]) +} + +pub(crate) async fn load_catalog(store: &ObjectStore, path: &str) -> Option +where + T: DeserializeOwned, +{ + let bytes = match store.read(path).await { + Ok(bytes) => bytes.to_bytes(), + Err(error) if error.kind() == ErrorKind::NotFound => return None, + Err(error) => { + warn!(error; "Failed to load series-index catalog, path: {path}"); + return None; + } + }; + match serde_json::from_slice(&bytes) { + Ok(catalog) => Some(catalog), + Err(error) => { + warn!(error; "Invalid series-index catalog, path: {path}, phase: load"); + None + } + } +} + +pub(crate) async fn store_catalog(store: &ObjectStore, path: &str, catalog: &T) -> Result<()> +where + T: Serialize, +{ + let bytes = serde_json::to_vec_pretty(catalog).context(SerdeJsonSnafu)?; + store + .write(path, bytes) + .await + .map(|_| ()) + .context(OpenDalSnafu) +} + +/// Best-effort removal of both catalogs when dropping a region. +pub(crate) async fn delete_catalogs(store: &ObjectStore, region_id: RegionId) { + for path in [ + series_catalog_path(region_id), + range_catalog_path(region_id), + ] { + if let Err(error) = store.delete(&path).await + && error.kind() != ErrorKind::NotFound + { + warn!(error; "Failed to delete index catalog, path: {path}"); + } + } +} + +/// Restores the in-memory snapshot once when opening a region. +pub(crate) async fn load_version_control( + store: &ObjectStore, + region_id: RegionId, + purger: &IndexFilePurger, +) -> SeriesIndexVersionControl { + let range = load_catalog::(store, &range_catalog_path(region_id)) + .await + .unwrap_or_default(); + let series = load_catalog::(store, &series_catalog_path(region_id)) + .await + .unwrap_or_default(); + // TODO: Handle catalog entries whose index files are missing from storage. + let version = SeriesIndexVersion { + range_indexes: range.indexes.into_iter().collect(), + series_indexes: series + .indexes + .into_iter() + .map(|entry| { + ( + entry.index_uuid, + SeriesIndexFileHandle::new(region_id, entry, purger.clone()), + ) + }) + .collect(), + }; + let control = SeriesIndexVersionControl::default(); + control.publish(std::sync::Arc::new(version)); + control +} + +#[cfg(test)] +mod tests { + use std::sync::Arc; + + use common_time::Timestamp; + use object_store::ObjectStore; + use object_store::layers::mock::{self, MockLayerBuilder}; + use object_store::services::Memory; + use store_api::storage::{FileId, RegionId}; + + use crate::series_index::catalog::{ + SeriesIndexCatalog, SeriesIndexEntry, load_catalog, load_version_control, + series_catalog_path, series_metadata, store_catalog, + }; + use crate::series_index::purger::series_index_channel; + + struct FailingCatalogReader; + + impl mock::Read for FailingCatalogReader { + async fn read( + &self, + _range: mock::BytesRange, + ) -> mock::Result<(mock::RpRead, mock::Buffer)> { + Err(mock::Error::new( + mock::ErrorKind::Unexpected, + "injected catalog read failure", + )) + } + + async fn open( + &self, + _range: mock::BytesRange, + ) -> mock::Result<(mock::RpRead, Box)> { + Err(mock::Error::new( + mock::ErrorKind::Unexpected, + "injected catalog read failure", + )) + } + } + + #[tokio::test] + async fn test_load_catalog_returns_none_on_error() { + let store = ObjectStore::new(Memory::default()).unwrap(); + let path = series_catalog_path(RegionId::new(1, 1)); + // Missing catalog. + assert!( + load_catalog::(&store, &path) + .await + .is_none() + ); + store.write(&path, "invalid").await.unwrap(); + assert!( + load_catalog::(&store, &path) + .await + .is_none() + ); + let layer = MockLayerBuilder::default() + .reader_factory(Arc::new(|_, _, _| Box::new(FailingCatalogReader))) + .build() + .unwrap(); + let store = store.layer(layer); + assert!( + load_catalog::(&store, &path) + .await + .is_none() + ); + } + + #[tokio::test] + async fn test_catalog_roundtrip() { + let store = ObjectStore::new(Memory::default()).unwrap(); + let region_id = RegionId::new(1, 1); + let entry = SeriesIndexEntry { + index_uuid: FileId::random(), + bucket_start: Timestamp::new_second(0), + bucket_end: Timestamp::new_second(100), + source_file_ids: vec![FileId::random()], + min_file_sequence: 1, + max_file_sequence: 2, + }; + store_catalog( + &store, + &series_catalog_path(region_id), + &SeriesIndexCatalog { + indexes: vec![entry.clone()], + }, + ) + .await + .unwrap(); + let (purger, _receiver) = series_index_channel(store.clone()); + let current = load_version_control(&store, region_id, &purger) + .await + .current(); + assert!(current.range_indexes.is_empty()); + assert_eq!(&entry, current.series_indexes[&entry.index_uuid].entry()); + + let metadata = series_metadata(&entry).unwrap(); + let decoded: SeriesIndexEntry = + serde_json::from_str(metadata[0].value.as_ref().unwrap()).unwrap(); + assert_eq!(entry, decoded); + } +} diff --git a/src/mito2/src/series_index/purger.rs b/src/mito2/src/series_index/purger.rs new file mode 100644 index 0000000000..039ca66ccc --- /dev/null +++ b/src/mito2/src/series_index/purger.rs @@ -0,0 +1,110 @@ +// Copyright 2023 Greptime Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! Deferred deletion of aggregate series-index files. + +use std::fmt::{self, Debug, Formatter}; + +use common_telemetry::{info, warn}; +use object_store::{ErrorKind, ObjectStore}; +use tokio::sync::mpsc::{UnboundedReceiver, UnboundedSender, unbounded_channel}; + +use crate::metrics::SERIES_INDEX_FILE_OPERATION_TOTAL; +use crate::series_index::catalog::series_index_path; +use crate::sst::file::RegionFileId; + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +pub(crate) enum IndexFileType { + Range, + Series, +} + +impl IndexFileType { + fn as_str(self) -> &'static str { + match self { + Self::Range => "range", + Self::Series => "series", + } + } +} + +#[derive(Debug, Clone, Copy)] +pub(crate) struct PurgeRequest { + pub(crate) file_id: RegionFileId, +} + +#[derive(Clone)] +pub(crate) struct IndexFilePurger { + store: ObjectStore, + sender: UnboundedSender, +} + +impl Debug for IndexFilePurger { + fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result { + f.debug_struct("IndexFilePurger").finish_non_exhaustive() + } +} + +impl IndexFilePurger { + pub(crate) fn purge(&self, request: PurgeRequest) { + if let Err(error) = self.sender.send(request) { + let store = self.store.clone(); + common_runtime::spawn_global(async move { + purge_file(&store, error.0).await; + }); + } + } +} + +pub(crate) fn file_operation(index_type: IndexFileType, operation: &str, result: &str) { + SERIES_INDEX_FILE_OPERATION_TOTAL + .with_label_values(&[index_type.as_str(), operation, result]) + .inc(); +} + +/// Processes queued deletions once each, independently of periodic maintenance. +pub(crate) async fn run_index_purge_task( + worker_id: u32, + store: ObjectStore, + mut receiver: UnboundedReceiver, +) { + info!("Start series-index purge task, worker: {worker_id}"); + while let Some(request) = receiver.recv().await { + purge_file(&store, request).await; + } + info!("Stop series-index purge task, worker: {worker_id}"); +} + +async fn purge_file(store: &ObjectStore, request: PurgeRequest) { + let path = series_index_path(request.file_id.region_id(), request.file_id.file_id()); + match store.delete(&path).await { + Ok(()) => { + file_operation(IndexFileType::Series, "delete", "success"); + } + Err(error) if error.kind() == ErrorKind::NotFound => { + file_operation(IndexFileType::Series, "delete", "success"); + } + Err(error) => { + file_operation(IndexFileType::Series, "delete", "failure"); + warn!(error; "Failed to delete series index, index_type: {}, path: {}, phase: deletion", IndexFileType::Series.as_str(), path); + } + } +} + +pub(crate) fn series_index_channel( + store: ObjectStore, +) -> (IndexFilePurger, UnboundedReceiver) { + let (sender, receiver) = unbounded_channel(); + (IndexFilePurger { store, sender }, receiver) +} diff --git a/src/mito2/src/series_index/searcher.rs b/src/mito2/src/series_index/searcher.rs index 87bc4eac9d..9b88a5b826 100644 --- a/src/mito2/src/series_index/searcher.rs +++ b/src/mito2/src/series_index/searcher.rs @@ -306,6 +306,7 @@ mod tests { object_store, path, SeriesIndexWriterOptions { row_group_size }, + None, ) .await .unwrap(); diff --git a/src/mito2/src/series_index/task.rs b/src/mito2/src/series_index/task.rs new file mode 100644 index 0000000000..76b42dd62f --- /dev/null +++ b/src/mito2/src/series_index/task.rs @@ -0,0 +1,113 @@ +// Copyright 2023 Greptime Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! Worker-owned background maintenance for series indexes. + +use std::sync::Arc; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::time::Duration; + +use common_telemetry::info; +use object_store::ObjectStore; +use tokio::sync::Notify; +use tokio::sync::mpsc::UnboundedReceiver; +use tokio::task::JoinHandle; +use tokio::time::{Instant, MissedTickBehavior}; + +use crate::series_index::purger::{PurgeRequest, run_index_purge_task}; + +/// Shared lifecycle state for a worker's series-index task. +#[derive(Debug)] +pub(crate) struct SeriesIndexTaskState { + running: AtomicBool, + notify: Notify, +} + +impl SeriesIndexTaskState { + pub(crate) fn new() -> Self { + Self { + running: AtomicBool::new(true), + notify: Notify::new(), + } + } + + pub(crate) fn is_running(&self) -> bool { + self.running.load(Ordering::Acquire) + } + + pub(crate) fn stop(&self) { + self.running.store(false, Ordering::Release); + // Retain a permit if maintenance has not started waiting yet. + self.notify.notify_one(); + } + + pub(crate) async fn notified(&self) { + self.notify.notified().await; + } +} + +/// Starts both tasks, detaching purge and returning the maintenance handle. +pub(crate) fn spawn_series_index_tasks( + worker_id: u32, + store: ObjectStore, + state: Arc, + purge_receiver: UnboundedReceiver, + interval: Duration, +) -> JoinHandle<()> { + // Snapshots may retain senders after the worker stops; purge until all senders drop. + common_runtime::spawn_global(run_index_purge_task(worker_id, store, purge_receiver)); + common_runtime::spawn_global(async move { + SeriesIndexTask { + worker_id, + state, + interval, + } + .run() + .await; + }) +} + +/// Periodic series-index maintenance for one region worker. +struct SeriesIndexTask { + worker_id: u32, + state: Arc, + interval: Duration, +} + +impl SeriesIndexTask { + /// Runs periodic maintenance until the worker stops. + async fn run(mut self) { + let worker_id = self.worker_id; + info!("Start series-index background task, worker: {worker_id}"); + let interval = self.interval; + let mut timer = tokio::time::interval_at(Instant::now() + interval, interval); + timer.set_missed_tick_behavior(MissedTickBehavior::Skip); + while self.state.is_running() { + tokio::select! { + _ = self.state.notified() => {} + _ = timer.tick() => { + if self.state.is_running() { + self.maintain().await; + } + } + } + } + info!("Stop series-index background task, worker: {worker_id}"); + } + + /// Runs periodic maintenance independently of incoming deletion requests. + async fn maintain(&mut self) { + // TODO: Reconcile indexes and perform other periodic maintenance here. + } +} diff --git a/src/mito2/src/series_index/version.rs b/src/mito2/src/series_index/version.rs new file mode 100644 index 0000000000..7db23da2f9 --- /dev/null +++ b/src/mito2/src/series_index/version.rs @@ -0,0 +1,119 @@ +// Copyright 2023 Greptime Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! Immutable index snapshots and aggregate series-file handles. + +use std::collections::{HashMap, HashSet}; +use std::fmt::{self, Debug, Formatter}; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::{Arc, RwLock}; + +use store_api::storage::{FileId, RegionId}; + +use crate::series_index::catalog::SeriesIndexEntry; +use crate::series_index::purger::{IndexFilePurger, PurgeRequest}; +use crate::sst::file::RegionFileId; + +/// A reference-counted series-index file with deferred deletion semantics. +#[derive(Clone)] +pub(crate) struct SeriesIndexFileHandle { + inner: Arc, +} + +impl Debug for SeriesIndexFileHandle { + fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result { + f.debug_struct("SeriesIndexFileHandle") + .field("file_id", &self.inner.file_id) + .field("deleted", &self.inner.deleted.load(Ordering::Relaxed)) + .finish() + } +} + +impl SeriesIndexFileHandle { + pub(crate) fn new( + region_id: RegionId, + entry: SeriesIndexEntry, + purger: IndexFilePurger, + ) -> Self { + Self { + inner: Arc::new(SeriesIndexFileHandleInner { + file_id: RegionFileId::new(region_id, entry.index_uuid), + entry, + deleted: AtomicBool::new(false), + purger, + }), + } + } + + pub(crate) fn entry(&self) -> &SeriesIndexEntry { + &self.inner.entry + } + + pub(crate) fn mark_deleted(&self) { + self.inner.deleted.store(true, Ordering::Release); + } +} + +struct SeriesIndexFileHandleInner { + file_id: RegionFileId, + entry: SeriesIndexEntry, + deleted: AtomicBool, + purger: IndexFilePurger, +} + +impl Drop for SeriesIndexFileHandleInner { + fn drop(&mut self) { + if self.deleted.load(Ordering::Acquire) { + self.purger.purge(PurgeRequest { + file_id: self.file_id, + }); + } + } +} + +/// Immutable series-index snapshot for one region. +#[derive(Debug, Default)] +pub(crate) struct SeriesIndexVersion { + pub(crate) range_indexes: HashSet, + pub(crate) series_indexes: HashMap, +} + +impl SeriesIndexVersion { + fn mark_all_deleted(&self) { + self.series_indexes + .values() + .for_each(SeriesIndexFileHandle::mark_deleted); + } +} + +/// Copy-on-write series-index snapshots owned by a region. +#[derive(Debug, Default)] +pub(crate) struct SeriesIndexVersionControl { + current: RwLock>, +} + +impl SeriesIndexVersionControl { + pub(crate) fn current(&self) -> Arc { + self.current.read().unwrap().clone() + } + + pub(crate) fn publish(&self, next: Arc) -> Arc { + std::mem::replace(&mut *self.current.write().unwrap(), next) + } + + pub(crate) fn mark_dropped(&self) { + self.publish(Arc::new(SeriesIndexVersion::default())) + .mark_all_deleted(); + } +} diff --git a/src/mito2/src/series_index/writer.rs b/src/mito2/src/series_index/writer.rs index 5c0dc3be89..db8e105ed4 100644 --- a/src/mito2/src/series_index/writer.rs +++ b/src/mito2/src/series_index/writer.rs @@ -27,6 +27,7 @@ use datatypes::timestamp::timestamp_array_to_primitive; use datatypes::value::Value; use mito_codec::row_converter::{CompositeValues, PrimaryKeyCodec, build_primary_key_codec}; use object_store::ObjectStore; +use parquet::file::metadata::KeyValue; use snafu::{OptionExt, ResultExt, ensure}; use store_api::codec::PrimaryKeyEncoding; use store_api::metadata::RegionMetadataRef; @@ -129,6 +130,7 @@ impl SeriesIndexWriter { object_store: ObjectStore, path: &str, options: SeriesIndexWriterOptions, + key_value_metadata: Option>, ) -> Result { let open_start = Instant::now(); ensure!( @@ -145,6 +147,7 @@ impl SeriesIndexWriter { path, &schema, options.row_group_size, + key_value_metadata, ) .await?; let codec = build_primary_key_codec(&metadata); @@ -763,6 +766,7 @@ mod tests { store.clone(), &path, SeriesIndexWriterOptions::default(), + None, ) .await .err() @@ -814,6 +818,7 @@ mod tests { store.clone(), "series.parquet", SeriesIndexWriterOptions { row_group_size: 2 }, + None, ) .await .unwrap(); @@ -921,6 +926,7 @@ mod tests { store.clone(), "groups.parquet", SeriesIndexWriterOptions { row_group_size: 2 }, + None, ) .await .unwrap(); @@ -941,6 +947,7 @@ mod tests { store.clone(), "empty.parquet", SeriesIndexWriterOptions::default(), + None, ) .await .unwrap() @@ -967,6 +974,7 @@ mod tests { store.clone(), "abort.parquet", SeriesIndexWriterOptions { row_group_size: 1 }, + None, ) .await .unwrap(); @@ -995,6 +1003,7 @@ mod tests { store.clone(), "dictionary-abort.parquet", SeriesIndexWriterOptions { row_group_size: 1 }, + None, ) .await .unwrap(); @@ -1050,6 +1059,7 @@ mod tests { store.clone(), "nullable.parquet", SeriesIndexWriterOptions::default(), + None, ) .await .unwrap(); @@ -1072,6 +1082,7 @@ mod tests { store, "invalid.parquet", SeriesIndexWriterOptions { row_group_size: 0 }, + None, ) .await .is_err() diff --git a/src/mito2/src/sst/file_purger.rs b/src/mito2/src/sst/file_purger.rs index 5f0617751d..935c027258 100644 --- a/src/mito2/src/sst/file_purger.rs +++ b/src/mito2/src/sst/file_purger.rs @@ -17,6 +17,7 @@ use std::sync::Arc; use common_telemetry::error; use store_api::region_request::PathType; +use store_api::storage::FileId; use crate::access_layer::AccessLayerRef; use crate::cache::CacheManagerRef; @@ -24,6 +25,7 @@ use crate::error::Result; use crate::schedule::scheduler::SchedulerRef; use crate::sst::file::{FileMeta, delete_files, delete_index}; use crate::sst::file_ref::FileReferenceManagerRef; +use crate::sst::range_index::RangeIndexDeleter; /// A worker to delete files in background. pub trait FilePurger: Send + Sync + fmt::Debug { @@ -58,6 +60,7 @@ pub struct LocalFilePurger { scheduler: SchedulerRef, sst_layer: AccessLayerRef, cache_manager: Option, + range_index_deleter: Option, } impl fmt::Debug for LocalFilePurger { @@ -88,7 +91,8 @@ pub fn should_enable_gc(global_gc_enabled: bool, _object_store_scheme: &'static /// the files from both the storage and the cache. /// /// If the storage is an object store, an `ObjectStoreFilePurger` is created, which -/// only manages the file references without deleting the actual files. +/// only manages SST file references. Companion range indexes are deleted directly in either +/// mode on final handle release when the SST is marked deleted and a deleter is provided. /// pub fn create_file_purger( gc_enabled: bool, @@ -97,6 +101,7 @@ pub fn create_file_purger( sst_layer: AccessLayerRef, cache_manager: Option, file_ref_manager: FileReferenceManagerRef, + range_index_deleter: Option, ) -> FilePurgerRef { // Only enable GC for: // - object store based storage @@ -104,9 +109,16 @@ pub fn create_file_purger( if should_enable_gc(gc_enabled, sst_layer.object_store().info().scheme()) && matches!(path_type, PathType::Data | PathType::Bare) { - Arc::new(ObjectStoreFilePurger { file_ref_manager }) + Arc::new(ObjectStoreFilePurger { + file_ref_manager, + scheduler, + range_index_deleter, + }) } else { - Arc::new(LocalFilePurger::new(scheduler, sst_layer, cache_manager)) + Arc::new( + LocalFilePurger::new(scheduler, sst_layer, cache_manager) + .with_range_index_deleter(range_index_deleter), + ) } } @@ -116,8 +128,12 @@ pub fn create_local_file_purger( sst_layer: AccessLayerRef, cache_manager: Option, _file_ref_manager: FileReferenceManagerRef, + range_index_deleter: Option, ) -> FilePurgerRef { - Arc::new(LocalFilePurger::new(scheduler, sst_layer, cache_manager)) + Arc::new( + LocalFilePurger::new(scheduler, sst_layer, cache_manager) + .with_range_index_deleter(range_index_deleter), + ) } impl LocalFilePurger { @@ -131,9 +147,16 @@ impl LocalFilePurger { scheduler, sst_layer, cache_manager, + range_index_deleter: None, } } + /// Attaches deletion of companion range indexes. + pub fn with_range_index_deleter(mut self, deleter: Option) -> Self { + self.range_index_deleter = deleter; + self + } + /// Stop the scheduler of the file purger. pub async fn stop_scheduler(&self) -> Result<()> { self.scheduler.stop(true).await @@ -177,6 +200,11 @@ impl LocalFilePurger { impl FilePurger for LocalFilePurger { fn remove_file(&self, file_meta: FileMeta, is_delete: bool, index_outdated: bool) { if is_delete { + schedule_range_index_deletion( + &self.scheduler, + self.range_index_deleter.as_ref(), + file_meta.file_id, + ); self.delete_file(file_meta); } else if index_outdated { self.delete_index(file_meta); @@ -184,18 +212,53 @@ impl FilePurger for LocalFilePurger { } } -#[derive(Debug)] pub struct ObjectStoreFilePurger { file_ref_manager: FileReferenceManagerRef, + scheduler: SchedulerRef, + range_index_deleter: Option, +} + +impl fmt::Debug for ObjectStoreFilePurger { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.debug_struct("ObjectStoreFilePurger") + .field("file_ref_manager", &self.file_ref_manager) + .field("range_index_deleter", &self.range_index_deleter) + .finish_non_exhaustive() + } +} + +/// Range indexes are deleted directly even when SST deletion is delegated to GC. +fn schedule_range_index_deletion( + scheduler: &SchedulerRef, + deleter: Option<&RangeIndexDeleter>, + file_id: FileId, +) { + let Some(deleter) = deleter.cloned() else { + return; + }; + if let Err(error) = scheduler.schedule(Box::pin(async move { + if let Err(error) = deleter.delete(file_id).await { + error!(error; "Failed to delete range index, file_id: {file_id}"); + } + })) { + error!(error; "Failed to schedule range-index deletion, file_id: {file_id}"); + } } impl FilePurger for ObjectStoreFilePurger { - fn remove_file(&self, file_meta: FileMeta, _is_delete: bool, _index_outdated: bool) { + fn remove_file(&self, file_meta: FileMeta, is_delete: bool, _index_outdated: bool) { // if not on local file system, instead inform the global file purger to remove the file reference. // notice that no matter whether the file is deleted or not, we need to remove the reference // because the file is no longer in use nonetheless. // for same reason, we don't care about index_outdated here. self.file_ref_manager.remove_file(&file_meta); + if is_delete { + schedule_range_index_deletion( + &self.scheduler, + self.range_index_deleter.as_ref(), + file_meta.file_id, + ); + } } fn new_file(&self, file_meta: &FileMeta) { diff --git a/src/mito2/src/sst/parquet/index_writer.rs b/src/mito2/src/sst/parquet/index_writer.rs index 44dfbc44f1..6d6414b9cd 100644 --- a/src/mito2/src/sst/parquet/index_writer.rs +++ b/src/mito2/src/sst/parquet/index_writer.rs @@ -21,6 +21,7 @@ use parquet::arrow::AsyncArrowWriter; use parquet::arrow::async_writer::AsyncFileWriter; use parquet::basic::{Compression, Encoding, ZstdLevel}; use parquet::errors::ParquetError; +use parquet::file::metadata::KeyValue; use parquet::file::properties::WriterProperties; use snafu::{OptionExt, ResultExt}; @@ -94,6 +95,7 @@ impl ParquetIndexWriter { path: &str, schema: &SchemaRef, row_group_size: usize, + key_value_metadata: Option>, ) -> Result { let file_name = path.rsplit('/').next().unwrap_or(path).to_string(); let output = object_store @@ -108,6 +110,7 @@ impl ParquetIndexWriter { .set_max_row_group_row_count(Some(row_group_size)) .set_column_index_truncate_length(None) .set_statistics_truncate_length(None) + .set_key_value_metadata(key_value_metadata) .build(); let writer = AsyncArrowWriter::try_new(AsyncWriter::new(output), schema.clone(), Some(properties)) diff --git a/src/mito2/src/sst/range_index.rs b/src/mito2/src/sst/range_index.rs index fa766acacb..c4e576b6cd 100644 --- a/src/mito2/src/sst/range_index.rs +++ b/src/mito2/src/sst/range_index.rs @@ -14,6 +14,7 @@ //! Per-SST series row-range index. +mod deleter; mod searcher; mod writer; @@ -26,6 +27,11 @@ pub use writer::{ SstRangeIndexWriter, SstRangeIndexWriterMetrics, SstRangeIndexWriterOptions, range_index_schema, }; +pub use crate::sst::range_index::deleter::RangeIndexDeleter; +// Used by the upcoming query and index-building integration. +#[allow(unused_imports)] +pub(crate) use crate::sst::range_index::deleter::range_index_path; + const ROW_GROUP_ID_COLUMN: &str = "row_group_id"; const START_COLUMN: &str = "start"; const END_COLUMN: &str = "end"; diff --git a/src/mito2/src/sst/range_index/deleter.rs b/src/mito2/src/sst/range_index/deleter.rs new file mode 100644 index 0000000000..5865f0b0b9 --- /dev/null +++ b/src/mito2/src/sst/range_index/deleter.rs @@ -0,0 +1,58 @@ +// Copyright 2023 Greptime Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! Direct deletion of per-SST range index files. + +use object_store::{ErrorKind, ObjectStore}; +use snafu::ResultExt; +use store_api::storage::{FileId, RegionId}; + +use crate::error::{OpenDalSnafu, Result}; +use crate::metrics::SERIES_INDEX_FILE_OPERATION_TOTAL; + +/// Deletes range indexes belonging to one region, independently of SST garbage collection. +#[derive(Debug, Clone)] +pub struct RangeIndexDeleter { + store: ObjectStore, + region_id: RegionId, +} + +impl RangeIndexDeleter { + /// Creates a deleter using the owning region ID, including for imported SSTs. + pub fn new(store: ObjectStore, region_id: RegionId) -> Self { + Self { store, region_id } + } + + /// Deletes the range index directly from the index store. + pub async fn delete(&self, file_id: FileId) -> Result<()> { + let path = range_index_path(self.region_id, file_id); + let result = match self.store.delete(&path).await { + Ok(()) => Ok(()), + Err(error) if error.kind() == ErrorKind::NotFound => Ok(()), + Err(error) => Err(error).context(OpenDalSnafu), + }; + SERIES_INDEX_FILE_OPERATION_TOTAL + .with_label_values(&[ + "range", + "delete", + if result.is_ok() { "success" } else { "failure" }, + ]) + .inc(); + result + } +} + +pub(crate) fn range_index_path(region_id: RegionId, file_id: FileId) -> String { + format!("{}/range/{file_id}.parquet", region_id.as_u64()) +} diff --git a/src/mito2/src/sst/range_index/writer.rs b/src/mito2/src/sst/range_index/writer.rs index aa28cf1d30..ac1058b5a5 100644 --- a/src/mito2/src/sst/range_index/writer.rs +++ b/src/mito2/src/sst/range_index/writer.rs @@ -147,6 +147,7 @@ impl SstRangeIndexWriter { path, &schema, options.index_row_group_size, + None, ) .await?; let codec = SparsePrimaryKeyCodec::new(&metadata); diff --git a/src/mito2/src/worker.rs b/src/mito2/src/worker.rs index 0c23affaf5..8efcf47ab6 100644 --- a/src/mito2/src/worker.rs +++ b/src/mito2/src/worker.rs @@ -45,6 +45,7 @@ use common_runtime::JoinHandle; use common_stat::get_total_memory_bytes; use common_telemetry::{error, info, warn}; use futures::future::try_join_all; +use object_store::ObjectStore; use object_store::manager::ObjectStoreManagerRef; use prometheus::{Histogram, IntGauge}; use rand::{Rng, rng}; @@ -57,6 +58,7 @@ use store_api::storage::{FileId, RegionId}; use tokio::sync::mpsc::{Receiver, Sender}; use tokio::sync::{Mutex, Semaphore, mpsc, oneshot, watch}; +use crate::access_layer::new_fs_cache_store; use crate::cache::write_cache::{WriteCache, WriteCacheRef}; use crate::cache::{CacheManager, CacheManagerRef, WriteCacheUploadStoreWrapperRef}; use crate::compaction::CompactionScheduler; @@ -77,6 +79,9 @@ use crate::request::{ SenderDdlRequest, SenderWriteRequest, WorkerRequest, WorkerRequestWithTime, }; use crate::schedule::scheduler::{LocalScheduler, SchedulerRef}; +use crate::series_index::{ + IndexFilePurger, SeriesIndexTaskState, series_index_channel, spawn_series_index_tasks, +}; use crate::sst::file::RegionFileId; use crate::sst::file_ref::FileReferenceManagerRef; use crate::sst::index::IndexBuildScheduler; @@ -190,6 +195,11 @@ impl WorkerGroup { .with_buffer_size(Some(config.index.write_buffer_size.as_bytes() as _)); let index_build_job_pool = Arc::new(LocalScheduler::new(config.max_background_index_builds)); + let series_index_store = if config.experimental_series_index_root.trim().is_empty() { + None + } else { + Some(new_fs_cache_store(&config.experimental_series_index_root).await?) + }; let flush_job_pool = Arc::new(LocalScheduler::new(config.max_background_flushes)); let compact_job_pool = Arc::new(LocalScheduler::new(config.max_background_compactions)); let flush_semaphore = Arc::new(Semaphore::new(config.max_background_flushes)); @@ -242,6 +252,7 @@ impl WorkerGroup { object_store_manager: object_store_manager.clone(), write_buffer_manager: write_buffer_manager.clone(), index_build_job_pool: index_build_job_pool.clone(), + series_index_store: series_index_store.clone(), flush_job_pool: flush_job_pool.clone(), compact_job_pool: compact_job_pool.clone(), purge_scheduler: purge_scheduler.clone(), @@ -399,6 +410,11 @@ impl WorkerGroup { }); let index_build_job_pool = Arc::new(LocalScheduler::new(config.max_background_index_builds)); + let series_index_store = if config.experimental_series_index_root.trim().is_empty() { + None + } else { + Some(new_fs_cache_store(&config.experimental_series_index_root).await?) + }; let flush_job_pool = Arc::new(LocalScheduler::new(config.max_background_flushes)); let compact_job_pool = Arc::new(LocalScheduler::new(config.max_background_compactions)); let flush_semaphore = Arc::new(Semaphore::new(config.max_background_flushes)); @@ -452,6 +468,7 @@ impl WorkerGroup { object_store_manager: object_store_manager.clone(), write_buffer_manager: write_buffer_manager.clone(), index_build_job_pool: index_build_job_pool.clone(), + series_index_store: series_index_store.clone(), flush_job_pool: flush_job_pool.clone(), compact_job_pool: compact_job_pool.clone(), purge_scheduler: purge_scheduler.clone(), @@ -546,6 +563,7 @@ struct WorkerStarter { write_buffer_manager: WriteBufferManagerRef, compact_job_pool: SchedulerRef, index_build_job_pool: SchedulerRef, + series_index_store: Option, flush_job_pool: SchedulerRef, purge_scheduler: SchedulerRef, listener: WorkerListener, @@ -574,6 +592,26 @@ impl WorkerStarter { let (sender, receiver) = mpsc::channel(self.config.worker_channel_size); let running = Arc::new(AtomicBool::new(true)); + let series_index_task_state = self + .series_index_store + .as_ref() + .map(|_| Arc::new(SeriesIndexTaskState::new())); + let mut series_index_purger = None; + let series_index_handle = self + .series_index_store + .clone() + .zip(series_index_task_state.clone()) + .map(|(store, state)| { + let (purger, purge_receiver) = series_index_channel(store.clone()); + series_index_purger = Some(purger); + spawn_series_index_tasks( + self.id, + store, + state, + purge_receiver, + self.config.experimental_series_index_maintenance_interval, + ) + }); let now = self.time_provider.current_time_millis(); let id_string = self.id.to_string(); let mut worker_thread = RegionWorkerLoop { @@ -598,6 +636,8 @@ impl WorkerStarter { self.index_build_job_pool, self.config.max_background_index_builds, ), + series_index_store: self.series_index_store, + series_index_purger, flush_scheduler: FlushScheduler::new(self.flush_job_pool), compaction_scheduler: CompactionScheduler::new( self.compact_job_pool, @@ -639,6 +679,8 @@ impl WorkerStarter { catchup_regions, sender, handle: Mutex::new(Some(handle)), + series_index_handle: Mutex::new(series_index_handle), + series_index_task_state, running, }) } @@ -658,6 +700,10 @@ pub(crate) struct RegionWorker { sender: Sender, /// Handle to the worker thread. handle: Mutex>>, + /// Handle to the series-index maintenance task. + series_index_handle: Mutex>>, + /// Controls the worker-owned series-index maintenance task. + series_index_task_state: Option>, /// Whether to run the worker thread. running: Arc, } @@ -674,6 +720,9 @@ impl RegionWorker { ); // Manually set the running flag to false to avoid printing more warning logs. self.set_running(false); + if let Some(state) = &self.series_index_task_state { + state.stop(); + } return WorkerStoppedSnafu { id: self.id }.fail(); } @@ -685,10 +734,15 @@ impl RegionWorker { /// This method waits until the worker thread exists. async fn stop(&self) -> Result<()> { let handle = self.handle.lock().await.take(); + self.set_running(false); + if let Some(state) = &self.series_index_task_state { + state.stop(); + } + + let mut worker_result = Ok(()); if let Some(handle) = handle { info!("Stop region worker {}", self.id); - self.set_running(false); if self .sender .send(WorkerRequestWithTime::new(WorkerRequest::Stop)) @@ -698,8 +752,16 @@ impl RegionWorker { warn!("Worker {} is already exited before stop", self.id); } - handle.await.context(JoinSnafu)?; + worker_result = handle.await.context(JoinSnafu); } + let series_index_result = if let Some(handle) = self.series_index_handle.lock().await.take() + { + handle.await.context(JoinSnafu) + } else { + Ok(()) + }; + worker_result?; + series_index_result?; Ok(()) } @@ -749,6 +811,9 @@ impl RegionWorker { impl Drop for RegionWorker { fn drop(&mut self) { + if let Some(state) = &self.series_index_task_state { + state.stop(); + } if self.is_running() { self.set_running(false); // Once we drop the sender, the worker thread will receive a disconnected error. @@ -870,6 +935,9 @@ struct RegionWorkerLoop { write_buffer_manager: WriteBufferManagerRef, /// Scheduler for index build task. index_build_scheduler: IndexBuildScheduler, + /// Store for companion range indexes deleted by the region SST purger. + series_index_store: Option, + series_index_purger: Option, /// Schedules background flush requests. flush_scheduler: FlushScheduler, /// Scheduler for compaction tasks. @@ -1556,10 +1624,14 @@ mod tests { let group = env .create_worker_group(MitoConfig { num_workers: 4, + experimental_series_index_root: "series-index".to_string(), ..Default::default() }) .await; - group.stop().await.unwrap(); + tokio::time::timeout(Duration::from_secs(5), group.stop()) + .await + .expect("series-index tasks should stop without waiting for their interval") + .unwrap(); } } diff --git a/src/mito2/src/worker/handle_catchup.rs b/src/mito2/src/worker/handle_catchup.rs index 71d3f05c69..5ef45f82f7 100644 --- a/src/mito2/src/worker/handle_catchup.rs +++ b/src/mito2/src/worker/handle_catchup.rs @@ -125,6 +125,8 @@ impl RegionWorkerLoop { self.partition_expr_fetcher.clone(), ) .cache(Some(self.cache_manager.clone())) + .series_index_store(self.series_index_store.clone()) + .series_index_purger(self.series_index_purger.clone()) .hook(self.plugins.get()) .options(region.version().options.clone())? .skip_wal_replay(true) diff --git a/src/mito2/src/worker/handle_create.rs b/src/mito2/src/worker/handle_create.rs index bf254dc1e6..3acb5da76d 100644 --- a/src/mito2/src/worker/handle_create.rs +++ b/src/mito2/src/worker/handle_create.rs @@ -71,6 +71,8 @@ impl RegionWorkerLoop { .metadata_builder(builder) .parse_options(request.options)? .cache(Some(self.cache_manager.clone())) + .series_index_store(self.series_index_store.clone()) + .series_index_purger(self.series_index_purger.clone()) .hook(self.plugins.get()); opener.ensure_region_requirements(requirements)?; diff --git a/src/mito2/src/worker/handle_drop.rs b/src/mito2/src/worker/handle_drop.rs index a012927edc..8112d6a159 100644 --- a/src/mito2/src/worker/handle_drop.rs +++ b/src/mito2/src/worker/handle_drop.rs @@ -126,8 +126,13 @@ where .await?; self.cleanup_dropped_region_runtime_state(region_id).await; + if let Some(store) = &self.series_index_store { + crate::series_index::delete_catalogs(store, region_id).await; + } + // Marks region version as dropped region.version_control.mark_dropped(); + region.series_index_version_control.mark_dropped(); info!( "Region {} is dropped logically, but some files are not deleted yet", region_id diff --git a/src/mito2/src/worker/handle_open.rs b/src/mito2/src/worker/handle_open.rs index 8193acdfc1..825a677b76 100644 --- a/src/mito2/src/worker/handle_open.rs +++ b/src/mito2/src/worker/handle_open.rs @@ -196,6 +196,8 @@ impl RegionWorkerLoop { ) .skip_wal_replay(request.skip_wal_replay) .cache(Some(self.cache_manager.clone())) + .series_index_store(self.series_index_store.clone()) + .series_index_purger(self.series_index_purger.clone()) .hook(self.plugins.get()) .wal_entry_reader(wal_entry_receiver.map(|receiver| Box::new(receiver) as _)) .replay_checkpoint(request.checkpoint.map(|checkpoint| checkpoint.entry_id)) diff --git a/tests-integration/tests/http.rs b/tests-integration/tests/http.rs index bf9c5ea488..c696bd40ac 100644 --- a/tests-integration/tests/http.rs +++ b/tests-integration/tests/http.rs @@ -2332,6 +2332,8 @@ manifest_checkpoint_distance = 10 experimental_manifest_keep_removed_file_count = 256 experimental_manifest_keep_removed_file_ttl = "1h" compress_manifest = false +experimental_series_index_root = "" +experimental_series_index_maintenance_interval = "5m" experimental_compaction_memory_limit = "unlimited" experimental_compaction_on_exhausted = "wait" auto_flush_interval = "10m"