diff --git a/config/config.md b/config/config.md index 6544a28f93d..e380f075548 100644 --- a/config/config.md +++ b/config/config.md @@ -177,6 +177,7 @@ | `region_engine.mito.manifest_checkpoint_distance` | Integer | `10` | Number of meta action updated to trigger a new checkpoint for the manifest. | | `region_engine.mito.compress_manifest` | Bool | `false` | Whether to compress manifest and checkpoint file by gzip (default false). | | `region_engine.mito.experimental_enable_series_index` | Bool | `false` | Under development; do not enable. Whether to enable series indexes.
Indexes are stored on the local filesystem under `{data_home}/series_index`. | +| `region_engine.mito.experimental_series_index_max_size` | String | `5GiB` | Approximate series and range index size limit in open regions, shared across workers.
Workers periodically refresh usage and skip maintenance when full. In-flight reconciliation
can exceed the limit. Closed-region files, temporary output, catalogs, and old snapshots
retained by readers are not counted.
Minimum: 1KiB. Takes effect on restart. | | `region_engine.mito.experimental_enable_range_index` | Bool | `false` | Whether to build and query range indexes when series indexes are enabled.
Obsolete range-index metadata and files are still cleaned up when disabled. | | `region_engine.mito.experimental_series_index_maintenance_interval` | String | `5m` | Interval between series-index maintenance runs. Zero uses the default of 5 min. | | `region_engine.mito.experimental_series_index_bucket_width` | String | `5days` | Requested minimum series-index bucket width (default: 5 days), rounded up to
an exact multiple of each region's compaction time window. | @@ -639,6 +640,7 @@ | `region_engine.mito.experimental_manifest_keep_removed_file_ttl` | String | `1h` | How long to keep removed files in the `removed_files` field of manifest
after they are removed from manifest.
files will only be removed from `removed_files` field
if both `keep_removed_file_count` and `keep_removed_file_ttl` is reached. | | `region_engine.mito.compress_manifest` | Bool | `false` | Whether to compress manifest and checkpoint file by gzip (default false). | | `region_engine.mito.experimental_enable_series_index` | Bool | `false` | Under development; do not enable. Whether to enable series indexes.
Indexes are stored on the local filesystem under `{data_home}/series_index`. | +| `region_engine.mito.experimental_series_index_max_size` | String | `5GiB` | Approximate series and range index size limit in open regions, shared across workers.
Workers periodically refresh usage and skip maintenance when full. In-flight reconciliation
can exceed the limit. Closed-region files, temporary output, catalogs, and old snapshots
retained by readers are not counted.
Minimum: 1KiB. Takes effect on restart. | | `region_engine.mito.experimental_enable_range_index` | Bool | `false` | Whether to build and query range indexes when series indexes are enabled.
Obsolete range-index metadata and files are still cleaned up when disabled. | | `region_engine.mito.experimental_series_index_maintenance_interval` | String | `5m` | Interval between series-index maintenance runs. Zero uses the default of 5 min. | | `region_engine.mito.experimental_series_index_bucket_width` | String | `5days` | Requested minimum series-index bucket width (default: 5 days), rounded up to
an exact multiple of each region's compaction time window. | diff --git a/config/datanode.example.toml b/config/datanode.example.toml index 45baaefbdfb..da2d353fef8 100644 --- a/config/datanode.example.toml +++ b/config/datanode.example.toml @@ -502,6 +502,13 @@ compress_manifest = false ## Indexes are stored on the local filesystem under `{data_home}/series_index`. #+ experimental_enable_series_index = false +## Approximate series and range index size limit in open regions, shared across workers. +## Workers periodically refresh usage and skip maintenance when full. In-flight reconciliation +## can exceed the limit. Closed-region files, temporary output, catalogs, and old snapshots +## retained by readers are not counted. +## Minimum: 1KiB. Takes effect on restart. +#+ experimental_series_index_max_size = "5GiB" + ## Whether to build and query range indexes when series indexes are enabled. ## Obsolete range-index metadata and files are still cleaned up when disabled. #+ experimental_enable_range_index = false diff --git a/config/standalone.example.toml b/config/standalone.example.toml index ef9dd561f6a..e0a8deb9c22 100644 --- a/config/standalone.example.toml +++ b/config/standalone.example.toml @@ -707,6 +707,13 @@ compress_manifest = false ## Indexes are stored on the local filesystem under `{data_home}/series_index`. #+ experimental_enable_series_index = false +## Approximate series and range index size limit in open regions, shared across workers. +## Workers periodically refresh usage and skip maintenance when full. In-flight reconciliation +## can exceed the limit. Closed-region files, temporary output, catalogs, and old snapshots +## retained by readers are not counted. +## Minimum: 1KiB. Takes effect on restart. +#+ experimental_series_index_max_size = "5GiB" + ## Whether to build and query range indexes when series indexes are enabled. ## Obsolete range-index metadata and files are still cleaned up when disabled. #+ experimental_enable_range_index = false diff --git a/src/cmd/tests/load_config_test.rs b/src/cmd/tests/load_config_test.rs index 34c5bc8c9ba..d4b19812df2 100644 --- a/src/cmd/tests/load_config_test.rs +++ b/src/cmd/tests/load_config_test.rs @@ -159,6 +159,7 @@ fn test_load_datanode_example_config() { }, region_engine: vec![ RegionEngineConfig::Mito(MitoConfig { + experimental_series_index_max_size: ReadableSize::gb(5), auto_flush_interval: Duration::from_secs(10 * 60), default_region_write_buffer_size: ReadableSize::mb(0), write_cache_ttl: Some(Duration::from_secs(60 * 60 * 8)), @@ -397,6 +398,7 @@ fn test_load_standalone_example_config() { }), region_engine: vec![ RegionEngineConfig::Mito(MitoConfig { + experimental_series_index_max_size: ReadableSize::gb(5), auto_flush_interval: Duration::from_secs(10 * 60), default_region_write_buffer_size: ReadableSize::mb(0), write_cache_ttl: Some(Duration::from_secs(60 * 60 * 8)), diff --git a/src/mito2/AGENTS.md b/src/mito2/AGENTS.md index 583ec31bd4a..b8fd2f8892c 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/` | Series index writer/searcher, bucket planning, catalogs, worker-owned maintenance and purge tasks, and index lifecycle | +| `series_index` | `src/mito2/src/series_index/` | Series index writer/searcher, bucket planning, catalogs, worker-owned reconciliation (`maintenance.rs`), shared approximate usage (`task.rs`), catalogs, and index lifecycle | | `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 | @@ -62,6 +62,10 @@ filtered `RecordBatch` stream. - **Manifest format** (`manifest/action.rs`): affects crash recovery and follower replay. Keep it backward compatible. - **SST/Parquet layout** (`sst/`): readers must stay compatible with existing files. +- **Series-index disk limit**: tasks estimate usage from entry sizes in open-region snapshots + and defer new builds when the shared estimate is full. Cleanup still prunes obsolete + range entries and expires series indexes to reclaim published usage. Concurrent builds may + overshoot; closed-region files are not counted. Series handles and the SST purger own file deletion. - **Series-index coverage** (`series_index/catalog.rs`): `SeriesIndexEntry` stores compaction-window width and SST summaries keyed by aligned start in both catalogs and Parquet footers. - **Request types** (`request.rs`): usually tied to proto definitions consumed by `datanode`. diff --git a/src/mito2/src/config.rs b/src/mito2/src/config.rs index 5164e2fd72b..035feb75a15 100644 --- a/src/mito2/src/config.rs +++ b/src/mito2/src/config.rs @@ -92,6 +92,10 @@ pub struct MitoConfig { /// Under development; do not enable. Whether to enable series indexes (default false). /// Indexes are stored on the local filesystem under `{data_home}/series_index`. pub experimental_enable_series_index: bool, + /// Approximate series and range index size limit in open regions (default: 5 GiB). + /// Workers share a periodically refreshed estimate and skip maintenance when full. + /// In-flight reconciliation can exceed the limit; closed-region files are not counted. + pub experimental_series_index_max_size: ReadableSize, /// Whether to build and query range indexes when series indexes are enabled (default false). /// Obsolete range-index metadata and files are still cleaned up when disabled. pub experimental_enable_range_index: bool, @@ -221,6 +225,7 @@ impl Default for MitoConfig { compress_manifest: false, max_background_index_builds: divide_num_cpus(8), experimental_enable_series_index: false, + experimental_series_index_max_size: ReadableSize::gb(5), experimental_enable_range_index: false, experimental_series_index_maintenance_interval: DEFAULT_SERIES_INDEX_MAINTENANCE_INTERVAL, @@ -280,6 +285,14 @@ impl MitoConfig { /// /// Returns an error if there is a configuration that unable to sanitize. pub fn sanitize(&mut self, data_home: &str) -> Result<()> { + if self.experimental_enable_series_index { + snafu::ensure!( + self.experimental_series_index_max_size.as_bytes() >= 1024, + crate::error::InvalidConfigSnafu { + reason: "experimental_series_index_max_size must be at least 1KiB" + } + ); + } // Use default value if `num_workers` is 0. if self.num_workers == 0 { self.num_workers = divide_num_cpus(2); @@ -442,12 +455,17 @@ mod tests { let mut config: MitoConfig = toml::from_str( "experimental_enable_series_index = true experimental_enable_range_index = false + experimental_series_index_max_size = '64MiB' experimental_series_index_maintenance_interval = '30s' experimental_series_index_bucket_width = '2days'", ) .unwrap(); config.sanitize("/data").unwrap(); assert!(config.experimental_enable_series_index); + assert_eq!( + ReadableSize::mb(64), + config.experimental_series_index_max_size + ); assert!(!config.experimental_enable_range_index); assert_eq!( config.experimental_series_index_maintenance_interval, @@ -459,6 +477,10 @@ mod tests { ); let restored: MitoConfig = toml::from_str(&toml::to_string(&config).unwrap()).unwrap(); assert_eq!(config, restored); + config.experimental_series_index_max_size = ReadableSize(1023); + assert!(config.sanitize("/data").is_err()); + config.experimental_series_index_max_size = ReadableSize(1024); + config.sanitize("/data").unwrap(); let mut config: MitoConfig = toml::from_str("experimental_series_index_maintenance_interval = '0s'").unwrap(); diff --git a/src/mito2/src/engine/index_build_test.rs b/src/mito2/src/engine/index_build_test.rs index ce950f8a246..d2d59b038ff 100644 --- a/src/mito2/src/engine/index_build_test.rs +++ b/src/mito2/src/engine/index_build_test.rs @@ -1512,46 +1512,6 @@ async fn test_index_build_type_manual_consistency() { assert_listener_counts(&listener, 2, 2); } -#[tokio::test] -async fn test_gate_index_build_listener_smoke() { - use store_api::storage::{FileId, RegionId}; - - use crate::engine::listener::{EventListener, GateIndexBuildListener}; - use crate::sst::file::RegionFileId; - - let gate = Arc::new(GateIndexBuildListener::default()); - - // Initial counts are zero. - assert_eq!(gate.begin_count(), 0); - assert_eq!(gate.finish_count(), 0); - assert_eq!(gate.abort_count(), 0); - - // Spawn a task that will block in on_index_build_begin. - let gate_clone = gate.clone(); - let handle = tokio::spawn(async move { - gate_clone - .on_index_build_begin(RegionFileId::new(RegionId::new(1, 1), FileId::random())) - .await; - }); - - // Wait for begin to arrive. - tokio::time::timeout(std::time::Duration::from_secs(5), gate.wait_begin(1)) - .await - .unwrap(); - assert_eq!(gate.begin_count(), 1); - assert_eq!(gate.finish_count(), 0); - assert_eq!(gate.abort_count(), 0); - - // Release the blocked begin. - gate.release_begin(); - - // The spawned task should now complete. - tokio::time::timeout(std::time::Duration::from_secs(5), handle) - .await - .unwrap() - .unwrap(); -} - #[tokio::test] async fn test_index_build_type_manual_duplicate_in_flight() { let mut env = TestEnv::with_prefix("test_index_build_type_manual_duplicate_in_flight_").await; diff --git a/src/mito2/src/engine/listener.rs b/src/mito2/src/engine/listener.rs index 6ef4a293a2a..432e3367603 100644 --- a/src/mito2/src/engine/listener.rs +++ b/src/mito2/src/engine/listener.rs @@ -583,12 +583,12 @@ impl EventListener for GateIndexBuildListener { info!("Region {} index build begin (gated)", region_file_id); self.begin_count.fetch_add(1, Ordering::Relaxed); self.begin_notify.notify_one(); - // Block until the test releases the gate. - let _permit = self - .begin_blocker + // Consume the release so subsequent builds remain blocked. + self.begin_blocker .acquire() .await - .expect("gate semaphore should not be closed"); + .expect("gate semaphore should not be closed") + .forget(); } async fn on_index_build_finish(&self, region_file_id: RegionFileId) { diff --git a/src/mito2/src/engine/scan_test.rs b/src/mito2/src/engine/scan_test.rs index bf089737f35..e10968d3a32 100644 --- a/src/mito2/src/engine/scan_test.rs +++ b/src/mito2/src/engine/scan_test.rs @@ -1340,6 +1340,7 @@ async fn check_two_phase_series_scan( let store = region.series_index_store.clone().unwrap(); let sequence = region.flushed_sequence(); let entry = SeriesIndexEntry { + file_size: 0, index_uuid: FileId::random(), bucket_start: Timestamp::new_millisecond(0), bucket_end: Timestamp::new_millisecond(2001), @@ -1412,7 +1413,13 @@ async fn check_two_phase_series_scan( .unwrap(); writer.write(0, &batch).await.unwrap(); writer.finish().await.unwrap(); - index_version.range_indexes.insert(file_id); + index_version.range_indexes.insert( + file_id, + crate::series_index::RangeIndexEntry { + file_id, + file_size: 0, + }, + ); } region .series_index_version_control @@ -1607,7 +1614,7 @@ async fn check_two_phase_series_scan( // succeed even when that file is unavailable. let region = engine.find_region(region_id).unwrap(); let version = region.series_index_version_control.current(); - let file_id = *version.range_indexes.iter().next().unwrap(); + let file_id = *version.range_indexes.keys().next().unwrap(); region .series_index_store .as_ref() diff --git a/src/mito2/src/metrics.rs b/src/mito2/src/metrics.rs index 49e91d93efd..4c275c248c9 100644 --- a/src/mito2/src/metrics.rs +++ b/src/mito2/src/metrics.rs @@ -318,6 +318,14 @@ lazy_static! { // Index metrics. lazy_static! { + /// Approximate published index bytes in open regions, refreshed by maintenance. + pub static ref SERIES_INDEX_DISK_BYTES: IntGauge = register_int_gauge!( + "greptime_mito_series_index_disk_bytes", "estimated series and range index bytes in open regions" + ).unwrap(); + /// Maintenance passes deferred by the estimated disk usage. + pub static ref SERIES_INDEX_CAPACITY_DEFERRED: IntCounter = register_int_counter!( + "greptime_mito_series_index_capacity_deferred_total", "series-index capacity deferrals" + ).unwrap(); // Index metrics. /// Outcomes of series-index reconciliation passes. pub static ref SERIES_INDEX_RECONCILE_TOTAL: IntCounterVec = register_int_counter_vec!( diff --git a/src/mito2/src/read/series_candidate.rs b/src/mito2/src/read/series_candidate.rs index f05f54e6662..3ba5512a7bd 100644 --- a/src/mito2/src/read/series_candidate.rs +++ b/src/mito2/src/read/series_candidate.rs @@ -778,6 +778,7 @@ mod tests { )); let store = env.access_layer.object_store().clone(); let entry = SeriesIndexEntry { + file_size: 0, index_uuid: FileId::random(), bucket_start: Timestamp::new_millisecond(0), bucket_end: Timestamp::new_millisecond(20), diff --git a/src/mito2/src/series_index.rs b/src/mito2/src/series_index.rs index 4a4133e1a9b..3aa70256709 100644 --- a/src/mito2/src/series_index.rs +++ b/src/mito2/src/series_index.rs @@ -39,7 +39,7 @@ use store_api::metric_engine_consts::{ use crate::error::Result; #[cfg(test)] -pub(crate) use crate::series_index::catalog::SeriesIndexEntry; +pub(crate) use crate::series_index::catalog::{RangeIndexEntry, SeriesIndexEntry}; pub(crate) use crate::series_index::catalog::{ delete_catalogs, load_version_control, series_index_path, }; diff --git a/src/mito2/src/series_index/bucket.rs b/src/mito2/src/series_index/bucket.rs index 7fecb253441..446948838f6 100644 --- a/src/mito2/src/series_index/bucket.rs +++ b/src/mito2/src/series_index/bucket.rs @@ -334,6 +334,7 @@ impl SeriesBucket { .filter_map(|file| file.meta_ref().sequence.map(|sequence| sequence.get())) .min()?; Some(SeriesIndexEntry { + file_size: 0, index_uuid: FileId::random(), bucket_start: self.start, bucket_end: self.end, diff --git a/src/mito2/src/series_index/builder.rs b/src/mito2/src/series_index/builder.rs index 8d361be364e..63a51badaeb 100644 --- a/src/mito2/src/series_index/builder.rs +++ b/src/mito2/src/series_index/builder.rs @@ -21,7 +21,6 @@ use common_telemetry::warn; use futures::TryStreamExt; use object_store::ObjectStore; use snafu::{OptionExt, ensure}; -use store_api::storage::FileId; use crate::error::{Result, UnexpectedSnafu}; use crate::read::BoxedRecordBatchStream; @@ -34,7 +33,7 @@ use crate::region::MitoRegionRef; use crate::region::version::VersionRef; use crate::series_index::bucket::SeriesBucket; use crate::series_index::catalog::{ - SeriesIndexEntry, range_index_path, series_index_path, series_metadata, + RangeIndexEntry, SeriesIndexEntry, range_index_path, series_index_path, series_metadata, }; use crate::series_index::purger::{IndexFilePurger, IndexFileType, file_operation}; use crate::series_index::version::SeriesIndexFileHandle; @@ -68,7 +67,7 @@ pub(crate) async fn build_range_index( region: &MitoRegionRef, version: &VersionRef, file: FileHandle, -) -> Result> { +) -> Result> { let file_id = file.file_id().file_id(); let Some((context, mut selection)) = reader_input(region, file).await? else { return Ok(None); @@ -116,9 +115,12 @@ pub(crate) async fn build_range_index( } return Err(error); } - writer.finish().await?; + let metrics = writer.finish().await?; file_operation(IndexFileType::Range, "build", "success"); - Ok(Some(file_id)) + Ok(Some(RangeIndexEntry { + file_id, + file_size: metrics.output_bytes, + })) } /// Builds only the series index. Callers build needed range indexes separately. @@ -200,11 +202,14 @@ pub(crate) async fn build_series_index( } return Err(error); } - writer.finish().await?; + let metrics = writer.finish().await?; file_operation(IndexFileType::Series, "build", "success"); Ok(SeriesIndexFileHandle::new( region.region_id, - entry.clone(), + SeriesIndexEntry { + file_size: metrics.output_bytes, + ..entry.clone() + }, purger.clone(), )) } @@ -275,10 +280,19 @@ mod tests { let mut range_bytes = HashMap::new(); if build_ranges { for file in files { - let id = build_range_index(&store, ®ion, &version, file.clone()) + let completed = build_range_index(&store, ®ion, &version, file.clone()) .await .unwrap() .unwrap(); + let id = completed.file_id; + assert_eq!( + completed.file_size, + store + .stat(&range_index_path(region.region_id, id)) + .await + .unwrap() + .content_length() + ); range_bytes.insert( id, store @@ -314,7 +328,21 @@ mod tests { assert!(result.is_err()); } else { let handle = result.unwrap(); - assert_eq!(handle.entry(), &entry); + assert_eq!( + handle.entry().file_size, + store + .stat(&series_index_path(region.region_id, entry.index_uuid)) + .await + .unwrap() + .content_length() + ); + assert_eq!( + handle.entry(), + &SeriesIndexEntry { + file_size: handle.entry().file_size, + ..entry.clone() + } + ); assert!( store .exists(&series_index_path(region.region_id, entry.index_uuid)) @@ -504,12 +532,13 @@ mod tests { let mut completed = Vec::new(); let result: Result> = async { for file in files { - let Some(id) = build_range_index(&store, ®ion, &version, file.clone()).await? + let Some(entry) = + build_range_index(&store, ®ion, &version, file.clone()).await? else { // Defer series construction if a needed range is not ready. return Ok(None); }; - completed.push(id); + completed.push(entry.file_id); } build_series_index(&store, ®ion, &version, &bucket, &entry, &purger) .await diff --git a/src/mito2/src/series_index/catalog.rs b/src/mito2/src/series_index/catalog.rs index d56e10698fd..95cf87f83d3 100644 --- a/src/mito2/src/series_index/catalog.rs +++ b/src/mito2/src/series_index/catalog.rs @@ -60,6 +60,8 @@ pub(crate) struct WindowSequence { /// relies on this complete-coverage contract, not on `source_file_ids`. #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub(crate) struct SeriesIndexEntry { + /// Completed size in the catalog; zero in the footer written before completion. + pub(crate) file_size: u64, pub(crate) index_uuid: FileId, /// Inclusive bucket start. pub(crate) bucket_start: Timestamp, @@ -97,6 +99,13 @@ impl SeriesIndexEntry { } } +/// A completed per-SST range index. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +pub(crate) struct RangeIndexEntry { + pub(crate) file_id: FileId, + pub(crate) file_size: u64, +} + #[derive(Debug, Default, Serialize, Deserialize)] pub(crate) struct SeriesIndexCatalog { pub(crate) indexes: Vec, @@ -104,7 +113,7 @@ pub(crate) struct SeriesIndexCatalog { #[derive(Debug, Default, Serialize, Deserialize)] pub(crate) struct RangeIndexCatalog { - pub(crate) indexes: Vec, + pub(crate) indexes: Vec, } pub(crate) fn range_catalog_path(region_id: RegionId) -> String { @@ -185,9 +194,12 @@ pub(crate) async fn load_version_control( 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::new( - range.indexes.into_iter().collect(), + range + .indexes + .into_iter() + .map(|entry| (entry.file_id, entry)) + .collect(), series .indexes .into_iter() @@ -216,9 +228,9 @@ mod tests { use store_api::storage::{FileId, RegionId}; use crate::series_index::catalog::{ - RangeIndexCatalog, SeriesIndexCatalog, SeriesIndexEntry, WindowSequence, load_catalog, - load_version_control, range_catalog_path, series_catalog_path, series_metadata, - store_catalog, + RangeIndexCatalog, RangeIndexEntry, SeriesIndexCatalog, SeriesIndexEntry, WindowSequence, + load_catalog, load_version_control, range_catalog_path, series_catalog_path, + series_metadata, store_catalog, }; use crate::series_index::purger::series_index_channel; @@ -250,6 +262,7 @@ mod tests { fn coverage_uses_exclusive_time_end_and_inclusive_file_sequences() { let region_id = RegionId::new(1, 1); let entry = SeriesIndexEntry { + file_size: 0, index_uuid: FileId::random(), bucket_start: Timestamp::new_second(1), bucket_end: Timestamp::new_second(2), @@ -304,7 +317,10 @@ mod tests { .write( &range_catalog_path(region_id), serde_json::to_vec(&RangeIndexCatalog { - indexes: vec![file_id], + indexes: vec![RangeIndexEntry { + file_id, + file_size: 1, + }], }) .unwrap(), ) @@ -315,7 +331,8 @@ mod tests { .await .unwrap(); let control = load_version_control(&store, region_id, &purger).await; - assert!(control.current().range_indexes.contains(&file_id)); + assert_eq!(1, control.current().range_indexes.len()); + assert_eq!(1, control.current().range_indexes[&file_id].file_size); assert!(control.current().series_indexes.is_empty()); let layer = MockLayerBuilder::default() .reader_factory(Arc::new(|_, _, _| Box::new(FailingCatalogReader))) @@ -342,6 +359,7 @@ mod tests { let store = ObjectStore::new(Memory::default()).unwrap(); let region_id = RegionId::new(1, 1); let entry = SeriesIndexEntry { + file_size: 0, index_uuid: FileId::random(), bucket_start: Timestamp::new_second(0), bucket_end: Timestamp::new_second(100), diff --git a/src/mito2/src/series_index/maintenance.rs b/src/mito2/src/series_index/maintenance.rs index d782cf5a2c1..f5221bb1ebe 100644 --- a/src/mito2/src/series_index/maintenance.rs +++ b/src/mito2/src/series_index/maintenance.rs @@ -20,7 +20,7 @@ use std::time::{Duration, Instant}; use common_telemetry::{debug, info}; use object_store::ObjectStore; -use store_api::storage::RegionId; +use store_api::storage::{FileId, RegionId}; use crate::error::Result; use crate::metrics::{SERIES_INDEX_RECONCILE_ELAPSED, SERIES_INDEX_RECONCILE_TOTAL}; @@ -74,6 +74,7 @@ impl ReconcileStats { } /// Reconciles indexes for one region snapshot, persists catalogs, then atomically publishes it. +#[allow(clippy::too_many_arguments)] pub(crate) async fn reconcile_series_indexes( worker_id: u32, store: ObjectStore, @@ -82,6 +83,7 @@ pub(crate) async fn reconcile_series_indexes( now_ms: i64, purger: IndexFilePurger, enable_range_index: bool, + allow_builds: bool, ) -> Result { let total_start = Instant::now(); // Use this snapshot throughout reconciliation, even if the region version advances. @@ -104,6 +106,7 @@ pub(crate) async fn reconcile_series_indexes( &purger, &mut unpublished, enable_range_index, + allow_builds, ) .await?; SERIES_INDEX_RECONCILE_ELAPSED @@ -164,6 +167,7 @@ async fn build_index_version( purger: &IndexFilePurger, unpublished: &mut UnpublishedSeriesFiles, enable_range_index: bool, + allow_builds: bool, ) -> Result<(Option, ReconcileStats)> { let mut stats = ReconcileStats::default(); let files = version @@ -179,25 +183,23 @@ async fn build_index_version( .map(|file| file.file_id().file_id()) .collect::>(); let current = region.series_index_version(); - // The SST purger deletes companion range files after final handle release. Prune - // metadata here using the captured SST snapshot, independently of physical deletion; - // a later region-version change is picked up by the next reconciliation. - stats.removed_range = current.range_indexes.difference(&visible).count(); - let buckets = match version.compaction_time_window { - Some(window) => rounded_bucket_width(requested_bucket_width, window) + // Empty build inputs still let the planner expire established bucket coverage. + let buckets = match (allow_builds, version.compaction_time_window) { + (true, Some(window)) => rounded_bucket_width(requested_bucket_width, window) // Successful rounding guarantees the window fits in i64; subsecond windows // use the same one-second minimum as rounded_bucket_width. .map(|width| { group_files_into_series_buckets(&files, width, (window.as_secs() as i64).max(1)) }) .unwrap_or_default(), - None => { + (true, None) => { debug!( "Deferring series indexes without compaction window, worker: {worker_id}, region: {}", region.region_id ); Vec::new() } + (false, _) => Vec::new(), }; let plan = plan_series_indexes( buckets, @@ -207,33 +209,41 @@ async fn build_index_version( ); stats.computed_buckets = plan.computed_buckets; stats.skipped_buckets = plan.skipped_buckets; + let mut next = prune_index_version(¤t, &visible, &plan.expired_index_ids, &mut stats); + if !allow_builds { + if let Some(next) = &mut next { + next.index_buckets = plan.index_buckets; + } + return Ok((next, stats)); + } if plan.builds.is_empty() - && plan.expired_index_ids.is_empty() - && stats.removed_range == 0 + && next.is_none() && (!enable_range_index || current.range_indexes.len() == visible.len()) { return Ok((None, stats)); } - let mut range_indexes = current.range_indexes.clone(); - range_indexes.retain(|file_id| visible.contains(file_id)); - let mut series_indexes = current.series_indexes.clone(); - for id in plan - .expired_index_ids - .iter() - .chain(&plan.superseded_index_ids) - { + let SeriesIndexVersion { + mut range_indexes, + mut series_indexes, + .. + } = next.unwrap_or_else(|| SeriesIndexVersion { + range_indexes: current.range_indexes.clone(), + series_indexes: current.series_indexes.clone(), + index_buckets: Default::default(), + }); + let index_buckets = plan.index_buckets; + for id in &plan.superseded_index_ids { series_indexes.remove(id); } for (bucket, expected) in plan.builds { // Complete companion indexes independently so a failed series build preserves them. for file in &bucket.files { if enable_range_index - && !range_indexes.contains(&file.file_id().file_id()) - && let Some(file_id) = - build_range_index(store, region, version, file.clone()).await? + && !range_indexes.contains_key(&file.file_id().file_id()) + && let Some(entry) = build_range_index(store, region, version, file.clone()).await? { stats.built_range += 1; - range_indexes.insert(file_id); + range_indexes.insert(entry.file_id, entry); } } let series_handle = @@ -246,12 +256,12 @@ async fn build_index_version( if enable_range_index { for file in files { let file_id = file.file_id().file_id(); - if range_indexes.contains(&file_id) { + if range_indexes.contains_key(&file_id) { continue; } - if let Some(file_id) = build_range_index(store, region, version, file).await? { + if let Some(entry) = build_range_index(store, region, version, file).await? { stats.built_range += 1; - range_indexes.insert(file_id); + range_indexes.insert(entry.file_id, entry); } } } @@ -268,11 +278,43 @@ async fn build_index_version( let next = SeriesIndexVersion { range_indexes, series_indexes, - index_buckets: plan.index_buckets, + index_buckets, }; Ok((Some(next), stats)) } +/// Prunes obsolete metadata without building indexes or retiring published handles. +/// Returns `None` without cloning the current version when no entries need removal. +/// Physical range deletion belongs to the SST purger; series handles are retired only +/// after the cleaned catalogs and snapshot have been published successfully. +fn prune_index_version( + current: &SeriesIndexVersion, + visible: &HashSet, + expired_index_ids: &[FileId], + stats: &mut ReconcileStats, +) -> Option { + if current.range_indexes.keys().all(|id| visible.contains(id)) + && expired_index_ids + .iter() + .all(|id| !current.series_indexes.contains_key(id)) + { + return None; + } + let mut next = SeriesIndexVersion { + range_indexes: current.range_indexes.clone(), + series_indexes: current.series_indexes.clone(), + index_buckets: current.index_buckets.clone(), + }; + next.range_indexes + .retain(|file_id, _| visible.contains(file_id)); + for id in expired_index_ids { + next.series_indexes.remove(id); + } + stats.removed_range = current.range_indexes.len() - next.range_indexes.len(); + stats.removed_series = current.series_indexes.len() - next.series_indexes.len(); + Some(next) +} + /// Writes changed catalogs in a stable order; the two writes are not atomic together. async fn persist_index_catalogs( store: &ObjectStore, @@ -281,8 +323,9 @@ async fn persist_index_catalogs( stats: &ReconcileStats, ) -> Result<()> { if stats.built_range + stats.removed_range > 0 { - let mut range_entries = next.range_indexes.iter().copied().collect::>(); - range_entries.sort_unstable_by(|left, right| left.as_bytes().cmp(right.as_bytes())); + let mut range_entries = next.range_indexes.values().copied().collect::>(); + range_entries + .sort_unstable_by(|left, right| left.file_id.as_bytes().cmp(right.file_id.as_bytes())); store_catalog( store, &range_catalog_path(region_id), @@ -328,3 +371,77 @@ fn publish_index_version(region: &MitoRegionRef, next: Arc) } } } + +#[cfg(test)] +mod tests { + use std::collections::HashMap; + + use common_time::Timestamp; + use object_store::services::Memory; + + use super::*; + use crate::series_index::catalog::{RangeIndexEntry, SeriesIndexEntry}; + use crate::series_index::purger::series_index_channel; + + #[rstest::rstest] + fn test_prune_index_version( + #[values(false, true)] remove_range: bool, + #[values(false, true)] remove_series: bool, + ) { + let range_id = FileId::random(); + let series_id = FileId::random(); + let store = ObjectStore::new(Memory::default()).unwrap(); + let (purger, mut receiver) = series_index_channel(store); + let current = SeriesIndexVersion::new( + HashMap::from([( + range_id, + RangeIndexEntry { + file_id: range_id, + file_size: 10, + }, + )]), + HashMap::from([( + series_id, + SeriesIndexFileHandle::new( + RegionId::new(1, 1), + SeriesIndexEntry { + index_uuid: series_id, + file_size: 20, + bucket_start: Timestamp::new_second(0), + bucket_end: Timestamp::new_second(100), + source_file_ids: vec![range_id], + min_file_sequence: 1, + max_file_sequence: 1, + compaction_window_secs: 100, + window_sequences: Default::default(), + }, + purger, + ), + )]), + ); + let visible = if remove_range { + HashSet::new() + } else { + HashSet::from([range_id]) + }; + // Unknown and repeated IDs must not inflate removal counts. + let mut expired = vec![FileId::random()]; + if remove_series { + expired.extend([series_id, series_id]); + } + let mut stats = ReconcileStats::default(); + let next = prune_index_version(¤t, &visible, &expired, &mut stats); + assert_eq!(remove_range || remove_series, next.is_some()); + assert_eq!(usize::from(remove_range), stats.removed_range); + assert_eq!(usize::from(remove_series), stats.removed_series); + if let Some(next) = next { + assert_eq!(!remove_range, next.range_indexes.contains_key(&range_id)); + assert_eq!(!remove_series, next.series_indexes.contains_key(&series_id)); + } + assert_eq!(1, current.range_indexes.len()); + assert_eq!(1, current.series_indexes.len()); + drop(current); + // Pruning alone must never retire published files. + assert!(receiver.try_recv().is_err()); + } +} diff --git a/src/mito2/src/series_index/searcher.rs b/src/mito2/src/series_index/searcher.rs index 4613edae913..0b9befe4b8a 100644 --- a/src/mito2/src/series_index/searcher.rs +++ b/src/mito2/src/series_index/searcher.rs @@ -312,6 +312,7 @@ mod tests { ) -> (SeriesIndexFileHandle, UnboundedReceiver) { let (purger, receiver) = series_index_channel(store.clone()); let entry = SeriesIndexEntry { + file_size: 0, index_uuid: FileId::random(), bucket_start: common_time::Timestamp::new_second(0), bucket_end: common_time::Timestamp::new_second(60), diff --git a/src/mito2/src/series_index/task.rs b/src/mito2/src/series_index/task.rs index 2afbb245014..5579cf87ce4 100644 --- a/src/mito2/src/series_index/task.rs +++ b/src/mito2/src/series_index/task.rs @@ -15,7 +15,7 @@ //! Worker-owned background maintenance for series indexes. use std::sync::Arc; -use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::time::Duration; use common_telemetry::{info, warn}; @@ -25,7 +25,9 @@ use tokio::sync::mpsc::UnboundedReceiver; use tokio::task::JoinHandle; use tokio::time::{Instant, MissedTickBehavior}; -use crate::metrics::SERIES_INDEX_RECONCILE_TOTAL; +use crate::metrics::{ + SERIES_INDEX_CAPACITY_DEFERRED, SERIES_INDEX_DISK_BYTES, SERIES_INDEX_RECONCILE_TOTAL, +}; use crate::region::{RegionLeaderState, RegionMapRef, RegionRoleState}; use crate::series_index::maintenance::reconcile_series_indexes; use crate::series_index::purger::{IndexFilePurger, PurgeRequest, run_index_purge_task}; @@ -75,6 +77,8 @@ pub(crate) fn spawn_series_index_tasks( bucket_width: Duration, purger: IndexFilePurger, purge_receiver: UnboundedReceiver, + disk_usage: Arc, + max_size: u64, interval: Duration, time_provider: TimeProviderRef, enable_range_index: bool, @@ -92,6 +96,9 @@ pub(crate) fn spawn_series_index_tasks( regions, bucket_width, purger, + disk_usage, + max_size, + reported_usage: 0, state, interval, time_provider, @@ -108,6 +115,9 @@ struct SeriesIndexTask { regions: RegionMapRef, bucket_width: Duration, purger: IndexFilePurger, + disk_usage: Arc, + max_size: u64, + reported_usage: u64, worker_id: u32, state: Arc, interval: Duration, @@ -136,8 +146,33 @@ impl SeriesIndexTask { info!("Stop series-index background task, worker: {worker_id}"); } + /// Each worker contributes only its open regions, refreshed at maintenance boundaries. + fn refresh_usage(&mut self) { + let usage = self + .regions + .list_regions() + .iter() + .map(|region| region.series_index_version().disk_usage()) + .sum(); + self.report_usage(usage); + } + + fn report_usage(&mut self, usage: u64) { + if usage >= self.reported_usage { + let delta = usage - self.reported_usage; + self.disk_usage.fetch_add(delta, Ordering::Relaxed); + SERIES_INDEX_DISK_BYTES.add(delta as i64); + } else { + let delta = self.reported_usage - usage; + self.disk_usage.fetch_sub(delta, Ordering::Relaxed); + SERIES_INDEX_DISK_BYTES.sub(delta as i64); + } + self.reported_usage = usage; + } + /// Runs periodic maintenance independently of incoming deletion requests. async fn maintain(&mut self) { + self.refresh_usage(); for region in self.regions.list_regions() { if !self.state.is_running() { break; @@ -150,6 +185,11 @@ impl SeriesIndexTask { ) { continue; } + // Full capacity defers builds, but cleanup must still reclaim published usage. + let allow_builds = self.disk_usage.load(Ordering::Relaxed) < self.max_size; + if !allow_builds { + SERIES_INDEX_CAPACITY_DEFERRED.inc(); + } if let Err(error) = reconcile_series_indexes( self.worker_id, self.store.clone(), @@ -158,6 +198,7 @@ impl SeriesIndexTask { self.time_provider.current_time_millis(), self.purger.clone(), self.enable_range_index, + allow_builds, ) .await { @@ -166,10 +207,17 @@ impl SeriesIndexTask { .inc(); warn!(error; "Failed to reconcile series indexes, worker: {}, region: {}", self.worker_id, region.region_id); } + self.refresh_usage(); } } } +impl Drop for SeriesIndexTask { + fn drop(&mut self) { + self.report_usage(0); + } +} + #[cfg(test)] mod tests { use object_store::services::Memory; @@ -211,6 +259,9 @@ mod tests { regions, bucket_width: Duration::from_secs(100), purger, + disk_usage: Arc::default(), + max_size: u64::MAX, + reported_usage: 0, worker_id: 0, state: Arc::new(SeriesIndexTaskState::new()), interval: Duration::from_secs(3600), @@ -230,4 +281,236 @@ mod tests { } engine.stop().await.unwrap(); } + + #[rstest::rstest] + #[case::at_capacity(0, true)] + #[case::above_capacity(1, true)] + #[case::at_capacity_without_range(0, false)] + #[case::above_capacity_without_range(1, false)] + #[tokio::test] + async fn test_full_capacity_cleanup_and_recovery( + #[case] excess: u64, + #[case] enable_range_index: bool, + ) { + use std::sync::Mutex; + + use object_store::layers::mock::MockLayerBuilder; + + use crate::series_index::catalog::load_version_control; + use crate::series_index::tests::prepare_region_with_timestamps; + use crate::time_provider::mock::MockTimeProvider; + + let mut env = TestEnv::with_prefix("series-capacity-cleanup").await; + let (engine, region) = + prepare_region_with_timestamps(&mut env, &[1000, 2000, 3000, 4000, 5000]).await; + assert_eq!( + 5, + region + .version() + .ssts + .levels() + .iter() + .flat_map(|level| level.files()) + .count() + ); + let mut options = region.version().options.clone(); + options.ttl = Some(common_time::TimeToLive::Duration(Duration::from_secs(100))); + region.version_control.alter_options(options); + let writes = Arc::new(Mutex::new(Vec::new())); + let captured = writes.clone(); + let store = ObjectStore::new(Memory::default()).unwrap().layer( + MockLayerBuilder::default() + .writer_factory(Arc::new(move |path, _, inner| { + captured.lock().unwrap().push(path.to_string()); + inner + })) + .build() + .unwrap(), + ); + let (purger, mut receiver) = series_index_channel(store.clone()); + let regions = Arc::new(RegionMap::default()); + regions.insert_region(region.clone()); + let clock = Arc::new(MockTimeProvider::new(0)); + let usage = Arc::new(AtomicU64::new(0)); + let mut task = SeriesIndexTask { + store: store.clone(), + regions, + bucket_width: Duration::from_secs(100), + purger, + disk_usage: usage.clone(), + max_size: u64::MAX, + reported_usage: 0, + worker_id: 0, + state: Arc::new(SeriesIndexTaskState::new()), + interval: Duration::from_secs(3600), + time_provider: clock.clone(), + enable_range_index, + }; + task.maintain().await; + let previous = region.series_index_version(); + assert_eq!(1, previous.series_indexes.len()); + let old_id = *previous.series_indexes.keys().next().unwrap(); + task.max_size = previous.disk_usage() - excess; + writes.lock().unwrap().clear(); + task.maintain().await; + assert!(Arc::ptr_eq(&previous, ®ion.series_index_version())); + assert!(writes.lock().unwrap().is_empty()); + + // Removing the newest SST makes range coverage obsolete and would normally + // trigger a series replacement over the four remaining SSTs. + let sources = region.version(); + let newest = sources + .ssts + .levels() + .iter() + .flat_map(|level| level.files()) + .max_by_key(|file| file.meta_ref().sequence) + .unwrap() + .meta_ref() + .clone(); + region.version_control.apply_edit( + Some(crate::manifest::action::RegionEdit { + files_to_remove: vec![newest.clone()], + files_to_add: Vec::new(), + timestamp_ms: None, + compaction_time_window: None, + flushed_entry_id: None, + flushed_sequence: None, + committed_sequence: None, + }), + &[], + crate::test_util::new_noop_file_purger(), + ); + task.maintain().await; + let cleaned = region.series_index_version(); + assert!(!cleaned.range_indexes.contains_key(&newest.file_id)); + assert_eq!( + usize::from(enable_range_index) * 4, + cleaned.range_indexes.len() + ); + assert_eq!(previous.index_buckets, cleaned.index_buckets); + assert_eq!(1, cleaned.series_indexes.len()); + assert!(cleaned.series_indexes.contains_key(&old_id)); + assert_eq!(cleaned.disk_usage(), usage.load(Ordering::Relaxed)); + assert_eq!( + previous.disk_usage() + - previous + .range_indexes + .get(&newest.file_id) + .map_or(0, |e| e.file_size), + cleaned.disk_usage() + ); + assert!( + writes + .lock() + .unwrap() + .iter() + .all(|path| path == &range_catalog_path(region.region_id)) + ); + let restored = load_version_control(&store, region.region_id, &task.purger).await; + assert_eq!(cleaned.range_indexes, restored.current().range_indexes); + assert_eq!(cleaned.index_buckets, restored.current().index_buckets); + drop(restored); + + // Expiration must work even when the shared estimate remains full. + task.max_size = cleaned.disk_usage() - excess; + clock.set_now(201_000); + writes.lock().unwrap().clear(); + task.maintain().await; + let expired = region.series_index_version(); + assert!(expired.series_indexes.is_empty()); + assert!(expired.index_buckets.is_empty()); + assert_eq!(cleaned.range_indexes, expired.range_indexes); + assert_eq!(expired.disk_usage(), usage.load(Ordering::Relaxed)); + assert!(usage.load(Ordering::Relaxed) < task.max_size); + assert_eq!( + *writes.lock().unwrap(), + vec![series_catalog_path(region.region_id)] + ); + let restored = load_version_control(&store, region.region_id, &task.purger).await; + assert!(restored.current().series_indexes.is_empty()); + assert_eq!(expired.range_indexes, restored.current().range_indexes); + assert!(receiver.try_recv().is_err()); + drop(previous); + assert!(receiver.try_recv().is_err()); + drop(cleaned); + assert_eq!(old_id, receiver.try_recv().unwrap().file_id.file_id()); + assert!(receiver.try_recv().is_err()); + + // Make the remaining SSTs eligible again; reclaimed capacity admits a build + // without closing the region or increasing the configured limit. + clock.set_now(0); + task.maintain().await; + let rebuilt = region.series_index_version(); + assert_eq!(1, rebuilt.series_indexes.len()); + assert!(!rebuilt.series_indexes.contains_key(&old_id)); + assert_eq!( + 4, + rebuilt + .series_indexes + .values() + .next() + .unwrap() + .entry() + .source_file_ids + .len() + ); + assert_eq!(rebuilt.disk_usage(), usage.load(Ordering::Relaxed)); + engine.stop().await.unwrap(); + } + + #[tokio::test] + async fn test_shared_usage_defers_builds_and_allows_overshoot() { + let mut env = TestEnv::with_prefix("series-approximate-usage").await; + let (engine, region) = prepare_region(&mut env).await; + let usage = Arc::new(AtomicU64::new(0)); + let store = ObjectStore::new(Memory::default()).unwrap(); + let make_task = || SeriesIndexTask { + store: store.clone(), + regions: Arc::new(RegionMap::default()), + bucket_width: Duration::from_secs(100), + purger: series_index_channel(store.clone()).0, + disk_usage: usage.clone(), + max_size: 1024, + reported_usage: 0, + worker_id: 0, + state: Arc::new(SeriesIndexTaskState::new()), + interval: Duration::from_secs(3600), + time_provider: Arc::new(crate::time_provider::StdTimeProvider), + enable_range_index: true, + }; + let mut task = make_task(); + task.regions.insert_region(region.clone()); + task.maintain().await; + let bytes = region.series_index_version().disk_usage(); + assert!( + bytes > task.max_size, + "a started reconciliation may exceed the limit" + ); + assert_eq!(bytes, usage.load(Ordering::Relaxed)); + // Another worker sees the same estimate and defers builds without creating coverage. + let mut other_env = TestEnv::with_prefix("series-approximate-other").await; + let (other_engine, other_region) = prepare_region(&mut other_env).await; + let mut other = make_task(); + other.regions.insert_region(other_region.clone()); + other.max_size = bytes; // Exact equality also skips. + other.maintain().await; + assert_eq!(0, other_region.series_index_version().disk_usage()); + assert_eq!(bytes, usage.load(Ordering::Relaxed)); + + // Closing a region reduces the estimate on the next pass, even while full. + task.regions.remove_region(region.region_id); + task.maintain().await; + assert_eq!(0, usage.load(Ordering::Relaxed)); + other.maintain().await; + assert!(other_region.series_index_version().disk_usage() > 0); + assert_eq!( + other_region.series_index_version().disk_usage(), + usage.load(Ordering::Relaxed) + ); + drop(other); + assert_eq!(0, usage.load(Ordering::Relaxed)); + engine.stop().await.unwrap(); + other_engine.stop().await.unwrap(); + } } diff --git a/src/mito2/src/series_index/tests.rs b/src/mito2/src/series_index/tests.rs index c1efe3afdcf..8ee898277bc 100644 --- a/src/mito2/src/series_index/tests.rs +++ b/src/mito2/src/series_index/tests.rs @@ -48,11 +48,16 @@ pub(super) async fn prepare_region(env: &mut TestEnv) -> (MitoEngine, MitoRegion prepare_region_with_timestamps(env, &[1000, 2000, 3000, 4000]).await } -async fn prepare_region_with_timestamps( +pub(crate) async fn prepare_region_with_timestamps( env: &mut TestEnv, timestamps: &[i64], ) -> (MitoEngine, MitoRegionRef) { - let engine = env.create_engine(MitoConfig::default()).await; + let engine = env + .create_engine(MitoConfig { + min_compaction_interval: Duration::from_secs(3600), + ..Default::default() + }) + .await; let metadata = Arc::new(sst_region_metadata_with_encoding( PrimaryKeyEncoding::Sparse, )); @@ -93,6 +98,10 @@ async fn prepare_region_with_timestamps( .handle_request(region_id, RegionRequest::Create(request)) .await .unwrap(); + // Flush completion can precede compaction scheduling. Keep background TTL + // cleanup from racing tests that alter options after preparing the SSTs. + let region = engine.get_region(region_id).unwrap(); + region.update_schedule_compaction_millis(); for &ts in timestamps { engine .handle_request( @@ -122,7 +131,6 @@ async fn prepare_region_with_timestamps( .unwrap(); flush_region(&engine, region_id, None).await; } - let region = engine.get_region(region_id).unwrap(); (engine, region) } @@ -171,6 +179,7 @@ impl IndexTest { 0, self.purger.clone(), enable_range_index, + true, ) .await } @@ -311,7 +320,7 @@ async fn test_reconcile_restores_and_reuses_indexes() { let missing_paths = [ range_index_path( region.region_id, - *first.range_indexes.iter().next().unwrap(), + *first.range_indexes.keys().next().unwrap(), ), series_index_path(region.region_id, first_id), ]; @@ -325,6 +334,7 @@ async fn test_reconcile_restores_and_reuses_indexes() { .current(); assert_eq!(first.range_indexes, restored.range_indexes); assert_eq!(first.index_buckets, restored.index_buckets); + assert_eq!(first.disk_usage(), restored.disk_usage()); assert_eq!( first.series_indexes[&first_id].entry(), restored.series_indexes[&first_id].entry() @@ -697,13 +707,13 @@ async fn test_failed_reconcile_retires_only_unpublished_series() { ); assert_eq!(if failure == "range" { 2 } else { 3 }, plan.builds.len()); let (bucket, entry) = &plan.builds[0]; - let mut ranges = HashSet::new(); + let mut ranges = HashMap::new(); for file in &bucket.files { - let id = build_range_index(store, ®ion, &version, file.clone()) + let entry = build_range_index(store, ®ion, &version, file.clone()) .await .unwrap() .unwrap(); - ranges.insert(id); + ranges.insert(entry.file_id, entry); } let handle = build_series_index(store, ®ion, &version, bucket, entry, purger) .await @@ -838,6 +848,8 @@ async fn test_maintenance_wakeup_and_timer(#[case] enable_range_index: bool) { Duration::from_secs(100), purger, receiver, + Arc::default(), + u64::MAX, Duration::from_secs(3600), clock.clone(), enable_range_index, @@ -866,6 +878,7 @@ async fn test_drop_catalogs_and_retained_snapshot_after_task_stop() { let store = ObjectStore::new(Memory::default()).unwrap(); let region_id = RegionId::new(1, 1); let entry = SeriesIndexEntry { + file_size: 0, index_uuid: FileId::random(), bucket_start: Timestamp::new_second(0), bucket_end: Timestamp::new_second(100), @@ -914,6 +927,8 @@ async fn test_drop_catalogs_and_retained_snapshot_after_task_stop() { Duration::from_secs(100), purger, receiver, + Arc::default(), + u64::MAX, Duration::from_secs(3600), Arc::new(StdTimeProvider), true, diff --git a/src/mito2/src/series_index/version.rs b/src/mito2/src/series_index/version.rs index c31dfde906b..fa2115fc1f2 100644 --- a/src/mito2/src/series_index/version.rs +++ b/src/mito2/src/series_index/version.rs @@ -14,7 +14,7 @@ //! Immutable index snapshots and aggregate series-file handles. -use std::collections::{BTreeMap, HashMap, HashSet}; +use std::collections::{BTreeMap, HashMap}; use std::fmt::{self, Debug, Formatter}; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, RwLock}; @@ -23,7 +23,7 @@ use common_time::Timestamp; use store_api::storage::{FileId, RegionId}; use crate::series_index::bucket::IndexBucket; -use crate::series_index::catalog::SeriesIndexEntry; +use crate::series_index::catalog::{RangeIndexEntry, SeriesIndexEntry}; use crate::series_index::purger::{IndexFilePurger, PurgeRequest}; use crate::sst::file::RegionFileId; @@ -94,7 +94,7 @@ impl Drop for SeriesIndexFileHandleInner { pub(crate) struct SeriesIndexVersion { /// Range indexes for visible SSTs; reconciliation removes IDs absent from its SST snapshot. /// Physical deletion is independently handled by the SST file purger. - pub(crate) range_indexes: HashSet, + pub(crate) range_indexes: HashMap, pub(crate) series_indexes: HashMap, pub(crate) index_buckets: BTreeMap, } @@ -102,7 +102,7 @@ pub(crate) struct SeriesIndexVersion { impl SeriesIndexVersion { /// Restores bucket lookup from immutable index coverage stored in the catalog. pub(crate) fn new( - range_indexes: HashSet, + range_indexes: HashMap, series_indexes: HashMap, ) -> Self { let mut index_buckets = BTreeMap::new(); @@ -116,6 +116,19 @@ impl SeriesIndexVersion { } } + /// Approximate installed usage; old snapshots and unpublished outputs are excluded. + pub(crate) fn disk_usage(&self) -> u64 { + self.range_indexes + .values() + .map(|entry| entry.file_size) + .sum::() + + self + .series_indexes + .values() + .map(|handle| handle.entry().file_size) + .sum::() + } + fn mark_all_deleted(&self) { self.series_indexes .values() diff --git a/src/mito2/src/sst/parquet/reader.rs b/src/mito2/src/sst/parquet/reader.rs index 9584ff0a6af..714e62ae4af 100644 --- a/src/mito2/src/sst/parquet/reader.rs +++ b/src/mito2/src/sst/parquet/reader.rs @@ -713,7 +713,7 @@ impl ParquetReaderBuilder { && context .version .range_indexes - .contains(&self.file_handle.file_id().file_id())) + .contains_key(&self.file_handle.file_id().file_id())) .then(|| context.store.clone()) }); let context = FileRangeContext::new( diff --git a/src/mito2/src/worker.rs b/src/mito2/src/worker.rs index 24eec5dc2b3..0f67d8d2a1b 100644 --- a/src/mito2/src/worker.rs +++ b/src/mito2/src/worker.rs @@ -35,7 +35,7 @@ mod handle_write; use std::collections::HashMap; use std::path::Path; use std::sync::Arc; -use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::time::Duration; use common_base::Plugins; @@ -198,6 +198,7 @@ impl WorkerGroup { let index_build_job_pool = Arc::new(LocalScheduler::new(config.max_background_index_builds)); let series_index_store = series_index_store_from_config(&config, data_home).await?; + let series_index_disk_usage = Arc::new(AtomicU64::new(0)); 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)); @@ -251,6 +252,7 @@ impl WorkerGroup { write_buffer_manager: write_buffer_manager.clone(), index_build_job_pool: index_build_job_pool.clone(), series_index_store: series_index_store.clone(), + series_index_disk_usage: series_index_disk_usage.clone(), flush_job_pool: flush_job_pool.clone(), compact_job_pool: compact_job_pool.clone(), purge_scheduler: purge_scheduler.clone(), @@ -410,6 +412,7 @@ impl WorkerGroup { let index_build_job_pool = Arc::new(LocalScheduler::new(config.max_background_index_builds)); let series_index_store = series_index_store_from_config(&config, data_home).await?; + let series_index_disk_usage = Arc::new(AtomicU64::new(0)); 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)); @@ -464,6 +467,7 @@ impl WorkerGroup { write_buffer_manager: write_buffer_manager.clone(), index_build_job_pool: index_build_job_pool.clone(), series_index_store: series_index_store.clone(), + series_index_disk_usage: series_index_disk_usage.clone(), flush_job_pool: flush_job_pool.clone(), compact_job_pool: compact_job_pool.clone(), purge_scheduler: purge_scheduler.clone(), @@ -572,6 +576,7 @@ struct WorkerStarter { compact_job_pool: SchedulerRef, index_build_job_pool: SchedulerRef, series_index_store: Option, + series_index_disk_usage: Arc, flush_job_pool: SchedulerRef, purge_scheduler: SchedulerRef, listener: WorkerListener, @@ -620,6 +625,8 @@ impl WorkerStarter { self.config.experimental_series_index_bucket_width, purger, purge_receiver, + self.series_index_disk_usage.clone(), + self.config.experimental_series_index_max_size.as_bytes(), self.config.experimental_series_index_maintenance_interval, self.time_provider.clone(), self.config.experimental_enable_range_index, @@ -951,7 +958,7 @@ struct RegionWorkerLoop { index_build_scheduler: IndexBuildScheduler, /// Controls the worker-owned series-index task. series_index_task_state: Option>, - /// Store for companion range indexes deleted by the region SST purger. + /// Local store for series and range indexes managed by reconciliation. series_index_store: Option, series_index_purger: Option, /// Schedules background flush requests. diff --git a/tests-integration/tests/http.rs b/tests-integration/tests/http.rs index 3e64a6d74e3..a4841afad93 100644 --- a/tests-integration/tests/http.rs +++ b/tests-integration/tests/http.rs @@ -2561,6 +2561,7 @@ experimental_manifest_keep_removed_file_count = 256 experimental_manifest_keep_removed_file_ttl = "1h" compress_manifest = false experimental_enable_series_index = false +experimental_series_index_max_size = "5GiB" experimental_enable_range_index = false experimental_series_index_maintenance_interval = "5m" experimental_series_index_bucket_width = "5days"