mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-27 15:37:02 +00:00
fix(ecstore): cap mmap-copy shard reads to stop large single-part GET OOM (#5140)
This commit is contained in:
@@ -56,3 +56,33 @@ pub const DEFAULT_OBJECT_MMAP_READ_ENABLE: bool = true;
|
|||||||
///
|
///
|
||||||
/// Prefer [`DEFAULT_OBJECT_MMAP_READ_ENABLE`].
|
/// Prefer [`DEFAULT_OBJECT_MMAP_READ_ENABLE`].
|
||||||
pub const DEFAULT_OBJECT_ZERO_COPY_ENABLE: bool = DEFAULT_OBJECT_MMAP_READ_ENABLE;
|
pub const DEFAULT_OBJECT_ZERO_COPY_ENABLE: bool = DEFAULT_OBJECT_MMAP_READ_ENABLE;
|
||||||
|
|
||||||
|
/// Environment variable capping the byte length a single mmap-copy read may
|
||||||
|
/// materialize in memory.
|
||||||
|
///
|
||||||
|
/// The mmap-copy read path returns the whole requested range as one owned
|
||||||
|
/// allocation before the first byte is served. GET/heal shard reads request
|
||||||
|
/// the entire part span in one call, so for a large single-part object
|
||||||
|
/// (e.g. a multi-gigabyte non-multipart upload) an uncapped mmap-copy read
|
||||||
|
/// allocates the whole shard in memory — stalling first-byte latency past the
|
||||||
|
/// disk-read timeout and OOM-killing memory-limited deployments
|
||||||
|
/// (<https://github.com/rustfs/rustfs/issues/5123>). Reads longer than this
|
||||||
|
/// cap fall back to the bounded streaming reader instead.
|
||||||
|
///
|
||||||
|
/// - Purpose: Bound per-shard-read memory for mmap-based reads
|
||||||
|
/// - Acceptable values: byte count as an unsigned integer; `0` disables
|
||||||
|
/// mmap-copy for all non-empty reads (every read streams)
|
||||||
|
/// - Example: `export RUSTFS_OBJECT_MMAP_READ_MAX_LENGTH=8388608`
|
||||||
|
pub const ENV_OBJECT_MMAP_READ_MAX_LENGTH: &str = "RUSTFS_OBJECT_MMAP_READ_MAX_LENGTH";
|
||||||
|
|
||||||
|
/// Default mmap-copy read length cap: 32 MiB per shard read.
|
||||||
|
///
|
||||||
|
/// Large enough that typical multipart part shards (parts up to a few hundred
|
||||||
|
/// megabytes across the erasure set) keep the mmap fast path, small enough
|
||||||
|
/// that whole-part reads of huge single-part objects stream instead of
|
||||||
|
/// materializing gigabytes per shard.
|
||||||
|
///
|
||||||
|
/// The cap bounds memory per shard reader, so a single part read can still
|
||||||
|
/// materialize up to `data_shards x cap` bytes; raising the cap raises that
|
||||||
|
/// per-request bound proportionally.
|
||||||
|
pub const DEFAULT_OBJECT_MMAP_READ_MAX_LENGTH: usize = 32 * 1024 * 1024;
|
||||||
|
|||||||
@@ -22,7 +22,10 @@ use crate::diagnostics::get::{
|
|||||||
use crate::disk::{self, DiskAPI as _, DiskStore, FileReader, MmapCopyStageMetrics, error::DiskError};
|
use crate::disk::{self, DiskAPI as _, DiskStore, FileReader, MmapCopyStageMetrics, error::DiskError};
|
||||||
use crate::erasure::coding::{BitrotReader, BitrotWriterWrapper, CustomWriter};
|
use crate::erasure::coding::{BitrotReader, BitrotWriterWrapper, CustomWriter};
|
||||||
use bytes::Bytes;
|
use bytes::Bytes;
|
||||||
use rustfs_config::{DEFAULT_OBJECT_MMAP_READ_ENABLE, ENV_OBJECT_MMAP_READ_ENABLE, ENV_OBJECT_ZERO_COPY_ENABLE};
|
use rustfs_config::{
|
||||||
|
DEFAULT_OBJECT_MMAP_READ_ENABLE, DEFAULT_OBJECT_MMAP_READ_MAX_LENGTH, ENV_OBJECT_MMAP_READ_ENABLE,
|
||||||
|
ENV_OBJECT_MMAP_READ_MAX_LENGTH, ENV_OBJECT_ZERO_COPY_ENABLE,
|
||||||
|
};
|
||||||
use rustfs_utils::HashAlgorithm;
|
use rustfs_utils::HashAlgorithm;
|
||||||
use std::future::Future;
|
use std::future::Future;
|
||||||
use std::io::{self, Cursor};
|
use std::io::{self, Cursor};
|
||||||
@@ -81,6 +84,21 @@ pub(crate) fn object_mmap_read_enabled() -> bool {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Mmap-copy read length cap. Cached: this is consulted on every shard open
|
||||||
|
/// and `std::env::var` takes a process-global lock. In test builds the env
|
||||||
|
/// var is read directly so `temp_env` overrides take effect.
|
||||||
|
pub(crate) fn object_mmap_read_max_length() -> usize {
|
||||||
|
#[cfg(test)]
|
||||||
|
{
|
||||||
|
rustfs_utils::get_env_usize(ENV_OBJECT_MMAP_READ_MAX_LENGTH, DEFAULT_OBJECT_MMAP_READ_MAX_LENGTH)
|
||||||
|
}
|
||||||
|
#[cfg(not(test))]
|
||||||
|
{
|
||||||
|
static CACHED: std::sync::OnceLock<usize> = std::sync::OnceLock::new();
|
||||||
|
*CACHED.get_or_init(|| rustfs_utils::get_env_usize(ENV_OBJECT_MMAP_READ_MAX_LENGTH, DEFAULT_OBJECT_MMAP_READ_MAX_LENGTH))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[derive(Clone)]
|
#[derive(Clone)]
|
||||||
struct BitrotReaderSource {
|
struct BitrotReaderSource {
|
||||||
inline_data: Option<Bytes>,
|
inline_data: Option<Bytes>,
|
||||||
@@ -305,7 +323,28 @@ async fn open_disk_reader(
|
|||||||
let metrics_path = metrics_path.filter(|_| rustfs_io_metrics::get_stage_metrics_enabled());
|
let metrics_path = metrics_path.filter(|_| rustfs_io_metrics::get_stage_metrics_enabled());
|
||||||
let stage_metrics_enabled = metrics_path.is_some();
|
let stage_metrics_enabled = metrics_path.is_some();
|
||||||
|
|
||||||
if use_mmap_read && disk.is_local() {
|
// Mmap-copy materializes the whole `offset..offset+length` range as one
|
||||||
|
// owned allocation before any byte is served, and GET/heal shard reads
|
||||||
|
// request the entire part span in one call. Over-cap reads (e.g. a huge
|
||||||
|
// single-part object, rustfs#5123) take the bounded streaming path below.
|
||||||
|
let use_mmap_copy = if use_mmap_read && disk.is_local() {
|
||||||
|
let mmap_read_cap = object_mmap_read_max_length();
|
||||||
|
let within_cap = length <= mmap_read_cap;
|
||||||
|
if !within_cap {
|
||||||
|
rustfs_io_metrics::record_zero_copy_fallback("read_length_exceeds_cap");
|
||||||
|
debug!(
|
||||||
|
length,
|
||||||
|
mmap_read_cap,
|
||||||
|
path = %path,
|
||||||
|
"zero_copy_read_skipped_over_cap"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
within_cap
|
||||||
|
} else {
|
||||||
|
false
|
||||||
|
};
|
||||||
|
|
||||||
|
if use_mmap_copy {
|
||||||
let start = stage_metrics_enabled.then(Instant::now);
|
let start = stage_metrics_enabled.then(Instant::now);
|
||||||
let zero_copy_start = Instant::now();
|
let zero_copy_start = Instant::now();
|
||||||
let mmap_metrics = metrics_path.map(|metrics_path| MmapCopyStageMetrics {
|
let mmap_metrics = metrics_path.map(|metrics_path| MmapCopyStageMetrics {
|
||||||
@@ -685,6 +724,96 @@ mod tests {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn object_mmap_read_max_length_defaults_and_env_override() {
|
||||||
|
temp_env::with_var(ENV_OBJECT_MMAP_READ_MAX_LENGTH, None::<&str>, || {
|
||||||
|
assert_eq!(object_mmap_read_max_length(), DEFAULT_OBJECT_MMAP_READ_MAX_LENGTH);
|
||||||
|
});
|
||||||
|
temp_env::with_var(ENV_OBJECT_MMAP_READ_MAX_LENGTH, Some("1024"), || {
|
||||||
|
assert_eq!(object_mmap_read_max_length(), 1024);
|
||||||
|
});
|
||||||
|
temp_env::with_var(ENV_OBJECT_MMAP_READ_MAX_LENGTH, Some("0"), || {
|
||||||
|
assert_eq!(object_mmap_read_max_length(), 0);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
// rustfs#5123: whole-part shard reads of large single-part objects must not
|
||||||
|
// be materialized in memory by the mmap-copy path; over-cap reads stream.
|
||||||
|
#[tokio::test]
|
||||||
|
async fn open_disk_reader_streams_when_length_exceeds_mmap_cap() {
|
||||||
|
use crate::disk::endpoint::Endpoint;
|
||||||
|
use crate::disk::{DiskOption, new_disk};
|
||||||
|
use tokio::io::AsyncReadExt;
|
||||||
|
|
||||||
|
let dir = tempfile::tempdir().expect("tempdir should be created");
|
||||||
|
let mut endpoint =
|
||||||
|
Endpoint::try_from(dir.path().to_str().expect("tempdir path should be utf8")).expect("endpoint should parse");
|
||||||
|
endpoint.set_pool_index(0);
|
||||||
|
endpoint.set_set_index(0);
|
||||||
|
endpoint.set_disk_index(0);
|
||||||
|
let disk = new_disk(
|
||||||
|
&endpoint,
|
||||||
|
&DiskOption {
|
||||||
|
cleanup: false,
|
||||||
|
health_check: false,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("local disk should be created");
|
||||||
|
|
||||||
|
let payload = vec![7u8; 4096];
|
||||||
|
disk.make_volume("test-bucket").await.expect("volume should be created");
|
||||||
|
disk.write_all("test-bucket", "obj/part.1", Bytes::from(payload.clone()))
|
||||||
|
.await
|
||||||
|
.expect("shard file should be written");
|
||||||
|
|
||||||
|
temp_env::async_with_vars([(ENV_OBJECT_MMAP_READ_MAX_LENGTH, Some("1024"))], async {
|
||||||
|
let over_cap = open_disk_reader(&disk, "test-bucket", "obj/part.1", 0, payload.len(), true, None)
|
||||||
|
.await
|
||||||
|
.expect("over-cap read should open");
|
||||||
|
assert!(
|
||||||
|
matches!(over_cap, ShardReader::Stream(_)),
|
||||||
|
"read longer than the mmap cap must take the streaming path"
|
||||||
|
);
|
||||||
|
let mut over_cap = over_cap;
|
||||||
|
let mut streamed = Vec::new();
|
||||||
|
over_cap
|
||||||
|
.read_to_end(&mut streamed)
|
||||||
|
.await
|
||||||
|
.expect("streaming fallback should read the range");
|
||||||
|
assert_eq!(streamed, payload, "streaming fallback must return the same bytes");
|
||||||
|
|
||||||
|
let under_cap = open_disk_reader(&disk, "test-bucket", "obj/part.1", 0, 512, true, None)
|
||||||
|
.await
|
||||||
|
.expect("under-cap read should open");
|
||||||
|
assert!(
|
||||||
|
matches!(under_cap, ShardReader::InMemory(_)),
|
||||||
|
"read within the mmap cap keeps the mmap-copy fast path"
|
||||||
|
);
|
||||||
|
|
||||||
|
let at_cap = open_disk_reader(&disk, "test-bucket", "obj/part.1", 0, 1024, true, None)
|
||||||
|
.await
|
||||||
|
.expect("at-cap read should open");
|
||||||
|
assert!(
|
||||||
|
matches!(at_cap, ShardReader::InMemory(_)),
|
||||||
|
"read of exactly the mmap cap keeps the mmap-copy fast path"
|
||||||
|
);
|
||||||
|
})
|
||||||
|
.await;
|
||||||
|
|
||||||
|
// Cap of 0 disables mmap-copy for every non-empty read.
|
||||||
|
temp_env::async_with_vars([(ENV_OBJECT_MMAP_READ_MAX_LENGTH, Some("0"))], async {
|
||||||
|
let disabled = open_disk_reader(&disk, "test-bucket", "obj/part.1", 0, 512, true, None)
|
||||||
|
.await
|
||||||
|
.expect("read with cap 0 should open");
|
||||||
|
assert!(
|
||||||
|
matches!(disabled, ShardReader::Stream(_)),
|
||||||
|
"cap 0 must route every non-empty read to the streaming path"
|
||||||
|
);
|
||||||
|
})
|
||||||
|
.await;
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_create_bitrot_reader_with_inline_data() {
|
async fn test_create_bitrot_reader_with_inline_data() {
|
||||||
let test_data = b"hello world test data";
|
let test_data = b"hello world test data";
|
||||||
|
|||||||
Reference in New Issue
Block a user