mirror of
https://github.com/GreptimeTeam/greptimedb.git
synced 2026-10-03 02:25:35 +00:00
fix: pause GC during maintenance mode (#8450)
Skip scheduled meta GC while cluster maintenance mode is enabled and reject manual GC requests explicitly instead of returning an empty success report. Also increase mito GC's default lingering time to 1h and update generated config docs and config API expectations. Signed-off-by: discord9 <discord9@163.com>
This commit is contained in:
+1
-1
@@ -608,7 +608,7 @@
|
||||
| `region_engine.mito.bloom_filter_index.mem_threshold_on_create` | String | `auto` | Memory threshold for the index creation.<br/>- `auto`: automatically determine the threshold based on the system memory size (default)<br/>- `unlimited`: no memory limit<br/>- `[size]` e.g. `64MB`: fixed memory threshold |
|
||||
| `region_engine.mito.gc` | -- | -- | -- |
|
||||
| `region_engine.mito.gc.enable` | Bool | `false` | Whether GC is enabled. Need to be the same with metasrv's `gc.enable` or unexpected behavior will occur |
|
||||
| `region_engine.mito.gc.lingering_time` | String | `1m` | Lingering time before deleting files.<br/>Should be long enough to allow long running queries to finish.<br/>If set to None, then unused files will be deleted immediately. |
|
||||
| `region_engine.mito.gc.lingering_time` | String | `1h` | Lingering time before deleting files.<br/>Should be long enough to allow long running queries to finish.<br/>If set to None, then unused files will be deleted immediately. |
|
||||
| `region_engine.mito.gc.unknown_file_lingering_time` | String | `1d` | Lingering time before deleting unknown files (files with undetermined expel time).<br/>Only applies during full file listing GC.<br/>This uses the object's last-modified timestamp as a heuristic (strict less-than comparison);<br/>do not configure this value too small in production to avoid deleting pre-manifest files<br/>from in-progress compaction or flush.<br/>For active/open regions, an unknown file is deleted only if its object last-modified time exceeds this TTL.<br/>If the object store does not provide a last-modified timestamp, the file is conservatively kept.<br/>For dropped regions, unknown files are deleted immediately. |
|
||||
| `region_engine.file` | -- | -- | Enable the file engine. |
|
||||
| `region_engine.metric` | -- | -- | Metric engine options. |
|
||||
|
||||
@@ -678,7 +678,7 @@ enable = false
|
||||
## Lingering time before deleting files.
|
||||
## Should be long enough to allow long running queries to finish.
|
||||
## If set to None, then unused files will be deleted immediately.
|
||||
lingering_time = "1m"
|
||||
lingering_time = "1h"
|
||||
|
||||
## Lingering time before deleting unknown files (files with undetermined expel time).
|
||||
## Only applies during full file listing GC.
|
||||
|
||||
@@ -426,6 +426,12 @@ pub enum Error {
|
||||
location: Location,
|
||||
},
|
||||
|
||||
#[snafu(display("Manual GC is rejected because maintenance mode is enabled"))]
|
||||
ManualGcRejectedByMaintenanceMode {
|
||||
#[snafu(implicit)]
|
||||
location: Location,
|
||||
},
|
||||
|
||||
#[cfg(feature = "mysql_kvbackend")]
|
||||
#[snafu(display("Failed to parse mysql url: {}", mysql_url))]
|
||||
ParseMySqlUrl {
|
||||
@@ -1194,6 +1200,7 @@ impl ErrorExt for Error {
|
||||
| Error::ParseAddr { .. }
|
||||
| Error::UnsupportedSelectorType { .. }
|
||||
| Error::InvalidArguments { .. }
|
||||
| Error::ManualGcRejectedByMaintenanceMode { .. }
|
||||
| Error::ProcedureNotFound { .. }
|
||||
| Error::TooManyPartitions { .. }
|
||||
| Error::TomlFormat { .. }
|
||||
|
||||
@@ -344,6 +344,7 @@ impl TestEnv {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: rx,
|
||||
config,
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||
|
||||
@@ -151,6 +151,7 @@ async fn test_handle_tick() {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config: GcSchedulerOptions::default(),
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||
|
||||
@@ -70,6 +70,7 @@ async fn test_gc_candidate_filtering_by_role() {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config: GcSchedulerOptions::default(),
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||
@@ -149,6 +150,7 @@ async fn test_gc_candidate_size_threshold() {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config,
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||
@@ -241,6 +243,7 @@ async fn test_gc_candidate_scoring() {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config,
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||
@@ -327,6 +330,7 @@ async fn test_gc_candidate_regions_per_table_threshold() {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config,
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||
|
||||
@@ -77,6 +77,7 @@ async fn test_concurrent_table_processing_limits() {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config,
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||
@@ -159,6 +160,7 @@ async fn test_datanode_processes_tables_with_partial_gc_failures() {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config: GcSchedulerOptions::default(),
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||
@@ -269,6 +271,7 @@ async fn test_region_gc_concurrency_limit() {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config,
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||
@@ -381,6 +384,7 @@ async fn test_region_gc_concurrency_with_partial_failures() {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config,
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||
@@ -520,6 +524,7 @@ async fn test_region_gc_concurrency_with_retryable_errors() {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config,
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||
|
||||
@@ -59,6 +59,7 @@ async fn test_different_gc_weights() {
|
||||
|
||||
let scheduler1 = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config: config1,
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||
@@ -83,6 +84,7 @@ async fn test_different_gc_weights() {
|
||||
|
||||
let scheduler2 = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config: config2,
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||
@@ -161,6 +163,7 @@ async fn test_regions_per_table_threshold() {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config,
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||
|
||||
@@ -83,6 +83,7 @@ async fn test_gc_regions_failure_handling() {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config: GcSchedulerOptions::default(),
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||
@@ -178,6 +179,7 @@ async fn test_get_file_references_failure() {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config: GcSchedulerOptions {
|
||||
retry_backoff_duration: Duration::from_millis(10), // shorten for test
|
||||
@@ -256,6 +258,7 @@ async fn test_get_table_route_failure() {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config: GcSchedulerOptions::default(),
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||
|
||||
@@ -68,6 +68,7 @@ async fn test_full_file_listing_first_time_gc() {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config,
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||
@@ -140,6 +141,7 @@ async fn test_full_file_listing_interval_enforcement() {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config,
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||
@@ -240,6 +242,7 @@ async fn test_full_file_listing_no_interval_passed() {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config,
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||
|
||||
@@ -76,6 +76,7 @@ async fn test_full_gc_workflow() {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config: GcSchedulerOptions::default(),
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||
@@ -205,6 +206,7 @@ async fn test_tracker_cleanup() {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config,
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(old_region_gc_tracker)),
|
||||
|
||||
@@ -51,6 +51,7 @@ async fn test_empty_file_refs_manifest() {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config: GcSchedulerOptions::default(),
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||
@@ -139,6 +140,7 @@ async fn test_multiple_regions_per_table() {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager: crate::gc::scheduler::new_test_runtime_switch_manager(),
|
||||
receiver: GcScheduler::channel().1,
|
||||
config: GcSchedulerOptions::default(),
|
||||
region_gc_tracker: Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||
|
||||
@@ -18,16 +18,18 @@ use std::time::{Duration, Instant};
|
||||
|
||||
use common_meta::DatanodeId;
|
||||
use common_meta::key::TableMetadataManagerRef;
|
||||
use common_meta::key::runtime_switch::RuntimeSwitchManagerRef;
|
||||
use common_procedure::ProcedureManagerRef;
|
||||
use common_telemetry::tracing::Instrument as _;
|
||||
use common_telemetry::{error, info};
|
||||
use snafu::ResultExt;
|
||||
use store_api::storage::{GcReport, RegionId};
|
||||
use tokio::sync::mpsc::{Receiver, Sender};
|
||||
use tokio::sync::{Mutex, oneshot};
|
||||
|
||||
use crate::cluster::MetaPeerClientRef;
|
||||
use crate::define_ticker;
|
||||
use crate::error::{Error, Result};
|
||||
use crate::error::{self, Error, Result};
|
||||
use crate::gc::Region2Peers;
|
||||
use crate::gc::ctx::{DefaultGcSchedulerCtx, SchedulerCtx};
|
||||
use crate::gc::dropped::DroppedRegionCollector;
|
||||
@@ -115,6 +117,8 @@ define_ticker!(
|
||||
/// [`GcScheduler`] is used to periodically trigger garbage collection on datanodes.
|
||||
pub struct GcScheduler {
|
||||
pub(crate) ctx: Arc<dyn SchedulerCtx>,
|
||||
/// Runtime switch manager to check maintenance mode.
|
||||
pub(crate) runtime_switch_manager: RuntimeSwitchManagerRef,
|
||||
/// The receiver of events.
|
||||
pub(crate) receiver: Receiver<Event>,
|
||||
/// GC configuration.
|
||||
@@ -133,6 +137,7 @@ impl GcScheduler {
|
||||
meta_peer_client: MetaPeerClientRef,
|
||||
mailbox: MailboxRef,
|
||||
server_addr: String,
|
||||
runtime_switch_manager: RuntimeSwitchManagerRef,
|
||||
config: GcSchedulerOptions,
|
||||
) -> Result<(Self, GcTicker)> {
|
||||
// Validate configuration before creating the scheduler
|
||||
@@ -148,6 +153,7 @@ impl GcScheduler {
|
||||
mailbox,
|
||||
server_addr,
|
||||
)?),
|
||||
runtime_switch_manager,
|
||||
receiver: rx,
|
||||
config,
|
||||
region_gc_tracker: Arc::new(Mutex::new(HashMap::new())),
|
||||
@@ -192,7 +198,11 @@ impl GcScheduler {
|
||||
.instrument(span)
|
||||
.await;
|
||||
if let Err(e) = &result {
|
||||
error!(e; "Failed to handle manual gc");
|
||||
if matches!(e, Error::ManualGcRejectedByMaintenanceMode { .. }) {
|
||||
info!("Rejected manual gc request: {}", e);
|
||||
} else {
|
||||
error!(e; "Failed to handle manual gc");
|
||||
}
|
||||
}
|
||||
let _ = sender.send(result);
|
||||
}
|
||||
@@ -204,6 +214,10 @@ impl GcScheduler {
|
||||
METRIC_META_GC_SCHEDULER_CYCLES_TOTAL.inc();
|
||||
let _timer = METRIC_META_GC_SCHEDULER_DURATION_SECONDS.start_timer();
|
||||
info!("Start to trigger gc");
|
||||
if self.is_maintenance_mode_enabled().await? {
|
||||
info!("Skip gc trigger because maintenance mode is enabled");
|
||||
return Ok(GcJobReport::default());
|
||||
}
|
||||
let span = common_telemetry::tracing::info_span!("meta_gc_handle_tick");
|
||||
let report = self.trigger_gc().instrument(span).await?;
|
||||
|
||||
@@ -227,6 +241,11 @@ impl GcScheduler {
|
||||
) -> Result<GcJobReport> {
|
||||
info!("Start to handle manual gc request");
|
||||
|
||||
if self.is_maintenance_mode_enabled().await? {
|
||||
info!("Skip manual gc request because maintenance mode is enabled");
|
||||
return error::ManualGcRejectedByMaintenanceModeSnafu {}.fail();
|
||||
}
|
||||
|
||||
// No specific regions, use default tick behavior
|
||||
let Some(regions) = region_ids else {
|
||||
let report = self.trigger_gc().await?;
|
||||
@@ -294,11 +313,26 @@ impl GcScheduler {
|
||||
info!("Finished manual gc request");
|
||||
Ok(report)
|
||||
}
|
||||
|
||||
pub(crate) async fn is_maintenance_mode_enabled(&self) -> Result<bool> {
|
||||
self.runtime_switch_manager
|
||||
.maintenance_mode()
|
||||
.await
|
||||
.context(error::RuntimeSwitchManagerSnafu)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) fn new_test_runtime_switch_manager() -> RuntimeSwitchManagerRef {
|
||||
Arc::new(common_meta::key::runtime_switch::RuntimeSwitchManager::new(
|
||||
Arc::new(common_meta::kv_backend::memory::MemoryKvBackend::new()),
|
||||
))
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::collections::HashMap;
|
||||
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||
use std::time::Duration;
|
||||
|
||||
use common_meta::datanode::RegionStat;
|
||||
@@ -309,6 +343,72 @@ mod tests {
|
||||
|
||||
use super::*;
|
||||
|
||||
#[derive(Default)]
|
||||
struct CountingSchedulerCtx {
|
||||
get_table_to_region_stats_calls: AtomicUsize,
|
||||
get_table_reparts_calls: AtomicUsize,
|
||||
gc_regions_calls: AtomicUsize,
|
||||
}
|
||||
|
||||
impl CountingSchedulerCtx {
|
||||
fn assert_no_scheduler_work(&self) {
|
||||
assert_eq!(
|
||||
0,
|
||||
self.get_table_to_region_stats_calls.load(Ordering::Relaxed),
|
||||
"get_table_to_region_stats should not be called"
|
||||
);
|
||||
assert_eq!(
|
||||
0,
|
||||
self.get_table_reparts_calls.load(Ordering::Relaxed),
|
||||
"get_table_reparts should not be called"
|
||||
);
|
||||
assert_eq!(
|
||||
0,
|
||||
self.gc_regions_calls.load(Ordering::Relaxed),
|
||||
"gc_regions should not be called"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl SchedulerCtx for CountingSchedulerCtx {
|
||||
async fn get_table_to_region_stats(&self) -> Result<HashMap<TableId, Vec<RegionStat>>> {
|
||||
self.get_table_to_region_stats_calls
|
||||
.fetch_add(1, Ordering::Relaxed);
|
||||
panic!("get_table_to_region_stats should not be called in maintenance mode")
|
||||
}
|
||||
|
||||
async fn get_table_reparts(&self) -> Result<Vec<(TableId, TableRepartValue)>> {
|
||||
self.get_table_reparts_calls.fetch_add(1, Ordering::Relaxed);
|
||||
panic!("get_table_reparts should not be called in maintenance mode")
|
||||
}
|
||||
|
||||
async fn get_table_route(
|
||||
&self,
|
||||
_table_id: TableId,
|
||||
) -> Result<(TableId, PhysicalTableRouteValue)> {
|
||||
unreachable!("get_table_route should not be called in this test")
|
||||
}
|
||||
|
||||
async fn batch_get_table_route(
|
||||
&self,
|
||||
_table_ids: &[TableId],
|
||||
) -> Result<HashMap<TableId, PhysicalTableRouteValue>> {
|
||||
unreachable!("batch_get_table_route should not be called in this test")
|
||||
}
|
||||
|
||||
async fn gc_regions(
|
||||
&self,
|
||||
_region_ids: &[RegionId],
|
||||
_full_file_listing: bool,
|
||||
_timeout: Duration,
|
||||
_region_routes_override: Region2Peers,
|
||||
) -> Result<GcReport> {
|
||||
self.gc_regions_calls.fetch_add(1, Ordering::Relaxed);
|
||||
panic!("gc_regions should not be called in maintenance mode")
|
||||
}
|
||||
}
|
||||
|
||||
struct ErrorMockSchedulerCtx;
|
||||
|
||||
#[async_trait::async_trait]
|
||||
@@ -356,6 +456,7 @@ mod tests {
|
||||
|
||||
let scheduler = GcScheduler {
|
||||
ctx: Arc::new(ErrorMockSchedulerCtx),
|
||||
runtime_switch_manager: new_test_runtime_switch_manager(),
|
||||
receiver: rx,
|
||||
config: GcSchedulerOptions::default(),
|
||||
region_gc_tracker: Arc::new(Mutex::new(HashMap::new())),
|
||||
@@ -372,4 +473,61 @@ mod tests {
|
||||
|
||||
assert!(result.is_err());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_maintenance_mode_skips_manual_gc() {
|
||||
let (tx, rx) = GcScheduler::channel();
|
||||
drop(tx);
|
||||
let runtime_switch_manager = new_test_runtime_switch_manager();
|
||||
runtime_switch_manager.set_maintenance_mode().await.unwrap();
|
||||
|
||||
let ctx = Arc::new(CountingSchedulerCtx::default());
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager,
|
||||
receiver: rx,
|
||||
config: GcSchedulerOptions::default(),
|
||||
region_gc_tracker: Arc::new(Mutex::new(HashMap::new())),
|
||||
last_tracker_cleanup: Arc::new(Mutex::new(Instant::now())),
|
||||
};
|
||||
|
||||
let result = scheduler
|
||||
.handle_manual_gc(
|
||||
Some(vec![RegionId::new(1, 0)]),
|
||||
Some(false),
|
||||
Some(Duration::from_secs(1)),
|
||||
)
|
||||
.await;
|
||||
|
||||
let err = result.unwrap_err();
|
||||
assert!(matches!(
|
||||
err,
|
||||
error::Error::ManualGcRejectedByMaintenanceMode { .. }
|
||||
));
|
||||
assert!(err.to_string().contains("maintenance mode is enabled"));
|
||||
ctx.assert_no_scheduler_work();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_maintenance_mode_skips_tick_gc() {
|
||||
let (tx, rx) = GcScheduler::channel();
|
||||
drop(tx);
|
||||
let runtime_switch_manager = new_test_runtime_switch_manager();
|
||||
runtime_switch_manager.set_maintenance_mode().await.unwrap();
|
||||
|
||||
let ctx = Arc::new(CountingSchedulerCtx::default());
|
||||
let scheduler = GcScheduler {
|
||||
ctx: ctx.clone(),
|
||||
runtime_switch_manager,
|
||||
receiver: rx,
|
||||
config: GcSchedulerOptions::default(),
|
||||
region_gc_tracker: Arc::new(Mutex::new(HashMap::new())),
|
||||
last_tracker_cleanup: Arc::new(Mutex::new(Instant::now())),
|
||||
};
|
||||
|
||||
let result = scheduler.handle_tick().await;
|
||||
|
||||
assert!(result.is_ok());
|
||||
ctx.assert_no_scheduler_work();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -516,6 +516,7 @@ impl MetasrvBuilder {
|
||||
meta_peer_client.clone(),
|
||||
mailbox.clone(),
|
||||
options.grpc.server_addr.clone(),
|
||||
runtime_switch_manager.clone(),
|
||||
options.gc.clone(),
|
||||
)?;
|
||||
gc_scheduler.try_start()?;
|
||||
|
||||
+9
-1
@@ -173,7 +173,7 @@ impl Default for GcConfig {
|
||||
Self {
|
||||
enable: false,
|
||||
// expect long running queries to be finished(or at least be able to notify it's using a deleted file) within a reasonable time
|
||||
lingering_time: Some(Duration::from_secs(60)),
|
||||
lingering_time: Some(Duration::from_secs(60 * 60)),
|
||||
// 1 day, for unknown expel time, which is when this file get removed from manifest.
|
||||
// Only applies to full-listing GC for active/open regions. A long default avoids
|
||||
// accidentally deleting pre-manifest files (e.g. compaction/flush still in progress).
|
||||
@@ -1135,4 +1135,12 @@ mod tests {
|
||||
]
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_gc_config_default_lingering_time() {
|
||||
assert_eq!(
|
||||
GcConfig::default().lingering_time,
|
||||
Some(Duration::from_secs(60 * 60))
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1982,7 +1982,7 @@ apply_on_query = "auto"
|
||||
mem_threshold_on_create = "auto"
|
||||
{vector_index_config}[region_engine.mito.gc]
|
||||
enable = false
|
||||
lingering_time = "1m"
|
||||
lingering_time = "1h"
|
||||
unknown_file_lingering_time = "1day"
|
||||
max_concurrent_lister_per_gc_job = 32
|
||||
max_concurrent_gc_job = 4
|
||||
|
||||
Reference in New Issue
Block a user