mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-04 12:31:36 +00:00
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 <housemecn@gmail.com> Co-authored-by: Chris <anzhengchao@gmail.com>
This commit is contained in:
@@ -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);
|
||||
|
||||
@@ -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<Arc<FdCacheEntry>> = cached_fd.is_none().then(|| {
|
||||
let opened_fd: Option<Arc<FdCacheEntry>> = 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");
|
||||
|
||||
@@ -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<BytesMut> = 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<oneshot::Sender<()>>,
|
||||
dropped: Option<oneshot::Sender<()>>,
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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 |
|
||||
|
||||
Reference in New Issue
Block a user