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
This commit is contained in:
GatewayJ
2026-09-23 22:19:23 +08:00
committed by GitHub
parent 311e4306e6
commit f219a0aba8
6 changed files with 363 additions and 109 deletions
+4
View File
@@ -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
@@ -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);
@@ -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};
}
+131 -65
View File
@@ -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<ObjectPartInfo> {
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::<usize>().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<ListingMetadataRead> {
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<Vec<ObjectPartInfo>> {
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::<usize>().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<CheckPartsResp> {
@@ -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;
@@ -2490,15 +2490,15 @@ impl SetDisks {
part_numbers: &[usize],
read_quorum: usize,
) -> disk::error::Result<Vec<ObjectPartInfo>> {
let bucket = bucket.to_string();
let part_meta_paths = part_meta_paths.to_vec();
let bucket: Arc<str> = 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 {
+128 -40
View File
@@ -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<DiskStore>],
path: &str,
length: i64,
shard_size: usize,
) -> (Vec<Option<BitrotWriterWrapper>>, Vec<Option<DiskError>>) {
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<Mutex<HashMap<String, Arc<tokio::sync::Semaphore>>>> = 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<tokio::sync::Barrier>,
fail: bool,
}
#[async_trait::async_trait]
impl InternodeDataTransport for BarrierTransport {
async fn open_read(&self, _: ReadStreamRequest) -> disk::error::Result<FileReader> {
unreachable!("writer setup must not read")
}
async fn open_walk_dir(&self, _: WalkDirStreamRequest) -> disk::error::Result<FileReader> {
unreachable!("writer setup must not list")
}
async fn open_write(&self, request: WriteStreamRequest) -> disk::error::Result<FileWriter> {
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::<Vec<_>>(), [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();