From f219a0aba8eb3a4ff643e158b0b4d39c92b8d506 Mon Sep 17 00:00:00 2001 From: GatewayJ <835269233@qq.com> Date: Wed, 23 Sep 2026 22:19:23 +0800 Subject: [PATCH] perf(ecstore): parallelize multipart I/O setup and metadata reads (#8085) * perf(ecstore): parallelize multipart I/O setup and metadata reads * perf(ecstore): share multipart paths and increase read concurrency * fix(ecstore): route multipart benchmark through storage API facade --- crates/ecstore/Cargo.toml | 4 + .../benches/multipart_read_parts_benchmark.rs | 92 ++++++++ crates/ecstore/benches/storage_api/mod.rs | 4 + crates/ecstore/src/disk/local.rs | 196 ++++++++++++------ .../src/set_disk/core/io_primitives.rs | 8 +- crates/ecstore/src/set_disk/ops/multipart.rs | 168 +++++++++++---- 6 files changed, 363 insertions(+), 109 deletions(-) create mode 100644 crates/ecstore/benches/multipart_read_parts_benchmark.rs diff --git a/crates/ecstore/Cargo.toml b/crates/ecstore/Cargo.toml index 4a8ec662a..a1b2d4a6c 100644 --- a/crates/ecstore/Cargo.toml +++ b/crates/ecstore/Cargo.toml @@ -282,5 +282,9 @@ harness = false name = "single_block_non_inline_benchmark" harness = false +[[bench]] +name = "multipart_read_parts_benchmark" +harness = false + [lib] doctest = false diff --git a/crates/ecstore/benches/multipart_read_parts_benchmark.rs b/crates/ecstore/benches/multipart_read_parts_benchmark.rs new file mode 100644 index 000000000..4bfaa8871 --- /dev/null +++ b/crates/ecstore/benches/multipart_read_parts_benchmark.rs @@ -0,0 +1,92 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use criterion::{BenchmarkId, Criterion, Throughput, criterion_group, criterion_main}; +use rustfs_filemeta::ObjectPartInfo; +use std::hint::black_box; +use std::time::Duration; + +mod storage_api; +use storage_api::multipart::{DiskAPI, DiskOption, Endpoint, new_disk}; + +fn bench_multipart_read_parts(c: &mut Criterion) { + let runtime = tokio::runtime::Builder::new_multi_thread() + .worker_threads(4) + .enable_all() + .build() + .expect("benchmark runtime"); + let root = tempfile::tempdir().expect("benchmark disk"); + let mut endpoint = Endpoint::try_from(root.path().to_str().expect("UTF-8 path")).expect("endpoint"); + endpoint.set_pool_index(0); + endpoint.set_set_index(0); + endpoint.set_disk_index(0); + let disk = runtime + .block_on(new_disk( + &endpoint, + &DiskOption { + cleanup: false, + health_check: false, + }, + )) + .expect("local disk"); + let bucket = "multipart-bench"; + runtime.block_on(disk.make_volume(bucket)).expect("benchmark volume"); + let upload = root.path().join(bucket).join("upload"); + std::fs::create_dir_all(&upload).expect("upload directory"); + let mut paths = Vec::with_capacity(1024); + for number in 1..=1024 { + let part = ObjectPartInfo { + number, + etag: format!("{number:032x}"), + size: 1024, + actual_size: 1024, + ..Default::default() + }; + std::fs::write(upload.join(format!("part.{number}")), b"data").expect("part data"); + std::fs::write(upload.join(format!("part.{number}.meta")), part.marshal_msg().expect("metadata")).expect("part metadata"); + paths.push(format!("upload/part.{number}.meta")); + } + + // Measures the real DiskAPI path against warm local files. Fixture creation + // and correctness checks stay outside the timed region; this is not a cold-disk benchmark. + let mut group = c.benchmark_group("multipart_read_parts"); + group.sample_size(20); + group.warm_up_time(Duration::from_secs(1)); + group.measurement_time(Duration::from_secs(2)); + for count in [1, 32, 256, 1024] { + let paths = &paths[..count]; + let parts = runtime.block_on(disk.read_parts(bucket, paths)).expect("read parts"); + assert_eq!(parts.len(), count); + assert!( + parts + .iter() + .enumerate() + .all(|(i, part)| part.number == i + 1 && part.error.is_none()) + ); + group.throughput(Throughput::Elements(u64::try_from(count).expect("part count"))); + group.bench_with_input(BenchmarkId::from_parameter(count), &paths, |b, paths| { + b.iter(|| { + black_box( + runtime + .block_on(disk.read_parts(bucket, black_box(paths))) + .expect("read parts"), + ) + }); + }); + } + group.finish(); +} + +criterion_group!(benches, bench_multipart_read_parts); +criterion_main!(benches); diff --git a/crates/ecstore/benches/storage_api/mod.rs b/crates/ecstore/benches/storage_api/mod.rs index 238f5132e..ae80a1d42 100644 --- a/crates/ecstore/benches/storage_api/mod.rs +++ b/crates/ecstore/benches/storage_api/mod.rs @@ -29,3 +29,7 @@ pub(crate) mod erasure { pub(crate) mod single_block_non_inline { pub(crate) use super::{BitrotWriterWrapper, CustomWriter, Erasure}; } + +pub(crate) mod multipart { + pub(crate) use rustfs_ecstore::api::disk::{DiskAPI, DiskOption, Endpoint, new_disk}; +} diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index 6c279059b..0af9065a1 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -53,6 +53,7 @@ use crate::disk::{ use crate::erasure::coding::{self, bitrot_verify}; use crate::runtime::sources as runtime_sources; use bytes::Bytes; +use futures::{StreamExt, TryStreamExt, stream}; use metrics::counter; #[cfg(target_os = "linux")] use metrics::gauge; @@ -88,6 +89,9 @@ use tokio::time::{Instant, Sleep, interval_at, timeout}; use tracing::{debug, error, info, warn}; use uuid::Uuid; +// Bound outstanding filesystem jobs and metadata buffers per disk request. +const PART_METADATA_READ_CONCURRENCY: usize = 8; + const DELETED_OBJECTS_CLEANUP_INTERVAL: Duration = Duration::from_secs(60 * 5); const STALE_TMP_OBJECT_EXPIRY: Duration = Duration::from_secs(24 * 60 * 60); @@ -6679,6 +6683,53 @@ impl LocalDisk { Ok(data) } + async fn read_part_metadata(&self, bucket: &str, volume_dir: &Path, path_str: &str) -> Result { + let path = Path::new(path_str); + let num = path + .file_name() + .and_then(|v| v.to_str()) + .unwrap_or_default() + .strip_prefix("part.") + .and_then(|v| v.strip_suffix(".meta")) + .and_then(|v| v.parse::().ok()) + .unwrap_or_default(); + let data_path = self.io_get_object_path( + bucket, + &path_join_buf(&[ + path.parent().unwrap_or_else(|| Path::new("")).to_string_lossy().as_ref(), + &format!("part.{num}"), + ]), + )?; + let metadata_path = self.io_get_object_path(bucket, path.to_string_lossy().as_ref()); + // A part's existence check, metadata read and decode share one dispatch. + // Keep open errors unmapped for the existing missing-volume fallback. + let result = tokio::task::spawn_blocking(move || -> Result<_> { + let part_error = |error: String| ObjectPartInfo { + number: num, + error: Some(error), + ..Default::default() + }; + if let Err(err) = std::fs::metadata(data_path) { + return Ok(Ok(part_error(err.to_string()))); + } + // Invalid metadata paths remain request errors, but missing data wins + // first, as it does in the serial reader. + let metadata_path = metadata_path?; + Ok(read_all_data_std(&metadata_path) + .map(|(data, _)| ObjectPartInfo::unmarshal(&data).unwrap_or_else(|err| part_error(err.to_string())))) + }) + .await; + let result = match result { + Ok(result) => self.resolve_read_all_result(bucket, volume_dir, result?).await, + Err(err) => Err(DiskError::from(err)), + }; + Ok(result.unwrap_or_else(|err| ObjectPartInfo { + number: num, + error: Some(err.to_string()), + ..Default::default() + })) + } + async fn read_listing_metadata(&self, volume: &str, object_name: &str) -> Result { let object_dir = self.io_get_object_path(volume, object_name)?; let metadata_path = object_dir.join(STORAGE_FORMAT_FILE); @@ -9340,72 +9391,14 @@ impl DiskAPI for LocalDisk { #[tracing::instrument(level = "trace", skip_all)] async fn read_parts(&self, bucket: &str, paths: &[String]) -> Result> { let volume_dir = self.io_get_bucket_path(bucket)?; - - let mut ret = vec![ObjectPartInfo::default(); paths.len()]; - - for (i, path_str) in paths.iter().enumerate() { - let path = Path::new(path_str); - let file_name = path.file_name().and_then(|v| v.to_str()).unwrap_or_default(); - let num = file_name - .strip_prefix("part.") - .and_then(|v| v.strip_suffix(".meta")) - .and_then(|v| v.parse::().ok()) - .unwrap_or_default(); - - if let Err(err) = access( - self.io_get_object_path( - bucket, - path_join_buf(&[ - path.parent().unwrap_or_else(|| Path::new("")).to_string_lossy().as_ref(), - &format!("part.{num}"), - ]) - .as_str(), - )?, - ) + stream::iter(0..paths.len()) + .map(|index| self.read_part_metadata(bucket, &volume_dir, &paths[index])) + .buffered(PART_METADATA_READ_CONCURRENCY) + .try_fold(Vec::with_capacity(paths.len()), |mut parts, part| async move { + parts.push(part); + Ok(parts) + }) .await - { - ret[i] = ObjectPartInfo { - number: num, - error: Some(err.to_string()), - ..Default::default() - }; - continue; - } - - let data = match self - .read_all_data( - bucket, - volume_dir.clone(), - self.io_get_object_path(bucket, path.to_string_lossy().as_ref())?, - ) - .await - { - Ok(data) => data, - Err(err) => { - ret[i] = ObjectPartInfo { - number: num, - error: Some(err.to_string()), - ..Default::default() - }; - continue; - } - }; - - match ObjectPartInfo::unmarshal(&data) { - Ok(meta) => { - ret[i] = meta; - } - Err(err) => { - ret[i] = ObjectPartInfo { - number: num, - error: Some(err.to_string()), - ..Default::default() - }; - } - }; - } - - Ok(ret) } #[tracing::instrument(level = "trace", skip_all)] async fn check_parts(&self, volume: &str, path: &str, fi: &FileInfo) -> Result { @@ -12460,6 +12453,79 @@ mod test { ); } + #[tokio::test] + async fn test_read_parts_preserves_order_across_concurrency_windows() { + let dir = tempfile::tempdir().expect("temporary disk"); + let endpoint = Endpoint::try_from(dir.path().to_str().expect("UTF-8 path")).expect("endpoint"); + let disk = LocalDisk::new(&endpoint, false).await.expect("local disk"); + let bucket = "bucket"; + ensure_test_volume(&disk, bucket).await; + let count = PART_METADATA_READ_CONCURRENCY * 3 + 1; + let mut paths = Vec::with_capacity(count); + let mut expected = Vec::with_capacity(count); + for number in (1..=count).rev() { + let part = ObjectPartInfo { + number, + etag: format!("etag-{number}"), + size: number, + actual_size: i64::try_from(number).expect("small part"), + ..Default::default() + }; + let data_path = format!("upload/part.{number}"); + let meta_path = format!("{data_path}.meta"); + disk.write_all(bucket, &data_path, Bytes::from_static(b"data")) + .await + .expect("part data"); + disk.write_all(bucket, &meta_path, Bytes::from(part.marshal_msg().expect("metadata"))) + .await + .expect("part metadata"); + paths.push(meta_path); + expected.push(part); + } + assert_eq!(disk.read_parts(bucket, &paths).await.expect("read parts"), expected); + assert!(disk.read_parts(bucket, &[]).await.expect("empty batch").is_empty()); + let missing_meta = "upload/part.99.meta".to_owned(); + disk.write_all(bucket, "upload/part.99", Bytes::from_static(b"data")) + .await + .expect("data without metadata"); + paths.insert(PART_METADATA_READ_CONCURRENCY, missing_meta.clone()); + let parts = disk.read_parts(bucket, &paths).await.expect("per-part failures"); + let error = disk + .read_all_data( + bucket, + disk.io_get_bucket_path(bucket).expect("bucket path"), + disk.io_get_object_path(bucket, &missing_meta).expect("metadata path"), + ) + .await + .expect_err("missing metadata"); + assert_eq!(parts[PART_METADATA_READ_CONCURRENCY].number, 99); + assert_eq!(parts[PART_METADATA_READ_CONCURRENCY].error.as_deref(), Some(error.to_string().as_str())); + let successful: Vec<_> = parts.into_iter().filter(|part| part.error.is_none()).collect(); + assert_eq!(successful, expected, "later windows must survive an earlier part error"); + + #[cfg(unix)] + { + let meta_path = disk + .io_get_object_path(bucket, "upload") + .expect("upload path") + .join("part.100.meta"); + std::os::unix::fs::symlink(dir.path(), meta_path).expect("invalid metadata symlink"); + let paths = ["upload/part.100.meta".to_owned()]; + let parts = disk + .read_parts(bucket, &paths) + .await + .expect("missing data remains a per-part error"); + assert!(parts[0].error.is_some()); + disk.write_all(bucket, "upload/part.100", Bytes::from_static(b"data")) + .await + .expect("data for invalid metadata path"); + assert_eq!( + disk.read_parts(bucket, &paths).await.expect_err("reject metadata symlink"), + DiskError::InvalidPath + ); + } + } + #[tokio::test] async fn test_read_parts_reports_bad_metadata_and_missing_data_part() { use tempfile::tempdir; diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index 0d9618d83..b9ea895b6 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -2490,15 +2490,15 @@ impl SetDisks { part_numbers: &[usize], read_quorum: usize, ) -> disk::error::Result> { - let bucket = bucket.to_string(); - let part_meta_paths = part_meta_paths.to_vec(); + let bucket: Arc = Arc::from(bucket); + let part_meta_paths: Arc<[String]> = Arc::from(part_meta_paths); let tasks: Vec<_> = disks .iter() .map(|disk| { let disk = disk.clone(); - let bucket = bucket.clone(); - let part_meta_paths = part_meta_paths.clone(); + let bucket = Arc::clone(&bucket); + let part_meta_paths = Arc::clone(&part_meta_paths); async move { if let Some(disk) = disk { diff --git a/crates/ecstore/src/set_disk/ops/multipart.rs b/crates/ecstore/src/set_disk/ops/multipart.rs index 15a4603f2..eb100eaec 100644 --- a/crates/ecstore/src/set_disk/ops/multipart.rs +++ b/crates/ecstore/src/set_disk/ops/multipart.rs @@ -66,6 +66,7 @@ use crate::disk::DiskOption; use crate::disk::STORAGE_FORMAT_FILE; #[cfg(test)] use crate::disk::new_disk; +use crate::erasure::coding::BitrotWriterWrapper; use crate::multipart_listing::paginate_multipart_listing; #[cfg(test)] use crate::object_api::ObjectLockConfigSnapshot; @@ -77,7 +78,7 @@ use crate::storage_api_contracts::multipart::MultipartOperations; #[cfg(test)] use crate::storage_api_contracts::object::HTTPPreconditions; use crate::storage_api_contracts::object::ObjectOperations; -use futures::{StreamExt, stream}; +use futures::{StreamExt, future::join_all, stream}; #[cfg(test)] use http::HeaderMap; use rustfs_filemeta::metadata_keys; @@ -100,6 +101,49 @@ use std::time::Duration; use tokio::io::AsyncReadExt; use tokio::task::JoinSet; +// The erasure-set width bounds fan-out. Await every opener so errors retain +// their disk slots and every successful writer remains owned until quorum is checked. +async fn create_part_writers( + disks: &[Option], + path: &str, + length: i64, + shard_size: usize, +) -> (Vec>, Vec>) { + join_all(disks.iter().map(|disk| async move { + let Some(disk) = disk else { + return (None, Some(DiskError::DiskNotFound)); + }; + match create_bitrot_writer( + false, + Some(disk), + RUSTFS_META_TMP_BUCKET, + path, + length, + shard_size, + HashAlgorithm::HighwayHash256S, + ) + .await + { + Ok(writer) => (Some(writer), None), + Err(err) => { + warn!( + event = EVENT_SET_DISK_MULTIPART, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_SET_DISK, + disk = ?disk, + state = "bitrot_writer_skipped", + error = ?err, + "Set disk multipart bitrot writer skipped" + ); + (None, Some(err)) + } + } + })) + .await + .into_iter() + .unzip() +} + const MULTIPART_LIST_IO_CONCURRENCY: usize = 16; static CAPPED_MULTIPART_STAGING: OnceLock>>> = OnceLock::new(); @@ -1542,45 +1586,13 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { .map_err(Error::from)?); let writer_setup_stage_start = rustfs_io_metrics::put_stage_metrics_enabled().then(Instant::now); - let mut writers = Vec::with_capacity(shuffle_disks.len()); - let mut errors = Vec::with_capacity(shuffle_disks.len()); - for disk_op in shuffle_disks.iter() { - if let Some(disk) = disk_op { - let writer = match create_bitrot_writer( - false, - Some(disk), - RUSTFS_META_TMP_BUCKET, - &tmp_part_path, - erasure.shard_file_size(data.size()), - erasure.shard_size(), - HashAlgorithm::HighwayHash256S, - ) - .await - { - Ok(writer) => writer, - Err(err) => { - warn!( - event = EVENT_SET_DISK_MULTIPART, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_SET_DISK, - disk = ?disk, - state = "bitrot_writer_skipped", - error = ?err, - "Set disk multipart bitrot writer skipped" - ); - errors.push(Some(err)); - writers.push(None); - continue; - } - }; - - writers.push(Some(writer)); - errors.push(None); - } else { - errors.push(Some(DiskError::DiskNotFound)); - writers.push(None); - } - } + let (mut writers, errors) = create_part_writers( + &shuffle_disks, + &tmp_part_path, + erasure.shard_file_size(data.size()), + erasure.shard_size(), + ) + .await; if let Some(stage_start) = writer_setup_stage_start { rustfs_io_metrics::record_put_object_stage_duration( @@ -3594,6 +3606,82 @@ mod tests { use tempfile::TempDir; use tokio::sync::{Notify, RwLock}; + #[tokio::test] + async fn multipart_writer_setup_opens_disks_concurrently_and_preserves_error_slots() { + use crate::cluster::rpc::internode_data_transport::{ + InternodeDataTransport, InternodeDataTransportCapabilities, ReadStreamRequest, WalkDirStreamRequest, + WriteStreamRequest, + }; + use crate::cluster::rpc::remote_disk::RemoteDisk; + use crate::disk::{Disk, FileReader, FileWriter}; + + #[derive(Debug)] + struct BarrierTransport { + barrier: Arc, + fail: bool, + } + #[async_trait::async_trait] + impl InternodeDataTransport for BarrierTransport { + async fn open_read(&self, _: ReadStreamRequest) -> disk::error::Result { + unreachable!("writer setup must not read") + } + async fn open_walk_dir(&self, _: WalkDirStreamRequest) -> disk::error::Result { + unreachable!("writer setup must not list") + } + async fn open_write(&self, request: WriteStreamRequest) -> disk::error::Result { + assert_eq!(request.volume, RUSTFS_META_TMP_BUCKET); + assert_eq!(request.path, "upload/part.1"); + self.barrier.wait().await; + if self.fail { + Err(DiskError::FileAccessDenied) + } else { + Ok(Box::new(tokio::io::sink())) + } + } + fn name(&self) -> &'static str { + "multipart-writer-test" + } + fn capabilities(&self) -> InternodeDataTransportCapabilities { + InternodeDataTransportCapabilities::tcp_http() + } + } + + let barrier = Arc::new(tokio::sync::Barrier::new(3)); + let mut disks = Vec::new(); + for i in 0..3 { + let endpoint = Endpoint { + url: url::Url::parse(&format!("http://multipart-test.invalid:9000/disk{i}")).expect("endpoint"), + is_local: false, + pool_idx: 0, + set_idx: 0, + disk_idx: i, + }; + let disk = RemoteDisk::new( + &endpoint, + &DiskOption { + cleanup: false, + health_check: false, + }, + Arc::new(BarrierTransport { + barrier: Arc::clone(&barrier), + fail: i == 1, + }), + ) + .await + .expect("remote disk"); + disks.push(Some(Arc::new(Disk::Remote(Box::new(disk))))); + } + disks.insert(1, None); + // A serial opener cannot cross the barrier. The timeout only bounds failures; + // the assertion depends on all three independent openers making progress. + let (writers, errors) = + tokio::time::timeout(Duration::from_secs(10), create_part_writers(&disks, "upload/part.1", 1024, 256)) + .await + .expect("all disk openers must be polled concurrently"); + assert_eq!(writers.iter().map(Option::is_some).collect::>(), [true, false, false, true]); + assert_eq!(errors, [None, Some(DiskError::DiskNotFound), Some(DiskError::FileAccessDenied), None]); + } + #[test] fn multipart_bucket_incarnation_metadata_is_consistent_and_non_nil() { let incarnation = Uuid::new_v4();