mirror of
https://github.com/GreptimeTeam/greptimedb.git
synced 2026-08-18 12:08:22 +00:00
feat: prepare soft-drop WAL retirement (#8475)
* fix: flush soft-dropped regions on close Signed-off-by: Lei, HUANG <ratuthomm@gmail.com> * feat(mito2): handle flush-on-close race with concurrent in-flight flush When a region close with `flush_on_close: true` races with an already-running flush, pass the actual close request (including the flush_on_close flag) to the DDL handler instead of a default request so the pending flush is correctly awaited. Files: `src/mito2/src/worker/handle_close.rs` Also adds a test verifying that closing with flush-on-close while a flush is in progress still persists all written data correctly. Files: `src/mito2/src/engine/close_test.rs` Signed-off-by: Lei, HUANG <ratuthomm@gmail.com> * feat: support full WAL retirement Signed-off-by: Lei, HUANG <ratuthomm@gmail.com> * fix: complete close request migration Signed-off-by: Lei, HUANG <ratuthomm@gmail.com> * fix: finish close request callsites Signed-off-by: Lei, HUANG <ratuthomm@gmail.com> * refactor: guard Kafka provider setup behind index collector check Move Kafka provider initialization and `get_or_insert` inside the existing `if let Some(collector)` block so these operations are skipped when no global index collector is configured. Affected file: - `src/log-store/src/kafka/log_store.rs` Signed-off-by: Lei, HUANG <ratuthomm@gmail.com> * chore: avoid to_vec Signed-off-by: Lei, HUANG <ratuthomm@gmail.com> * refactor: replace imperative close-region loop with functional combinators Transform the region close dispatch in `DropTableExecutor` from mutable `Vec` and `push` loops to iterator chains with `join_all`, improving idiomatic Rust style and readability. - `src/common/meta/src/ddl/drop_table/executor.rs` — rewired datanode region-close logic to use `peers.map()` and nested `join_all`, moving `node_manager.datanode()` inside the closure to align with the new structure Signed-off-by: Lei, HUANG <ratuthomm@gmail.com> * fix: decouple Kafka client from WAL checkpoint Signed-off-by: Lei, HUANG <ratuthomm@gmail.com> * fix: merge Kafka WAL index checkpoints Signed-off-by: Lei, HUANG <ratuthomm@gmail.com> * fix: delegate Kafka WAL retirement to metasrv Signed-off-by: Lei, HUANG <ratuthomm@gmail.com> * chore: rebase main and resolve conflicts Signed-off-by: Lei, HUANG <ratuthomm@gmail.com> * fix: license header Signed-off-by: Lei, HUANG <ratuthomm@gmail.com> * chore: bump proto to commits on main Signed-off-by: Lei, HUANG <ratuthomm@gmail.com> * refactor: remove Kafka obsolete-all index changes Signed-off-by: Lei, HUANG <ratuthomm@gmail.com> --------- Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>
This commit is contained in:
Generated
+1
-1
@@ -5938,7 +5938,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "greptime-proto"
|
||||
version = "0.1.0"
|
||||
source = "git+https://github.com/GreptimeTeam/greptime-proto.git?rev=4a69c9e75240a76afeb0796798085813669288f8#4a69c9e75240a76afeb0796798085813669288f8"
|
||||
source = "git+https://github.com/GreptimeTeam/greptime-proto.git?rev=6fc4d88851c5ad16846698e4dff6e74c1f7a48ba#6fc4d88851c5ad16846698e4dff6e74c1f7a48ba"
|
||||
dependencies = [
|
||||
"prost 0.14.1",
|
||||
"prost-types 0.14.1",
|
||||
|
||||
+1
-1
@@ -158,7 +158,7 @@ fs2 = "0.4"
|
||||
fst = "0.4.7"
|
||||
futures = "0.3"
|
||||
futures-util = "0.3"
|
||||
greptime-proto = { git = "https://github.com/GreptimeTeam/greptime-proto.git", rev = "4a69c9e75240a76afeb0796798085813669288f8" }
|
||||
greptime-proto = { git = "https://github.com/GreptimeTeam/greptime-proto.git", rev = "6fc4d88851c5ad16846698e4dff6e74c1f7a48ba" }
|
||||
hex = "0.4"
|
||||
http = "1"
|
||||
humantime = "2.1"
|
||||
|
||||
@@ -179,6 +179,7 @@ impl DropTableProcedure {
|
||||
&self.context.node_manager,
|
||||
&self.context.leader_region_registry,
|
||||
&self.data.physical_region_routes,
|
||||
true,
|
||||
)
|
||||
.await?;
|
||||
self.context
|
||||
|
||||
@@ -313,6 +313,7 @@ impl DropTableExecutor {
|
||||
}),
|
||||
body: Some(region_request::Body::Close(PbCloseRegionRequest {
|
||||
region_id: region_id.as_u64(),
|
||||
flush_on_close: false,
|
||||
})),
|
||||
};
|
||||
|
||||
@@ -348,11 +349,13 @@ impl DropTableExecutor {
|
||||
}
|
||||
|
||||
/// Closes all table regions on datanodes without deleting region files or metadata tombstones.
|
||||
/// When `flush_leaders_on_close` is set, only leader regions are flushed before close.
|
||||
pub async fn on_close_regions(
|
||||
&self,
|
||||
node_manager: &NodeManagerRef,
|
||||
leader_region_registry: &LeaderRegionRegistryRef,
|
||||
region_routes: &[RegionRoute],
|
||||
flush_leaders_on_close: bool,
|
||||
) -> Result<()> {
|
||||
let table_id = self.table_id;
|
||||
let mut seen_peer_ids = HashSet::new();
|
||||
@@ -360,39 +363,57 @@ impl DropTableExecutor {
|
||||
.into_iter()
|
||||
.chain(find_followers(region_routes))
|
||||
.filter(|peer| seen_peer_ids.insert(peer.id));
|
||||
let mut close_region_tasks = Vec::new();
|
||||
|
||||
for datanode in peers {
|
||||
let requester = node_manager.datanode(&datanode).await;
|
||||
let close_region_tasks = peers.map(|datanode| {
|
||||
let region_ids = find_leader_regions(region_routes, &datanode)
|
||||
.into_iter()
|
||||
.chain(find_follower_regions(region_routes, &datanode))
|
||||
.map(|region_number| RegionId::new(table_id, region_number));
|
||||
.map(|region_number| {
|
||||
(
|
||||
RegionId::new(table_id, region_number),
|
||||
flush_leaders_on_close,
|
||||
)
|
||||
})
|
||||
.chain(
|
||||
find_follower_regions(region_routes, &datanode)
|
||||
.into_iter()
|
||||
.map(|region_number| (RegionId::new(table_id, region_number), false)),
|
||||
)
|
||||
.collect::<Vec<_>>();
|
||||
|
||||
for region_id in region_ids {
|
||||
debug!("Closing region {region_id} on Datanode {datanode:?}");
|
||||
let request = RegionRequest {
|
||||
header: Some(RegionRequestHeader {
|
||||
tracing_context: TracingContext::from_current_span().to_w3c(),
|
||||
..Default::default()
|
||||
}),
|
||||
body: Some(region_request::Body::Close(PbCloseRegionRequest {
|
||||
region_id: region_id.as_u64(),
|
||||
})),
|
||||
};
|
||||
async move {
|
||||
let requester = node_manager.datanode(&datanode).await;
|
||||
let close_region_tasks =
|
||||
region_ids.into_iter().map(|(region_id, flush_on_close)| {
|
||||
debug!("Closing region {region_id} on Datanode {datanode:?}");
|
||||
let request = RegionRequest {
|
||||
header: Some(RegionRequestHeader {
|
||||
tracing_context: TracingContext::from_current_span().to_w3c(),
|
||||
..Default::default()
|
||||
}),
|
||||
body: Some(region_request::Body::Close(PbCloseRegionRequest {
|
||||
region_id: region_id.as_u64(),
|
||||
flush_on_close,
|
||||
})),
|
||||
};
|
||||
|
||||
let datanode = datanode.clone();
|
||||
let requester = requester.clone();
|
||||
close_region_tasks.push(async move {
|
||||
if let Err(err) = requester.handle(request).await
|
||||
&& err.status_code() != StatusCode::RegionNotFound
|
||||
{
|
||||
return Err(add_peer_context_if_needed(datanode)(err));
|
||||
}
|
||||
Ok(())
|
||||
});
|
||||
let datanode = datanode.clone();
|
||||
let requester = requester.clone();
|
||||
async move {
|
||||
if let Err(err) = requester.handle(request).await
|
||||
&& err.status_code() != StatusCode::RegionNotFound
|
||||
{
|
||||
return Err(add_peer_context_if_needed(datanode)(err));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
});
|
||||
|
||||
join_all(close_region_tasks)
|
||||
.await
|
||||
.into_iter()
|
||||
.collect::<Result<Vec<_>>>()?;
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
join_all(close_region_tasks)
|
||||
.await
|
||||
|
||||
@@ -293,16 +293,16 @@ async fn test_soft_drop_closes_regions_and_keeps_tombstone() {
|
||||
let Some(region_request::Body::Close(req)) = request.body else {
|
||||
unreachable!();
|
||||
};
|
||||
requests.push((peer.id, req.region_id));
|
||||
requests.push((peer.id, req.region_id, req.flush_on_close));
|
||||
}
|
||||
requests.sort_unstable();
|
||||
assert_eq!(
|
||||
requests,
|
||||
vec![
|
||||
(1, RegionId::new(table_id, 1).as_u64()),
|
||||
(1, RegionId::new(table_id, 2).as_u64()),
|
||||
(2, RegionId::new(table_id, 1).as_u64()),
|
||||
(2, RegionId::new(table_id, 2).as_u64()),
|
||||
(1, RegionId::new(table_id, 1).as_u64(), true),
|
||||
(1, RegionId::new(table_id, 2).as_u64(), false),
|
||||
(2, RegionId::new(table_id, 1).as_u64(), false),
|
||||
(2, RegionId::new(table_id, 2).as_u64(), true),
|
||||
]
|
||||
);
|
||||
assert!(rx.try_recv().is_err());
|
||||
|
||||
@@ -153,7 +153,7 @@ impl RegionAliveKeeper {
|
||||
|
||||
async fn close_staled_region(&self, region_id: RegionId) {
|
||||
info!("Closing staled region: {region_id}");
|
||||
let request = RegionRequest::Close(RegionCloseRequest {});
|
||||
let request = RegionRequest::Close(RegionCloseRequest::default());
|
||||
if let Err(e) = self.region_server.handle_request(region_id, request).await
|
||||
&& e.status_code() != StatusCode::RegionNotFound
|
||||
{
|
||||
|
||||
@@ -522,7 +522,10 @@ mod tests {
|
||||
.unwrap();
|
||||
|
||||
region_server
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let mut heartbeat_env = HeartbeatResponseTestEnv::new();
|
||||
|
||||
@@ -40,8 +40,10 @@ impl InstructionHandler for CloseRegionsHandler {
|
||||
.collect::<Vec<_>>();
|
||||
|
||||
let futs = region_ids.iter().map(|region_id| {
|
||||
ctx.region_server
|
||||
.handle_request(*region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
ctx.region_server.handle_request(
|
||||
*region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
});
|
||||
|
||||
let results = join_all(futs).await;
|
||||
|
||||
@@ -169,11 +169,17 @@ mod tests {
|
||||
.await
|
||||
.unwrap();
|
||||
region_server
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
region_server
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
|
||||
@@ -1852,7 +1852,10 @@ impl RegionServerInner {
|
||||
|
||||
for (region_id, engine) in regions {
|
||||
let closed = engine
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await;
|
||||
match closed {
|
||||
Ok(_) => debug!("Region {region_id} is closed"),
|
||||
@@ -2218,7 +2221,10 @@ mod tests {
|
||||
);
|
||||
|
||||
let response = mock_region_server
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(response.affected_rows, 0);
|
||||
|
||||
@@ -511,6 +511,10 @@ impl LogStore for KafkaLogStore {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn obsolete_all(&self, _provider: &Provider, _region_id: RegionId) -> Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Returns the highest entry id of the specified topic in remote WAL.
|
||||
fn latest_entry_id(&self, provider: &Provider) -> Result<EntryId> {
|
||||
let provider = provider
|
||||
|
||||
@@ -86,6 +86,10 @@ impl LogStore for NoopLogStore {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn obsolete_all(&self, _provider: &Provider, _region_id: RegionId) -> Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn latest_entry_id(&self, _provider: &Provider) -> Result<EntryId> {
|
||||
Ok(0)
|
||||
}
|
||||
|
||||
@@ -484,6 +484,12 @@ impl LogStore for RaftEngineLogStore {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn obsolete_all(&self, provider: &Provider, region_id: RegionId) -> Result<()> {
|
||||
let latest_entry_id = self.latest_entry_id(provider)?;
|
||||
self.obsolete(provider, region_id, latest_entry_id).await?;
|
||||
self.delete_namespace(provider).await
|
||||
}
|
||||
|
||||
fn latest_entry_id(&self, provider: &Provider) -> Result<EntryId> {
|
||||
let ns = provider
|
||||
.as_raft_engine_provider()
|
||||
@@ -576,6 +582,38 @@ mod tests {
|
||||
assert!(logstore.list_namespaces().await.unwrap().is_empty());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_obsolete_all_removes_entries_and_namespace() {
|
||||
let dir = create_temp_dir("raft-engine-logstore-test");
|
||||
let logstore = RaftEngineLogStore::try_new(
|
||||
dir.path().to_str().unwrap().to_string(),
|
||||
&RaftEngineConfig::default(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let region_id = RegionId::new(1, 1);
|
||||
let provider = Provider::raft_engine_provider(region_id.as_u64());
|
||||
logstore.create_namespace(&provider).await.unwrap();
|
||||
for entry_id in 1..=3 {
|
||||
logstore
|
||||
.append(
|
||||
EntryImpl::create(
|
||||
entry_id,
|
||||
region_id.as_u64(),
|
||||
entry_id.to_string().into_bytes(),
|
||||
)
|
||||
.into(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
logstore.obsolete_all(&provider, region_id).await.unwrap();
|
||||
|
||||
assert_eq!(0, logstore.latest_entry_id(&provider).unwrap());
|
||||
assert!(logstore.list_namespaces().await.unwrap().is_empty());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_append_and_read() {
|
||||
let dir = create_temp_dir("raft-engine-logstore-test");
|
||||
|
||||
@@ -604,7 +604,7 @@ mod test {
|
||||
engine
|
||||
.handle_request(
|
||||
physical_region_id,
|
||||
RegionRequest::Close(RegionCloseRequest {}),
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
@@ -632,7 +632,7 @@ mod test {
|
||||
engine
|
||||
.handle_request(
|
||||
nonexistent_region_id,
|
||||
RegionRequest::Close(RegionCloseRequest {}),
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
@@ -736,7 +736,7 @@ mod test {
|
||||
metric_engine
|
||||
.handle_request(
|
||||
physical_region_id,
|
||||
RegionRequest::Close(RegionCloseRequest {}),
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
@@ -828,7 +828,10 @@ mod test {
|
||||
// Closes all regions
|
||||
for region_id in logical_region_ids.iter().chain(physical_region_ids.iter()) {
|
||||
metric_engine
|
||||
.handle_request(*region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
*region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
@@ -29,7 +29,7 @@ impl MetricEngineInner {
|
||||
pub async fn close_region(
|
||||
&self,
|
||||
region_id: RegionId,
|
||||
_req: RegionCloseRequest,
|
||||
req: RegionCloseRequest,
|
||||
) -> Result<AffectedRows> {
|
||||
let data_region_id = utils::to_data_region_id(region_id);
|
||||
if self
|
||||
@@ -38,7 +38,8 @@ impl MetricEngineInner {
|
||||
.unwrap()
|
||||
.exist_physical_region(data_region_id)
|
||||
{
|
||||
self.close_physical_region(data_region_id).await?;
|
||||
self.close_physical_region(data_region_id, req.flush_on_close)
|
||||
.await?;
|
||||
self.state
|
||||
.write()
|
||||
.unwrap()
|
||||
@@ -59,18 +60,25 @@ impl MetricEngineInner {
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) async fn close_physical_region(&self, region_id: RegionId) -> Result<AffectedRows> {
|
||||
pub(crate) async fn close_physical_region(
|
||||
&self,
|
||||
region_id: RegionId,
|
||||
flush_on_close: bool,
|
||||
) -> Result<AffectedRows> {
|
||||
let data_region_id = utils::to_data_region_id(region_id);
|
||||
let metadata_region_id = utils::to_metadata_region_id(region_id);
|
||||
|
||||
self.mito
|
||||
.handle_request(data_region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
data_region_id,
|
||||
RegionRequest::Close(RegionCloseRequest { flush_on_close }),
|
||||
)
|
||||
.await
|
||||
.with_context(|_| CloseMitoRegionSnafu { region_id })?;
|
||||
self.mito
|
||||
.handle_request(
|
||||
metadata_region_id,
|
||||
RegionRequest::Close(RegionCloseRequest {}),
|
||||
RegionRequest::Close(RegionCloseRequest { flush_on_close }),
|
||||
)
|
||||
.await
|
||||
.with_context(|_| CloseMitoRegionSnafu { region_id })?;
|
||||
|
||||
@@ -105,7 +105,7 @@ impl MetricEngineInner {
|
||||
utils::to_metadata_region_id(physical_region_id),
|
||||
utils::to_data_region_id(physical_region_id)
|
||||
);
|
||||
if let Err(err) = self.close_physical_region(physical_region_id).await {
|
||||
if let Err(err) = self.close_physical_region(physical_region_id, false).await {
|
||||
error!(err; "Failed to close physical region {}", physical_region_id);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -304,7 +304,7 @@ mod tests {
|
||||
metric_engine
|
||||
.handle_request(
|
||||
target_physical_region_id,
|
||||
RegionRequest::Close(RegionCloseRequest {}),
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
@@ -498,7 +498,10 @@ async fn test_catchup_with_manifest_update(factory: Option<LogStoreFactory>) {
|
||||
|
||||
async fn close_region(engine: &MitoEngine, region_id: RegionId) {
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
@@ -15,15 +15,20 @@
|
||||
use std::sync::Arc;
|
||||
use std::sync::atomic::Ordering;
|
||||
|
||||
use api::v1::Rows;
|
||||
use common_base::Plugins;
|
||||
use common_recordbatch::RecordBatches;
|
||||
use store_api::region_engine::RegionEngine;
|
||||
use store_api::region_request::{RegionCloseRequest, RegionRequest};
|
||||
use store_api::storage::RegionId;
|
||||
use store_api::storage::{RegionId, ScanRequest};
|
||||
|
||||
use crate::config::MitoConfig;
|
||||
use crate::engine::flush_test::MockRegionHook;
|
||||
use crate::engine::listener::AlterFlushListener;
|
||||
use crate::engine::region_hook::RegionHookRef;
|
||||
use crate::test_util::{CreateRequestBuilder, TestEnv};
|
||||
use crate::test_util::{
|
||||
CreateRequestBuilder, TestEnv, build_rows, flush_region, put_rows, rows_schema,
|
||||
};
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_engine_close_region() {
|
||||
@@ -43,7 +48,10 @@ async fn test_engine_close_region_with_format(flat_format: bool) {
|
||||
let region_id = RegionId::new(1, 1);
|
||||
// It's okay to close a region doesn't exist.
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
@@ -55,14 +63,20 @@ async fn test_engine_close_region_with_format(flat_format: bool) {
|
||||
|
||||
// Close the created region.
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(!engine.is_region_exists(region_id));
|
||||
|
||||
// It's okay to close this region again.
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
@@ -93,7 +107,10 @@ async fn test_region_hook_on_close() {
|
||||
assert_eq!(hook.dropped_count.load(Ordering::Relaxed), 0);
|
||||
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
@@ -103,3 +120,160 @@ async fn test_region_hook_on_close() {
|
||||
assert_eq!(hook.dropped_count.load(Ordering::Relaxed), 0);
|
||||
assert_eq!(hook.files_removed_count.load(Ordering::Relaxed), 0);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_engine_close_region_flush_on_close() {
|
||||
let mut env = TestEnv::with_prefix("close-flush-on-close").await;
|
||||
let engine = env.create_engine(MitoConfig::default()).await;
|
||||
|
||||
let region_id = RegionId::new(1, 1);
|
||||
let request = CreateRequestBuilder::new().build();
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Create(request.clone()))
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
put_rows(
|
||||
&engine,
|
||||
region_id,
|
||||
Rows {
|
||||
schema: rows_schema(&request),
|
||||
rows: build_rows(0, 3),
|
||||
},
|
||||
)
|
||||
.await;
|
||||
assert!(
|
||||
!engine
|
||||
.get_region(region_id)
|
||||
.unwrap()
|
||||
.version()
|
||||
.memtables
|
||||
.is_empty()
|
||||
);
|
||||
|
||||
engine
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest {
|
||||
flush_on_close: true,
|
||||
}),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(!engine.is_region_exists(region_id));
|
||||
|
||||
engine
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Open(store_api::region_request::RegionOpenRequest {
|
||||
engine: String::new(),
|
||||
table_dir: request.table_dir.clone(),
|
||||
path_type: store_api::region_request::PathType::Bare,
|
||||
options: request.options.clone(),
|
||||
skip_wal_replay: true,
|
||||
checkpoint: None,
|
||||
requirements: Default::default(),
|
||||
}),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let stream = engine
|
||||
.scan_to_stream(region_id, ScanRequest::default())
|
||||
.await
|
||||
.unwrap();
|
||||
let batches = RecordBatches::try_collect(stream).await.unwrap();
|
||||
assert_eq!(3, batches.iter().map(|b| b.num_rows()).sum::<usize>());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_engine_close_region_flush_on_close_while_flushing() {
|
||||
let mut env = TestEnv::with_prefix("close-flush-on-close-while-flushing").await;
|
||||
let listener = Arc::new(AlterFlushListener::default());
|
||||
let engine = env
|
||||
.create_engine_with(MitoConfig::default(), None, Some(listener.clone()), None)
|
||||
.await;
|
||||
|
||||
let region_id = RegionId::new(1, 1);
|
||||
let request = CreateRequestBuilder::new().build();
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Create(request.clone()))
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
put_rows(
|
||||
&engine,
|
||||
region_id,
|
||||
Rows {
|
||||
schema: rows_schema(&request),
|
||||
rows: build_rows(0, 3),
|
||||
},
|
||||
)
|
||||
.await;
|
||||
|
||||
let flush_engine = engine.clone();
|
||||
let flush_task = tokio::spawn(async move {
|
||||
flush_region(&flush_engine, region_id, None).await;
|
||||
});
|
||||
listener.wait_flush_begin().await;
|
||||
|
||||
put_rows(
|
||||
&engine,
|
||||
region_id,
|
||||
Rows {
|
||||
schema: rows_schema(&request),
|
||||
rows: build_rows(3, 6),
|
||||
},
|
||||
)
|
||||
.await;
|
||||
|
||||
let request_count = listener.request_count();
|
||||
let close_engine = engine.clone();
|
||||
let close_task = tokio::spawn(async move {
|
||||
close_engine
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest {
|
||||
flush_on_close: true,
|
||||
}),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
});
|
||||
listener.wait_request_count(request_count + 1).await;
|
||||
|
||||
let second_flush_listener = listener.clone();
|
||||
let second_flush_task = tokio::spawn(async move {
|
||||
second_flush_listener.wait_flush_begin().await;
|
||||
second_flush_listener.wake_flush();
|
||||
});
|
||||
listener.wake_flush();
|
||||
|
||||
flush_task.await.unwrap();
|
||||
close_task.await.unwrap();
|
||||
second_flush_task.abort();
|
||||
assert!(!engine.is_region_exists(region_id));
|
||||
|
||||
engine
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Open(store_api::region_request::RegionOpenRequest {
|
||||
engine: String::new(),
|
||||
table_dir: request.table_dir,
|
||||
path_type: store_api::region_request::PathType::Bare,
|
||||
options: request.options,
|
||||
skip_wal_replay: true,
|
||||
checkpoint: None,
|
||||
requirements: Default::default(),
|
||||
}),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let stream = engine
|
||||
.scan_to_stream(region_id, ScanRequest::default())
|
||||
.await
|
||||
.unwrap();
|
||||
let batches = RecordBatches::try_collect(stream).await.unwrap();
|
||||
assert_eq!(6, batches.iter().map(|b| b.num_rows()).sum::<usize>());
|
||||
}
|
||||
|
||||
@@ -105,7 +105,10 @@ async fn test_engine_create_close_create_region_with_format(flat_format: bool) {
|
||||
.unwrap();
|
||||
// Close the region.
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
// Create the same region id again.
|
||||
|
||||
@@ -419,7 +419,10 @@ async fn test_stalled_write_fails_fast_if_region_closed_during_editing() {
|
||||
let close_engine = engine.clone();
|
||||
let close_task = tokio::spawn(async move {
|
||||
close_engine
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
});
|
||||
|
||||
|
||||
@@ -447,7 +447,10 @@ async fn test_skip_remote_wal_replay_sets_topic_latest_entry_id(factory: Option<
|
||||
assert_eq!(1, region.topic_latest_entry_id.load(Ordering::Relaxed));
|
||||
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
@@ -589,7 +592,10 @@ async fn test_remote_wal_open_without_replayed_memtable_sets_topic_latest_entry_
|
||||
// The empty flush updates `topic_latest_entry_id` from KafkaLogStore's topic stats.
|
||||
flush_region(&engine, region_id, None).await;
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
|
||||
@@ -245,7 +245,10 @@ async fn test_engine_region_open_with_options_with_format(flat_format: bool) {
|
||||
|
||||
// Close the region.
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
@@ -306,7 +309,10 @@ async fn test_engine_region_open_with_custom_store_with_format(flat_format: bool
|
||||
|
||||
// Close the custom region.
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
@@ -628,7 +634,10 @@ async fn test_open_compaction_region_with_format(flat_format: bool) {
|
||||
|
||||
// Close the region.
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
@@ -708,7 +717,10 @@ async fn test_open_backfills_partition_expr_with_fetcher() {
|
||||
|
||||
// close and reopen to trigger backfill in opener
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
engine
|
||||
@@ -742,7 +754,10 @@ async fn test_open_backfills_partition_expr_with_fetcher() {
|
||||
|
||||
// reopen again to ensure no further changes and still Some
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let engine = env.reopen_engine(engine, MitoConfig::default()).await;
|
||||
@@ -784,7 +799,10 @@ async fn test_open_keeps_none_without_fetcher() {
|
||||
assert!(meta.partition_expr.is_none());
|
||||
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let engine = env.reopen_engine(engine, MitoConfig::default()).await;
|
||||
|
||||
@@ -78,7 +78,10 @@ async fn test_close_region_skip_wal(insert: bool) {
|
||||
|
||||
// Close the region. This should trigger a flush.
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
@@ -146,7 +149,10 @@ async fn test_close_follower_region_skip_wal() {
|
||||
|
||||
// Close the region. This should trigger a flush.
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
@@ -214,7 +220,10 @@ async fn test_close_follower_region_skip_wal_with_pending_data() {
|
||||
assert!(region.is_follower());
|
||||
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
@@ -268,7 +277,10 @@ async fn test_close_region_skip_wal_while_flush_in_flight_closes_region() {
|
||||
let engine_cloned = engine.clone();
|
||||
let close_job = tokio::spawn(async move {
|
||||
engine_cloned
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
});
|
||||
@@ -344,7 +356,10 @@ async fn test_close_region_skip_wal_rejects_writes_queued_after_close() {
|
||||
let engine_cloned = engine.clone();
|
||||
let close_job = tokio::spawn(async move {
|
||||
engine_cloned
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
});
|
||||
@@ -445,7 +460,10 @@ async fn test_concurrent_close_region_skip_wal_while_flush_in_flight_succeeds()
|
||||
let engine_cloned = engine.clone();
|
||||
let first_close = tokio::spawn(async move {
|
||||
engine_cloned
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
});
|
||||
@@ -455,7 +473,10 @@ async fn test_concurrent_close_region_skip_wal_while_flush_in_flight_succeeds()
|
||||
let engine_cloned = engine.clone();
|
||||
let second_close = tokio::spawn(async move {
|
||||
engine_cloned
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
});
|
||||
@@ -518,7 +539,10 @@ async fn test_close_region_after_truncate_skip_wal() {
|
||||
assert!(!region.version().memtables.is_empty());
|
||||
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
|
||||
@@ -2107,7 +2107,7 @@ mod tests {
|
||||
scheduler.add_ddl_request_to_pending(SenderDdlRequest {
|
||||
sender: OptionOutputTx::from(sender),
|
||||
region_id: builder.region_id(),
|
||||
request: DdlRequest::Close(store_api::region_request::RegionCloseRequest {}),
|
||||
request: DdlRequest::Close(store_api::region_request::RegionCloseRequest::default()),
|
||||
});
|
||||
|
||||
let version_data = version_control.current();
|
||||
|
||||
@@ -1307,7 +1307,10 @@ pub async fn reopen_region(
|
||||
) {
|
||||
// Close the region.
|
||||
engine
|
||||
.handle_request(region_id, RegionRequest::Close(RegionCloseRequest {}))
|
||||
.handle_request(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest::default()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
|
||||
@@ -159,6 +159,15 @@ impl<S: LogStore> Wal<S> {
|
||||
.map_err(BoxedError::new)
|
||||
.context(DeleteWalSnafu { region_id })
|
||||
}
|
||||
/// Marks all WAL entries of a region as obsolete and removes its dedicated namespace when
|
||||
/// supported by the backend.
|
||||
pub async fn obsolete_all(&self, region_id: RegionId, provider: &Provider) -> Result<()> {
|
||||
self.store
|
||||
.obsolete_all(provider, region_id)
|
||||
.await
|
||||
.map_err(BoxedError::new)
|
||||
.context(DeleteWalSnafu { region_id })
|
||||
}
|
||||
}
|
||||
|
||||
/// WAL batch writer.
|
||||
|
||||
@@ -187,6 +187,14 @@ mod tests {
|
||||
unreachable!()
|
||||
}
|
||||
|
||||
async fn obsolete_all(
|
||||
&self,
|
||||
_provider: &Provider,
|
||||
_region_id: RegionId,
|
||||
) -> Result<(), Self::Error> {
|
||||
unreachable!()
|
||||
}
|
||||
|
||||
fn entry(
|
||||
&self,
|
||||
_data: Vec<u8>,
|
||||
|
||||
@@ -1148,8 +1148,9 @@ impl<S: LogStore> RegionWorkerLoop<S> {
|
||||
.await;
|
||||
continue;
|
||||
}
|
||||
DdlRequest::Close(_) => {
|
||||
self.handle_close_request(ddl.region_id, ddl.sender).await;
|
||||
DdlRequest::Close(req) => {
|
||||
self.handle_close_request(ddl.region_id, req, ddl.sender)
|
||||
.await;
|
||||
continue;
|
||||
}
|
||||
DdlRequest::Alter(req) => {
|
||||
|
||||
@@ -27,6 +27,7 @@ impl<S: LogStore> RegionWorkerLoop<S> {
|
||||
pub(crate) async fn handle_close_request(
|
||||
&mut self,
|
||||
region_id: RegionId,
|
||||
request: RegionCloseRequest,
|
||||
sender: OptionOutputTx,
|
||||
) {
|
||||
let Some(region) = self.regions.get_region(region_id) else {
|
||||
@@ -36,9 +37,10 @@ impl<S: LogStore> RegionWorkerLoop<S> {
|
||||
|
||||
info!("Try to close region {}, worker: {}", region_id, self.id);
|
||||
|
||||
// If the region is using Noop WAL and has data in memtable and region is flushable (like,
|
||||
// not in follower state), we should flush it before closing to ensure durability.
|
||||
if region.provider == Provider::Noop
|
||||
// If the close request asks for a flush, or the region is using Noop WAL,
|
||||
// and has data in memtable and region is flushable (like, not in follower state),
|
||||
// we should flush it before closing to ensure durability.
|
||||
if (request.flush_on_close || region.provider == Provider::Noop)
|
||||
&& !region
|
||||
.version_control
|
||||
.current()
|
||||
@@ -53,7 +55,7 @@ impl<S: LogStore> RegionWorkerLoop<S> {
|
||||
.add_ddl_request_to_pending(SenderDdlRequest {
|
||||
region_id,
|
||||
sender,
|
||||
request: DdlRequest::Close(RegionCloseRequest {}),
|
||||
request: DdlRequest::Close(request),
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -86,6 +86,14 @@ pub trait LogStore: Send + Sync + 'static + std::fmt::Debug {
|
||||
entry_id: EntryId,
|
||||
) -> Result<(), Self::Error>;
|
||||
|
||||
/// Marks all entries of a region as obsolete and removes its dedicated namespace when
|
||||
/// supported by the backend.
|
||||
async fn obsolete_all(
|
||||
&self,
|
||||
provider: &Provider,
|
||||
region_id: RegionId,
|
||||
) -> Result<(), Self::Error>;
|
||||
|
||||
/// Makes an entry instance of the associated Entry type
|
||||
fn entry(
|
||||
&self,
|
||||
|
||||
@@ -189,6 +189,10 @@ impl RegionRequest {
|
||||
reason: "RemoteDynFilter request should be handled separately by RegionServer",
|
||||
}
|
||||
.fail(),
|
||||
region_request::Body::CleanUp(_) => UnexpectedSnafu {
|
||||
reason: "CleanUp request should be handled separately by RegionServer",
|
||||
}
|
||||
.fail(),
|
||||
region_request::Body::ApplyStagingManifest(apply) => {
|
||||
make_region_apply_staging_manifest(apply)
|
||||
}
|
||||
@@ -329,7 +333,9 @@ fn make_region_close(close: CloseRequest) -> Result<Vec<(RegionId, RegionRequest
|
||||
let region_id = close.region_id.into();
|
||||
Ok(vec![(
|
||||
region_id,
|
||||
RegionRequest::Close(RegionCloseRequest {}),
|
||||
RegionRequest::Close(RegionCloseRequest {
|
||||
flush_on_close: close.flush_on_close,
|
||||
}),
|
||||
)])
|
||||
}
|
||||
|
||||
@@ -647,8 +653,11 @@ impl RegionOpenRequest {
|
||||
}
|
||||
|
||||
/// Close region request.
|
||||
#[derive(Debug)]
|
||||
pub struct RegionCloseRequest {}
|
||||
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
|
||||
pub struct RegionCloseRequest {
|
||||
/// Whether to flush the region before closing it.
|
||||
pub flush_on_close: bool,
|
||||
}
|
||||
|
||||
/// Alter metadata of a region.
|
||||
#[derive(Debug, PartialEq, Eq, Clone)]
|
||||
|
||||
Reference in New Issue
Block a user