From 753e8aa657e181e26dced0ea0800cf0001604a86 Mon Sep 17 00:00:00 2001 From: GatewayJ <835269233@qq.com> Date: Sat, 3 Oct 2026 14:53:00 +0800 Subject: [PATCH] perf(ecstore): reduce multipart and object I/O overhead (#8095) * perf(ecstore): share multipart cleanup paths across disks * perf(ecstore): avoid encode copies and cached descriptor duplication * test(ecstore): cover default ingest selection and overrides * test(ecstore): cover cached positioned reads in both modes --------- Co-authored-by: Hauser Co-authored-by: Chris --- crates/ecstore/benches/erasure_benchmark.rs | 26 ++++++++ crates/ecstore/src/disk/local.rs | 51 ++++++++++------ crates/ecstore/src/erasure/coding/encode.rs | 59 ++++++++++++++----- .../src/set_disk/core/io_primitives.rs | 4 +- docs/operations/object-io-tuning-ab-matrix.md | 1 + 5 files changed, 108 insertions(+), 33 deletions(-) diff --git a/crates/ecstore/benches/erasure_benchmark.rs b/crates/ecstore/benches/erasure_benchmark.rs index c9f7adb52..b7c5bc9bd 100644 --- a/crates/ecstore/benches/erasure_benchmark.rs +++ b/crates/ecstore/benches/erasure_benchmark.rs @@ -307,6 +307,32 @@ fn bench_memory_patterns(c: &mut Criterion) { group.sample_size(10); group.measurement_time(Duration::from_secs(5)); + // The reader owns an EC-sized buffer before encoding. Keep ingest outside + // the timed region to isolate the borrowed-copy versus ownership transfer. + let erasure = Erasure::new(data_shards, parity_shards, block_size); + let payload = generate_test_data(block_size); + for owned in [false, true] { + group.bench_function(if owned { "owned_ingest" } else { "borrowed_ingest" }, |b| { + b.iter_batched( + || { + let mut buffer = bytes::BytesMut::with_capacity(erasure.shard_size() * (data_shards + parity_shards)); + buffer.extend_from_slice(&payload); + buffer + }, + |buffer| { + let shards = if owned { + erasure.encode_data_bytes_mut(buffer, block_size) + } else { + erasure.encode_data(&buffer) + } + .expect("encode benchmark block"); + black_box(shards) + }, + criterion::BatchSize::PerIteration, + ); + }); + } + // Test reusing the same Erasure instance group.bench_function("reuse_erasure_instance", |b| { let erasure = Erasure::new(data_shards, parity_shards, block_size); diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index a0d9712dd..5a3994250 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -3320,7 +3320,7 @@ impl LocalIoBackend for StdBackend { let end_offset_u64 = u64::try_from(end_offset).map_err(|_| DiskError::FileCorrupt)?; // Descriptor cache (rustfs/backlog#1801): on a hit the read reuses an - // already-open descriptor (via dup below) and skips `access` + + // already-open descriptor by reference and skips `access` + // `File::open`. Linux-only — on other Unix `cached_fd` is None and the // read opens per call exactly as before. `fd_lookup` snapshots the // invalidation generation BEFORE the open so a heal/delete that lands @@ -3357,22 +3357,19 @@ impl LocalIoBackend for StdBackend { let file_open_start = metrics_enabled.then(StdInstant::now); // Acquire the read handle (rustfs/backlog#1801). On a descriptor-cache - // hit this reuses the cached descriptor via `dup` (one syscall, no path - // resolution or permission re-check) and skips the volume access probe; + // hit this borrows the cached descriptor without duplicating it + // and skips the volume access probe; // on a miss it resolves the volume, access-checks, and opens the file. - // `File::try_clone` shares the cached descriptor's open-file offset, so + // Cached readers share the same descriptor, so // the read below is positioned (mmap offset argument / `read_exact_at`) // and never depends on the descriptor's current offset. `cached_fd` being // None also marks this call as a miss for the cache-insert side-channel. // The cached length is the metadata snapshot captured at open time; // all in-place/replacement writers invalidate this entry before // publishing a mutation, so cache hits avoid a redundant fstat. + let mut opened_file = None; let (file, cached_len, access_check_duration) = if let Some(cached) = cached_fd.as_ref() { - ( - cached.file.as_ref().try_clone().map_err(DiskError::from)?, - Some(cached.len), - StdDuration::ZERO, - ) + (cached.file.as_ref(), Some(cached.len), StdDuration::ZERO) } else { // Measure the volume access probe only — the part-path resolution // above is accounted in `path_resolve_duration` (rustfs/backlog#1801). @@ -3383,7 +3380,8 @@ impl LocalIoBackend for StdBackend { .map_err(|e| DiskError::from(to_access_error(e, DiskError::VolumeAccessDenied)))?; } let access_check_duration = access_check_start.map_or(StdDuration::ZERO, |started_at| started_at.elapsed()); - (std::fs::File::open(&file_path).map_err(DiskError::from)?, None, access_check_duration) + let file = opened_file.insert(std::fs::File::open(&file_path).map_err(DiskError::from)?); + (&*file, None, access_check_duration) }; let file_open_duration = file_open_start.map_or(StdDuration::ZERO, |started_at| started_at.elapsed()); @@ -3410,7 +3408,7 @@ impl LocalIoBackend for StdBackend { #[cfg(target_os = "macos")] if should_reclaim_after_read { - let _ = set_std_fd_nocache(&file); + let _ = set_std_fd_nocache(file); } let mut mmap_map_duration = StdDuration::ZERO; @@ -3488,7 +3486,7 @@ impl LocalIoBackend for StdBackend { if should_populate_mmap_read { mmap_options.populate(); } - let mmap = unsafe { mmap_options.map(&file) }.map_err(DiskError::other)?; + let mmap = unsafe { mmap_options.map(file) }.map_err(DiskError::other)?; let mmap_map_faults_after = read_mmap_page_fault_counts(metrics_enabled); mmap_map_duration = mmap_map_start.map_or(StdDuration::ZERO, |started_at| started_at.elapsed()); mmap_map_fault_delta = mmap_page_fault_delta(mmap_map_faults_before, mmap_map_faults_after); @@ -3513,8 +3511,8 @@ impl LocalIoBackend for StdBackend { let direct_read_copy_start = metrics_enabled.then(StdInstant::now); let direct_read_copy_faults_before = read_mmap_page_fault_counts(metrics_enabled); let mut buffer = vec![0; length]; - // Positioned read: a cache hit reads through a `dup`'d handle - // that shares the cached descriptor's offset, so this must not + // Positioned read: cache hits share the cached descriptor, + // so concurrent readers must not // touch the descriptor offset (rustfs/backlog#1801). file.read_exact_at(&mut buffer, offset_u64).map_err(DiskError::from)?; let direct_read_copy_faults_after = read_mmap_page_fault_counts(metrics_enabled); @@ -3536,7 +3534,7 @@ impl LocalIoBackend for StdBackend { u64::try_from(_reclaim_len).map_err(|_| DiskError::other("read reclaim length overflow"))?, ) .ok_or_else(|| DiskError::other("read reclaim length overflow"))?; - fadvise(&file, _reclaim_offset, Some(reclaim_len), Advice::DontNeed) + fadvise(file, _reclaim_offset, Some(reclaim_len), Advice::DontNeed) .map_err(std::io::Error::from) .map_err(DiskError::from)?; } @@ -3546,10 +3544,10 @@ impl LocalIoBackend for StdBackend { // Hand the freshly opened descriptor back so the async caller can index // the cache — None on a hit (the cache already holds it). mmap/reclaim // above only borrowed `file`, so it is still owned here and moves into the - // Arc; `cached_fd.is_none()` is true exactly when this call did the open. + // Arc; `opened_file` is populated only when this call did the open. // Non-Linux has no fd cache, so skip the Arc allocation there. #[cfg(target_os = "linux")] - let opened_fd: Option> = cached_fd.is_none().then(|| { + let opened_fd: Option> = opened_file.map(|file| { Arc::new(FdCacheEntry { file: Arc::new(file), len: metadata_len, @@ -23850,6 +23848,25 @@ mod test { .expect("operation should succeed"); assert_eq!(second, Bytes::from_static(payload)); + // Both read methods must preserve offsets on the shared cached descriptor. + for method in [ + RUSTFS_OBJECT_MMAP_READ_METHOD_MMAP_COPY, + RUSTFS_OBJECT_MMAP_READ_METHOD_DIRECT_READ_COPY, + ] { + let reads = temp_env::async_with_vars([(ENV_RUSTFS_OBJECT_MMAP_READ_METHOD, Some(method))], async { + futures::future::join_all((0..payload.len()).map(|offset| backend.pread_bytes(volume, object, offset, 1, None))) + .await + }) + .await; + for (offset, read) in reads.into_iter().enumerate() { + assert_eq!( + read.expect("cached offset read"), + &payload[offset..offset + 1], + "{method} at offset {offset}" + ); + } + } + // Invalidating by the object prefix drops the cached descriptor. backend.invalidate_cached_fds_under(volume, "obj/abc"); assert_eq!(cache.entry_count().await, 0, "prefix invalidation must drop the cached descriptor"); diff --git a/crates/ecstore/src/erasure/coding/encode.rs b/crates/ecstore/src/erasure/coding/encode.rs index 946e665f0..9b328aa18 100644 --- a/crates/ecstore/src/erasure/coding/encode.rs +++ b/crates/ecstore/src/erasure/coding/encode.rs @@ -43,7 +43,7 @@ const ENV_RUSTFS_ERASURE_ENCODE_BYTESMUT_INGEST: &str = "RUSTFS_ERASURE_ENCODE_B const DEFAULT_RUSTFS_ERASURE_ENCODE_MAX_INFLIGHT_BYTES: usize = 32 * 1024 * 1024; const DEFAULT_RUSTFS_ERASURE_ENCODE_MAX_INFLIGHT_BLOCKS: usize = 32; const DEFAULT_RUSTFS_ERASURE_ENCODE_BATCH_BLOCKS: usize = 4; -const DEFAULT_RUSTFS_ERASURE_ENCODE_BYTESMUT_INGEST: bool = false; +const DEFAULT_RUSTFS_ERASURE_ENCODE_BYTESMUT_INGEST: bool = true; pub(crate) enum IntegrityEncodeMode { Inline(usize), @@ -890,8 +890,6 @@ impl Erasure { // capacity and never reallocates. Reading into uninitialized spare capacity // (instead of resize + slice read) also skips zero-filling each fresh buffer. let ingest_capacity = expanded_block_bytes.max(block_size); - // Pre-allocate buffer pool for this encoding session - let mut buf_pool: Vec = Vec::with_capacity(4); let mut buf = BytesMut::with_capacity(ingest_capacity); loop { match read_full_buf_or_eof(&mut reader, &mut buf, block_size).await { @@ -901,8 +899,6 @@ impl Erasure { total += n; let encode_buf = buf; let res = self.clone().encode_block_bytes_mut(encode_buf, n).await?; - // Try to reuse buffer from pool, or allocate new one - buf = buf_pool.pop().unwrap_or_else(|| BytesMut::with_capacity(ingest_capacity)); let queued_bytes = res.queued_bytes(); let _producer_stage = rustfs_io_metrics::track_ec_encode_producer_bytes(queued_bytes); let send_wait_stage_start = stage_timer_if_enabled(); @@ -910,11 +906,8 @@ impl Erasure { return Err(std::io::Error::other(format!("Failed to send encoded data : {err}"))); } record_internal_stage_if_enabled("erasure_encode_send_wait", send_wait_stage_start); - // Return buffer to pool if it has sufficient capacity - if buf.capacity() >= ingest_capacity && buf_pool.len() < 4 { - buf_pool.push(buf); - buf = BytesMut::with_capacity(ingest_capacity); - } + // Encoded shards own the previous allocation until writing completes. + buf = BytesMut::with_capacity(ingest_capacity); } Ok(None) => break, Err(e) if e.kind() == std::io::ErrorKind::UnexpectedEof => { @@ -1054,18 +1047,18 @@ impl Erasure { let mut task = AbortOnDropTask::new(tokio::spawn(async move { let block_size = self.block_size; let mut total = 0; - let mut buf = vec![0u8; block_size]; + let ingest_capacity = expanded_block_bytes.max(block_size); + let mut buf = BytesMut::with_capacity(ingest_capacity); let mut pending_batch = Vec::with_capacity(batch_blocks); let mut pending_batch_bytes = 0usize; let mut pending_batch_stage = None; loop { - match rustfs_utils::read_full_or_eof(&mut reader, &mut buf).await { + match read_full_buf_or_eof(&mut reader, &mut buf, block_size).await { Ok(Some(n)) => { debug_assert!(n > 0, "non-zero block_size prevents zero-length reads"); total += n; - let encode_buf = std::mem::take(&mut buf); - let (res, returned_buf) = self.clone().encode_block(encode_buf, n).await?; - buf = returned_buf; + let res = self.clone().encode_block_bytes_mut(buf, n).await?; + buf = BytesMut::with_capacity(ingest_capacity); let queued_bytes = res.queued_bytes(); pending_batch_bytes = pending_batch_bytes.saturating_add(queued_bytes); pending_batch.push(res); @@ -1242,6 +1235,42 @@ mod tests { use tokio::io::{AsyncWrite, AsyncWriteExt, ReadBuf}; use tokio::sync::oneshot; + #[test] + fn bytesmut_ingest_selector_defaults_to_owned_and_honors_overrides() { + const CHILD_CASE: &str = "RUSTFS_TEST_BYTESMUT_INGEST_SELECTOR_CASE"; + if let Ok(case) = std::env::var(CHILD_CASE) { + let expected = match case.as_str() { + "default" | "true" => true, + "false" => false, + _ => panic!("unexpected ingest selector case: {case}"), + }; + assert_eq!(use_bytesmut_ingest(), expected, "ingest selector case: {case}"); + return; + } + + // Each selector probe needs a fresh OnceLock and an isolated environment. + for case in ["default", "true", "false"] { + let mut child = std::process::Command::new(std::env::current_exe().expect("ingest selector test executable")); + child.args([ + "--exact", + "erasure::coding::encode::tests::bytesmut_ingest_selector_defaults_to_owned_and_honors_overrides", + "--nocapture", + ]); + child.env(CHILD_CASE, case); + child.env_remove(ENV_RUSTFS_ERASURE_ENCODE_BYTESMUT_INGEST); + if case != "default" { + child.env(ENV_RUSTFS_ERASURE_ENCODE_BYTESMUT_INGEST, case); + } + let output = child.output().expect("start isolated ingest selector test"); + assert!( + output.status.success(), + "ingest selector case {case} failed:\n{}\n{}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr), + ); + } + } + struct PendingReader { entered: Option>, dropped: Option>, diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index e73ee7332..0d025b45d 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -5481,11 +5481,13 @@ impl SetDisks { // Use improved simple batch processor instead of join_all for better performance let processor = runtime_sources::batch_processors().write_processor(); + // Batch tasks require owned inputs; share the immutable paths across disks. + let paths: Arc<[String]> = Arc::from(paths); let tasks: Vec<_> = disks .iter() .map(|disk| { let disk = disk.clone(); - let paths = paths.to_vec(); + let paths = Arc::clone(&paths); async move { if let Some(disk) = disk { diff --git a/docs/operations/object-io-tuning-ab-matrix.md b/docs/operations/object-io-tuning-ab-matrix.md index 09ebf3cbe..62871e549 100644 --- a/docs/operations/object-io-tuning-ab-matrix.md +++ b/docs/operations/object-io-tuning-ab-matrix.md @@ -69,6 +69,7 @@ iostat -xz 5 > telemetry/iostat.txt & | Knob | Default | Controls | Validating stage | Risk if widened | | --- | --- | --- | --- | --- | | `RUSTFS_ERASURE_ENCODE_MAX_INFLIGHT_BYTES` | 32MiB (`crates/ecstore/src/erasure/coding/encode.rs`) | EC encode producer/consumer memory budget (blocks queued between encode and shard write) | `set_disk_encode` P95 + `rustfs_ec_encode_inflight_bytes_current` | RSS growth under high concurrency | +| `RUSTFS_ERASURE_ENCODE_BYTESMUT_INGEST` | `true` | Streaming encode reads into an owned, EC-sized buffer, avoiding the input-to-encoder block copy. Set `false` to select the previous streaming Vec ingest path; batched encode always uses owned buffers. Cached at first use. | `erasure_encode_cpu`, PUT throughput and RSS | No disk-format change; queue budget does not bound total process memory | | `RUSTFS_OBJECT_IO_BUFFER_SIZE` | 128KiB (`crates/config/src/constants/object.rs`) | Streaming read-in / write-out block size | `ingress_prepare`, `set_disk_encode` | Larger buffers = fewer polls, more resident memory | | `RUSTFS_OBJECT_DUPLEX_BUFFER_SIZE` | 4MiB (`crates/config/src/constants/object.rs`) | Duplex pipe capacity (shared; PUT uses it less than GET) | `set_disk_encode` feed smoothness | Memory per in-flight request | | `RUSTFS_DURABILITY_MODE` / `RUSTFS_DRIVE_SYNC_ENABLE` | mode-dependent | Per-shard fsync/sync discipline on commit | `set_disk_rename` P99 | Weakening it changes the durability contract — a deliberate tradeoff, never a free win |