perf(ecstore): gate bounded GET metadata fanout (#5917)

Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
houseme
2026-08-10 13:20:21 +08:00
committed by GitHub
parent 785ee719e7
commit 88e285c523
3 changed files with 299 additions and 23 deletions
+236 -18
View File
@@ -461,6 +461,34 @@ impl MetadataQuorumAccumulator {
None None
} }
pub(in crate::set_disk) fn can_still_reach_early_stop_with_pending(&self, pending: usize) -> bool {
if !self.allow_early_stop {
return false;
}
if self.delete_marker_votes.saturating_add(pending) >= self.default_write_quorum() {
return true;
}
if self.conflicting_metadata
|| self.delete_marker_seen
|| self.not_found_responses > 0
|| self.version_not_found_responses > 0
|| self.hard_errors > 0
{
return false;
}
if !self.requested_version_id.is_empty()
&& self.matching_version_votes.saturating_add(pending) >= self.read_quorum_for_version()
{
return true;
}
match &self.candidate {
Some(candidate) => self
.candidate_latest_quorum(candidate)
.is_some_and(|latest_quorum| self.candidate_votes.saturating_add(pending) >= latest_quorum),
None => pending >= self.default_write_quorum(),
}
}
/// Compute the read quorum threshold for version-aware early-stop. /// Compute the read quorum threshold for version-aware early-stop.
/// Uses `total_disks / 2` (like `missing_response_quorum`) when /// Uses `total_disks / 2` (like `missing_response_quorum`) when
/// `default_parity_count` is set, otherwise requires all disks. /// `default_parity_count` is set, otherwise requires all disks.
@@ -1982,7 +2010,7 @@ pub(in crate::set_disk) fn should_allow_metadata_early_stop(
healing: bool, healing: bool,
incl_free_versions: bool, incl_free_versions: bool,
) -> bool { ) -> bool {
if read_data { if read_data && !is_get_metadata_data_read_early_stop_enabled() {
return false; return false;
} }
@@ -2289,23 +2317,39 @@ impl SetDisks {
let object = Arc::new(object.to_string()); let object = Arc::new(object.to_string());
let version_id = Arc::new(version_id.to_string()); let version_id = Arc::new(version_id.to_string());
let mut join_set = JoinSet::new(); let mut join_set = JoinSet::new();
let bounded_fanout = is_get_metadata_early_stop_bounded_fanout_enabled();
let mut next_disk_index = 0usize;
let spawn_read_version =
|join_set: &mut JoinSet<(usize, disk::error::Result<FileInfo>, Duration)>, index: usize, disk: Option<DiskStore>| {
let opts = opts.clone();
let org_bucket = org_bucket.clone();
let bucket = bucket.clone();
let object = object.clone();
let version_id = version_id.clone();
join_set.spawn(async move {
let response_start = Instant::now();
let result = if let Some(disk) = disk {
Self::record_read_version_call(&object, index);
disk.read_version(&org_bucket, &bucket, &object, &version_id, &opts).await
} else {
Err(DiskError::DiskNotFound)
};
(index, result, response_start.elapsed())
});
};
for (index, disk) in disks.iter().cloned().enumerate() { if bounded_fanout {
let opts = opts.clone(); let initial_target = accumulator.default_write_quorum().min(disks.len());
let org_bucket = org_bucket.clone(); while next_disk_index < initial_target {
let bucket = bucket.clone(); if let Some(disk) = disks.get(next_disk_index).cloned() {
let object = object.clone(); spawn_read_version(&mut join_set, next_disk_index, disk);
let version_id = version_id.clone(); }
join_set.spawn(async move { next_disk_index = next_disk_index.saturating_add(1);
let response_start = Instant::now(); }
let result = if let Some(disk) = disk { } else {
Self::record_read_version_call(&object, index); for (index, disk) in disks.iter().cloned().enumerate() {
disk.read_version(&org_bucket, &bucket, &object, &version_id, &opts).await spawn_read_version(&mut join_set, index, disk);
} else { }
Err(DiskError::DiskNotFound)
};
(index, result, response_start.elapsed())
});
} }
while let Some(result) = join_set.join_next().await { while let Some(result) = join_set.join_next().await {
@@ -2337,7 +2381,11 @@ impl SetDisks {
.early_stop_decision() .early_stop_decision()
.or_else(|| accumulator.version_early_stop_decision()) .or_else(|| accumulator.version_early_stop_decision())
{ {
let saved_responses = join_set.len(); let saved_responses = if bounded_fanout {
disks.len().saturating_sub(observations.len())
} else {
join_set.len()
};
join_set.abort_all(); join_set.abort_all();
rustfs_io_metrics::record_get_object_metadata_early_stop_hit(GET_OBJECT_PATH_LEGACY_DUPLEX, decision.reason); rustfs_io_metrics::record_get_object_metadata_early_stop_hit(GET_OBJECT_PATH_LEGACY_DUPLEX, decision.reason);
rustfs_io_metrics::record_get_object_metadata_early_stop_saved_responses( rustfs_io_metrics::record_get_object_metadata_early_stop_saved_responses(
@@ -2348,6 +2396,16 @@ impl SetDisks {
let diagnostics = MetadataFanoutDiagnostics::new(fanout_start.elapsed(), observations); let diagnostics = MetadataFanoutDiagnostics::new(fanout_start.elapsed(), observations);
return Ok((ress, errors, diagnostics)); return Ok((ress, errors, diagnostics));
} }
if bounded_fanout
&& next_disk_index < disks.len()
&& !accumulator.can_still_reach_early_stop_with_pending(join_set.len())
{
if let Some(disk) = disks.get(next_disk_index).cloned() {
spawn_read_version(&mut join_set, next_disk_index, disk);
}
next_disk_index = next_disk_index.saturating_add(1);
}
} }
rustfs_io_metrics::record_get_object_metadata_early_stop_miss( rustfs_io_metrics::record_get_object_metadata_early_stop_miss(
@@ -5254,6 +5312,166 @@ mod tests {
drop(dirs); drop(dirs);
} }
fn valid_metadata_fanout_fileinfo(
bucket: &str,
object: &str,
version_id: Uuid,
data_dir: Uuid,
mod_time: OffsetDateTime,
) -> FileInfo {
let mut fi = FileInfo::new(object, 2, 2);
fi.volume = bucket.to_string();
fi.name = object.to_string();
fi.size = 1;
fi.erasure.index = 1;
fi.version_id = Some(version_id);
fi.is_latest = true;
fi.data_dir = Some(data_dir);
fi.mod_time = Some(mod_time);
fi.metadata.insert("etag".to_string(), "etag-1".to_string());
fi.add_object_part(1, "part-etag".to_string(), 1, fi.mod_time, 1, None, None);
fi
}
async fn install_metadata_fanout_fileinfo(
disks: &[Option<DiskStore>],
bucket: &str,
object: &str,
missing_part_disk: Option<usize>,
) {
let version_id = Uuid::new_v4();
let data_dir = Uuid::new_v4();
let mod_time = OffsetDateTime::now_utc();
for (index, disk) in disks
.iter()
.enumerate()
.filter_map(|(index, disk)| disk.as_ref().map(|disk| (index, disk)))
{
if missing_part_disk != Some(index) {
disk.write_all(bucket, &format!("{object}/{data_dir}/part.1"), Bytes::from_static(b"x"))
.await
.expect("part data should be installed on every disk");
}
disk.write_metadata(
bucket,
bucket,
object,
valid_metadata_fanout_fileinfo(bucket, object, version_id, data_dir, mod_time),
)
.await
.expect("metadata should be installed on every disk");
}
}
#[tokio::test]
async fn bounded_metadata_early_stop_ab_limits_data_get_read_version_fanout() {
const DISKS: usize = 4;
let bucket = "bounded-data-get-fanout-bucket";
let control_object = "bounded-data-get-control-object";
let treatment_object = "bounded-data-get-treatment-object";
let (dirs, disks) = call_counter_local_disks(bucket, DISKS).await;
install_metadata_fanout_fileinfo(&disks, bucket, control_object, None).await;
install_metadata_fanout_fileinfo(&disks, bucket, treatment_object, None).await;
temp_env::async_with_vars(
[
("RUSTFS_GET_METADATA_EARLY_STOP_ENABLE", Some("true")),
("RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE", None),
("RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT", Some("true")),
],
async {
let calls = disk_call_counters::observe(control_object);
let (_, _, diagnostics) =
SetDisks::read_all_fileinfo_observed(&disks, bucket, bucket, control_object, "", true, false, false, true, 2)
.await
.expect("control metadata should resolve");
assert_eq!(
calls.total(disk_call_counters::KIND_READ_VERSION),
DISKS as u64,
"control path should keep the default data-read full fanout"
);
assert_eq!(diagnostics.total_responses(), DISKS);
},
)
.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 calls = disk_call_counters::observe(treatment_object);
let (parts_metadata, errs, diagnostics) = SetDisks::read_all_fileinfo_observed(
&disks,
bucket,
bucket,
treatment_object,
"",
true,
false,
false,
true,
2,
)
.await
.expect("healthy object metadata should reach early-stop quorum");
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"
);
assert_eq!(diagnostics.total_responses(), 3);
assert_eq!(parts_metadata.iter().filter(|fi| fi.name == treatment_object).count(), 3);
assert!(errs.iter().all(Option::is_none));
},
)
.await;
drop(dirs);
}
#[tokio::test]
async fn bounded_metadata_early_stop_falls_back_to_full_fanout_on_data_read_error() {
const DISKS: usize = 4;
let bucket = "bounded-data-get-error-bucket";
let object = "bounded-data-get-error-object";
let (dirs, disks) = call_counter_local_disks(bucket, DISKS).await;
install_metadata_fanout_fileinfo(&disks, bucket, object, Some(0)).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 calls = disk_call_counters::observe(object);
let (_, errs, diagnostics) =
SetDisks::read_all_fileinfo_observed(&disks, bucket, bucket, object, "", true, false, false, true, 2)
.await
.expect("metadata fanout should complete after falling back to all disks");
assert_eq!(
calls.total(disk_call_counters::KIND_READ_VERSION),
DISKS as u64,
"a data-read error must force bounded fanout to schedule every disk before returning"
);
assert_eq!(diagnostics.total_responses(), DISKS);
assert!(
errs.iter()
.any(|err| err.as_ref().is_some_and(|err| matches!(err, DiskError::FileNotFound)))
);
},
)
.await;
drop(dirs);
}
/// Bound for the pause handshake. This is a hang-guard, not a timing /// Bound for the pause handshake. This is a hang-guard, not a timing
/// dependency: under a working barrier `wait_until_paused` returns via the /// dependency: under a working barrier `wait_until_paused` returns via the
/// `Notify` handshake far below this bound regardless of IO pressure, so the /// `Notify` handshake far below this bound regardless of IO pressure, so the
+49 -3
View File
@@ -673,9 +673,9 @@ const DEFAULT_RUSTFS_GET_SMALL_OBJECT_DIRECT_MEMORY_THRESHOLD: usize = 128 * 102
const ENV_RUSTFS_GET_METADATA_EARLY_STOP_ENABLE: &str = "RUSTFS_GET_METADATA_EARLY_STOP_ENABLE"; const ENV_RUSTFS_GET_METADATA_EARLY_STOP_ENABLE: &str = "RUSTFS_GET_METADATA_EARLY_STOP_ENABLE";
// Enabled by default (backlog#872): the early-stop path only engages for // Enabled by default (backlog#872): the early-stop path only engages for
// requests `should_allow_metadata_early_stop` classifies as safe (metadata-only // requests `should_allow_metadata_early_stop` classifies as safe (metadata-only
// reads without version_id / healing / free-version needs) and still requires // reads by default, without version_id / healing / free-version needs) and
// a full read-quorum agreement before stopping. Set the env var to `false` to // still requires a full read-quorum agreement before stopping. Set the env var
// fall back to full-wait metadata fanout. // to `false` to fall back to full-wait metadata fanout.
const DEFAULT_RUSTFS_GET_METADATA_EARLY_STOP_ENABLE: bool = true; const DEFAULT_RUSTFS_GET_METADATA_EARLY_STOP_ENABLE: bool = true;
const ENV_RUSTFS_GET_METADATA_EARLY_STOP_ROLLOUT_PCT: &str = "RUSTFS_GET_METADATA_EARLY_STOP_ROLLOUT_PCT"; const ENV_RUSTFS_GET_METADATA_EARLY_STOP_ROLLOUT_PCT: &str = "RUSTFS_GET_METADATA_EARLY_STOP_ROLLOUT_PCT";
@@ -684,6 +684,12 @@ const DEFAULT_RUSTFS_GET_METADATA_EARLY_STOP_ROLLOUT_PCT: u32 = 100;
const ENV_RUSTFS_GET_METADATA_VERSION_EARLY_STOP_ENABLE: &str = "RUSTFS_GET_METADATA_VERSION_EARLY_STOP_ENABLE"; const ENV_RUSTFS_GET_METADATA_VERSION_EARLY_STOP_ENABLE: &str = "RUSTFS_GET_METADATA_VERSION_EARLY_STOP_ENABLE";
const DEFAULT_RUSTFS_GET_METADATA_VERSION_EARLY_STOP_ENABLE: bool = false; const DEFAULT_RUSTFS_GET_METADATA_VERSION_EARLY_STOP_ENABLE: bool = false;
const ENV_RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE: &str = "RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE";
const DEFAULT_RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE: bool = false;
const ENV_RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT: &str = "RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT";
const DEFAULT_RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT: bool = false;
// --- Multipart Reader-Setup Prefetch Configuration (backlog#870) --- // --- Multipart Reader-Setup Prefetch Configuration (backlog#870) ---
const ENV_RUSTFS_GET_MULTIPART_READER_SETUP_PREFETCH: &str = "RUSTFS_GET_MULTIPART_READER_SETUP_PREFETCH"; const ENV_RUSTFS_GET_MULTIPART_READER_SETUP_PREFETCH: &str = "RUSTFS_GET_MULTIPART_READER_SETUP_PREFETCH";
@@ -1194,6 +1200,46 @@ fn is_version_early_stop_enabled() -> bool {
} }
} }
fn is_get_metadata_data_read_early_stop_enabled() -> bool {
#[cfg(test)]
{
rustfs_utils::get_env_bool(
ENV_RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE,
DEFAULT_RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE,
)
}
#[cfg(not(test))]
{
static CACHED: OnceLock<bool> = OnceLock::new();
*CACHED.get_or_init(|| {
rustfs_utils::get_env_bool(
ENV_RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE,
DEFAULT_RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE,
)
})
}
}
fn is_get_metadata_early_stop_bounded_fanout_enabled() -> bool {
#[cfg(test)]
{
rustfs_utils::get_env_bool(
ENV_RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT,
DEFAULT_RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT,
)
}
#[cfg(not(test))]
{
static CACHED: OnceLock<bool> = OnceLock::new();
*CACHED.get_or_init(|| {
rustfs_utils::get_env_bool(
ENV_RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT,
DEFAULT_RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT,
)
})
}
}
/// Check if multipart reads prefetch the next part's bitrot reader setup /// Check if multipart reads prefetch the next part's bitrot reader setup
/// while the current part decodes (backlog#870). /// while the current part decodes (backlog#870).
/// ///
+14 -2
View File
@@ -3834,18 +3834,19 @@ mod tests {
assert!(metadata_early_stop_permitted(true, true, false, "", false, false)); assert!(metadata_early_stop_permitted(true, true, false, "", false, false));
// observe=false (non-observed fanout) also disables early-stop. // observe=false (non-observed fanout) also disables early-stop.
assert!(!metadata_early_stop_permitted(true, false, false, "", false, false)); assert!(!metadata_early_stop_permitted(true, false, false, "", false, false));
// Data reads are never eligible regardless of caller opt-in. // Data reads require their own explicit rollout gate.
assert!(!metadata_early_stop_permitted(true, true, true, "", false, false)); assert!(!metadata_early_stop_permitted(true, true, true, "", false, false));
}, },
); );
} }
#[test] #[test]
fn metadata_early_stop_rejects_data_reads() { fn metadata_early_stop_requires_explicit_data_read_opt_in() {
temp_env::with_vars( temp_env::with_vars(
[ [
(ENV_RUSTFS_GET_METADATA_EARLY_STOP_ENABLE, Some("true")), (ENV_RUSTFS_GET_METADATA_EARLY_STOP_ENABLE, Some("true")),
(ENV_RUSTFS_GET_METADATA_VERSION_EARLY_STOP_ENABLE, Some("true")), (ENV_RUSTFS_GET_METADATA_VERSION_EARLY_STOP_ENABLE, Some("true")),
(ENV_RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE, None),
], ],
|| { || {
assert!(!should_allow_metadata_early_stop(true, "", false, false)); assert!(!should_allow_metadata_early_stop(true, "", false, false));
@@ -3854,6 +3855,17 @@ mod tests {
assert!(should_allow_metadata_early_stop(false, "version-id", false, false)); assert!(should_allow_metadata_early_stop(false, "version-id", false, false));
}, },
); );
temp_env::with_vars(
[
(ENV_RUSTFS_GET_METADATA_EARLY_STOP_ENABLE, Some("true")),
(ENV_RUSTFS_GET_METADATA_VERSION_EARLY_STOP_ENABLE, Some("true")),
(ENV_RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE, Some("true")),
],
|| {
assert!(should_allow_metadata_early_stop(true, "", false, false));
assert!(should_allow_metadata_early_stop(true, "version-id", false, false));
},
);
} }
#[test] #[test]