Merge pull request #393 from rustfs/pref

improve multi put speed
This commit is contained in:
loverustfs
2025-05-09 22:40:03 +08:00
committed by GitHub
3 changed files with 87 additions and 54 deletions
+5 -5
View File
@@ -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));
+78 -49
View File
@@ -510,24 +510,33 @@ impl SetDisks {
meta: Vec<u8>,
write_quorum: usize,
) -> Result<Vec<Option<DiskStore>>> {
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);
}
@@ -539,7 +548,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);
}
@@ -941,7 +950,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) {
@@ -1111,39 +1120,47 @@ impl SetDisks {
version_id: &str,
read_data: bool,
healing: bool,
) -> (Vec<FileInfo>, Vec<Option<Error>>) {
let mut futures = Vec::with_capacity(disks.len());
) -> Result<(Vec<FileInfo>, Vec<Option<Error>>)> {
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);
@@ -1154,7 +1171,7 @@ impl SetDisks {
}
}
}
(ress, errors)
Ok((ress, errors))
}
async fn read_all_xl(
@@ -1767,7 +1784,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);
@@ -2174,7 +2191,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: {}",
@@ -3955,7 +3972,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
}
@@ -4249,7 +4266,7 @@ impl StorageAPI for SetDisks {
false,
false,
)
.await
.await?
} else {
Self::read_all_xl(&disks, bucket, object, false, false).await
}
@@ -4389,30 +4406,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);
@@ -5041,7 +5070,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: {}",
+4
View File
@@ -115,6 +115,10 @@ pub async fn lstat(path: impl AsRef<Path>) -> io::Result<Metadata> {
fs::metadata(path).await
}
pub fn lstat_std(path: impl AsRef<Path>) -> io::Result<Metadata> {
std::fs::metadata(path)
}
pub async fn make_dir_all(path: impl AsRef<Path>) -> io::Result<()> {
fs::create_dir_all(path.as_ref()).await
}