feat(mito2): reconcile series indexes in background

Signed-off-by: evenyag <realevenyag@gmail.com>
This commit is contained in:
evenyag
2026-09-11 23:53:32 +08:00
parent 555485c40e
commit 1d5b7bc033
22 changed files with 1302 additions and 86 deletions
+6
View File
@@ -165,6 +165,9 @@
| `region_engine.mito.worker_request_batch_size` | Integer | `64` | Max batch size for a worker to handle requests. |
| `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_series_index_root` | String | `""` | Under development; do not enable. Root directory for local series indexes.<br/>Empty disables the feature. Relative paths resolve under `data_home`. |
| `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<br/>an exact multiple of each region's compaction time window. |
| `region_engine.mito.max_background_flushes` | Integer | Auto | Max number of running background flush jobs (default: 1/2 of cpu cores). |
| `region_engine.mito.max_background_compactions` | Integer | Auto | Max number of running background compaction jobs (default: 1/4 of cpu cores). |
| `region_engine.mito.max_background_purges` | Integer | Auto | Max number of running background purge jobs (default: number of cpu cores). |
@@ -612,6 +615,9 @@
| `region_engine.mito.experimental_manifest_keep_removed_file_count` | Integer | `256` | Number of removed files to keep in manifest's `removed_files` field before also<br/>remove them from `removed_files`. Mostly for debugging purpose.<br/>If set to 0, it will only use `keep_removed_file_ttl` to decide when to remove files<br/>from `removed_files` field. |
| `region_engine.mito.experimental_manifest_keep_removed_file_ttl` | String | `1h` | How long to keep removed files in the `removed_files` field of manifest<br/>after they are removed from manifest.<br/>files will only be removed from `removed_files` field<br/>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_series_index_root` | String | `""` | Under development; do not enable. Root directory for local series indexes.<br/>Empty disables the feature. Relative paths resolve under `data_home`. |
| `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<br/>an exact multiple of each region's compaction time window. |
| `region_engine.mito.max_background_flushes` | Integer | Auto | Max number of running background flush jobs (default: 1/2 of cpu cores). |
| `region_engine.mito.max_background_compactions` | Integer | Auto | Max number of running background compaction jobs (default: 1/4 of cpu cores). |
| `region_engine.mito.max_background_purges` | Integer | Auto | Max number of running background purge jobs (default: number of cpu cores). |
+11
View File
@@ -474,6 +474,17 @@ experimental_manifest_keep_removed_file_ttl = "1h"
## Whether to compress manifest and checkpoint file by gzip (default false).
compress_manifest = false
## Under development; do not enable. Root directory for local series indexes.
## Empty disables the feature. Relative paths resolve under `data_home`.
#+ experimental_series_index_root = ""
## Interval between series-index maintenance runs. Zero uses the default of 5 min.
#+ experimental_series_index_maintenance_interval = "5m"
## Requested minimum series-index bucket width (default: 5 days), rounded up to
## an exact multiple of each region's compaction time window.
#+ experimental_series_index_bucket_width = "5days"
## Max number of running background flush jobs (default: 1/2 of cpu cores).
## @toml2docs:none-default="Auto"
#+ max_background_flushes = 4
+11
View File
@@ -635,6 +635,17 @@ manifest_checkpoint_distance = 10
## Whether to compress manifest and checkpoint file by gzip (default false).
compress_manifest = false
## Under development; do not enable. Root directory for local series indexes.
## Empty disables the feature. Relative paths resolve under `data_home`.
#+ experimental_series_index_root = ""
## Interval between series-index maintenance runs. Zero uses the default of 5 min.
#+ experimental_series_index_maintenance_interval = "5m"
## Requested minimum series-index bucket width (default: 5 days), rounded up to
## an exact multiple of each region's compaction time window.
#+ experimental_series_index_bucket_width = "5days"
## Max number of running background flush jobs (default: 1/2 of cpu cores).
## @toml2docs:none-default="Auto"
#+ max_background_flushes = 4
+1 -1
View File
@@ -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 and SST builders, catalogs, immutable snapshots, file lifecycle, and background maintenance |
| `series_index` | `src/mito2/src/series_index/` | Series index writer/searcher, bucket planning, catalogs, worker-owned maintenance and purge tasks, 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 |
+15 -3
View File
@@ -89,14 +89,16 @@ pub struct MitoConfig {
// Background job configs:
/// Max number of running background index build jobs (default: 1/8 of cpu cores).
pub max_background_index_builds: usize,
// TODO: Document both series-index settings in the example configs and regenerate configuration
// docs before exposing the feature.
/// Under development; do not enable. Root directory for loading series indexes, currently stored
/// on the local filesystem. Empty disables the feature. Relative paths resolve under `data_home`.
pub experimental_series_index_root: String,
/// Interval between series-index maintenance runs (default 5 min). Zero uses the default.
#[serde(with = "humantime_serde")]
pub experimental_series_index_maintenance_interval: Duration,
/// Under development; do not enable. Requested minimum bucket width for series indexes.
/// It is rounded up to an exact multiple of each region's compaction time window.
#[serde(with = "humantime_serde")]
pub experimental_series_index_bucket_width: Duration,
/// Max number of running background flush jobs (default: 1/2 of cpu cores).
pub max_background_flushes: usize,
/// Max number of running background compaction jobs (default: 1/4 of cpu cores).
@@ -217,6 +219,7 @@ impl Default for MitoConfig {
experimental_series_index_root: String::new(),
experimental_series_index_maintenance_interval:
DEFAULT_SERIES_INDEX_MAINTENANCE_INTERVAL,
experimental_series_index_bucket_width: Duration::from_secs(5 * 24 * 60 * 60),
max_background_flushes: divide_num_cpus(2),
max_background_compactions: divide_num_cpus(4),
max_background_purges: get_total_cpu_cores(),
@@ -439,9 +442,14 @@ mod tests {
.experimental_series_index_root
.is_empty()
);
assert_eq!(
MitoConfig::default().experimental_series_index_bucket_width,
Duration::from_secs(5 * 24 * 60 * 60)
);
let mut config: MitoConfig = toml::from_str(
"experimental_series_index_root = 'indexes'
experimental_series_index_maintenance_interval = '30s'",
experimental_series_index_maintenance_interval = '30s'
experimental_series_index_bucket_width = '2days'",
)
.unwrap();
config.sanitize("/data").unwrap();
@@ -450,6 +458,10 @@ mod tests {
config.experimental_series_index_maintenance_interval,
Duration::from_secs(30)
);
assert_eq!(
config.experimental_series_index_bucket_width,
Duration::from_secs(2 * 24 * 60 * 60)
);
let restored: MitoConfig = toml::from_str(&toml::to_string(&config).unwrap()).unwrap();
assert_eq!(config, restored);
+1 -1
View File
@@ -38,11 +38,11 @@ use tokio::sync::{Notify, Semaphore};
use crate::config::MitoConfig;
use crate::engine::MitoEngine;
use crate::engine::flush_test::MockTimeProvider;
use crate::engine::listener::{CompactionListener, EventListener};
use crate::test_util::{
CreateRequestBuilder, TestEnv, build_rows_for_key, column_metadata_to_column_schema, put_rows,
};
use crate::time_provider::mock::MockTimeProvider;
pub(crate) async fn put_and_flush(
engine: &MitoEngine,
+1 -1
View File
@@ -35,12 +35,12 @@ use tokio::sync::{Barrier, mpsc, oneshot};
use crate::config::MitoConfig;
use crate::engine::MitoEngine;
use crate::engine::flush_test::MockTimeProvider;
use crate::engine::listener::EventListener;
use crate::manifest::action::RegionEdit;
use crate::region::{MitoRegionRef, RegionLeaderState, RegionRoleState};
use crate::sst::file::FileMeta;
use crate::test_util::{CreateRequestBuilder, TestEnv, build_rows, rows_schema};
use crate::time_provider::mock::MockTimeProvider;
#[tokio::test]
async fn test_edit_region_schedule_compaction() {
+2 -39
View File
@@ -17,7 +17,7 @@
use std::assert_matches;
use std::collections::HashMap;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicI64, AtomicUsize, Ordering};
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::time::Duration;
use api::v1::Rows;
@@ -54,7 +54,7 @@ use crate::test_util::{
prepare_test_for_kafka_log_store, put_rows, raft_engine_log_store_factory, reopen_region,
rows_schema, single_kafka_log_store_factory,
};
use crate::time_provider::TimeProvider;
use crate::time_provider::mock::MockTimeProvider;
use crate::worker::MAX_INITIAL_CHECK_DELAY_SECS;
async fn set_write_buffer_size_to_current_usage(engine: &MitoEngine, region_id: RegionId) {
@@ -1334,43 +1334,6 @@ fn kafka_wal_options(topic: &Option<String>) -> HashMap<String, String> {
.unwrap_or_default()
}
#[derive(Debug)]
pub(crate) struct MockTimeProvider {
now: AtomicI64,
elapsed: AtomicI64,
}
impl TimeProvider for MockTimeProvider {
fn current_time_millis(&self) -> i64 {
self.now.load(Ordering::Relaxed)
}
fn elapsed_since(&self, _current_millis: i64) -> i64 {
self.elapsed.load(Ordering::Relaxed)
}
fn wait_duration(&self, _duration: Duration) -> Duration {
Duration::from_millis(20)
}
}
impl MockTimeProvider {
pub(crate) fn new(now: i64) -> Self {
Self {
now: AtomicI64::new(now),
elapsed: AtomicI64::new(0),
}
}
pub(crate) fn set_now(&self, now: i64) {
self.now.store(now, Ordering::Relaxed);
}
fn set_elapsed(&self, elapsed: i64) {
self.elapsed.store(elapsed, Ordering::Relaxed);
}
}
#[tokio::test]
async fn test_auto_flush_engine() {
test_auto_flush_engine_with_format(false).await;
+13
View File
@@ -319,6 +319,19 @@ lazy_static! {
// Index metrics.
lazy_static! {
// Index metrics.
/// Outcomes of series-index reconciliation passes.
pub static ref SERIES_INDEX_RECONCILE_TOTAL: IntCounterVec = register_int_counter_vec!(
"greptime_mito_series_index_reconcile_total",
"series-index reconciliation passes",
&["result"],
).unwrap();
/// Elapsed time of series-index reconciliation phases.
pub static ref SERIES_INDEX_RECONCILE_ELAPSED: HistogramVec = register_histogram_vec!(
"greptime_mito_series_index_reconcile_elapsed",
"series-index reconciliation elapsed time",
&["phase"],
exponential_buckets(0.01, 10.0, 7).unwrap(),
).unwrap();
/// Series and range index file operations.
pub static ref SERIES_INDEX_FILE_OPERATION_TOTAL: IntCounterVec = register_int_counter_vec!(
"greptime_mito_series_index_file_operation_total",
-1
View File
@@ -242,7 +242,6 @@ impl StagingPartitionInfo {
impl MitoRegion {
/// Returns the current immutable series-index snapshot.
#[allow(dead_code)] // Used by the upcoming query integration.
pub(crate) fn series_index_version(&self) -> Arc<SeriesIndexVersion> {
self.series_index_version_control.current()
}
+9 -14
View File
@@ -12,42 +12,37 @@
// See the License for the specific language governing permissions and
// limitations under the License.
//! Series index writer and searcher.
//! Series index construction, search, and maintenance.
//!
//! Under development. Index files are currently stored on the local filesystem.
// These components are consumed by the upcoming query and maintenance integration.
#[allow(dead_code)]
mod catalog;
// Consumed by the follow-up background-maintenance integration.
#[allow(dead_code)]
mod bucket;
// Consumed by the follow-up background-maintenance integration.
#[allow(dead_code)]
mod builder;
#[allow(dead_code)]
mod catalog;
mod maintenance;
mod purger;
mod searcher;
mod task;
#[cfg(test)]
mod tests;
#[allow(dead_code)]
mod version;
mod writer;
use futures::stream::BoxStream;
pub use searcher::SeriesIndexSearcher;
use store_api::metric_engine_consts::{
DATA_SCHEMA_TABLE_ID_COLUMN_NAME as TABLE_ID_COLUMN,
DATA_SCHEMA_TSID_COLUMN_NAME as TSID_COLUMN,
};
pub use writer::{
SeriesIndexWriter, SeriesIndexWriterMetrics, SeriesIndexWriterOptions, series_index_schema,
};
use crate::error::Result;
pub(crate) use crate::series_index::catalog::{delete_catalogs, load_version_control};
pub(crate) use crate::series_index::purger::{IndexFilePurger, series_index_channel};
pub use crate::series_index::searcher::SeriesIndexSearcher;
pub(crate) use crate::series_index::task::{SeriesIndexTaskState, spawn_series_index_tasks};
pub(crate) use crate::series_index::version::{SeriesIndexVersion, SeriesIndexVersionControl};
pub use crate::series_index::writer::{
SeriesIndexWriter, SeriesIndexWriterMetrics, SeriesIndexWriterOptions, series_index_schema,
};
pub(crate) const MIN_TS_COLUMN: &str = "__series_min_ts";
pub(crate) const MAX_TS_COLUMN: &str = "__series_max_ts";
+305
View File
@@ -0,0 +1,305 @@
// Copyright 2023 Greptime Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//! Worker-owned series-index reconciliation and publication.
use std::collections::HashSet;
use std::sync::Arc;
use std::time::{Duration, Instant};
use common_telemetry::{debug, info};
use object_store::ObjectStore;
use store_api::storage::RegionId;
use crate::error::Result;
use crate::metrics::{SERIES_INDEX_RECONCILE_ELAPSED, SERIES_INDEX_RECONCILE_TOTAL};
use crate::read::series_candidate::is_sparse_metric_metadata;
use crate::region::MitoRegionRef;
use crate::region::version::VersionRef;
use crate::series_index::bucket::{
group_files_into_series_buckets, plan_series_indexes, rounded_bucket_width,
};
use crate::series_index::builder::{build_range_index, build_series_index};
use crate::series_index::catalog::{
RangeIndexCatalog, SeriesIndexCatalog, range_catalog_path, series_catalog_path, store_catalog,
};
use crate::series_index::purger::IndexFilePurger;
use crate::series_index::version::{SeriesIndexFileHandle, SeriesIndexVersion};
/// Retires newly completed series files unless their snapshot is published.
#[derive(Default)]
struct UnpublishedSeriesFiles(Vec<SeriesIndexFileHandle>);
impl UnpublishedSeriesFiles {
fn disarm(&mut self) {
self.0.clear();
}
}
impl Drop for UnpublishedSeriesFiles {
fn drop(&mut self) {
for handle in &self.0 {
handle.mark_deleted();
}
}
}
#[derive(Debug, Default)]
pub(crate) struct ReconcileStats {
pub(crate) source_files: usize,
pub(crate) built_range: usize,
pub(crate) built_series: usize,
pub(crate) removed_range: usize,
pub(crate) removed_series: usize,
pub(crate) computed_buckets: usize,
pub(crate) skipped_buckets: usize,
}
impl ReconcileStats {
fn changed(&self) -> bool {
self.built_range + self.built_series + self.removed_range + self.removed_series > 0
}
}
/// Reconciles indexes for one region snapshot, persists catalogs, then atomically publishes it.
pub(crate) async fn reconcile_series_indexes(
worker_id: u32,
store: ObjectStore,
region: MitoRegionRef,
requested_bucket_width: Duration,
now_ms: i64,
purger: IndexFilePurger,
) -> Result<ReconcileStats> {
let total_start = Instant::now();
// Use this snapshot throughout reconciliation, even if the region version advances.
let version = region.version_control.current().version;
if !is_sparse_metric_metadata(&version.metadata) {
SERIES_INDEX_RECONCILE_TOTAL
.with_label_values(&["noop"])
.inc();
return Ok(ReconcileStats::default());
}
let build_start = Instant::now();
let mut unpublished = UnpublishedSeriesFiles::default();
let (next, stats) = build_index_version(
worker_id,
&store,
&region,
&version,
requested_bucket_width,
now_ms,
&purger,
&mut unpublished,
)
.await?;
SERIES_INDEX_RECONCILE_ELAPSED
.with_label_values(&["build"])
.observe(build_start.elapsed().as_secs_f64());
// Persist both catalogs before making the new snapshot visible to readers.
if let Some(next) = next {
persist_index_catalogs(&store, region.region_id, &next).await?;
publish_index_version(&region, Arc::new(next));
}
unpublished.disarm();
let result = if stats.changed() { "changed" } else { "noop" };
SERIES_INDEX_RECONCILE_TOTAL
.with_label_values(&[result])
.inc();
SERIES_INDEX_RECONCILE_ELAPSED
.with_label_values(&["total"])
.observe(total_start.elapsed().as_secs_f64());
if stats.changed() {
info!(
"Reconciled series-index snapshot, worker: {worker_id}, region: {}, elapsed: {:?}, stats: {:?}",
region.region_id,
total_start.elapsed(),
stats
);
} else {
debug!(
"Series-index reconciliation made no changes, worker: {worker_id}, region: {}",
region.region_id
);
}
Ok(stats)
}
/// Builds a changed snapshot, retaining reusable indexes and removing obsolete coverage.
/// Returns `None` when no range or series indexes were added or removed.
#[allow(clippy::too_many_arguments)]
async fn build_index_version(
worker_id: u32,
store: &ObjectStore,
region: &MitoRegionRef,
version: &VersionRef,
requested_bucket_width: Duration,
now_ms: i64,
purger: &IndexFilePurger,
unpublished: &mut UnpublishedSeriesFiles,
) -> Result<(Option<SeriesIndexVersion>, ReconcileStats)> {
let mut stats = ReconcileStats::default();
let files = version
.ssts
.levels()
.iter()
.flat_map(|level| level.files())
.cloned()
.collect::<Vec<_>>();
stats.source_files = files.len();
let visible = files
.iter()
.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)
// 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 => {
debug!(
"Deferring series indexes without compaction window, worker: {worker_id}, region: {}",
region.region_id
);
Vec::new()
}
};
let plan = plan_series_indexes(
buckets,
current.index_buckets.clone(),
version.options.ttl,
now_ms,
);
stats.computed_buckets = plan.computed_buckets;
stats.skipped_buckets = plan.skipped_buckets;
if plan.builds.is_empty()
&& plan.expired_index_ids.is_empty()
&& stats.removed_range == 0
&& 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)
{
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 !range_indexes.contains(&file.file_id().file_id())
&& let Some(file_id) =
build_range_index(store, region, version, file.clone()).await?
{
stats.built_range += 1;
range_indexes.insert(file_id);
}
}
let series_handle =
build_series_index(store, region, version, &bucket, &expected, purger).await?;
unpublished.0.push(series_handle.clone());
stats.built_series += 1;
series_indexes.insert(expected.index_uuid, series_handle);
}
// Cover SSTs outside planned aggregate builds, including skipped buckets.
for file in files {
let file_id = file.file_id().file_id();
if range_indexes.contains(&file_id) {
continue;
}
if let Some(file_id) = build_range_index(store, region, version, file).await? {
stats.built_range += 1;
range_indexes.insert(file_id);
}
}
stats.removed_series = current
.series_indexes
.keys()
.filter(|id| !series_indexes.contains_key(id))
.count();
// Bucket reconciliation is speculative until indexes change. Recompute its coverage
// next time rather than replacing the published snapshot on a no-op pass.
if !stats.changed() {
return Ok((None, stats));
}
let next = SeriesIndexVersion {
range_indexes,
series_indexes,
index_buckets: plan.index_buckets,
};
Ok((Some(next), stats))
}
/// Writes catalogs in a stable order; the two writes are not atomic together.
async fn persist_index_catalogs(
store: &ObjectStore,
region_id: RegionId,
next: &SeriesIndexVersion,
) -> Result<()> {
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 series_entries = next
.series_indexes
.values()
.map(|handle| handle.entry().clone())
.collect::<Vec<_>>();
series_entries.sort_unstable_by_key(|entry| {
(
entry.bucket_start,
entry.bucket_end,
entry.min_file_sequence,
entry.max_file_sequence,
)
});
store_catalog(
store,
&range_catalog_path(region_id),
&RangeIndexCatalog {
indexes: range_entries,
},
)
.await?;
store_catalog(
store,
&series_catalog_path(region_id),
&SeriesIndexCatalog {
indexes: series_entries,
},
)
.await?;
Ok(())
}
/// Publishes the snapshot and retires series files absent from the new version.
fn publish_index_version(region: &MitoRegionRef, next: Arc<SeriesIndexVersion>) {
let previous = region.series_index_version_control.publish(next.clone());
for (id, handle) in &previous.series_indexes {
if !next.series_indexes.contains_key(id) {
// Purge only after readers release their retained handles.
handle.mark_deleted();
}
}
}
+53
View File
@@ -108,3 +108,56 @@ pub(crate) fn series_index_channel(
let (sender, receiver) = unbounded_channel();
(IndexFilePurger { store, sender }, receiver)
}
#[cfg(test)]
mod tests {
use std::time::Duration;
use object_store::ObjectStore;
use object_store::services::Memory;
use store_api::storage::{FileId, RegionId};
use super::*;
#[tokio::test]
async fn test_purge_drains_queue_after_last_sender_drops() {
let store = ObjectStore::new(Memory::default()).unwrap();
let (purger, receiver) = series_index_channel(store.clone());
let mut paths = Vec::new();
for _ in 0..3 {
let file_id = RegionFileId::new(RegionId::new(1, 1), FileId::random());
let path = series_index_path(file_id.region_id(), file_id.file_id());
store.write(&path, "index").await.unwrap();
paths.push(path);
purger.purge(PurgeRequest { file_id });
}
drop(purger);
tokio::time::timeout(
Duration::from_secs(10),
run_index_purge_task(0, store.clone(), receiver),
)
.await
.unwrap();
for path in paths {
assert!(!store.exists(&path).await.unwrap());
}
}
#[tokio::test]
async fn test_purge_falls_back_after_receiver_drops() {
let store = ObjectStore::new(Memory::default()).unwrap();
let (purger, receiver) = series_index_channel(store.clone());
drop(receiver);
let file_id = RegionFileId::new(RegionId::new(1, 1), FileId::random());
let path = series_index_path(file_id.region_id(), file_id.file_id());
store.write(&path, "index").await.unwrap();
purger.purge(PurgeRequest { file_id });
tokio::time::timeout(Duration::from_secs(10), async {
while store.exists(&path).await.unwrap() {
tokio::task::yield_now().await;
}
})
.await
.unwrap();
}
}
+126 -11
View File
@@ -18,14 +18,18 @@ use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use common_telemetry::info;
use common_telemetry::{info, warn};
use object_store::ObjectStore;
use tokio::sync::Notify;
use tokio::sync::mpsc::UnboundedReceiver;
use tokio::task::JoinHandle;
use tokio::time::{Instant, MissedTickBehavior};
use crate::series_index::purger::{PurgeRequest, run_index_purge_task};
use crate::metrics::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};
use crate::time_provider::TimeProviderRef;
/// Shared lifecycle state for a worker's series-index task.
#[derive(Debug)]
@@ -46,6 +50,10 @@ impl SeriesIndexTaskState {
self.running.load(Ordering::Acquire)
}
pub(crate) fn wake(&self) {
self.notify.notify_one();
}
pub(crate) fn stop(&self) {
self.running.store(false, Ordering::Release);
// Retain a permit if maintenance has not started waiting yet.
@@ -58,20 +66,34 @@ impl SeriesIndexTaskState {
}
/// Starts both tasks, detaching purge and returning the maintenance handle.
#[allow(clippy::too_many_arguments)]
pub(crate) fn spawn_series_index_tasks(
worker_id: u32,
store: ObjectStore,
regions: RegionMapRef,
state: Arc<SeriesIndexTaskState>,
bucket_width: Duration,
purger: IndexFilePurger,
purge_receiver: UnboundedReceiver<PurgeRequest>,
interval: Duration,
time_provider: TimeProviderRef,
) -> JoinHandle<()> {
// Snapshots may retain senders after the worker stops; purge until all senders drop.
common_runtime::spawn_global(run_index_purge_task(worker_id, store, purge_receiver));
common_runtime::spawn_global(run_index_purge_task(
worker_id,
store.clone(),
purge_receiver,
));
common_runtime::spawn_global(async move {
SeriesIndexTask {
worker_id,
store,
regions,
bucket_width,
purger,
state,
interval,
time_provider,
}
.run()
.await;
@@ -80,9 +102,14 @@ pub(crate) fn spawn_series_index_tasks(
/// Periodic series-index maintenance for one region worker.
struct SeriesIndexTask {
store: ObjectStore,
regions: RegionMapRef,
bucket_width: Duration,
purger: IndexFilePurger,
worker_id: u32,
state: Arc<SeriesIndexTaskState>,
interval: Duration,
time_provider: TimeProviderRef,
}
impl SeriesIndexTask {
@@ -90,17 +117,17 @@ impl SeriesIndexTask {
async fn run(mut self) {
let worker_id = self.worker_id;
info!("Start series-index background task, worker: {worker_id}");
let interval = self.interval;
let interval = self.time_provider.wait_duration(self.interval);
let mut timer = tokio::time::interval_at(Instant::now() + interval, interval);
timer.set_missed_tick_behavior(MissedTickBehavior::Skip);
// Schedule future ticks from a late tick rather than the original cadence.
timer.set_missed_tick_behavior(MissedTickBehavior::Delay);
while self.state.is_running() {
tokio::select! {
_ = self.state.notified() => {}
_ = timer.tick() => {
if self.state.is_running() {
self.maintain().await;
}
}
_ = timer.tick() => {}
}
if self.state.is_running() {
self.maintain().await;
}
}
info!("Stop series-index background task, worker: {worker_id}");
@@ -108,6 +135,94 @@ impl SeriesIndexTask {
/// Runs periodic maintenance independently of incoming deletion requests.
async fn maintain(&mut self) {
// TODO: Reconcile indexes and perform other periodic maintenance here.
for region in self.regions.list_regions() {
if !self.state.is_running() {
break;
}
// Best effort: the region can still change state during reconciliation.
// Local indexes can be built on followers as well as writable leaders.
if !matches!(
region.state(),
RegionRoleState::Follower | RegionRoleState::Leader(RegionLeaderState::Writable)
) {
continue;
}
if let Err(error) = reconcile_series_indexes(
self.worker_id,
self.store.clone(),
region.clone(),
self.bucket_width,
self.time_provider.current_time_millis(),
self.purger.clone(),
)
.await
{
SERIES_INDEX_RECONCILE_TOTAL
.with_label_values(&["failure"])
.inc();
warn!(error; "Failed to reconcile series indexes, worker: {}, region: {}", self.worker_id, region.region_id);
}
}
}
}
#[cfg(test)]
mod tests {
use object_store::services::Memory;
use store_api::region_engine::RegionEngine;
use super::*;
use crate::region::RegionMap;
use crate::series_index::catalog::{range_catalog_path, series_catalog_path};
use crate::series_index::purger::series_index_channel;
use crate::series_index::tests::prepare_region;
use crate::test_util::TestEnv;
#[rstest::rstest]
#[case::follower(RegionRoleState::Follower, true)]
#[case::writable(RegionRoleState::Leader(RegionLeaderState::Writable), true)]
#[case::staging(RegionRoleState::Leader(RegionLeaderState::Staging), false)]
#[case::entering_staging(RegionRoleState::Leader(RegionLeaderState::EnteringStaging), false)]
#[case::altering(RegionRoleState::Leader(RegionLeaderState::Altering), false)]
#[case::dropping(RegionRoleState::Leader(RegionLeaderState::Dropping), false)]
#[case::truncating(RegionRoleState::Leader(RegionLeaderState::Truncating), false)]
#[case::editing(RegionRoleState::Leader(RegionLeaderState::Editing), false)]
#[case::downgrading(RegionRoleState::Leader(RegionLeaderState::Downgrading), false)]
#[tokio::test]
async fn test_maintenance_region_states(#[case] role: RegionRoleState, #[case] builds: bool) {
let mut env = TestEnv::with_prefix("series-maintenance-state").await;
let (engine, region) = prepare_region(&mut env).await;
// Install the desired state without triggering the corresponding DDL.
region.switch_state_to_staging(RegionLeaderState::Writable);
region
.manifest_ctx
.exit_staging(region.region_id, role)
.unwrap();
let store = ObjectStore::new(Memory::default()).unwrap();
let (purger, _receiver) = series_index_channel(store.clone());
let regions = Arc::new(RegionMap::default());
regions.insert_region(region.clone());
let mut task = SeriesIndexTask {
store: store.clone(),
regions,
bucket_width: Duration::from_secs(100),
purger,
worker_id: 0,
state: Arc::new(SeriesIndexTaskState::new()),
interval: Duration::from_secs(3600),
time_provider: Arc::new(crate::time_provider::StdTimeProvider),
};
task.maintain().await;
assert_eq!(
builds,
!region.series_index_version().series_indexes.is_empty()
);
for path in [
range_catalog_path(region.region_id),
series_catalog_path(region.region_id),
] {
assert_eq!(builds, store.exists(&path).await.unwrap());
}
engine.stop().await.unwrap();
}
}
+648 -7
View File
@@ -12,25 +12,36 @@
// See the License for the specific language governing permissions and
// limitations under the License.
//! Shared fixtures for series-index construction tests.
//! Behavioral coverage for series-index reconciliation.
use std::sync::Arc;
use std::time::Duration;
use api::v1::helper::row;
use api::v1::value::ValueData;
use api::v1::{ColumnDataType, Rows, SemanticType, WriteHint};
use object_store::ObjectStore;
use object_store::layers::mock::{self, MockLayerBuilder, oio};
use object_store::services::Memory;
use store_api::codec::PrimaryKeyEncoding;
use store_api::metric_engine_consts::PRIMARY_KEY_ENCODING;
use store_api::region_engine::RegionEngine;
use store_api::region_request::{RegionPutRequest, RegionRequest};
use store_api::storage::RegionId;
use store_api::storage::consts::PRIMARY_KEY_COLUMN_NAME;
use store_api::storage::{FileId, RegionId};
use super::catalog::{
delete_catalogs, load_version_control, range_catalog_path, range_index_path,
series_catalog_path, series_index_path,
};
use super::maintenance::reconcile_series_indexes;
use super::purger::series_index_channel;
use crate::config::MitoConfig;
use crate::engine::MitoEngine;
use crate::region::MitoRegionRef;
use crate::region::{MitoRegionRef, RegionMap};
use crate::test_util::sst_util::{new_sparse_primary_key, sst_region_metadata_with_encoding};
use crate::test_util::{CreateRequestBuilder, TestEnv, flush_region, rows_schema};
use crate::time_provider::StdTimeProvider;
/// Builds real sparse SSTs; background maintenance is disabled so tests control publication.
pub(super) async fn prepare_region(env: &mut TestEnv) -> (MitoEngine, MitoRegionRef) {
@@ -66,10 +77,12 @@ async fn prepare_region_with_timestamps(
"100s".to_string(),
);
// Keep the source SSTs stable while tests reconcile multiple buckets.
request.options.insert(
"compaction.twcs.trigger_file_num".to_string(),
"100".to_string(),
);
for window in ["active_window", "inactive_window"] {
request.options.insert(
format!("compaction.twcs.{window}.trigger_file_num"),
"100".to_string(),
);
}
let full_schema = rows_schema(&request);
let mut pk_column = full_schema[0].clone();
pk_column.column_name = PRIMARY_KEY_COLUMN_NAME.to_string();
@@ -112,3 +125,631 @@ async fn prepare_region_with_timestamps(
let region = engine.get_region(region_id).unwrap();
(engine, region)
}
struct IndexTest {
region: MitoRegionRef,
store: ObjectStore,
purger: super::purger::IndexFilePurger,
}
impl IndexTest {
fn new(
region: MitoRegionRef,
) -> (
Self,
tokio::sync::mpsc::UnboundedReceiver<super::purger::PurgeRequest>,
) {
let store = ObjectStore::new(Memory::default()).unwrap();
let (purger, receiver) = series_index_channel(store.clone());
(
Self {
region,
store,
purger,
},
receiver,
)
}
async fn reconcile(
&self,
store: ObjectStore,
) -> crate::error::Result<super::maintenance::ReconcileStats> {
reconcile_series_indexes(
0,
store,
self.region.clone(),
Duration::from_secs(100),
0,
self.purger.clone(),
)
.await
}
}
#[tokio::test]
async fn test_reconcile_restores_and_reuses_indexes() {
let mut env = TestEnv::with_prefix("series-reconcile").await;
let (engine, region) = prepare_region(&mut env).await;
let (test, _receiver) = IndexTest::new(region.clone());
let store = &test.store;
let purger = &test.purger;
let stats = test.reconcile(store.clone()).await.unwrap();
assert_eq!((4, 1), (stats.built_range, stats.built_series));
let first = region.series_index_version();
let first_id = *first.series_indexes.keys().next().unwrap();
let missing_paths = [
range_index_path(
region.region_id,
*first.range_indexes.iter().next().unwrap(),
),
series_index_path(region.region_id, first_id),
];
for path in &missing_paths {
store.delete(path).await.unwrap();
}
// Restore catalog entries even when their index files are missing.
let restored = load_version_control(store, region.region_id, purger)
.await
.current();
assert_eq!(first.range_indexes, restored.range_indexes);
assert_eq!(first.index_buckets, restored.index_buckets);
assert_eq!(
first.series_indexes[&first_id].entry(),
restored.series_indexes[&first_id].entry()
);
// Catalog changes are not reloaded, and missing index files are not repaired.
for path in [
range_catalog_path(region.region_id),
series_catalog_path(region.region_id),
] {
store.write(&path, "{}").await.unwrap();
}
let stats = test.reconcile(store.clone()).await.unwrap();
assert_eq!((0, 0), (stats.built_range, stats.built_series));
assert!(Arc::ptr_eq(&first, &region.series_index_version()));
for path in [
range_catalog_path(region.region_id),
series_catalog_path(region.region_id),
] {
assert_eq!(
b"{}".as_slice(),
store.read(&path).await.unwrap().to_bytes().as_ref()
);
}
for path in &missing_paths {
assert!(!store.exists(path).await.unwrap());
}
engine.stop().await.unwrap();
}
#[tokio::test]
async fn test_reconcile_range_only_build_and_removal() {
let mut env = TestEnv::with_prefix("series-range-lifecycle").await;
let (engine, region) = prepare_region_with_timestamps(&mut env, &[1000]).await;
let (test, _receiver) = IndexTest::new(region.clone());
let initial = region.series_index_version();
let stats = test.reconcile(test.store.clone()).await.unwrap();
assert_eq!((1, 0), (stats.built_range, stats.built_series));
let built = region.series_index_version();
assert!(!Arc::ptr_eq(&initial, &built));
let files = region
.version()
.ssts
.levels()
.iter()
.flat_map(|level| level.files())
.map(|file| file.meta_ref().clone())
.collect();
region.version_control.apply_edit(
Some(crate::manifest::action::RegionEdit {
files_to_remove: files,
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(),
);
let stats = test.reconcile(test.store.clone()).await.unwrap();
assert_eq!(
(0, 0, 1),
(stats.built_range, stats.built_series, stats.removed_range)
);
let removed = region.series_index_version();
assert!(removed.range_indexes.is_empty());
assert!(!Arc::ptr_eq(&built, &removed));
let restored = load_version_control(&test.store, region.region_id, &test.purger).await;
assert!(restored.current().range_indexes.is_empty());
test.reconcile(test.store.clone()).await.unwrap();
assert!(Arc::ptr_eq(&removed, &region.series_index_version()));
engine.stop().await.unwrap();
}
struct FailingSeriesWriter {
inner: oio::Writer,
fail: bool,
}
#[rstest::rstest]
#[case::series("series")]
#[case::range_catalog("range-index.json")]
#[case::series_catalog("series-index.json")]
#[tokio::test]
async fn test_reconcile_replaces_whole_bucket_after_failure(#[case] failure: &str) {
let mut env = TestEnv::with_prefix("series-replacement").await;
let (engine, region) =
prepare_region_with_timestamps(&mut env, &[1000, 2000, 3000, 4000, 5000]).await;
// Retain all source handles while hiding the newest SST for the initial build.
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();
let edit = |add| crate::manifest::action::RegionEdit {
files_to_add: if add { vec![newest.clone()] } else { vec![] },
files_to_remove: if add { vec![] } else { vec![newest.clone()] },
timestamp_ms: None,
compaction_time_window: None,
flushed_entry_id: None,
flushed_sequence: None,
committed_sequence: None,
};
region.version_control.apply_edit(
Some(edit(false)),
&[],
crate::test_util::new_noop_file_purger(),
);
let (test, mut receiver) = IndexTest::new(region.clone());
test.reconcile(test.store.clone()).await.unwrap();
let previous = region.series_index_version();
let old_id = *previous.series_indexes.keys().next().unwrap();
assert_eq!(
4,
previous.series_indexes[&old_id]
.entry()
.source_file_ids
.len()
);
region.version_control.apply_edit(
Some(edit(true)),
&[],
crate::test_util::new_noop_file_purger(),
);
let failure = failure.to_string();
let layer = MockLayerBuilder::default()
.writer_factory(Arc::new(move |path, _, inner| {
Box::new(FailingSeriesWriter {
inner,
fail: if failure == "series" {
path.contains("/series/")
} else {
path.ends_with(&failure)
},
})
}))
.build()
.unwrap();
assert!(
test.reconcile(test.store.clone().layer(layer))
.await
.is_err()
);
assert!(Arc::ptr_eq(&previous, &region.series_index_version()));
// Failed catalog publication may retire a completed replacement, never its predecessor.
while let Ok(request) = receiver.try_recv() {
assert_ne!(old_id, request.file_id.file_id());
}
let stats = test.reconcile(test.store.clone()).await.unwrap();
assert_eq!((1, 1), (stats.built_series, stats.removed_series));
let replacement = region.series_index_version();
assert_eq!(1, replacement.series_indexes.len());
assert!(!replacement.series_indexes.contains_key(&old_id));
let entry = replacement.series_indexes.values().next().unwrap().entry();
assert_eq!(5, entry.source_file_ids.len());
assert!(entry.source_file_ids.contains(&newest.file_id));
let restored = load_version_control(&test.store, region.region_id, &test.purger).await;
assert_eq!(replacement.index_buckets, restored.current().index_buckets);
assert_eq!(1, restored.current().series_indexes.len());
test.reconcile(test.store.clone()).await.unwrap();
assert!(Arc::ptr_eq(&replacement, &region.series_index_version()));
// Publication retires the old file only after the last reader releases its snapshot.
let old_path = series_index_path(region.region_id, old_id);
assert!(test.store.exists(&old_path).await.unwrap());
assert!(receiver.try_recv().is_err());
drop(previous);
assert_eq!(old_id, receiver.try_recv().unwrap().file_id.file_id());
assert!(receiver.try_recv().is_err());
engine.stop().await.unwrap();
}
impl mock::Write for FailingSeriesWriter {
async fn write(&mut self, buffer: mock::Buffer) -> mock::Result<()> {
self.inner.write(buffer).await
}
async fn close(&mut self) -> mock::Result<mock::Metadata> {
if self.fail {
return Err(mock::Error::new(
mock::ErrorKind::Unexpected,
"injected series write failure",
));
}
self.inner.close().await
}
async fn abort(&mut self) -> mock::Result<()> {
self.inner.abort().await
}
}
#[tokio::test]
async fn test_failed_series_build_keeps_completed_sst_range_indexes() {
let mut env = TestEnv::with_prefix("series-build-failure").await;
let (engine, region) = prepare_region(&mut env).await;
let (test, mut receiver) = IndexTest::new(region.clone());
let store = &test.store;
let layer = MockLayerBuilder::default()
.writer_factory(Arc::new(|path, _, inner| {
Box::new(FailingSeriesWriter {
inner,
fail: path.contains("/series/"),
})
}))
.build()
.unwrap();
let failing_store = store.clone().layer(layer);
assert!(test.reconcile(failing_store).await.is_err());
assert!(region.series_index_version().range_indexes.is_empty());
assert!(receiver.try_recv().is_err());
for file in region
.version()
.ssts
.levels()
.iter()
.flat_map(|level| level.files())
{
let path = range_index_path(region.region_id, file.file_id().file_id());
assert!(store.exists(&path).await.unwrap());
}
// A retry can rebuild and publish the complete snapshot.
let stats = test.reconcile(store.clone()).await.unwrap();
assert_eq!((4, 1), (stats.built_range, stats.built_series));
engine.stop().await.unwrap();
}
#[tokio::test]
async fn test_reconcile_publishes_after_region_version_changes() {
for during_catalog_write in [false, true] {
let mut env = TestEnv::with_prefix("series-version-change").await;
let (engine, region) = prepare_region(&mut env).await;
let initial_version = region.version();
let (test, mut receiver) = IndexTest::new(region.clone());
let store = &test.store;
let purger = &test.purger;
let target_region = region.clone();
let catalog_path = range_catalog_path(region.region_id);
let layer = MockLayerBuilder::default()
.writer_factory(Arc::new(move |path, _, inner| {
let should_advance = if during_catalog_write {
path == catalog_path
} else {
path.contains("/series/")
};
if should_advance {
// Replace the version while reconciliation is writing its captured snapshot.
target_region
.version_control
.alter_options(target_region.version().options.clone());
}
inner
}))
.build()
.unwrap();
let stats = test.reconcile(store.clone().layer(layer)).await.unwrap();
assert!(!Arc::ptr_eq(&initial_version, &region.version()));
assert_eq!((4, 1), (stats.built_range, stats.built_series));
let published = region.series_index_version();
assert_eq!(4, published.range_indexes.len());
assert_eq!(1, published.series_indexes.len());
let restored = load_version_control(store, region.region_id, purger).await;
assert_eq!(published.range_indexes, restored.current().range_indexes);
let handle = published.series_indexes.values().next().unwrap();
assert_eq!(
handle.entry(),
restored.current().series_indexes[&handle.entry().index_uuid].entry()
);
assert!(receiver.try_recv().is_err());
engine.stop().await.unwrap();
}
}
#[tokio::test]
async fn test_failed_reconcile_retires_only_unpublished_series() {
use std::collections::{HashMap, HashSet};
use std::sync::Mutex;
use super::bucket::{group_files_into_series_buckets, plan_series_indexes};
use super::builder::{build_range_index, build_series_index};
use super::purger::run_index_purge_task;
use super::version::SeriesIndexVersion;
for failure in ["series", "range", "range-index.json", "series-index.json"] {
let mut env = TestEnv::with_prefix("series-unpublished-cleanup").await;
let mut timestamps = vec![
-99000, -98000, -97000, -96000, 1000, 2000, 3000, 4000, 101000,
];
if failure != "range" {
timestamps.extend([102000, 103000, 104000]);
}
let (engine, region) = prepare_region_with_timestamps(&mut env, &timestamps).await;
let (test, mut receiver) = IndexTest::new(region.clone());
let store = &test.store;
let purger = &test.purger;
// Publish the first bucket so failed attempts must preserve a reused handle.
let version = region.version();
let files = version
.ssts
.levels()
.iter()
.flat_map(|level| level.files())
.cloned()
.collect::<Vec<_>>();
let plan = plan_series_indexes(
group_files_into_series_buckets(
&files,
100,
version.compaction_time_window.unwrap().as_secs() as i64,
),
Default::default(),
None,
0,
);
assert_eq!(if failure == "range" { 2 } else { 3 }, plan.builds.len());
let (bucket, entry) = &plan.builds[0];
let mut ranges = HashSet::new();
for file in &bucket.files {
let id = build_range_index(store, &region, &version, file.clone())
.await
.unwrap()
.unwrap();
ranges.insert(id);
}
let handle = build_series_index(store, &region, &version, bucket, entry, purger)
.await
.unwrap();
let reused_path = series_index_path(region.region_id, entry.index_uuid);
region
.series_index_version_control
.publish(Arc::new(SeriesIndexVersion::new(
ranges,
HashMap::from([(entry.index_uuid, handle)]),
)));
let previous = region.series_index_version();
let standalone_path = files
.iter()
.find(|file| file.time_range().0 == common_time::Timestamp::new_millisecond(101000))
.map(|file| range_index_path(region.region_id, file.file_id().file_id()))
.unwrap();
// Every completed unpublished output must be retired.
let attempted = Arc::new(Mutex::new(Vec::<String>::new()));
let recorded = attempted.clone();
let layer = MockLayerBuilder::default()
.writer_factory(Arc::new(move |path, _, inner| {
let mut paths = recorded.lock().unwrap();
if path.contains("/series/") {
paths.push(path.to_string());
}
let fail = match failure {
"series" => path.contains("/series/") && paths.len() == 2,
"range" => path == standalone_path,
catalog => path.ends_with(catalog),
};
Box::new(FailingSeriesWriter { inner, fail })
}))
.build()
.unwrap();
let result = test.reconcile(store.clone().layer(layer)).await;
assert!(
result.is_err(),
"expected {failure} failure, got {result:?}; attempted series paths: {:?}",
attempted.lock().unwrap()
);
assert!(Arc::ptr_eq(&previous, &region.series_index_version()));
let completed = attempted
.lock()
.unwrap()
.iter()
.take(if failure.ends_with(".json") { 2 } else { 1 })
.cloned()
.collect::<HashSet<_>>();
assert_eq!(
if failure.ends_with(".json") { 2 } else { 1 },
completed.len()
);
let mut retired = HashSet::new();
// Drain through the real purge task using a finite channel.
let (drain_purger, drain_receiver) = series_index_channel(store.clone());
while let Ok(request) = receiver.try_recv() {
let path = series_index_path(request.file_id.region_id(), request.file_id.file_id());
assert!(store.exists(&path).await.unwrap());
assert!(retired.insert(path));
drain_purger.purge(request);
}
assert_eq!(completed, retired);
drop(drain_purger);
run_index_purge_task(0, store.clone(), drain_receiver).await;
let attempted_paths = attempted.lock().unwrap().clone();
for path in attempted_paths {
assert!(!store.exists(&path).await.unwrap());
}
assert!(store.exists(&reused_path).await.unwrap());
let stats = test.reconcile(store.clone()).await.unwrap();
assert_eq!(if failure == "range" { 1 } else { 2 }, stats.built_series);
let published = region.series_index_version();
for handle in published.series_indexes.values() {
assert!(
store
.exists(&series_index_path(
region.region_id,
handle.entry().index_uuid
))
.await
.unwrap()
);
}
// Release all snapshots without retiring them: neither reused nor published files
// should have been marked deleted by the cleanup guard.
region.series_index_version_control.publish(Arc::default());
drop(published);
drop(previous);
assert!(receiver.try_recv().is_err());
engine.stop().await.unwrap();
}
}
/// Waits for a background state change without relying on exact scheduling delays.
async fn wait_for(mut ready: impl FnMut() -> bool) {
tokio::time::timeout(Duration::from_secs(10), async {
while !ready() {
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.unwrap();
}
#[tokio::test]
async fn test_maintenance_wakeup_and_timer() {
let mut env = TestEnv::with_prefix("series-task").await;
let (engine, region) = prepare_region(&mut env).await;
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 store = ObjectStore::new(Memory::default()).unwrap();
let (purger, receiver) = series_index_channel(store.clone());
let state = Arc::new(super::task::SeriesIndexTaskState::new());
let clock = Arc::new(crate::time_provider::mock::MockTimeProvider::new(0));
// A notification issued before the task starts must also trigger maintenance.
state.wake();
let regions = Arc::new(RegionMap::default());
regions.insert_region(region.clone());
let task = super::task::spawn_series_index_tasks(
0,
store,
regions,
state.clone(),
Duration::from_secs(100),
purger,
receiver,
Duration::from_secs(3600),
clock.clone(),
);
wait_for(|| !region.series_index_version().series_indexes.is_empty()).await;
// Advance event time; periodic maintenance must expire coverage without a wakeup.
clock.set_now(201_000);
wait_for(|| region.series_index_version().series_indexes.is_empty()).await;
state.stop();
task.await.unwrap();
engine.stop().await.unwrap();
}
#[tokio::test]
async fn test_drop_catalogs_and_retained_snapshot_after_task_stop() {
use std::collections::BTreeMap;
use common_time::Timestamp;
use super::catalog::{SeriesIndexCatalog, SeriesIndexEntry, WindowSequence, store_catalog};
let store = ObjectStore::new(Memory::default()).unwrap();
let region_id = RegionId::new(1, 1);
let entry = SeriesIndexEntry {
index_uuid: FileId::random(),
bucket_start: Timestamp::new_second(0),
bucket_end: Timestamp::new_second(100),
source_file_ids: vec![FileId::random()],
min_file_sequence: 1,
max_file_sequence: 2,
compaction_window_secs: 10,
window_sequences: BTreeMap::from([(
0,
WindowSequence {
start: 0,
end: 100,
max_sequence: 2,
},
)]),
};
let path = series_index_path(region_id, entry.index_uuid);
store.write(&path, "index").await.unwrap();
store_catalog(
&store,
&series_catalog_path(region_id),
&SeriesIndexCatalog {
indexes: vec![entry.clone()],
},
)
.await
.unwrap();
store_catalog(
&store,
&range_catalog_path(region_id),
&super::catalog::RangeIndexCatalog::default(),
)
.await
.unwrap();
let (purger, receiver) = series_index_channel(store.clone());
let control = load_version_control(&store, region_id, &purger).await;
let snapshot = control.current();
assert_eq!(&entry, snapshot.series_indexes[&entry.index_uuid].entry());
let state = Arc::new(super::task::SeriesIndexTaskState::new());
state.stop();
super::task::spawn_series_index_tasks(
0,
store.clone(),
Arc::new(RegionMap::default()),
state,
Duration::from_secs(100),
purger,
receiver,
Duration::from_secs(3600),
Arc::new(StdTimeProvider),
)
.await
.unwrap();
delete_catalogs(&store, region_id).await;
// Removing already absent catalogs is harmless.
delete_catalogs(&store, region_id).await;
assert!(!store.exists(&range_catalog_path(region_id)).await.unwrap());
assert!(!store.exists(&series_catalog_path(region_id)).await.unwrap());
control.mark_dropped();
assert!(store.exists(&path).await.unwrap());
drop(snapshot);
tokio::time::timeout(Duration::from_secs(10), async {
while store.exists(&path).await.unwrap() {
tokio::task::yield_now().await;
}
})
.await
.unwrap();
}
+32 -5
View File
@@ -288,8 +288,14 @@ mod tests {
use crate::sst::index::puffin_manager::PuffinManagerFactory;
use crate::sst::location;
#[rstest::rstest]
#[case(false)]
#[case(true)]
#[tokio::test]
async fn test_file_purge() {
async fn test_range_index_purge_on_handle_release(
#[case] gc_enabled: bool,
#[values(false, true)] is_delete: bool,
) {
common_telemetry::init_default_ut_logging();
let dir = create_temp_dir("file-purge");
@@ -320,7 +326,19 @@ mod tests {
let scheduler = Arc::new(LocalScheduler::new(3));
let file_purger = Arc::new(LocalFilePurger::new(scheduler.clone(), layer, None));
let index_store = ObjectStore::new(object_store::services::Memory::default()).unwrap();
let owner = RegionId::new(9, 1);
let index_path = crate::sst::range_index::range_index_path(owner, sst_file_id.file_id());
index_store.write(&index_path, "range index").await.unwrap();
let file_purger = create_file_purger(
gc_enabled,
PathType::Bare,
scheduler.clone(),
layer,
None,
Arc::new(crate::sst::file_ref::FileReferenceManager::new(None)),
Some(RangeIndexDeleter::new(index_store.clone(), owner)),
);
{
let handle = FileHandle::new(
@@ -344,13 +362,22 @@ mod tests {
},
file_purger,
);
// mark file as deleted and drop the handle, we expect the file is deleted.
handle.mark_deleted();
if is_delete {
handle.mark_deleted();
}
let reader = handle.clone();
drop(handle);
assert!(index_store.exists(&index_path).await.unwrap());
drop(reader);
}
scheduler.stop(true).await.unwrap();
assert!(!object_store.exists(&path).await.unwrap());
assert_eq!(
object_store.exists(&path).await.unwrap(),
gc_enabled || !is_delete
);
assert_eq!(index_store.exists(&index_path).await.unwrap(), !is_delete);
}
#[tokio::test]
+46
View File
@@ -50,3 +50,49 @@ impl TimeProvider for StdTimeProvider {
current_time_millis() - current_millis
}
}
/// Controllable event time with short waits for background-task tests.
#[cfg(test)]
pub(crate) mod mock {
use std::sync::atomic::{AtomicI64, Ordering};
use std::time::Duration;
use crate::time_provider::TimeProvider;
#[derive(Debug)]
pub(crate) struct MockTimeProvider {
now: AtomicI64,
elapsed: AtomicI64,
}
impl TimeProvider for MockTimeProvider {
fn current_time_millis(&self) -> i64 {
self.now.load(Ordering::Relaxed)
}
fn elapsed_since(&self, _current_millis: i64) -> i64 {
self.elapsed.load(Ordering::Relaxed)
}
fn wait_duration(&self, _duration: Duration) -> Duration {
Duration::from_millis(20)
}
}
impl MockTimeProvider {
pub(crate) fn new(now: i64) -> Self {
Self {
now: AtomicI64::new(now),
elapsed: AtomicI64::new(0),
}
}
pub(crate) fn set_now(&self, now: i64) {
self.now.store(now, Ordering::Relaxed);
}
pub(crate) fn set_elapsed(&self, elapsed: i64) {
self.elapsed.store(elapsed, Ordering::Relaxed);
}
}
}
+11 -3
View File
@@ -603,13 +603,17 @@ impl<S: LogStore> WorkerStarter<S> {
.zip(series_index_task_state.clone())
.map(|(store, state)| {
let (purger, purge_receiver) = series_index_channel(store.clone());
series_index_purger = Some(purger);
series_index_purger = Some(purger.clone());
spawn_series_index_tasks(
self.id,
store,
regions.clone(),
state,
self.config.experimental_series_index_bucket_width,
purger,
purge_receiver,
self.config.experimental_series_index_maintenance_interval,
self.time_provider.clone(),
)
});
let now = self.time_provider.current_time_millis();
@@ -636,6 +640,7 @@ impl<S: LogStore> WorkerStarter<S> {
self.index_build_job_pool,
self.config.max_background_index_builds,
),
series_index_task_state: series_index_task_state.clone(),
series_index_store: self.series_index_store,
series_index_purger,
flush_scheduler: FlushScheduler::new(self.flush_job_pool),
@@ -700,9 +705,9 @@ pub(crate) struct RegionWorker {
sender: Sender<WorkerRequestWithTime>,
/// Handle to the worker thread.
handle: Mutex<Option<JoinHandle<()>>>,
/// Handle to the series-index maintenance task.
/// Handle to the sequential series-index task.
series_index_handle: Mutex<Option<JoinHandle<()>>>,
/// Controls the worker-owned series-index maintenance task.
/// Controls the worker-owned series-index task.
series_index_task_state: Option<Arc<SeriesIndexTaskState>>,
/// Whether to run the worker thread.
running: Arc<AtomicBool>,
@@ -935,6 +940,8 @@ struct RegionWorkerLoop<S> {
write_buffer_manager: WriteBufferManagerRef,
/// Scheduler for index build task.
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.
series_index_store: Option<ObjectStore>,
series_index_purger: Option<IndexFilePurger>,
@@ -1625,6 +1632,7 @@ mod tests {
.create_worker_group(MitoConfig {
num_workers: 4,
experimental_series_index_root: "series-index".to_string(),
experimental_series_index_maintenance_interval: Duration::from_secs(3600),
..Default::default()
})
.await;
+3
View File
@@ -134,6 +134,9 @@ impl<S: LogStore> RegionWorkerLoop<S> {
.await?;
debug_assert!(!reopened_region.is_writable());
self.regions.insert_region(reopened_region.clone());
if let Some(state) = &self.series_index_task_state {
state.wake();
}
Ok(reopened_region)
}
+3
View File
@@ -96,6 +96,9 @@ impl<S: LogStore> RegionWorkerLoop<S> {
// Insert the MitoRegion into the RegionMap.
self.regions.insert_region(region);
if let Some(state) = &self.series_index_task_state {
state.wake();
}
Ok(0)
}
+4
View File
@@ -227,6 +227,7 @@ impl<S: LogStore> RegionWorkerLoop<S> {
let opening_regions = self.opening_regions.clone();
let region_count = self.region_count.clone();
let worker_id = self.id;
let series_index_task_state = self.series_index_task_state.clone();
opening_regions.insert_sender(region_id, sender);
common_runtime::spawn_global(async move {
match opener.open(&config, &wal).await {
@@ -249,6 +250,9 @@ impl<S: LogStore> RegionWorkerLoop<S> {
// Insert the Region into the RegionMap.
regions.insert_region(region);
if let Some(state) = &series_index_task_state {
state.wake();
}
let senders = opening_regions.remove_sender(region_id);
for sender in senders {
+1
View File
@@ -2383,6 +2383,7 @@ experimental_manifest_keep_removed_file_ttl = "1h"
compress_manifest = false
experimental_series_index_root = ""
experimental_series_index_maintenance_interval = "5m"
experimental_series_index_bucket_width = "5days"
experimental_compaction_memory_limit = "unlimited"
experimental_compaction_on_exhausted = "wait"
auto_flush_interval = "10m"