From fd03ba54f3dc0d7664b57c896e8e949cd0fd16c4 Mon Sep 17 00:00:00 2001 From: junxiang Mu <1948535941@qq.com> Date: Fri, 9 May 2025 17:13:25 +0800 Subject: [PATCH] improve multi put speed Signed-off-by: junxiang Mu <1948535941@qq.com> --- ecstore/src/disk/local.rs | 10 +-- ecstore/src/set_disk.rs | 127 +++++++++++++++++++++++--------------- ecstore/src/utils/fs.rs | 4 ++ 3 files changed, 87 insertions(+), 54 deletions(-) diff --git a/ecstore/src/disk/local.rs b/ecstore/src/disk/local.rs index bd10ccf69..a814baca3 100644 --- a/ecstore/src/disk/local.rs +++ b/ecstore/src/disk/local.rs @@ -39,7 +39,7 @@ use crate::set_disk::{ }; use crate::store_api::{BitrotAlgorithm, StorageAPI}; use crate::utils::fs::{ - access, lstat, remove, remove_all, remove_all_std, remove_std, rename, O_APPEND, O_CREATE, O_RDONLY, O_WRONLY, + access, lstat, lstat_std, remove, remove_all, remove_all_std, remove_std, rename, O_APPEND, O_CREATE, O_RDONLY, O_WRONLY, }; use crate::utils::os::get_info; use crate::utils::path::{ @@ -1337,10 +1337,10 @@ impl DiskAPI for LocalDisk { let src_volume_dir = self.get_bucket_path(src_volume)?; let dst_volume_dir = self.get_bucket_path(dst_volume)?; if !skip_access_checks(src_volume) { - utils::fs::access(&src_volume_dir).await.map_err(map_err_not_exists)? + utils::fs::access_std(&src_volume_dir).map_err(map_err_not_exists)? } if !skip_access_checks(dst_volume) { - utils::fs::access(&dst_volume_dir).await.map_err(map_err_not_exists)? + utils::fs::access_std(&dst_volume_dir).map_err(map_err_not_exists)? } let src_is_dir = has_suffix(src_path, SLASH_SEPARATOR); @@ -1363,7 +1363,7 @@ impl DiskAPI for LocalDisk { check_path_length(dst_file_path.to_string_lossy().as_ref())?; if src_is_dir { - let meta_op = match lstat(&src_file_path).await { + let meta_op = match lstat_std(&src_file_path) { Ok(meta) => Some(meta), Err(e) => { if is_sys_err_io(&e) { @@ -1384,7 +1384,7 @@ impl DiskAPI for LocalDisk { } } - if let Err(e) = utils::fs::remove(&dst_file_path).await { + if let Err(e) = utils::fs::remove_std(&dst_file_path) { if is_sys_err_not_empty(&e) || is_sys_err_not_dir(&e) { warn!("rename_part remove dst failed {:?} err {:?}", &dst_file_path, e); return Err(Error::new(DiskError::FileAccessDenied)); diff --git a/ecstore/src/set_disk.rs b/ecstore/src/set_disk.rs index 6e7387834..38dbd9cf5 100644 --- a/ecstore/src/set_disk.rs +++ b/ecstore/src/set_disk.rs @@ -513,24 +513,33 @@ impl SetDisks { meta: Vec, write_quorum: usize, ) -> Result>> { - let mut futures = Vec::with_capacity(disks.len()); + let src_bucket = Arc::new(src_bucket.to_string()); + let src_object = Arc::new(src_object.to_string()); + let dst_bucket = Arc::new(dst_bucket.to_string()); + let dst_object = Arc::new(dst_object.to_string()); let mut errs = Vec::with_capacity(disks.len()); - for disk in disks.iter() { + let futures = disks.iter().map(|disk| { + let disk = disk.clone(); let meta = meta.clone(); - futures.push(async move { + let src_bucket = src_bucket.clone(); + let src_object = src_object.clone(); + let dst_bucket = dst_bucket.clone(); + let dst_object = dst_object.clone(); + tokio::spawn(async move { if let Some(disk) = disk { - disk.rename_part(src_bucket, src_object, dst_bucket, dst_object, meta).await + disk.rename_part(&src_bucket, &src_object, &dst_bucket, &dst_object, meta) + .await } else { Err(Error::new(DiskError::DiskNotFound)) } }) - } + }); let results = join_all(futures).await; for result in results { - match result { + match result? { Ok(_) => { errs.push(None); } @@ -542,7 +551,7 @@ impl SetDisks { if let Some(err) = reduce_write_quorum_errs(&errs, object_op_ignored_errs().as_ref(), write_quorum) { warn!("rename_part errs {:?}", &errs); - Self::cleanup_multipart_path(disks, &[dst_object.to_owned(), format!("{}.meta", dst_object)]).await; + Self::cleanup_multipart_path(disks, &[dst_object.to_string(), format!("{}.meta", dst_object)]).await; return Err(err); } @@ -944,7 +953,7 @@ impl SetDisks { let disks = disks.clone(); let (parts_metadata, errs) = - Self::read_all_fileinfo(&disks, bucket, RUSTFS_META_MULTIPART_BUCKET, &upload_id_path, "", false, false).await; + Self::read_all_fileinfo(&disks, bucket, RUSTFS_META_MULTIPART_BUCKET, &upload_id_path, "", false, false).await?; let map_err_notfound = |err: Error| { if is_err_object_not_found(&err) { @@ -1114,39 +1123,47 @@ impl SetDisks { version_id: &str, read_data: bool, healing: bool, - ) -> (Vec, Vec>) { - let mut futures = Vec::with_capacity(disks.len()); + ) -> Result<(Vec, Vec>)> { let mut ress = Vec::with_capacity(disks.len()); let mut errors = Vec::with_capacity(disks.len()); - - for disk in disks.iter() { - let opts = ReadOptions { - read_data, - healing, - ..Default::default() - }; - futures.push(async move { + let opts = Arc::new(ReadOptions { + read_data, + healing, + ..Default::default() + }); + 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 futures = disks.iter().map(|disk| { + let disk = disk.clone(); + let opts = opts.clone(); + let org_bucket = org_bucket.clone(); + let bucket = bucket.clone(); + let object = object.clone(); + let version_id = version_id.clone(); + tokio::spawn(async move { if let Some(disk) = disk { if version_id.is_empty() { - match disk.read_xl(bucket, object, read_data).await { + match disk.read_xl(&bucket, &object, read_data).await { Ok(info) => { - let fi = file_info_from_raw(info, bucket, object, read_data).await?; + let fi = file_info_from_raw(info, &bucket, &object, read_data).await?; Ok(fi) } Err(err) => Err(err), } } else { - disk.read_version(org_bucket, bucket, object, version_id, &opts).await + disk.read_version(&org_bucket, &bucket, &object, &version_id, &opts).await } } else { Err(Error::new(DiskError::DiskNotFound)) } }) - } + }); let results = join_all(futures).await; for result in results { - match result { + match result? { Ok(res) => { ress.push(res); errors.push(None); @@ -1157,7 +1174,7 @@ impl SetDisks { } } } - (ress, errors) + Ok((ress, errors)) } async fn read_all_xl( @@ -1770,7 +1787,7 @@ impl SetDisks { let vid = opts.version_id.clone().unwrap_or_default(); // TODO: 优化并发 可用数量中断 - let (parts_metadata, errs) = Self::read_all_fileinfo(&disks, "", bucket, object, vid.as_str(), read_data, false).await; + let (parts_metadata, errs) = Self::read_all_fileinfo(&disks, "", bucket, object, vid.as_str(), read_data, false).await?; // warn!("get_object_fileinfo parts_metadata {:?}", &parts_metadata); // warn!("get_object_fileinfo {}/{} errs {:?}", bucket, object, &errs); @@ -2177,7 +2194,7 @@ impl SetDisks { let disks = { self.disks.read().await.clone() }; - let (mut parts_metadata, errs) = Self::read_all_fileinfo(&disks, "", bucket, object, version_id, true, true).await; + let (mut parts_metadata, errs) = Self::read_all_fileinfo(&disks, "", bucket, object, version_id, true, true).await?; if is_all_not_found(&errs) { warn!( "heal_object failed, all obj part not found, bucket: {}, obj: {}, version_id: {}", @@ -3958,7 +3975,7 @@ impl StorageAPI for SetDisks { let (mut metas, errs) = { if let Some(vid) = &src_opts.version_id { - Self::read_all_fileinfo(&disks, "", src_bucket, src_object, vid, true, false).await + Self::read_all_fileinfo(&disks, "", src_bucket, src_object, vid, true, false).await? } else { Self::read_all_xl(&disks, src_bucket, src_object, true, false).await } @@ -4252,7 +4269,7 @@ impl StorageAPI for SetDisks { false, false, ) - .await + .await? } else { Self::read_all_xl(&disks, bucket, object, false, false).await } @@ -4392,30 +4409,42 @@ impl StorageAPI for SetDisks { let part_suffix = format!("part.{}", part_id); let tmp_part = format!("{}x{}", Uuid::new_v4(), OffsetDateTime::now_utc().unix_timestamp()); - let tmp_part_path = format!("{}/{}", tmp_part, part_suffix); + let tmp_part_path = Arc::new(format!("{}/{}", tmp_part, part_suffix)); let mut writers = Vec::with_capacity(disks.len()); let erasure = Erasure::new(fi.erasure.data_blocks, fi.erasure.parity_blocks, fi.erasure.block_size); + let shared_size = erasure.shard_size(erasure.block_size); - for disk in disks.iter() { - if let Some(disk) = disk { - // let writer = disk.append_file(RUSTFS_META_TMP_BUCKET, &tmp_part_path).await?; - // let filewriter = disk - // .create_file("", RUSTFS_META_TMP_BUCKET, &tmp_part_path, data.content_length) - // .await?; - let writer = new_bitrot_filewriter( - disk.clone(), - RUSTFS_META_TMP_BUCKET, - &tmp_part_path, - false, - DEFAULT_BITROT_ALGO, - erasure.shard_size(erasure.block_size), - ) - .await?; - writers.push(Some(writer)); - } else { - writers.push(None); - } + let futures = disks.iter().map(|disk| { + let disk = disk.clone(); + let tmp_part_path = tmp_part_path.clone(); + tokio::spawn(async move { + if let Some(disk) = disk { + // let writer = disk.append_file(RUSTFS_META_TMP_BUCKET, &tmp_part_path).await?; + // let filewriter = disk + // .create_file("", RUSTFS_META_TMP_BUCKET, &tmp_part_path, data.content_length) + // .await?; + match new_bitrot_filewriter( + disk.clone(), + RUSTFS_META_TMP_BUCKET, + &tmp_part_path, + false, + DEFAULT_BITROT_ALGO, + shared_size, + ) + .await + { + Ok(writer) => Ok(Some(writer)), + Err(e) => Err(e), + } + } else { + Ok(None) + } + }) + }); + for x in join_all(futures).await { + let x = x??; + writers.push(x); } let erasure = Erasure::new(fi.erasure.data_blocks, fi.erasure.parity_blocks, fi.erasure.block_size); @@ -5044,7 +5073,7 @@ impl StorageAPI for SetDisks { let disks = self.disks.read().await; let disks = disks.clone(); - let (_, errs) = Self::read_all_fileinfo(&disks, "", bucket, object, version_id, false, false).await; + let (_, errs) = Self::read_all_fileinfo(&disks, "", bucket, object, version_id, false, false).await?; if is_all_not_found(&errs) { warn!( "heal_object failed, all obj part not found, bucket: {}, obj: {}, version_id: {}", diff --git a/ecstore/src/utils/fs.rs b/ecstore/src/utils/fs.rs index f60e28f08..d8110ca63 100644 --- a/ecstore/src/utils/fs.rs +++ b/ecstore/src/utils/fs.rs @@ -115,6 +115,10 @@ pub async fn lstat(path: impl AsRef) -> io::Result { fs::metadata(path).await } +pub fn lstat_std(path: impl AsRef) -> io::Result { + std::fs::metadata(path) +} + pub async fn make_dir_all(path: impl AsRef) -> io::Result<()> { fs::create_dir_all(path.as_ref()).await }