mirror of
https://github.com/GreptimeTeam/greptimedb.git
synced 2026-10-03 10:35:35 +00:00
feat(mito): limit approximate series index disk usage (#9313)
* feat(mito): limit series index disk usage Signed-off-by: evenyag <realevenyag@gmail.com> * refactor(mito): enforce series index quota during reconciliation Signed-off-by: evenyag <realevenyag@gmail.com> * refactor(mito): remove series index disk budget layer Signed-off-by: evenyag <realevenyag@gmail.com> * refactor(mito): simplify series index limit to estimated usage Signed-off-by: evenyag <realevenyag@gmail.com> * refactor(mito): trust index catalogs when loading snapshots Signed-off-by: evenyag <realevenyag@gmail.com> * refactor(mito): minimize series index disk limit changes Signed-off-by: evenyag <realevenyag@gmail.com> * fix(mito2): keep series index cleanup running at capacity Signed-off-by: evenyag <realevenyag@gmail.com> * fix(mito2): avoid no-op index clones and stabilize capacity tests Signed-off-by: evenyag <realevenyag@gmail.com> * fix: update config API expectation and stabilize index build tests Signed-off-by: evenyag <realevenyag@gmail.com> --------- Signed-off-by: evenyag <realevenyag@gmail.com>
This commit is contained in:
@@ -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)),
|
||||
|
||||
+5
-1
@@ -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`.
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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!(
|
||||
|
||||
@@ -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),
|
||||
|
||||
@@ -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,
|
||||
};
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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<Option<FileId>> {
|
||||
) -> Result<Option<RangeIndexEntry>> {
|
||||
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<Option<SeriesIndexFileHandle>> = 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
|
||||
|
||||
@@ -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<SeriesIndexEntry>,
|
||||
@@ -104,7 +113,7 @@ pub(crate) struct SeriesIndexCatalog {
|
||||
|
||||
#[derive(Debug, Default, Serialize, Deserialize)]
|
||||
pub(crate) struct RangeIndexCatalog {
|
||||
pub(crate) indexes: Vec<FileId>,
|
||||
pub(crate) indexes: Vec<RangeIndexEntry>,
|
||||
}
|
||||
|
||||
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::<SeriesIndexCatalog>(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),
|
||||
|
||||
@@ -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<ReconcileStats> {
|
||||
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<SeriesIndexVersion>, 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::<HashSet<_>>();
|
||||
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<FileId>,
|
||||
expired_index_ids: &[FileId],
|
||||
stats: &mut ReconcileStats,
|
||||
) -> Option<SeriesIndexVersion> {
|
||||
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::<Vec<_>>();
|
||||
range_entries.sort_unstable_by(|left, right| left.as_bytes().cmp(right.as_bytes()));
|
||||
let mut range_entries = next.range_indexes.values().copied().collect::<Vec<_>>();
|
||||
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<SeriesIndexVersion>)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[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());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -312,6 +312,7 @@ mod tests {
|
||||
) -> (SeriesIndexFileHandle, UnboundedReceiver<PurgeRequest>) {
|
||||
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),
|
||||
|
||||
@@ -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<PurgeRequest>,
|
||||
disk_usage: Arc<AtomicU64>,
|
||||
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<AtomicU64>,
|
||||
max_size: u64,
|
||||
reported_usage: u64,
|
||||
worker_id: u32,
|
||||
state: Arc<SeriesIndexTaskState>,
|
||||
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();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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<FileId>,
|
||||
pub(crate) range_indexes: HashMap<FileId, RangeIndexEntry>,
|
||||
pub(crate) series_indexes: HashMap<FileId, SeriesIndexFileHandle>,
|
||||
pub(crate) index_buckets: BTreeMap<Timestamp, IndexBucket>,
|
||||
}
|
||||
@@ -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<FileId>,
|
||||
range_indexes: HashMap<FileId, RangeIndexEntry>,
|
||||
series_indexes: HashMap<FileId, SeriesIndexFileHandle>,
|
||||
) -> 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::<u64>()
|
||||
+ self
|
||||
.series_indexes
|
||||
.values()
|
||||
.map(|handle| handle.entry().file_size)
|
||||
.sum::<u64>()
|
||||
}
|
||||
|
||||
fn mark_all_deleted(&self) {
|
||||
self.series_indexes
|
||||
.values()
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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<S> {
|
||||
compact_job_pool: SchedulerRef,
|
||||
index_build_job_pool: SchedulerRef,
|
||||
series_index_store: Option<ObjectStore>,
|
||||
series_index_disk_usage: Arc<AtomicU64>,
|
||||
flush_job_pool: SchedulerRef,
|
||||
purge_scheduler: SchedulerRef,
|
||||
listener: WorkerListener,
|
||||
@@ -620,6 +625,8 @@ impl<S: LogStore> WorkerStarter<S> {
|
||||
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<S> {
|
||||
index_build_scheduler: IndexBuildScheduler,
|
||||
/// Controls the worker-owned series-index task.
|
||||
series_index_task_state: Option<Arc<SeriesIndexTaskState>>,
|
||||
/// 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<ObjectStore>,
|
||||
series_index_purger: Option<IndexFilePurger>,
|
||||
/// Schedules background flush requests.
|
||||
|
||||
Reference in New Issue
Block a user