From 924958bab5dea0a039b1df2ce5c0fe6ac575e792 Mon Sep 17 00:00:00 2001 From: houseme Date: Wed, 12 Aug 2026 10:38:13 +0800 Subject: [PATCH] perf(get): slim metadata fanout allocations (#1803) (#5968) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Every GET fans out a `read_version` across all disks to resolve xl.meta. Each fanout allocated an `Arc` (3 bools) plus four `Arc` (`Arc::new(x.to_string())` = two allocations each) and cloned them into every spawned task. This trims the per-fanout allocation footprint. - `ReadOptions` is three bools, so it is now `Copy`. The fanout drops the `Arc` and hands each spawned task a copy; the two pre-existing `ReadOptions::clone()` sites (set_disk/read.rs, set_disk/ops/heal.rs) stop cloning a `Copy` type. - The four request strings use `Arc::::from(&str)` (one allocation each) instead of `Arc::new(..to_string())` (string buffer + Arc = two each) — four fewer allocations per fanout, transparent to the `read_version(&str)` call. Behavior is unchanged: the fanout still spawns one task per disk (the spawn is deliberate — `read_version_call_counter_observes_spawned_fanout` verifies the process-global counter observes every per-disk increment across workers), quorum / early-stop / full-wait semantics are untouched, and no result ordering or error handling changed. Two larger items from the audit are intentionally NOT in this PR: - `tokio::spawn` -> `FuturesUnordered`: the spawn is a tested, deliberate design (cross-worker counter observation for #1309/#1314), and converting would also change panic isolation. Left as-is. - `vec![FileInfo::default(); N]`: `FileInfo`'s empty containers (String / HashMap / Vec) do not allocate, so this is one `Vec` allocation, not the per-element allocation the audit implied — not a real hot spot. `cargo fmt`, `cargo clippy -p rustfs-ecstore --lib` (0 warnings), `cargo check --lib --tests`, and the 26 fanout / call-counter unit tests pass on macOS (the change is fully cross-platform). Co-authored-by: heihutu Co-authored-by: zhi22915 --- crates/ecstore/src/disk/mod.rs | 2 +- .../src/set_disk/core/io_primitives.rs | 34 ++++++++++--------- crates/ecstore/src/set_disk/ops/heal.rs | 4 +-- crates/ecstore/src/set_disk/read.rs | 6 ++-- 4 files changed, 24 insertions(+), 22 deletions(-) diff --git a/crates/ecstore/src/disk/mod.rs b/crates/ecstore/src/disk/mod.rs index 8112aa194..0f9c28bee 100644 --- a/crates/ecstore/src/disk/mod.rs +++ b/crates/ecstore/src/disk/mod.rs @@ -1251,7 +1251,7 @@ pub struct VolumeInfo { pub created: Option, } -#[derive(Deserialize, Serialize, Debug, Default, Clone)] +#[derive(Deserialize, Serialize, Debug, Default, Clone, Copy)] pub struct ReadOptions { pub incl_free_versions: bool, pub read_data: bool, diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index a75e74218..c41385fae 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -2222,18 +2222,18 @@ impl SetDisks { let mut ress = Vec::with_capacity(disks.len()); let mut errors = Vec::with_capacity(disks.len()); let mut observations = observe.then(|| Vec::with_capacity(disks.len())); - let opts = Arc::new(ReadOptions { + let opts = ReadOptions { incl_free_versions, read_data, healing, - }); - let org_bucket = Arc::new(org_bucket.to_string()); - let bucket = Arc::new(bucket.to_string()); - let object = Arc::new(object.to_string()); - let version_id = Arc::new(version_id.to_string()); + }; + let org_bucket: Arc = Arc::from(org_bucket); + let bucket: Arc = Arc::from(bucket); + let object: Arc = Arc::from(object); + let version_id: Arc = Arc::from(version_id); let futures = disks.iter().enumerate().map(|(disk_index, disk)| { let disk = disk.clone(); - let opts = opts.clone(); + let task_opts = opts; let org_bucket = org_bucket.clone(); let bucket = bucket.clone(); let object = object.clone(); @@ -2242,7 +2242,8 @@ impl SetDisks { let response_start = observe.then(Instant::now); let result = if let Some(disk) = disk { Self::record_read_version_call(&object, disk_index); - disk.read_version(&org_bucket, &bucket, &object, &version_id, &opts).await + disk.read_version(&org_bucket, &bucket, &object, &version_id, &task_opts) + .await } else { Err(DiskError::DiskNotFound) }; @@ -2307,21 +2308,21 @@ impl SetDisks { let mut observations = Vec::with_capacity(disks.len()); let mut accumulator = MetadataQuorumAccumulator::new(disks.len(), default_parity_count, true).with_requested_version_id(version_id); - let opts = Arc::new(ReadOptions { + let opts = ReadOptions { incl_free_versions, read_data, healing, - }); - let org_bucket = Arc::new(org_bucket.to_string()); - let bucket = Arc::new(bucket.to_string()); - let object = Arc::new(object.to_string()); - let version_id = Arc::new(version_id.to_string()); + }; + let org_bucket: Arc = Arc::from(org_bucket); + let bucket: Arc = Arc::from(bucket); + let object: Arc = Arc::from(object); + let version_id: Arc = Arc::from(version_id); 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, Duration)>, index: usize, disk: Option| { - let opts = opts.clone(); + let task_opts = opts; let org_bucket = org_bucket.clone(); let bucket = bucket.clone(); let object = object.clone(); @@ -2332,7 +2333,8 @@ impl SetDisks { 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 + disk.read_version(&org_bucket, &bucket, &object, &version_id, &task_opts) + .await } else { Err(DiskError::DiskNotFound) }; diff --git a/crates/ecstore/src/set_disk/ops/heal.rs b/crates/ecstore/src/set_disk/ops/heal.rs index 10d159781..4e6c81d96 100644 --- a/crates/ecstore/src/set_disk/ops/heal.rs +++ b/crates/ecstore/src/set_disk/ops/heal.rs @@ -362,9 +362,9 @@ impl SetDisks { healing: true, }; let checks = target_disks.into_iter().map(|disk| { - let read_options = read_options.clone(); + let task_read_options = read_options; async move { - let file_info = match disk.read_version("", bucket, object, version_id, &read_options).await { + let file_info = match disk.read_version("", bucket, object, version_id, &task_read_options).await { Ok(file_info) => file_info, Err( DiskError::DiskNotFound diff --git a/crates/ecstore/src/set_disk/read.rs b/crates/ecstore/src/set_disk/read.rs index ae7461f75..2404a567a 100644 --- a/crates/ecstore/src/set_disk/read.rs +++ b/crates/ecstore/src/set_disk/read.rs @@ -224,7 +224,7 @@ impl SetDisks { let bucket = bucket.to_string(); let object = object.to_string(); let version_id = version_id.to_string(); - let opts = opts.clone(); + let opts = *opts; let processor = runtime_sources::batch_processors().read_processor(); let tasks: Vec<_> = disks @@ -235,9 +235,9 @@ impl SetDisks { let bucket = bucket.clone(); let object = object.clone(); let version_id = version_id.clone(); - let opts = opts.clone(); + let task_opts = opts; - async move { disk.read_version(&bucket, &bucket, &object, &version_id, &opts).await } + async move { disk.read_version(&bucket, &bucket, &object, &version_id, &task_opts).await } }) }) .collect();