mirror of
https://github.com/GreptimeTeam/greptimedb.git
synced 2026-08-21 21:48:43 +00:00
fix(mito2): fence async index builds by schema generation (#8697)
* fix(mito2): fence async index builds by schema generation Signed-off-by: Dennis Zhuang <killme2008@gmail.com> * fix(mito2): retry stale index builds after schema changes Signed-off-by: Dennis Zhuang <killme2008@gmail.com> --------- Signed-off-by: Dennis Zhuang <killme2008@gmail.com>
This commit is contained in:
@@ -19,16 +19,19 @@ use std::sync::Arc;
|
||||
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||
use std::time::Duration;
|
||||
|
||||
use api::v1::Rows;
|
||||
use api::v1::region::{StrictWindow, compact_request};
|
||||
use api::v1::{Rows, SemanticType};
|
||||
use async_trait::async_trait;
|
||||
use common_base::readable_size::ReadableSize;
|
||||
use common_recordbatch::RecordBatches;
|
||||
use datatypes::arrow::array::AsArray;
|
||||
use datatypes::arrow::datatypes::TimestampMillisecondType;
|
||||
use datatypes::prelude::ConcreteDataType;
|
||||
use datatypes::schema::{ColumnSchema, SkippingIndexOptions, SkippingIndexType};
|
||||
use store_api::metadata::ColumnMetadata;
|
||||
use store_api::region_engine::RegionEngine;
|
||||
use store_api::region_request::{
|
||||
AlterKind, RegionAlterRequest, RegionBuildIndexRequest, RegionCloseRequest,
|
||||
AddColumn, AlterKind, RegionAlterRequest, RegionBuildIndexRequest, RegionCloseRequest,
|
||||
RegionCompactRequest, RegionRequest, SetIndexOption,
|
||||
};
|
||||
use store_api::storage::{RegionId, ScanRequest};
|
||||
@@ -151,6 +154,12 @@ impl IndexPublicationGate {
|
||||
self.stopped_notify.notified().await;
|
||||
}
|
||||
}
|
||||
|
||||
async fn wait_finished(&self, count: usize) {
|
||||
while self.finish_count.load(Ordering::Relaxed) < count {
|
||||
self.stopped_notify.notified().await;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
@@ -821,6 +830,307 @@ async fn test_index_build_type_schema_change_multiple_files() {
|
||||
assert_eq!(num_of_index_files(&engine, &scanner, region_id).await, 3);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_consecutive_schema_changes_publish_latest_index_generation() {
|
||||
let mut env = TestEnv::with_prefix("test_consecutive_schema_change_index_builds_").await;
|
||||
let mut config = async_build_mode_config(false);
|
||||
config.max_background_index_builds = 2;
|
||||
let gate = Arc::new(IndexPublicationGate::new(
|
||||
IndexPublicationPhase::BeforeManifestCommit,
|
||||
));
|
||||
let engine = Arc::new(
|
||||
env.create_engine_with(config, None, Some(gate.clone()), None)
|
||||
.await,
|
||||
);
|
||||
|
||||
let region_id = RegionId::new(1, 1);
|
||||
env.get_schema_metadata_manager()
|
||||
.register_region_table_info(
|
||||
region_id.table_id(),
|
||||
"test_table",
|
||||
"test_catalog",
|
||||
"test_schema",
|
||||
None,
|
||||
env.get_kv_backend(),
|
||||
)
|
||||
.await;
|
||||
|
||||
let request = CreateRequestBuilder::new().build();
|
||||
let table_dir = request.table_dir.clone();
|
||||
let column_schemas = rows_schema(&request);
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Create(request))
|
||||
.await
|
||||
.unwrap();
|
||||
put_and_flush(&engine, region_id, &column_schemas, 0..20).await;
|
||||
|
||||
let set_index = |option| {
|
||||
RegionRequest::Alter(RegionAlterRequest {
|
||||
kind: AlterKind::SetIndexes {
|
||||
options: vec![option],
|
||||
},
|
||||
})
|
||||
};
|
||||
|
||||
engine
|
||||
.handle_request(
|
||||
region_id,
|
||||
set_index(SetIndexOption::Inverted {
|
||||
column_name: "tag_0".to_string(),
|
||||
}),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
tokio::time::timeout(Duration::from_secs(10), gate.wait_entered(1))
|
||||
.await
|
||||
.expect("first schema generation did not reach manifest publication");
|
||||
|
||||
engine
|
||||
.handle_request(
|
||||
region_id,
|
||||
set_index(SetIndexOption::Inverted {
|
||||
column_name: "field_0".to_string(),
|
||||
}),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
engine
|
||||
.handle_request(
|
||||
region_id,
|
||||
set_index(SetIndexOption::Skipping {
|
||||
column_name: "field_0".to_string(),
|
||||
options: SkippingIndexOptions::new_unchecked(
|
||||
1024,
|
||||
0.01,
|
||||
SkippingIndexType::BloomFilter,
|
||||
),
|
||||
}),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
tokio::time::timeout(Duration::from_secs(10), gate.wait_stopped(1))
|
||||
.await
|
||||
.expect("superseded schema-change build was not coalesced");
|
||||
assert_eq!(
|
||||
gate.entered.load(Ordering::Relaxed),
|
||||
1,
|
||||
"the same SST was built concurrently"
|
||||
);
|
||||
gate.release(1);
|
||||
|
||||
tokio::time::timeout(Duration::from_secs(10), gate.wait_entered(2))
|
||||
.await
|
||||
.expect("latest schema generation was not scheduled");
|
||||
gate.release(1);
|
||||
tokio::time::timeout(Duration::from_secs(10), gate.wait_finished(1))
|
||||
.await
|
||||
.expect("latest schema generation was not published");
|
||||
|
||||
assert_eq!(gate.finish_count.load(Ordering::Relaxed), 1);
|
||||
|
||||
let region = engine.get_region(region_id).unwrap();
|
||||
let version = region.version();
|
||||
let files = current_file_metas(&engine, region_id).await;
|
||||
assert_eq!(files.len(), 1);
|
||||
assert!(
|
||||
files[0].is_index_consistent_with_region(&version.metadata.column_metadatas),
|
||||
"the published index must match the latest schema generation, file: {:?}, metadata: {:?}",
|
||||
files[0],
|
||||
version.metadata
|
||||
);
|
||||
assert_eq!(
|
||||
region.manifest_ctx.manifest().await.metadata.schema_version,
|
||||
version.metadata.schema_version
|
||||
);
|
||||
let scanner = engine
|
||||
.scanner(region_id, ScanRequest::default())
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(num_of_index_files(&engine, &scanner, region_id).await, 1);
|
||||
|
||||
reopen_region(&engine, region_id, table_dir, true, HashMap::new()).await;
|
||||
let reopened = engine.get_region(region_id).unwrap();
|
||||
let reopened_files = current_file_metas(&engine, region_id).await;
|
||||
assert_eq!(reopened_files.len(), 1);
|
||||
assert!(
|
||||
reopened_files[0]
|
||||
.is_index_consistent_with_region(&reopened.version().metadata.column_metadatas)
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_unrelated_schema_change_retries_stale_flush_index_build() {
|
||||
let mut env =
|
||||
TestEnv::with_prefix("test_unrelated_schema_change_retries_stale_index_build_").await;
|
||||
let gate = Arc::new(IndexPublicationGate::new(
|
||||
IndexPublicationPhase::BeforeManifestCommit,
|
||||
));
|
||||
let engine = Arc::new(
|
||||
env.create_engine_with(
|
||||
async_build_mode_config(true),
|
||||
None,
|
||||
Some(gate.clone()),
|
||||
None,
|
||||
)
|
||||
.await,
|
||||
);
|
||||
|
||||
let region_id = RegionId::new(1, 1);
|
||||
env.get_schema_metadata_manager()
|
||||
.register_region_table_info(
|
||||
region_id.table_id(),
|
||||
"test_table",
|
||||
"test_catalog",
|
||||
"test_schema",
|
||||
None,
|
||||
env.get_kv_backend(),
|
||||
)
|
||||
.await;
|
||||
|
||||
let request = CreateRequestBuilder::new().build_with_index();
|
||||
let column_schemas = rows_schema(&request);
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Create(request))
|
||||
.await
|
||||
.unwrap();
|
||||
put_and_flush(&engine, region_id, &column_schemas, 0..20).await;
|
||||
tokio::time::timeout(Duration::from_secs(10), gate.wait_entered(1))
|
||||
.await
|
||||
.expect("flush index build did not reach manifest publication");
|
||||
|
||||
engine
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Alter(RegionAlterRequest {
|
||||
kind: AlterKind::AddColumns {
|
||||
columns: vec![AddColumn {
|
||||
column_metadata: ColumnMetadata {
|
||||
column_schema: ColumnSchema::new(
|
||||
"field_1",
|
||||
ConcreteDataType::float64_datatype(),
|
||||
true,
|
||||
),
|
||||
semantic_type: SemanticType::Field,
|
||||
column_id: 3,
|
||||
},
|
||||
location: None,
|
||||
}],
|
||||
},
|
||||
}),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
gate.release(1);
|
||||
tokio::time::timeout(Duration::from_secs(10), gate.wait_entered(2))
|
||||
.await
|
||||
.expect("stale flush index build was not retried");
|
||||
gate.release(1);
|
||||
tokio::time::timeout(Duration::from_secs(10), gate.wait_stopped(2))
|
||||
.await
|
||||
.expect("retried index build did not stop");
|
||||
|
||||
assert_eq!(gate.abort_count.load(Ordering::Relaxed), 1);
|
||||
assert_eq!(gate.finish_count.load(Ordering::Relaxed), 1);
|
||||
|
||||
let region = engine.get_region(region_id).unwrap();
|
||||
let files = current_file_metas(&engine, region_id).await;
|
||||
assert_eq!(files.len(), 1);
|
||||
assert!(files[0].is_index_consistent_with_region(®ion.version().metadata.column_metadatas));
|
||||
let scanner = engine
|
||||
.scanner(region_id, ScanRequest::default())
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(num_of_index_files(&engine, &scanner, region_id).await, 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_schema_change_uses_manifest_file_generation() {
|
||||
let mut env = TestEnv::with_prefix("test_schema_change_uses_manifest_file_generation_").await;
|
||||
let mut config = async_build_mode_config(false);
|
||||
config.max_background_index_builds = 2;
|
||||
let gate = Arc::new(IndexPublicationGate::new(
|
||||
IndexPublicationPhase::AfterManifestCommit,
|
||||
));
|
||||
let engine = Arc::new(
|
||||
env.create_engine_with(config, None, Some(gate.clone()), None)
|
||||
.await,
|
||||
);
|
||||
|
||||
let region_id = RegionId::new(1, 1);
|
||||
env.get_schema_metadata_manager()
|
||||
.register_region_table_info(
|
||||
region_id.table_id(),
|
||||
"test_table",
|
||||
"test_catalog",
|
||||
"test_schema",
|
||||
None,
|
||||
env.get_kv_backend(),
|
||||
)
|
||||
.await;
|
||||
|
||||
let request = CreateRequestBuilder::new().build();
|
||||
let column_schemas = rows_schema(&request);
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Create(request))
|
||||
.await
|
||||
.unwrap();
|
||||
put_and_flush(&engine, region_id, &column_schemas, 0..20).await;
|
||||
|
||||
engine
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Alter(RegionAlterRequest {
|
||||
kind: AlterKind::SetIndexes {
|
||||
options: vec![SetIndexOption::Inverted {
|
||||
column_name: "tag_0".to_string(),
|
||||
}],
|
||||
},
|
||||
}),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
tokio::time::timeout(Duration::from_secs(10), gate.wait_entered(1))
|
||||
.await
|
||||
.expect("first index generation did not commit");
|
||||
|
||||
// The first publication is committed to the manifest but its worker
|
||||
// notification is blocked, so version control still has the previous file
|
||||
// metadata when the next schema generation is created.
|
||||
engine
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Alter(RegionAlterRequest {
|
||||
kind: AlterKind::SetIndexes {
|
||||
options: vec![SetIndexOption::Inverted {
|
||||
column_name: "field_0".to_string(),
|
||||
}],
|
||||
},
|
||||
}),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
gate.release(1);
|
||||
tokio::time::timeout(Duration::from_secs(10), gate.wait_entered(2))
|
||||
.await
|
||||
.expect("latest schema generation did not commit");
|
||||
gate.release(1);
|
||||
tokio::time::timeout(Duration::from_secs(10), gate.wait_stopped(2))
|
||||
.await
|
||||
.expect("index builds did not stop");
|
||||
|
||||
let region = engine.get_region(region_id).unwrap();
|
||||
let version = region.version();
|
||||
let files = current_file_metas(&engine, region_id).await;
|
||||
assert_eq!(files.len(), 1);
|
||||
assert!(
|
||||
files[0].is_index_consistent_with_region(&version.metadata.column_metadatas),
|
||||
"the latest index definition was not published"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_index_build_type_manual_basic() {
|
||||
let mut env = TestEnv::with_prefix("test_index_build_type_manual_").await;
|
||||
|
||||
@@ -198,6 +198,15 @@ pub struct RegionManifestBuilder {
|
||||
}
|
||||
|
||||
impl RegionManifestBuilder {
|
||||
fn removed_file(&self, file: &FileMeta) -> RemovedFile {
|
||||
let index_version = self
|
||||
.files
|
||||
.get(&file.file_id)
|
||||
.and_then(FileMeta::index_version)
|
||||
.or_else(|| file.index_version());
|
||||
RemovedFile::File(file.file_id, index_version)
|
||||
}
|
||||
|
||||
/// Start with a checkpoint.
|
||||
pub fn with_checkpoint(checkpoint: Option<RegionManifest>) -> Self {
|
||||
if let Some(s) = checkpoint {
|
||||
@@ -263,14 +272,11 @@ impl RegionManifestBuilder {
|
||||
removed_files.push(RemovedFile::Index(old_file.file_id, old_index));
|
||||
}
|
||||
}
|
||||
removed_files.extend(edit.files_to_remove.iter().map(|file| {
|
||||
let index_version = self
|
||||
.files
|
||||
.get(&file.file_id)
|
||||
.and_then(FileMeta::index_version)
|
||||
.or_else(|| file.index_version());
|
||||
RemovedFile::File(file.file_id, index_version)
|
||||
}));
|
||||
removed_files.extend(
|
||||
edit.files_to_remove
|
||||
.iter()
|
||||
.map(|file| self.removed_file(file)),
|
||||
);
|
||||
let at = edit
|
||||
.timestamp_ms
|
||||
.unwrap_or_else(|| Utc::now().timestamp_millis());
|
||||
@@ -319,10 +325,15 @@ impl RegionManifestBuilder {
|
||||
self.files.clear();
|
||||
}
|
||||
TruncateKind::Partial { files_to_remove } => {
|
||||
// With GC disabled, VersionControl may still hold an older
|
||||
// FileMeta and LocalFilePurger can miss a just-committed newer
|
||||
// index generation. We accept this narrow local-mode orphan
|
||||
// window; object-store deployments enable GC and collect the
|
||||
// generation recorded here.
|
||||
self.removed_files.add_removed_files(
|
||||
files_to_remove
|
||||
.iter()
|
||||
.map(|f| RemovedFile::File(f.file_id, f.index_version()))
|
||||
.map(|file| self.removed_file(file))
|
||||
.collect(),
|
||||
truncate
|
||||
.timestamp_ms
|
||||
@@ -939,7 +950,7 @@ mod tests {
|
||||
builder.apply_edit(
|
||||
1,
|
||||
RegionEdit {
|
||||
files_to_add: vec![current],
|
||||
files_to_add: vec![current.clone()],
|
||||
files_to_remove: Vec::new(),
|
||||
timestamp_ms: None,
|
||||
compaction_time_window: None,
|
||||
@@ -952,7 +963,7 @@ mod tests {
|
||||
2,
|
||||
RegionEdit {
|
||||
files_to_add: Vec::new(),
|
||||
files_to_remove: vec![stale],
|
||||
files_to_remove: vec![stale.clone()],
|
||||
timestamp_ms: Some(42),
|
||||
compaction_time_window: None,
|
||||
flushed_entry_id: None,
|
||||
@@ -965,6 +976,35 @@ mod tests {
|
||||
builder.removed_files.removed_files[0].files,
|
||||
HashSet::from([RemovedFile::File(file_id, Some(2))])
|
||||
);
|
||||
|
||||
let mut builder = RegionManifestBuilder::default();
|
||||
builder.apply_edit(
|
||||
1,
|
||||
RegionEdit {
|
||||
files_to_add: vec![current],
|
||||
files_to_remove: Vec::new(),
|
||||
timestamp_ms: None,
|
||||
compaction_time_window: None,
|
||||
flushed_entry_id: None,
|
||||
flushed_sequence: None,
|
||||
committed_sequence: None,
|
||||
},
|
||||
);
|
||||
builder.apply_truncate(
|
||||
2,
|
||||
RegionTruncate {
|
||||
region_id: RegionId::new(1, 1),
|
||||
kind: TruncateKind::Partial {
|
||||
files_to_remove: vec![stale],
|
||||
},
|
||||
timestamp_ms: Some(42),
|
||||
},
|
||||
);
|
||||
|
||||
assert_eq!(
|
||||
builder.removed_files.removed_files[0].files,
|
||||
HashSet::from([RemovedFile::File(file_id, Some(2))])
|
||||
);
|
||||
}
|
||||
|
||||
/// Test if old version can still be deserialized then serialized to the new version.
|
||||
|
||||
+119
-18
@@ -1065,8 +1065,37 @@ pub(crate) enum IndexPublication {
|
||||
manifest_version: ManifestVersion,
|
||||
file_meta: FileMeta,
|
||||
},
|
||||
/// The build no longer matches the current manifest.
|
||||
Stale(IndexPublicationStale),
|
||||
}
|
||||
|
||||
/// Why an index publication became stale.
|
||||
#[derive(Debug, PartialEq, Eq)]
|
||||
pub(crate) enum IndexPublicationStale {
|
||||
/// The source SST or region incarnation is no longer publishable.
|
||||
Stale,
|
||||
SourceChanged,
|
||||
/// The SST still matches, but the schema generation changed.
|
||||
SchemaChanged,
|
||||
}
|
||||
|
||||
/// Manifest state an index build is based on.
|
||||
///
|
||||
/// Both fields must still match when the rebuilt index is published. The file
|
||||
/// metadata identifies the exact SST generation, while the schema version
|
||||
/// identifies the exact index definition used by the builder.
|
||||
#[derive(Clone, Debug, PartialEq, Eq)]
|
||||
pub(crate) struct IndexBuildSource {
|
||||
pub(crate) file_meta: FileMeta,
|
||||
pub(crate) schema_version: u64,
|
||||
}
|
||||
|
||||
impl IndexBuildSource {
|
||||
pub(crate) fn new(file_meta: FileMeta, schema_version: u64) -> Self {
|
||||
Self {
|
||||
file_meta,
|
||||
schema_version,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Context to update the region manifest.
|
||||
@@ -1259,12 +1288,12 @@ impl ManifestContext {
|
||||
|
||||
/// Conditionally publishes rebuilt index metadata for `source`.
|
||||
///
|
||||
/// The source-generation check and manifest update share the manifest write
|
||||
/// lock, so a concurrent compaction or another index build cannot commit
|
||||
/// between them.
|
||||
/// The source SST and schema-generation checks share the manifest write lock
|
||||
/// with the update, so a concurrent compaction, schema change, or another
|
||||
/// index build cannot commit between them.
|
||||
pub(crate) async fn update_manifest_for_index(
|
||||
&self,
|
||||
source: &FileMeta,
|
||||
source: &IndexBuildSource,
|
||||
updated: FileMeta,
|
||||
) -> Result<IndexPublication> {
|
||||
let manager = self.manifest_manager.write().await;
|
||||
@@ -1276,15 +1305,25 @@ impl ManifestContext {
|
||||
| RegionRoleState::Leader(RegionLeaderState::Downgrading)
|
||||
) || manager.is_stopped()
|
||||
{
|
||||
return Ok(IndexPublication::Stale);
|
||||
return Ok(IndexPublication::Stale(
|
||||
IndexPublicationStale::SourceChanged,
|
||||
));
|
||||
}
|
||||
|
||||
if manifest.files.get(&source.file_id) != Some(source) {
|
||||
return Ok(IndexPublication::Stale);
|
||||
if manifest.files.get(&source.file_meta.file_id) != Some(&source.file_meta) {
|
||||
return Ok(IndexPublication::Stale(
|
||||
IndexPublicationStale::SourceChanged,
|
||||
));
|
||||
}
|
||||
|
||||
if manifest.metadata.schema_version != source.schema_version {
|
||||
return Ok(IndexPublication::Stale(
|
||||
IndexPublicationStale::SchemaChanged,
|
||||
));
|
||||
}
|
||||
|
||||
// Only index metadata is allowed to change in this publication.
|
||||
let mut committed = source.clone();
|
||||
let mut committed = source.file_meta.clone();
|
||||
committed.available_indexes = updated.available_indexes;
|
||||
committed.indexes = updated.indexes;
|
||||
committed.index_file_size = updated.index_file_size;
|
||||
@@ -1914,8 +1953,8 @@ mod tests {
|
||||
};
|
||||
use crate::manifest::manager::{RegionManifestManager, RegionManifestOptions};
|
||||
use crate::region::{
|
||||
IndexPublication, ManifestContext, ManifestStats, MitoRegion, RegionLeaderState,
|
||||
RegionRoleState, RegionStats,
|
||||
IndexBuildSource, IndexPublication, IndexPublicationStale, ManifestContext, ManifestStats,
|
||||
MitoRegion, RegionLeaderState, RegionRoleState, RegionStats,
|
||||
};
|
||||
use crate::sst::FormatType;
|
||||
use crate::sst::index::intermediate::IntermediateManager;
|
||||
@@ -2067,22 +2106,24 @@ mod tests {
|
||||
let mut first_update = source.clone();
|
||||
first_update.index_version = 1;
|
||||
first_update.index_file_size = 128;
|
||||
let schema_version = region.version().metadata.schema_version;
|
||||
let initial_source = IndexBuildSource::new(source.clone(), schema_version);
|
||||
let first_committed = match region
|
||||
.manifest_ctx
|
||||
.update_manifest_for_index(&source, first_update)
|
||||
.update_manifest_for_index(&initial_source, first_update)
|
||||
.await
|
||||
.unwrap()
|
||||
{
|
||||
IndexPublication::Committed { file_meta, .. } => file_meta,
|
||||
IndexPublication::Stale => panic!("first index publication should commit"),
|
||||
IndexPublication::Stale(_) => panic!("first index publication should commit"),
|
||||
};
|
||||
|
||||
let (ready_tx, ready_rx) = tokio::sync::oneshot::channel();
|
||||
let (release_tx, release_rx) = tokio::sync::oneshot::channel();
|
||||
let delayed_manifest_ctx = region.manifest_ctx.clone();
|
||||
let delayed_source = first_committed.clone();
|
||||
let delayed_source = IndexBuildSource::new(first_committed.clone(), schema_version);
|
||||
let delayed_publication = tokio::spawn(async move {
|
||||
let mut delayed_update = delayed_source.clone();
|
||||
let mut delayed_update = delayed_source.file_meta.clone();
|
||||
delayed_update.index_version = 2;
|
||||
delayed_update.index_file_size = 128;
|
||||
ready_tx.send(()).unwrap();
|
||||
@@ -2096,20 +2137,21 @@ mod tests {
|
||||
let mut newer_update = first_committed.clone();
|
||||
newer_update.index_version = 3;
|
||||
newer_update.index_file_size = 256;
|
||||
let newer_source = IndexBuildSource::new(first_committed, schema_version);
|
||||
let newer_committed = match region
|
||||
.manifest_ctx
|
||||
.update_manifest_for_index(&first_committed, newer_update)
|
||||
.update_manifest_for_index(&newer_source, newer_update)
|
||||
.await
|
||||
.unwrap()
|
||||
{
|
||||
IndexPublication::Committed { file_meta, .. } => file_meta,
|
||||
IndexPublication::Stale => panic!("newer index publication should commit"),
|
||||
IndexPublication::Stale(_) => panic!("newer index publication should commit"),
|
||||
};
|
||||
|
||||
release_tx.send(()).unwrap();
|
||||
assert!(matches!(
|
||||
delayed_publication.await.unwrap().unwrap(),
|
||||
IndexPublication::Stale
|
||||
IndexPublication::Stale(IndexPublicationStale::SourceChanged)
|
||||
));
|
||||
assert_eq!(
|
||||
region
|
||||
@@ -2122,6 +2164,65 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_index_publication_rejects_older_schema_generation() {
|
||||
let env = SchedulerEnv::new().await;
|
||||
let region = build_test_region(&env).await;
|
||||
let source_meta = crate::sst::file::FileMeta {
|
||||
region_id: region.region_id,
|
||||
file_id: FileId::random(),
|
||||
level: 1,
|
||||
file_size: 1024,
|
||||
..Default::default()
|
||||
};
|
||||
region
|
||||
.manifest_ctx
|
||||
.update_manifest(
|
||||
RegionLeaderState::Writable,
|
||||
RegionMetaActionList::with_action(RegionMetaAction::Edit(RegionEdit {
|
||||
files_to_add: vec![source_meta.clone()],
|
||||
..empty_edit()
|
||||
})),
|
||||
false,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let old_schema_version = region.version().metadata.schema_version;
|
||||
let source = IndexBuildSource::new(source_meta.clone(), old_schema_version);
|
||||
let mut new_metadata = region.version().metadata.as_ref().clone();
|
||||
new_metadata.schema_version += 1;
|
||||
region
|
||||
.manifest_ctx
|
||||
.update_manifest(
|
||||
RegionLeaderState::Writable,
|
||||
RegionMetaActionList::with_action(RegionMetaAction::Change(RegionChange {
|
||||
metadata: Arc::new(new_metadata),
|
||||
sst_format: FormatType::PrimaryKey,
|
||||
append_mode: None,
|
||||
})),
|
||||
false,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let mut updated = source_meta.clone();
|
||||
updated.index_version = 1;
|
||||
updated.index_file_size = 128;
|
||||
assert!(matches!(
|
||||
region
|
||||
.manifest_ctx
|
||||
.update_manifest_for_index(&source, updated)
|
||||
.await
|
||||
.unwrap(),
|
||||
IndexPublication::Stale(IndexPublicationStale::SchemaChanged)
|
||||
));
|
||||
|
||||
let manifest = region.manifest_ctx.manifest().await;
|
||||
assert_eq!(manifest.metadata.schema_version, old_schema_version + 1);
|
||||
assert_eq!(manifest.files.get(&source_meta.file_id), Some(&source_meta));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_exit_staging_partition_expr_change_and_edit_success() {
|
||||
let env = SchedulerEnv::new().await;
|
||||
|
||||
@@ -908,6 +908,8 @@ pub(crate) enum BackgroundNotify {
|
||||
IndexBuildStopped(IndexBuildStopped),
|
||||
/// Index build has failed.
|
||||
IndexBuildFailed(IndexBuildFailed),
|
||||
/// An index build must be retried against the latest schema generation.
|
||||
IndexBuildRetry(BuildIndexRequest),
|
||||
/// Compaction has finished.
|
||||
CompactionFinished(CompactionFinished),
|
||||
/// Compaction has been cancelled cooperatively.
|
||||
|
||||
+369
-92
@@ -24,7 +24,7 @@ pub(crate) mod store;
|
||||
pub(crate) mod vector_index;
|
||||
|
||||
use std::cmp::Ordering;
|
||||
use std::collections::{BinaryHeap, HashMap, HashSet};
|
||||
use std::collections::{HashMap, HashSet};
|
||||
use std::num::NonZeroUsize;
|
||||
use std::sync::Arc;
|
||||
|
||||
@@ -63,14 +63,16 @@ use crate::metrics::{
|
||||
use crate::read::Batch;
|
||||
use crate::region::options::IndexOptions;
|
||||
use crate::region::version::VersionControlRef;
|
||||
use crate::region::{IndexPublication, ManifestContextRef};
|
||||
use crate::region::{
|
||||
IndexBuildSource, IndexPublication, IndexPublicationStale, ManifestContextRef,
|
||||
};
|
||||
use crate::request::{
|
||||
BackgroundNotify, IndexBuildFailed, IndexBuildFinished, IndexBuildStopped, WorkerRequest,
|
||||
WorkerRequestWithTime,
|
||||
BackgroundNotify, BuildIndexRequest, IndexBuildFailed, IndexBuildFinished, IndexBuildStopped,
|
||||
WorkerRequest, WorkerRequestWithTime,
|
||||
};
|
||||
use crate::schedule::scheduler::{Job, SchedulerRef};
|
||||
use crate::sst::file::{
|
||||
ColumnIndexMetadata, FileHandle, FileMeta, IndexType, IndexTypes, RegionFileId, RegionIndexId,
|
||||
ColumnIndexMetadata, FileHandle, IndexType, IndexTypes, RegionFileId, RegionIndexId,
|
||||
};
|
||||
use crate::sst::file_purger::FilePurgerRef;
|
||||
use crate::sst::index::fulltext_index::creator::FulltextIndexer;
|
||||
@@ -90,9 +92,8 @@ pub(crate) const TYPE_VECTOR_INDEX: &str = "vector_index";
|
||||
/// Evicts local state for a stale index artifact.
|
||||
///
|
||||
/// The remote artifact may already be referenced by another region manifest
|
||||
/// because repartitioned regions share physical SST and index paths. Remote
|
||||
/// deletion therefore stays in the normal FilePurger/GC lifecycle; object-store
|
||||
/// GC can account for cross-region references and the configured lingering time.
|
||||
/// because repartitioned regions share physical SST and index paths. With GC
|
||||
/// enabled, remote deletion therefore stays in the reference-aware GC lifecycle.
|
||||
pub(crate) async fn cleanup_stale_index_caches(
|
||||
index_id: RegionIndexId,
|
||||
access_layer: &AccessLayerRef,
|
||||
@@ -705,8 +706,8 @@ pub struct IndexBuildTask {
|
||||
pub region_id: RegionId,
|
||||
/// The SST file handle to build index for.
|
||||
pub file: FileHandle,
|
||||
/// The file meta to build index for.
|
||||
pub file_meta: FileMeta,
|
||||
/// The manifest state this build is based on.
|
||||
pub(crate) source: IndexBuildSource,
|
||||
pub reason: IndexBuildType,
|
||||
pub access_layer: AccessLayerRef,
|
||||
pub(crate) listener: WorkerListener,
|
||||
@@ -727,8 +728,9 @@ impl std::fmt::Debug for IndexBuildTask {
|
||||
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||
f.debug_struct("IndexBuildTask")
|
||||
.field("region_id", &self.region_id)
|
||||
.field("origin_region_id", &self.file_meta.region_id)
|
||||
.field("file_id", &self.file_meta.file_id)
|
||||
.field("origin_region_id", &self.source.file_meta.region_id)
|
||||
.field("file_id", &self.source.file_meta.file_id)
|
||||
.field("schema_version", &self.source.schema_version)
|
||||
.field("reason", &self.reason)
|
||||
.finish()
|
||||
}
|
||||
@@ -759,8 +761,8 @@ impl IndexBuildTask {
|
||||
async fn do_index_build(&mut self, version_control: VersionControlRef) {
|
||||
self.listener
|
||||
.on_index_build_begin(RegionFileId::new(
|
||||
self.file_meta.region_id,
|
||||
self.file_meta.file_id,
|
||||
self.source.file_meta.region_id,
|
||||
self.source.file_meta.file_id,
|
||||
))
|
||||
.await;
|
||||
match self.index_build(version_control).await {
|
||||
@@ -768,7 +770,7 @@ impl IndexBuildTask {
|
||||
Err(e) => {
|
||||
warn!(
|
||||
e; "Index build task failed, region: {}, file_id: {}",
|
||||
self.region_id, self.file_meta.file_id,
|
||||
self.region_id, self.source.file_meta.file_id,
|
||||
);
|
||||
self.on_failure(e.into()).await
|
||||
}
|
||||
@@ -776,7 +778,7 @@ impl IndexBuildTask {
|
||||
let worker_request = WorkerRequest::Background {
|
||||
region_id: self.region_id,
|
||||
notify: BackgroundNotify::IndexBuildStopped(IndexBuildStopped {
|
||||
file_id: self.file_meta.file_id,
|
||||
file_id: self.source.file_meta.file_id,
|
||||
}),
|
||||
};
|
||||
let _ = self
|
||||
@@ -787,8 +789,8 @@ impl IndexBuildTask {
|
||||
|
||||
// Checks if the SST file still exists in object store and version to avoid conflict with compaction.
|
||||
async fn check_sst_file_exists(&self, version_control: &VersionControlRef) -> bool {
|
||||
let file_id = self.file_meta.file_id;
|
||||
let level = self.file_meta.level;
|
||||
let file_id = self.source.file_meta.file_id;
|
||||
let level = self.source.file_meta.level;
|
||||
// We should check current version instead of the version when the job is created.
|
||||
let version = version_control.current().version;
|
||||
|
||||
@@ -822,6 +824,7 @@ impl IndexBuildTask {
|
||||
version_control: VersionControlRef,
|
||||
) -> Result<IndexBuildOutcome> {
|
||||
let new_index_version = self
|
||||
.source
|
||||
.file_meta
|
||||
.index_version()
|
||||
.map_or(0, |version| version + 1);
|
||||
@@ -830,13 +833,13 @@ impl IndexBuildTask {
|
||||
if !self.check_sst_file_exists(&version_control).await {
|
||||
self.listener
|
||||
.on_index_build_abort(RegionFileId::new(
|
||||
self.file_meta.region_id,
|
||||
self.file_meta.file_id,
|
||||
self.source.file_meta.region_id,
|
||||
self.source.file_meta.file_id,
|
||||
))
|
||||
.await;
|
||||
return Ok(IndexBuildOutcome::Aborted(format!(
|
||||
"SST file not found during index build, region: {}, file_id: {}",
|
||||
self.region_id, self.file_meta.file_id
|
||||
self.region_id, self.source.file_meta.file_id
|
||||
)));
|
||||
}
|
||||
|
||||
@@ -856,7 +859,10 @@ impl IndexBuildTask {
|
||||
});
|
||||
|
||||
// Use the same file_id but with new version for index file.
|
||||
let region_file_id = RegionFileId::new(self.file_meta.region_id, self.file_meta.file_id);
|
||||
let region_file_id = RegionFileId::new(
|
||||
self.source.file_meta.region_id,
|
||||
self.source.file_meta.file_id,
|
||||
);
|
||||
let index_file_id = region_file_id.file_id();
|
||||
let mut indexer = self
|
||||
.indexer_builder
|
||||
@@ -887,13 +893,13 @@ impl IndexBuildTask {
|
||||
indexer.abort().await;
|
||||
self.listener
|
||||
.on_index_build_abort(RegionFileId::new(
|
||||
self.file_meta.region_id,
|
||||
self.file_meta.file_id,
|
||||
self.source.file_meta.region_id,
|
||||
self.source.file_meta.file_id,
|
||||
))
|
||||
.await;
|
||||
return Ok(IndexBuildOutcome::Aborted(format!(
|
||||
"SST file not found during index build, region: {}, file_id: {}",
|
||||
self.region_id, self.file_meta.file_id
|
||||
self.region_id, self.source.file_meta.file_id
|
||||
)));
|
||||
}
|
||||
|
||||
@@ -922,11 +928,14 @@ impl IndexBuildTask {
|
||||
notify: BackgroundNotify::IndexBuildFinished(index_build_finished),
|
||||
}
|
||||
}
|
||||
Ok(IndexPublication::Stale) => {
|
||||
Ok(IndexPublication::Stale(stale)) => {
|
||||
INDEX_PUBLICATION_STALE_TOTAL
|
||||
.with_label_values(&["manifest_commit"])
|
||||
.inc();
|
||||
let index_id = RegionIndexId::new(region_file_id, new_index_version);
|
||||
// If no successor publishes the same version, GC collects the
|
||||
// remote artifact. Repartition requires GC, while local mode
|
||||
// accepts this narrow orphan window when GC is disabled.
|
||||
cleanup_stale_index_caches(
|
||||
index_id,
|
||||
&self.access_layer,
|
||||
@@ -935,9 +944,28 @@ impl IndexBuildTask {
|
||||
)
|
||||
.await;
|
||||
self.listener.on_index_build_abort(region_file_id).await;
|
||||
if stale == IndexPublicationStale::SchemaChanged {
|
||||
// An index-unrelated schema change also invalidates the
|
||||
// generation fence. Retry to avoid leaving the SST
|
||||
// unindexed; per-SST coalescing limits this to one active
|
||||
// and one pending task, so repeated changes cannot storm.
|
||||
let retry = BuildIndexRequest {
|
||||
region_id: self.region_id,
|
||||
build_type: IndexBuildType::SchemaChange,
|
||||
file_metas: vec![self.source.file_meta.clone()],
|
||||
};
|
||||
let worker_request = WorkerRequest::Background {
|
||||
region_id: self.region_id,
|
||||
notify: BackgroundNotify::IndexBuildRetry(retry),
|
||||
};
|
||||
let _ = self
|
||||
.request_sender
|
||||
.send(WorkerRequestWithTime::new(worker_request))
|
||||
.await;
|
||||
}
|
||||
return Ok(IndexBuildOutcome::Aborted(format!(
|
||||
"Source SST or region incarnation changed before index publication, region: {}, file_id: {}",
|
||||
self.region_id, self.file_meta.file_id
|
||||
"Index build source changed before publication, region: {}, file_id: {}",
|
||||
self.region_id, self.source.file_meta.file_id
|
||||
)));
|
||||
}
|
||||
Err(e) => {
|
||||
@@ -964,8 +992,8 @@ impl IndexBuildTask {
|
||||
index_version: u64,
|
||||
) -> Result<()> {
|
||||
if let Some(write_cache) = &self.write_cache {
|
||||
let file_id = self.file_meta.file_id;
|
||||
let region_id = self.file_meta.region_id;
|
||||
let file_id = self.source.file_meta.file_id;
|
||||
let region_id = self.source.file_meta.region_id;
|
||||
let remote_store = self.access_layer.object_store();
|
||||
let mut upload_tracker = UploadTracker::new(region_id);
|
||||
let mut err = None;
|
||||
@@ -1010,14 +1038,14 @@ impl IndexBuildTask {
|
||||
output: IndexOutput,
|
||||
new_index_version: u64,
|
||||
) -> Result<IndexPublication> {
|
||||
let mut updated = self.file_meta.clone();
|
||||
let mut updated = self.source.file_meta.clone();
|
||||
updated.available_indexes = output.build_available_indexes();
|
||||
updated.indexes = output.build_indexes();
|
||||
updated.index_file_size = output.file_size;
|
||||
updated.index_version = new_index_version;
|
||||
let publication = self
|
||||
.manifest_ctx
|
||||
.update_manifest_for_index(&self.file_meta, updated)
|
||||
.update_manifest_for_index(&self.source, updated)
|
||||
.await?;
|
||||
if let IndexPublication::Committed {
|
||||
manifest_version, ..
|
||||
@@ -1080,10 +1108,86 @@ impl Ord for PendingIndexBuild {
|
||||
}
|
||||
}
|
||||
|
||||
impl PendingIndexBuild {
|
||||
/// Returns whether this task should replace another pending task for the
|
||||
/// same SST.
|
||||
fn supersedes(&self, other: &Self) -> bool {
|
||||
match self
|
||||
.task
|
||||
.source
|
||||
.schema_version
|
||||
.cmp(&other.task.source.schema_version)
|
||||
{
|
||||
Ordering::Greater => true,
|
||||
Ordering::Less => false,
|
||||
Ordering::Equal => match self
|
||||
.task
|
||||
.source
|
||||
.file_meta
|
||||
.index_version()
|
||||
.cmp(&other.task.source.file_meta.index_version())
|
||||
{
|
||||
Ordering::Greater => true,
|
||||
Ordering::Less => false,
|
||||
Ordering::Equal => self.task.reason.priority() > other.task.reason.priority(),
|
||||
},
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Default)]
|
||||
struct PendingIndexBuilds {
|
||||
tasks: HashMap<FileId, PendingIndexBuild>,
|
||||
}
|
||||
|
||||
impl PendingIndexBuilds {
|
||||
fn insert(&mut self, pending: PendingIndexBuild) -> Option<IndexBuildTask> {
|
||||
let file_id = pending.task.source.file_meta.file_id;
|
||||
match self.tasks.entry(file_id) {
|
||||
std::collections::hash_map::Entry::Vacant(entry) => {
|
||||
entry.insert(pending);
|
||||
None
|
||||
}
|
||||
std::collections::hash_map::Entry::Occupied(mut entry)
|
||||
if pending.supersedes(entry.get()) =>
|
||||
{
|
||||
Some(entry.insert(pending).task)
|
||||
}
|
||||
std::collections::hash_map::Entry::Occupied(_) => Some(pending.task),
|
||||
}
|
||||
}
|
||||
|
||||
fn highest_ready(&self, building_files: &HashSet<FileId>) -> Option<PendingIndexBuild> {
|
||||
self.tasks
|
||||
.values()
|
||||
.filter(|pending| !building_files.contains(&pending.task.source.file_meta.file_id))
|
||||
.max()
|
||||
.cloned()
|
||||
}
|
||||
|
||||
fn remove(&mut self, file_id: FileId) {
|
||||
self.tasks.remove(&file_id);
|
||||
}
|
||||
|
||||
fn drain(&mut self) -> impl Iterator<Item = PendingIndexBuild> {
|
||||
std::mem::take(&mut self.tasks).into_values()
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
fn len(&self) -> usize {
|
||||
self.tasks.len()
|
||||
}
|
||||
|
||||
fn is_empty(&self) -> bool {
|
||||
self.tasks.is_empty()
|
||||
}
|
||||
}
|
||||
|
||||
/// Tracks the index build status of a region scheduled by the [IndexBuildScheduler].
|
||||
struct IndexBuildStatus {
|
||||
building_files: HashSet<FileId>,
|
||||
pending_tasks: BinaryHeap<PendingIndexBuild>,
|
||||
/// At most one coalesced pending task is kept for each SST.
|
||||
pending_tasks: PendingIndexBuilds,
|
||||
/// Whether the active builds belong to a region incarnation that has
|
||||
/// stopped accepting index publications.
|
||||
retiring: bool,
|
||||
@@ -1093,13 +1197,13 @@ impl IndexBuildStatus {
|
||||
fn new() -> Self {
|
||||
IndexBuildStatus {
|
||||
building_files: HashSet::new(),
|
||||
pending_tasks: BinaryHeap::new(),
|
||||
pending_tasks: PendingIndexBuilds::default(),
|
||||
retiring: false,
|
||||
}
|
||||
}
|
||||
|
||||
async fn fail_pending(&mut self, err: Arc<Error>) {
|
||||
for pending in std::mem::take(&mut self.pending_tasks) {
|
||||
for pending in self.pending_tasks.drain() {
|
||||
pending.task.on_failure(err.clone()).await;
|
||||
}
|
||||
}
|
||||
@@ -1139,36 +1243,56 @@ impl IndexBuildScheduler {
|
||||
.entry(task.region_id)
|
||||
.or_insert_with(IndexBuildStatus::new);
|
||||
|
||||
let duplicate = if status.retiring {
|
||||
status
|
||||
.pending_tasks
|
||||
.iter()
|
||||
.any(|pending| pending.task.file_meta.file_id == task.file_meta.file_id)
|
||||
} else {
|
||||
status.building_files.contains(&task.file_meta.file_id)
|
||||
};
|
||||
if duplicate {
|
||||
let region_file_id =
|
||||
RegionFileId::new(task.file_meta.region_id, task.file_meta.file_id);
|
||||
let file_id = task.source.file_meta.file_id;
|
||||
let region_file_id = RegionFileId::new(
|
||||
task.source.file_meta.region_id,
|
||||
task.source.file_meta.file_id,
|
||||
);
|
||||
let can_wait_for_active = status.retiring || task.reason == IndexBuildType::SchemaChange;
|
||||
let (rejected, coalesced) =
|
||||
if status.building_files.contains(&file_id) && !can_wait_for_active {
|
||||
(Some(task), None)
|
||||
} else {
|
||||
let pending = PendingIndexBuild {
|
||||
task,
|
||||
version_control: version_control.clone(),
|
||||
};
|
||||
(None, status.pending_tasks.insert(pending))
|
||||
};
|
||||
let should_schedule = !status.retiring;
|
||||
|
||||
if let Some(rejected) = rejected {
|
||||
debug!(
|
||||
"Aborting index build task since index is already being built for region file {:?}",
|
||||
"Rejecting index build because region file {:?} is already being built",
|
||||
region_file_id
|
||||
);
|
||||
task.on_success(IndexBuildOutcome::Aborted(format!(
|
||||
"Index is already being built for region file {:?}",
|
||||
region_file_id
|
||||
)))
|
||||
.await;
|
||||
task.listener.on_index_build_abort(region_file_id).await;
|
||||
return Ok(());
|
||||
rejected
|
||||
.on_success(IndexBuildOutcome::Aborted(format!(
|
||||
"Index is already being built for region file {:?}",
|
||||
region_file_id
|
||||
)))
|
||||
.await;
|
||||
rejected.listener.on_index_build_abort(region_file_id).await;
|
||||
}
|
||||
|
||||
status.pending_tasks.push(PendingIndexBuild {
|
||||
task,
|
||||
version_control: version_control.clone(),
|
||||
});
|
||||
if let Some(coalesced) = coalesced {
|
||||
debug!(
|
||||
"Coalescing redundant index build task for region file {:?}",
|
||||
region_file_id
|
||||
);
|
||||
coalesced
|
||||
.on_success(IndexBuildOutcome::Aborted(format!(
|
||||
"Index build was coalesced for region file {:?}",
|
||||
region_file_id
|
||||
)))
|
||||
.await;
|
||||
coalesced
|
||||
.listener
|
||||
.on_index_build_abort(region_file_id)
|
||||
.await;
|
||||
}
|
||||
|
||||
if !status.retiring {
|
||||
if should_schedule {
|
||||
self.schedule_next_build_batch();
|
||||
}
|
||||
Ok(())
|
||||
@@ -1176,52 +1300,63 @@ impl IndexBuildScheduler {
|
||||
|
||||
/// Schedule tasks until reaching the files limit or no more tasks.
|
||||
fn schedule_next_build_batch(&mut self) {
|
||||
let mut building_count = 0;
|
||||
for status in self.region_status.values() {
|
||||
building_count += status.building_files.len();
|
||||
}
|
||||
let mut building_count = self
|
||||
.region_status
|
||||
.values()
|
||||
.map(|status| status.building_files.len())
|
||||
.sum::<usize>();
|
||||
|
||||
while building_count < self.files_limit {
|
||||
if let Some(pending) = self.find_next_task() {
|
||||
let task = pending.task;
|
||||
let region_id = task.region_id;
|
||||
let file_id = task.file_meta.file_id;
|
||||
let job = task.into_index_build_job(pending.version_control);
|
||||
if self.scheduler.schedule(job).is_ok() {
|
||||
let Some(pending) = self.find_next_task() else {
|
||||
break;
|
||||
};
|
||||
|
||||
let task = pending.task;
|
||||
let region_id = task.region_id;
|
||||
let file_id = task.source.file_meta.file_id;
|
||||
let job = task.clone().into_index_build_job(pending.version_control);
|
||||
match self.scheduler.schedule(job) {
|
||||
Ok(()) => {
|
||||
if let Some(status) = self.region_status.get_mut(®ion_id) {
|
||||
status.pending_tasks.remove(file_id);
|
||||
status.building_files.insert(file_id);
|
||||
building_count += 1;
|
||||
status
|
||||
.pending_tasks
|
||||
.retain(|pending| pending.task.file_meta.file_id != file_id);
|
||||
} else {
|
||||
error!(
|
||||
"Region status not found when scheduling index build task, region: {}",
|
||||
region_id
|
||||
);
|
||||
}
|
||||
} else {
|
||||
error!(
|
||||
"Failed to schedule index build job, region: {}, file_id: {}",
|
||||
region_id, file_id
|
||||
);
|
||||
break;
|
||||
}
|
||||
} else {
|
||||
// No more tasks to schedule.
|
||||
break;
|
||||
Err(err) => {
|
||||
error!(
|
||||
err;
|
||||
"Failed to schedule index build job, region: {}, file_id: {}",
|
||||
region_id,
|
||||
file_id
|
||||
);
|
||||
if let Some(status) = self.region_status.get_mut(®ion_id) {
|
||||
status.pending_tasks.remove(file_id);
|
||||
}
|
||||
common_runtime::spawn_global(async move {
|
||||
task.on_failure(Arc::new(err)).await;
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
self.region_status.retain(|_, status| {
|
||||
!status.building_files.is_empty() || !status.pending_tasks.is_empty()
|
||||
});
|
||||
}
|
||||
|
||||
/// Find the next task which has the highest priority to run.
|
||||
/// Find the ready task with the highest priority.
|
||||
fn find_next_task(&self) -> Option<PendingIndexBuild> {
|
||||
self.region_status
|
||||
.values()
|
||||
.filter(|status| !status.retiring)
|
||||
.filter_map(|status| status.pending_tasks.peek())
|
||||
.filter_map(|status| status.pending_tasks.highest_ready(&status.building_files))
|
||||
.max()
|
||||
.cloned()
|
||||
}
|
||||
|
||||
pub(crate) fn on_task_stopped(&mut self, region_id: RegionId, file_id: FileId) {
|
||||
@@ -1373,12 +1508,13 @@ mod tests {
|
||||
use crate::memtable::time_partition::TimePartitions;
|
||||
use crate::region::RegionLeaderState;
|
||||
use crate::region::version::{VersionBuilder, VersionControl};
|
||||
use crate::sst::file::RegionFileId;
|
||||
use crate::schedule::scheduler::{LocalScheduler, Scheduler};
|
||||
use crate::sst::file::{FileMeta, RegionFileId};
|
||||
use crate::sst::file_purger::NoopFilePurger;
|
||||
use crate::sst::location;
|
||||
use crate::sst::parquet::WriteOptions;
|
||||
use crate::test_util::memtable_util::EmptyMemtableBuilder;
|
||||
use crate::test_util::scheduler_util::SchedulerEnv;
|
||||
use crate::test_util::scheduler_util::{SchedulerEnv, VecScheduler};
|
||||
use crate::test_util::sst_util::{
|
||||
new_flat_source_from_record_batches, new_record_batch_by_range, sst_region_metadata,
|
||||
};
|
||||
@@ -1952,7 +2088,10 @@ mod tests {
|
||||
let task = IndexBuildTask {
|
||||
region_id,
|
||||
file,
|
||||
file_meta,
|
||||
source: IndexBuildSource::new(
|
||||
file_meta,
|
||||
version_control.current().version.metadata.schema_version,
|
||||
),
|
||||
reason: IndexBuildType::Flush,
|
||||
access_layer: env.access_layer.clone(),
|
||||
listener: WorkerListener::default(),
|
||||
@@ -2016,7 +2155,10 @@ mod tests {
|
||||
let task = IndexBuildTask {
|
||||
region_id,
|
||||
file,
|
||||
file_meta: file_meta.clone(),
|
||||
source: IndexBuildSource::new(
|
||||
file_meta.clone(),
|
||||
version_control.current().version.metadata.schema_version,
|
||||
),
|
||||
reason: IndexBuildType::Flush,
|
||||
access_layer: env.access_layer.clone(),
|
||||
listener: WorkerListener::default(),
|
||||
@@ -2094,7 +2236,10 @@ mod tests {
|
||||
let task = IndexBuildTask {
|
||||
region_id,
|
||||
file,
|
||||
file_meta: file_meta.clone(),
|
||||
source: IndexBuildSource::new(
|
||||
file_meta.clone(),
|
||||
version_control.current().version.metadata.schema_version,
|
||||
),
|
||||
reason: IndexBuildType::Flush,
|
||||
access_layer: env.access_layer.clone(),
|
||||
listener: WorkerListener::default(),
|
||||
@@ -2201,7 +2346,10 @@ mod tests {
|
||||
let task = IndexBuildTask {
|
||||
region_id,
|
||||
file,
|
||||
file_meta: file_meta.clone(),
|
||||
source: IndexBuildSource::new(
|
||||
file_meta.clone(),
|
||||
version_control.current().version.metadata.schema_version,
|
||||
),
|
||||
reason: IndexBuildType::Flush,
|
||||
access_layer: env.access_layer.clone(),
|
||||
listener: WorkerListener::default(),
|
||||
@@ -2296,7 +2444,10 @@ mod tests {
|
||||
let task = IndexBuildTask {
|
||||
region_id,
|
||||
file,
|
||||
file_meta: file_meta.clone(),
|
||||
source: IndexBuildSource::new(
|
||||
file_meta.clone(),
|
||||
version_control.current().version.metadata.schema_version,
|
||||
),
|
||||
reason: IndexBuildType::Flush,
|
||||
access_layer: env.access_layer.clone(),
|
||||
listener: WorkerListener::default(),
|
||||
@@ -2351,6 +2502,7 @@ mod tests {
|
||||
reason: IndexBuildType,
|
||||
) -> (IndexBuildTask, mpsc::Receiver<Result<IndexBuildOutcome>>) {
|
||||
let metadata = Arc::new(sst_region_metadata());
|
||||
let schema_version = metadata.schema_version;
|
||||
let manifest_ctx = env.mock_manifest_context(metadata.clone()).await;
|
||||
let file_purger = Arc::new(NoopFilePurger {});
|
||||
let indexer_builder = mock_indexer_builder(metadata, env).await;
|
||||
@@ -2369,7 +2521,7 @@ mod tests {
|
||||
let task = IndexBuildTask {
|
||||
region_id,
|
||||
file,
|
||||
file_meta,
|
||||
source: IndexBuildSource::new(file_meta, schema_version),
|
||||
reason,
|
||||
access_layer: env.access_layer.clone(),
|
||||
listener: WorkerListener::default(),
|
||||
@@ -2384,6 +2536,131 @@ mod tests {
|
||||
(task, result_rx)
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_scheduler_coalesces_latest_schema_generation_per_sst() {
|
||||
let job_scheduler = Arc::new(VecScheduler::default());
|
||||
let env = SchedulerEnv::new().await.scheduler(job_scheduler.clone());
|
||||
let mut scheduler = env.mock_index_build_scheduler(2);
|
||||
let metadata = Arc::new(sst_region_metadata());
|
||||
let region_id = metadata.region_id;
|
||||
let file_id = FileId::random();
|
||||
let file_purger = Arc::new(NoopFilePurger {});
|
||||
let files = HashMap::from([(
|
||||
file_id,
|
||||
FileMeta {
|
||||
region_id,
|
||||
file_id,
|
||||
file_size: 100,
|
||||
..Default::default()
|
||||
},
|
||||
)]);
|
||||
let version_control = mock_version_control(metadata, file_purger, files).await;
|
||||
|
||||
let (active, _active_rx) = create_mock_task_for_schedule_with_result(
|
||||
&env,
|
||||
file_id,
|
||||
region_id,
|
||||
IndexBuildType::Manual,
|
||||
)
|
||||
.await;
|
||||
scheduler
|
||||
.schedule_build(&version_control, active)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(job_scheduler.num_jobs(), 1);
|
||||
|
||||
let (mut older, mut older_rx) = create_mock_task_for_schedule_with_result(
|
||||
&env,
|
||||
file_id,
|
||||
region_id,
|
||||
IndexBuildType::SchemaChange,
|
||||
)
|
||||
.await;
|
||||
older.source.schema_version = 1;
|
||||
scheduler
|
||||
.schedule_build(&version_control, older)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let (mut latest, mut latest_rx) = create_mock_task_for_schedule_with_result(
|
||||
&env,
|
||||
file_id,
|
||||
region_id,
|
||||
IndexBuildType::SchemaChange,
|
||||
)
|
||||
.await;
|
||||
latest.source.schema_version = 2;
|
||||
scheduler
|
||||
.schedule_build(&version_control, latest)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let replaced = tokio::time::timeout(std::time::Duration::from_secs(5), older_rx.recv())
|
||||
.await
|
||||
.expect("replaced pending task result sender was not completed")
|
||||
.expect("replaced pending task result channel closed");
|
||||
assert!(matches!(
|
||||
replaced,
|
||||
Ok(IndexBuildOutcome::Aborted(reason)) if reason.contains("coalesced")
|
||||
));
|
||||
assert!(matches!(
|
||||
latest_rx.try_recv(),
|
||||
Err(mpsc::error::TryRecvError::Empty)
|
||||
));
|
||||
|
||||
let status = &scheduler.region_status[®ion_id];
|
||||
assert_eq!(status.building_files.len(), 1);
|
||||
assert_eq!(status.pending_tasks.len(), 1);
|
||||
assert_eq!(job_scheduler.num_jobs(), 1);
|
||||
|
||||
scheduler.on_task_stopped(region_id, file_id);
|
||||
let status = &scheduler.region_status[®ion_id];
|
||||
assert_eq!(status.building_files.len(), 1);
|
||||
assert!(status.pending_tasks.is_empty());
|
||||
assert_eq!(job_scheduler.num_jobs(), 2);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_scheduler_completes_sender_when_job_is_rejected() {
|
||||
let job_scheduler = Arc::new(LocalScheduler::new(1));
|
||||
job_scheduler.stop(false).await.unwrap();
|
||||
let env = SchedulerEnv::new().await.scheduler(job_scheduler);
|
||||
let mut scheduler = env.mock_index_build_scheduler(1);
|
||||
let metadata = Arc::new(sst_region_metadata());
|
||||
let region_id = metadata.region_id;
|
||||
let file_id = FileId::random();
|
||||
let file_purger = Arc::new(NoopFilePurger {});
|
||||
let files = HashMap::from([(
|
||||
file_id,
|
||||
FileMeta {
|
||||
region_id,
|
||||
file_id,
|
||||
file_size: 100,
|
||||
..Default::default()
|
||||
},
|
||||
)]);
|
||||
let version_control = mock_version_control(metadata, file_purger, files).await;
|
||||
let (task, mut result_rx) = create_mock_task_for_schedule_with_result(
|
||||
&env,
|
||||
file_id,
|
||||
region_id,
|
||||
IndexBuildType::Flush,
|
||||
)
|
||||
.await;
|
||||
|
||||
scheduler
|
||||
.schedule_build(&version_control, task)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let result = tokio::time::timeout(std::time::Duration::from_secs(5), result_rx.recv())
|
||||
.await
|
||||
.expect("scheduler rejection did not complete the result sender")
|
||||
.expect("result channel closed without a result");
|
||||
assert!(result.is_err());
|
||||
assert!(!scheduler.region_status.contains_key(®ion_id));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_scheduler_comprehensive() {
|
||||
let env = SchedulerEnv::new().await;
|
||||
|
||||
@@ -73,8 +73,8 @@ use crate::region::{
|
||||
RegionMapRef,
|
||||
};
|
||||
use crate::request::{
|
||||
BackgroundNotify, BulkInsertRequest, DdlRequest, SenderBulkRequest, SenderDdlRequest,
|
||||
SenderWriteRequest, WorkerRequest, WorkerRequestWithTime,
|
||||
BackgroundNotify, BulkInsertRequest, DdlRequest, OptionOutputTx, SenderBulkRequest,
|
||||
SenderDdlRequest, SenderWriteRequest, WorkerRequest, WorkerRequestWithTime,
|
||||
};
|
||||
use crate::schedule::scheduler::{LocalScheduler, SchedulerRef};
|
||||
use crate::sst::file::RegionFileId;
|
||||
@@ -1251,6 +1251,10 @@ impl<S: LogStore> RegionWorkerLoop<S> {
|
||||
BackgroundNotify::IndexBuildFailed(req) => {
|
||||
self.handle_index_build_failed(region_id, req).await
|
||||
}
|
||||
BackgroundNotify::IndexBuildRetry(req) => {
|
||||
self.handle_rebuild_index(req, OptionOutputTx::new(None))
|
||||
.await
|
||||
}
|
||||
BackgroundNotify::CompactionFinished(req) => {
|
||||
self.handle_compaction_finished(region_id, req).await
|
||||
}
|
||||
|
||||
@@ -26,11 +26,12 @@ use crate::cache::CacheStrategy;
|
||||
use crate::error::Result;
|
||||
use crate::manifest::action::RegionEdit;
|
||||
use crate::metrics::INDEX_PUBLICATION_STALE_TOTAL;
|
||||
use crate::region::MitoRegionRef;
|
||||
use crate::region::version::VersionRef;
|
||||
use crate::region::{IndexBuildSource, MitoRegionRef};
|
||||
use crate::request::{
|
||||
BuildIndexRequest, IndexBuildFailed, IndexBuildFinished, IndexBuildStopped, OptionOutputTx,
|
||||
};
|
||||
use crate::sst::file::{FileHandle, RegionFileId, RegionIndexId};
|
||||
use crate::sst::file::{FileHandle, FileMeta, RegionFileId, RegionIndexId};
|
||||
use crate::sst::index::{
|
||||
IndexBuildOutcome, IndexBuildTask, IndexBuildType, IndexerBuilderImpl, ResultMpscSender,
|
||||
cleanup_stale_index_caches,
|
||||
@@ -41,11 +42,12 @@ impl<S> RegionWorkerLoop<S> {
|
||||
pub(crate) fn new_index_build_task(
|
||||
&self,
|
||||
region: &MitoRegionRef,
|
||||
version: &VersionRef,
|
||||
file: FileHandle,
|
||||
file_meta: FileMeta,
|
||||
build_type: IndexBuildType,
|
||||
result_sender: ResultMpscSender,
|
||||
) -> IndexBuildTask {
|
||||
let version = region.version();
|
||||
let access_layer = region.access_layer.clone();
|
||||
|
||||
let puffin_manager = if let Some(write_cache) = self.cache_manager.write_cache() {
|
||||
@@ -77,7 +79,7 @@ impl<S> RegionWorkerLoop<S> {
|
||||
IndexBuildTask {
|
||||
region_id: region.region_id,
|
||||
file: file.clone(),
|
||||
file_meta: file.meta_ref().clone(),
|
||||
source: IndexBuildSource::new(file_meta, version.metadata.schema_version),
|
||||
reason: build_type,
|
||||
access_layer: access_layer.clone(),
|
||||
listener: self.listener.clone(),
|
||||
@@ -121,6 +123,20 @@ impl<S> RegionWorkerLoop<S> {
|
||||
|
||||
let version_control = region.version_control.clone();
|
||||
let version = version_control.current().version;
|
||||
// A committed index publication may not have reached version control
|
||||
// yet. Use the manifest's FileMeta as the conditional publication
|
||||
// source while keeping the builder's schema generation from `version`.
|
||||
let manifest = region.manifest_ctx.manifest().await;
|
||||
let current_file_meta = |file_id| {
|
||||
let file_meta = manifest.files.get(&file_id);
|
||||
if file_meta.is_none() {
|
||||
debug!(
|
||||
"Skipping index build because file is absent from manifest, region: {}, file_id: {}",
|
||||
region_id, file_id
|
||||
);
|
||||
}
|
||||
file_meta.cloned()
|
||||
};
|
||||
|
||||
let all_files: HashMap<FileId, FileHandle> = version
|
||||
.ssts
|
||||
@@ -135,18 +151,21 @@ impl<S> RegionWorkerLoop<S> {
|
||||
// If no specific files are provided, find files whose index is inconsistent with the region metadata.
|
||||
all_files
|
||||
.values()
|
||||
.filter(|file| {
|
||||
!file
|
||||
.meta_ref()
|
||||
.is_index_consistent_with_region(&version.metadata.column_metadatas)
|
||||
.filter_map(|file| {
|
||||
let file_meta = current_file_meta(file.meta_ref().file_id)?;
|
||||
(!file_meta.is_index_consistent_with_region(&version.metadata.column_metadatas))
|
||||
.then(|| (file.clone(), file_meta))
|
||||
})
|
||||
.cloned()
|
||||
.collect::<Vec<_>>()
|
||||
} else {
|
||||
request
|
||||
.file_metas
|
||||
.iter()
|
||||
.filter_map(|meta| all_files.get(&meta.file_id).cloned())
|
||||
.filter_map(|meta| {
|
||||
let file = all_files.get(&meta.file_id)?;
|
||||
let file_meta = current_file_meta(meta.file_id)?;
|
||||
Some((file.clone(), file_meta))
|
||||
})
|
||||
.collect::<Vec<_>>()
|
||||
};
|
||||
|
||||
@@ -162,7 +181,7 @@ impl<S> RegionWorkerLoop<S> {
|
||||
let num_tasks = build_tasks.len();
|
||||
let (tx, mut rx) = mpsc::channel::<Result<IndexBuildOutcome>>(num_tasks);
|
||||
|
||||
for file_handle in build_tasks {
|
||||
for (file_handle, file_meta) in build_tasks {
|
||||
debug!(
|
||||
"Scheduling index build for region {}, file_id {}",
|
||||
region_id,
|
||||
@@ -181,7 +200,9 @@ impl<S> RegionWorkerLoop<S> {
|
||||
|
||||
let task = self.new_index_build_task(
|
||||
®ion,
|
||||
&version,
|
||||
file_handle.clone(),
|
||||
file_meta,
|
||||
request.build_type.clone(),
|
||||
tx.clone(),
|
||||
);
|
||||
|
||||
Reference in New Issue
Block a user