improve multi put speed

Signed-off-by: junxiang Mu <1948535941@qq.com>
This commit is contained in:
junxiang Mu
2025-05-09 17:13:25 +08:00
parent 5cb040f863
commit fd03ba54f3
3 changed files with 87 additions and 54 deletions
+78 -49
View File
@@ -513,24 +513,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);
}
@@ -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<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);
@@ -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: {}",