perf(ecstore): hedge bounded GET metadata fanout (#5935)

Keep opt-in bounded GET data-read fanout from waiting on a single pending ReadVersion response when an unscheduled spare disk can satisfy quorum. Add a deterministic 2+2 regression that pauses the third scheduled metadata read and verifies the spare is started before returning.

Co-authored-by: heihutu <heihutu@gmail.com>
Co-authored-by: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
houseme
2026-08-11 01:11:04 +08:00
committed by GitHub
parent 1148e76279
commit 3747d19ce5
@@ -2330,6 +2330,8 @@ impl SetDisks {
let response_start = Instant::now();
let result = if let Some(disk) = disk {
Self::record_read_version_call(&object, index);
#[cfg(test)]
Self::read_version_fanout_barrier(&object, index).await;
disk.read_version(&org_bucket, &bucket, &object, &version_id, &opts).await
} else {
Err(DiskError::DiskNotFound)
@@ -2397,9 +2399,13 @@ impl SetDisks {
return Ok((ress, errors, diagnostics));
}
let pending_responses = join_set.len();
let should_hedge_single_pending_data_read =
read_data && pending_responses == 1 && accumulator.can_still_reach_early_stop_with_pending(pending_responses);
if bounded_fanout
&& next_disk_index < disks.len()
&& !accumulator.can_still_reach_early_stop_with_pending(join_set.len())
&& (!accumulator.can_still_reach_early_stop_with_pending(pending_responses)
|| should_hedge_single_pending_data_read)
{
if let Some(disk) = disks.get(next_disk_index).cloned() {
spawn_read_version(&mut join_set, next_disk_index, disk);
@@ -3317,6 +3323,12 @@ impl SetDisks {
#[inline(always)]
fn record_read_version_call(_object: &str, _disk_index: usize) {}
#[cfg(test)]
#[inline]
async fn read_version_fanout_barrier(object: &str, disk_index: usize) {
rename_fanout_barrier::checkpoint(object, disk_index, rename_fanout_barrier::PHASE_READ_VERSION).await;
}
/// Test-only awaitable pause point for the rename/commit fan-out (backlog#1325,
/// serving the barrier-style acceptances of #1312 / #1319 / #1313). `phase` is
/// [`rename_fanout_barrier::PHASE_RENAME`] or `PHASE_CLEANUP`. When a test has
@@ -4803,6 +4815,8 @@ pub(in crate::set_disk) mod rename_fanout_barrier_phase {
pub const RENAME: &str = "rename";
/// The per-disk old-data-dir cleanup phase of the commit fan-out.
pub const CLEANUP: &str = "cleanup";
/// The per-disk `read_version` phase of metadata read fan-out.
pub const READ_VERSION: &str = "read_version";
}
/// Test-only awaitable pause barrier + background-task introspection for the
@@ -4846,7 +4860,9 @@ pub(in crate::set_disk) mod rename_fanout_barrier {
use std::sync::{Arc, Mutex, OnceLock};
use tokio::sync::Notify;
pub use super::rename_fanout_barrier_phase::{CLEANUP as PHASE_CLEANUP, RENAME as PHASE_RENAME};
pub use super::rename_fanout_barrier_phase::{
CLEANUP as PHASE_CLEANUP, READ_VERSION as PHASE_READ_VERSION, RENAME as PHASE_RENAME,
};
/// One armed barrier: the fan-out task matching `(disk_index, phase)` pauses.
struct Armed {
@@ -5364,7 +5380,7 @@ mod tests {
}
#[tokio::test]
async fn bounded_metadata_early_stop_ab_limits_data_get_read_version_fanout() {
async fn bounded_metadata_early_stop_ab_hedges_data_get_read_version_fanout() {
const DISKS: usize = 4;
let bucket = "bounded-data-get-fanout-bucket";
let control_object = "bounded-data-get-control-object";
@@ -5419,13 +5435,73 @@ mod tests {
.await
.expect("healthy object metadata should reach early-stop quorum");
assert!(
(3..=DISKS as u64).contains(&calls.total(disk_call_counters::KIND_READ_VERSION)),
"healthy 2+2 bounded data-read fanout may finish at quorum before a spare hedge is needed"
);
assert!(
(3..=DISKS).contains(&diagnostics.total_responses()),
"treatment path should return after reaching quorum, with at most the spare hedge response observed"
);
assert!(parts_metadata.iter().filter(|fi| fi.name == treatment_object).count() >= 3);
assert!(errs.iter().all(Option::is_none));
},
)
.await;
drop(dirs);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn bounded_data_get_hedges_single_pending_read_version() {
const DISKS: usize = 4;
let bucket = "bounded-data-get-hedge-bucket";
let object = "bounded-data-get-hedge-object";
let (dirs, disks) = call_counter_local_disks(bucket, DISKS).await;
install_metadata_fanout_fileinfo(&disks, bucket, object, None).await;
temp_env::async_with_vars(
[
("RUSTFS_GET_METADATA_EARLY_STOP_ENABLE", Some("true")),
("RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE", Some("true")),
("RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT", Some("true")),
],
async {
let barrier = rename_fanout_barrier::arm(object, 2, rename_fanout_barrier::PHASE_READ_VERSION);
let calls = disk_call_counters::observe(object);
let disks_for_read = disks.clone();
let mut read = tokio::spawn(async move {
SetDisks::read_all_fileinfo_observed(&disks_for_read, bucket, bucket, object, "", true, false, false, true, 2)
.await
});
tokio::time::timeout(BARRIER_PAUSE_GUARD, barrier.wait_until_paused())
.await
.expect("third scheduled read_version should pause at the deterministic barrier");
tokio::time::timeout(BARRIER_PAUSE_GUARD, async {
while calls.for_disk(disk_call_counters::KIND_READ_VERSION, 3) == 0 {
tokio::task::yield_now().await;
}
})
.await
.expect("bounded data-read fanout should hedge by starting the spare disk");
let completed = tokio::time::timeout(BARRIER_PAUSE_GUARD, &mut read).await;
if completed.is_err() {
barrier.release();
}
let (parts_metadata, errs, diagnostics) = completed
.expect("spare metadata should allow early-stop without waiting for the paused disk")
.expect("metadata read task should not panic")
.expect("healthy spare metadata should resolve");
assert_eq!(
calls.total(disk_call_counters::KIND_READ_VERSION),
3,
"treatment path should stop after the 2+2 read/write quorum instead of issuing every disk read"
DISKS as u64,
"bounded data-read fanout should issue the paused disk plus one spare hedge"
);
assert_eq!(diagnostics.total_responses(), 3);
assert_eq!(parts_metadata.iter().filter(|fi| fi.name == treatment_object).count(), 3);
assert_eq!(parts_metadata.iter().filter(|fi| fi.name == object).count(), 3);
assert!(errs.iter().all(Option::is_none));
},
)