feat(mito2): add series index catalog and lifecycle foundation (#9053)

* refactor(mito2): add series index catalog and lifecycle components

Signed-off-by: evenyag <realevenyag@gmail.com>

* feat(mito2): restore series index catalogs on region open

Signed-off-by: evenyag <realevenyag@gmail.com>

* refactor(mito2): simplify series index foundation and maintenance

Signed-off-by: evenyag <realevenyag@gmail.com>

* refactor(mito2): separate series index purge task and simplify tests

Signed-off-by: evenyag <realevenyag@gmail.com>

* docs: defer experimental series index configuration examples

Signed-off-by: evenyag <realevenyag@gmail.com>

* feat(mito): make series index maintenance interval configurable

Signed-off-by: evenyag <realevenyag@gmail.com>

* fix(mito2): correct series index cleanup on close and drop

Signed-off-by: evenyag <realevenyag@gmail.com>

* test(mito2): revert drop test changes

Signed-off-by: evenyag <realevenyag@gmail.com>

* test: update config API expectation for series index settings

Signed-off-by: evenyag <realevenyag@gmail.com>

* refactor: use tokio unbounded channel for series index purger

Signed-off-by: evenyag <realevenyag@gmail.com>

* chore(mito2): simplify review test scope and clarify index config

Signed-off-by: evenyag <realevenyag@gmail.com>

---------

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