perf(get): reduce metrics and streaming handoff overhead (#3907)

This commit is contained in:
houseme
2026-06-26 18:12:51 +08:00
committed by GitHub
parent 30957e2b51
commit 0321e4c6ca
8 changed files with 938 additions and 71 deletions
+278 -60
View File
@@ -143,6 +143,7 @@ use s3s::dto::{
ServerSideEncryption, StorageClass, StreamingBlob, TaggingHeader, Timestamp, TimestampFormat, WebsiteRedirectLocation,
};
use s3s::header::{X_AMZ_RESTORE, X_AMZ_RESTORE_OUTPUT_PATH};
use s3s::stream::{ByteStream, RemainingLength};
use s3s::{S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error};
const DEFAULT_PUT_LARGE_CONCURRENCY_TUNING_MIN_SIZE_BYTES: i64 = 32 * 1024 * 1024;
@@ -169,8 +170,13 @@ use super::storage_api::object_usecase::{
StorageObjectToDelete as ObjectToDelete, StoragePutObjReader as PutObjReader,
};
type S3StdError = Box<dyn std::error::Error + Send + Sync + 'static>;
const ACCEPT_RANGES_BYTES: &str = "bytes";
const MAX_GET_OBJECT_MEMORY_BUFFER_BYTES: i64 = 64 * 1024 * 1024;
const MEDIUM_CONCURRENCY_GET_OBJECT_MEMORY_BUFFER_BYTES: i64 = 8 * 1024 * 1024;
const HIGH_CONCURRENCY_GET_OBJECT_MEMORY_BUFFER_BYTES: i64 = 4 * 1024 * 1024;
const VERY_HIGH_CONCURRENCY_GET_OBJECT_MEMORY_BUFFER_BYTES: i64 = 1024 * 1024;
const LOG_COMPONENT_APP: &str = "app";
const LOG_SUBSYSTEM_OBJECT: &str = "object";
const EVENT_PUT_OBJECT_STORE_INFLIGHT_SLOW: &str = "put_object_store_inflight_slow";
@@ -395,6 +401,14 @@ pin_project! {
}
}
pin_project! {
struct GetObjectReaderStream<R> {
#[pin]
inner: ReaderStream<R>,
remaining: usize,
}
}
impl MemoryTrackedBytesStream {
fn new(bytes: Bytes, guard: Option<rustfs_io_metrics::MemoryGaugeGuard>) -> Self {
Self {
@@ -405,6 +419,18 @@ impl MemoryTrackedBytesStream {
}
}
impl<R> GetObjectReaderStream<R>
where
R: AsyncRead,
{
fn new(reader: R, capacity: usize, remaining: usize) -> Self {
Self {
inner: ReaderStream::with_capacity(reader, capacity),
remaining,
}
}
}
impl futures::Stream for MemoryTrackedBytesStream {
type Item = std::io::Result<Bytes>;
@@ -419,6 +445,50 @@ impl futures::Stream for MemoryTrackedBytesStream {
}
}
impl<R> futures::Stream for GetObjectReaderStream<R>
where
R: AsyncRead,
{
type Item = Result<Bytes, S3StdError>;
fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
let mut this = self.project();
if *this.remaining == 0 {
return Poll::Ready(None);
}
match this.inner.as_mut().poll_next(cx) {
Poll::Ready(Some(Ok(mut bytes))) => {
if bytes.len() > *this.remaining {
bytes.truncate(*this.remaining);
}
*this.remaining -= bytes.len();
if bytes.is_empty() {
Poll::Ready(None)
} else {
Poll::Ready(Some(Ok(bytes)))
}
}
Poll::Ready(Some(Err(err))) => Poll::Ready(Some(Err(Box::new(err)))),
Poll::Ready(None) => Poll::Ready(None),
Poll::Pending => Poll::Pending,
}
}
fn size_hint(&self) -> (usize, Option<usize>) {
self.inner.size_hint()
}
}
impl<R> ByteStream for GetObjectReaderStream<R>
where
R: AsyncRead,
{
fn remaining_length(&self) -> RemainingLength {
RemainingLength::new_exact(self.remaining)
}
}
impl<R> ExtractArchiveEtagReader<R> {
fn new(inner: R, etag: Arc<Mutex<Option<String>>>) -> Self {
Self {
@@ -744,14 +814,56 @@ fn object_seek_support_threshold() -> usize {
})
}
fn object_seek_support_concurrency_thresholds() -> (usize, usize) {
static OBJECT_SEEK_SUPPORT_CONCURRENCY_THRESHOLDS: OnceLock<(usize, usize)> = OnceLock::new();
*OBJECT_SEEK_SUPPORT_CONCURRENCY_THRESHOLDS.get_or_init(|| {
let medium = rustfs_utils::get_env_usize(
rustfs_config::ENV_OBJECT_MEDIUM_CONCURRENCY_THRESHOLD,
rustfs_config::DEFAULT_OBJECT_MEDIUM_CONCURRENCY_THRESHOLD,
)
.max(1);
let high = rustfs_utils::get_env_usize(
rustfs_config::ENV_OBJECT_HIGH_CONCURRENCY_THRESHOLD,
rustfs_config::DEFAULT_OBJECT_HIGH_CONCURRENCY_THRESHOLD,
)
.max(medium + 1);
(medium, high)
})
}
fn concurrency_aware_seek_support_threshold(configured_threshold: i64, concurrent_requests: usize) -> i64 {
let (medium_threshold, high_threshold) = object_seek_support_concurrency_thresholds();
let effective_threshold = configured_threshold.min(MAX_GET_OBJECT_MEMORY_BUFFER_BYTES);
if concurrent_requests >= high_threshold.saturating_mul(2) {
return effective_threshold.min(VERY_HIGH_CONCURRENCY_GET_OBJECT_MEMORY_BUFFER_BYTES);
}
if concurrent_requests >= high_threshold {
return effective_threshold.min(HIGH_CONCURRENCY_GET_OBJECT_MEMORY_BUFFER_BYTES);
}
if concurrent_requests >= medium_threshold {
return effective_threshold.min(MEDIUM_CONCURRENCY_GET_OBJECT_MEMORY_BUFFER_BYTES);
}
effective_threshold
}
fn should_buffer_get_object_in_memory(
info: &ObjectInfo,
response_content_length: i64,
part_number: Option<usize>,
has_range: bool,
concurrent_requests: usize,
) -> bool {
let configured_threshold = object_seek_support_threshold() as i64;
should_buffer_get_object_in_memory_with_threshold(info, response_content_length, part_number, has_range, configured_threshold)
should_buffer_get_object_in_memory_with_threshold(
info,
response_content_length,
part_number,
has_range,
configured_threshold,
concurrent_requests,
)
}
fn should_buffer_get_object_in_memory_with_threshold(
@@ -760,12 +872,13 @@ fn should_buffer_get_object_in_memory_with_threshold(
part_number: Option<usize>,
has_range: bool,
configured_threshold: i64,
concurrent_requests: usize,
) -> bool {
if part_number.is_some() || has_range || response_content_length <= 0 || configured_threshold <= 0 {
return false;
}
let effective_threshold = configured_threshold.min(MAX_GET_OBJECT_MEMORY_BUFFER_BYTES);
let effective_threshold = concurrency_aware_seek_support_threshold(configured_threshold, concurrent_requests);
if configured_threshold > MAX_GET_OBJECT_MEMORY_BUFFER_BYTES
&& GET_OBJECT_BUFFER_THRESHOLD_WARNED
.compare_exchange(false, true, Ordering::Relaxed, Ordering::Relaxed)
@@ -1609,16 +1722,19 @@ impl DefaultObjectUsecase {
let expected = usize::try_from(response_content_length.max(0)).unwrap_or(usize::MAX);
let (stream_buffer_size, buffer_source) =
resolve_reader_stream_buffer_size(stream_buffer_size, get_reader_stream_buffer_size_override());
rustfs_io_metrics::record_get_object_stream_strategy(
stream_strategy.as_str(),
stream_buffer_size,
response_content_length,
);
let handoff_start = std::time::Instant::now();
let get_stage_metrics_enabled = rustfs_io_metrics::get_stage_metrics_enabled();
if get_stage_metrics_enabled {
rustfs_io_metrics::record_get_object_stream_strategy(
stream_strategy.as_str(),
stream_buffer_size,
response_content_length,
);
}
let handoff_start = get_stage_metrics_enabled.then(std::time::Instant::now);
#[cfg(feature = "tracing-chunk-debug")]
let stream = {
let mut emitted = 0usize;
ReaderStream::with_capacity(reader, stream_buffer_size).inspect(move |item| match item {
GetObjectReaderStream::new(reader, stream_buffer_size, expected).inspect(move |item| match item {
Ok(bytes) => {
emitted += bytes.len();
tracing::debug!(emitted, expected, chunk_len = bytes.len(), "GetObject ReaderStream emitted bytes");
@@ -1634,15 +1750,17 @@ impl DefaultObjectUsecase {
})
};
#[cfg(not(feature = "tracing-chunk-debug"))]
let stream = ReaderStream::with_capacity(reader, stream_buffer_size);
let blob = StreamingBlob::wrap(bytes_stream(stream, expected));
rustfs_io_metrics::record_get_object_response_handoff(
stream_strategy.as_str(),
buffer_source,
stream_buffer_size,
response_content_length,
handoff_start.elapsed().as_secs_f64(),
);
let stream = GetObjectReaderStream::new(reader, stream_buffer_size, expected);
let blob = StreamingBlob::new(stream);
if let Some(handoff_start) = handoff_start {
rustfs_io_metrics::record_get_object_response_handoff(
stream_strategy.as_str(),
buffer_source,
stream_buffer_size,
response_content_length,
handoff_start.elapsed().as_secs_f64(),
);
}
Some(blob)
}
@@ -1793,17 +1911,32 @@ impl DefaultObjectUsecase {
) -> S3Result<GetObjectPreparedRead<'a>> {
let h = req.headers.clone();
let io_planning = Self::acquire_get_object_io_planning(manager, wrapper, timeout_config, bucket, key).await?;
let store_lookup_start = std::time::Instant::now();
let store_lookup_start = rustfs_io_metrics::get_stage_metrics_enabled().then(std::time::Instant::now);
let store = get_validated_store(bucket).await?;
rustfs_io_metrics::record_get_object_stage_duration(
"s3_handler",
"store_lookup",
store_lookup_start.elapsed().as_secs_f64(),
);
if let Some(store_lookup_start) = store_lookup_start {
rustfs_io_metrics::record_get_object_stage_duration(
"s3_handler",
"store_lookup",
store_lookup_start.elapsed().as_secs_f64(),
);
}
let read_start = std::time::Instant::now();
let read_setup =
Self::prepare_get_object_read(req, &store, manager, bucket, key, rs, h, opts, part_number, read_start).await?;
let read_stage_start = rustfs_io_metrics::get_stage_metrics_enabled().then_some(read_start);
let read_setup = Self::prepare_get_object_read(
req,
&store,
manager,
bucket,
key,
rs,
h,
opts,
part_number,
read_start,
read_stage_start,
)
.await?;
Ok(GetObjectPreparedRead { io_planning, read_setup })
}
@@ -1820,16 +1953,19 @@ impl DefaultObjectUsecase {
opts: &ObjectOptions,
part_number: Option<usize>,
read_start: std::time::Instant,
read_stage_start: Option<std::time::Instant>,
) -> S3Result<GetObjectReadSetup> {
let reader = store
.get_object_reader(bucket, key, rs.clone(), h, opts)
.await
.map_err(map_get_object_reader_error)?;
rustfs_io_metrics::record_get_object_stage_duration(
"s3_handler",
"store_reader_setup",
read_start.elapsed().as_secs_f64(),
);
if let Some(read_stage_start) = read_stage_start {
rustfs_io_metrics::record_get_object_stage_duration(
"s3_handler",
"store_reader_setup",
read_stage_start.elapsed().as_secs_f64(),
);
}
let info = reader.object_info;
@@ -2098,6 +2234,7 @@ impl DefaultObjectUsecase {
response_content_length: i64,
optimal_buffer_size: usize,
enable_readahead: bool,
concurrent_requests: usize,
part_number: Option<usize>,
has_range: bool,
encryption_applied: bool,
@@ -2107,7 +2244,7 @@ impl DefaultObjectUsecase {
{
if encryption_applied {
let should_buffer_encrypted_object =
should_buffer_get_object_in_memory(info, response_content_length, part_number, has_range);
should_buffer_get_object_in_memory(info, response_content_length, part_number, has_range, concurrent_requests);
if should_buffer_encrypted_object {
let mut buf = Vec::with_capacity(response_content_length as usize);
@@ -2139,7 +2276,7 @@ impl DefaultObjectUsecase {
}
let should_provide_seek_support =
should_buffer_get_object_in_memory(info, response_content_length, part_number, has_range);
should_buffer_get_object_in_memory(info, response_content_length, part_number, has_range, concurrent_requests);
if should_provide_seek_support {
let mut buf = Vec::with_capacity(response_content_length as usize);
@@ -2882,6 +3019,7 @@ impl DefaultObjectUsecase {
response_content_length,
optimal_buffer_size,
enable_readahead,
concurrent_requests,
part_number,
rs.is_some(),
encryption_applied,
@@ -2974,13 +3112,15 @@ impl DefaultObjectUsecase {
let helper = OperationHelper::new(&req, EventName::ObjectAccessedGet, S3Operation::GetObject).suppress_event();
// mc get 3
let request_context_start = std::time::Instant::now();
let request_context_start = rustfs_io_metrics::get_stage_metrics_enabled().then(std::time::Instant::now);
let request_context = Self::prepare_get_object_request_context(&req).await?;
rustfs_io_metrics::record_get_object_stage_duration(
"s3_handler",
"request_context",
request_context_start.elapsed().as_secs_f64(),
);
if let Some(request_context_start) = request_context_start {
rustfs_io_metrics::record_get_object_stage_duration(
"s3_handler",
"request_context",
request_context_start.elapsed().as_secs_f64(),
);
}
let GetObjectRequestContext {
bucket,
key,
@@ -3025,14 +3165,16 @@ impl DefaultObjectUsecase {
encryption_applied,
} = read_setup;
let versioning_start = std::time::Instant::now();
let versioning_start = rustfs_io_metrics::get_stage_metrics_enabled().then(std::time::Instant::now);
let versioned = BucketVersioningSys::prefix_enabled(&bucket, &key).await;
rustfs_io_metrics::record_get_object_stage_duration(
"s3_handler",
"versioning_lookup",
versioning_start.elapsed().as_secs_f64(),
);
let output_build_start = std::time::Instant::now();
if let Some(versioning_start) = versioning_start {
rustfs_io_metrics::record_get_object_stage_duration(
"s3_handler",
"versioning_lookup",
versioning_start.elapsed().as_secs_f64(),
);
}
let output_build_start = rustfs_io_metrics::get_stage_metrics_enabled().then(std::time::Instant::now);
let output_context = self
.build_get_object_output_context(
&req,
@@ -3060,11 +3202,13 @@ impl DefaultObjectUsecase {
versioned,
)
.await?;
rustfs_io_metrics::record_get_object_stage_duration(
"s3_handler",
"output_build",
output_build_start.elapsed().as_secs_f64(),
);
if let Some(output_build_start) = output_build_start {
rustfs_io_metrics::record_get_object_stage_duration(
"s3_handler",
"output_build",
output_build_start.elapsed().as_secs_f64(),
);
}
let GetObjectOutputContext {
output,
event_info,
@@ -5520,7 +5664,7 @@ mod tests {
let configured_threshold = 20_i64 * 1024 * 1024 * 1024;
let response_len = 80_i64 * 1024 * 1024;
let should_buffer =
should_buffer_get_object_in_memory_with_threshold(&info, response_len, None, false, configured_threshold);
should_buffer_get_object_in_memory_with_threshold(&info, response_len, None, false, configured_threshold, 1);
assert!(
!should_buffer,
@@ -5538,21 +5682,24 @@ mod tests {
1024 * 1024,
None,
false,
configured_threshold
configured_threshold,
1
));
assert!(!should_buffer_get_object_in_memory_with_threshold(
&info,
1024 * 1024,
Some(1),
false,
configured_threshold
configured_threshold,
1
));
assert!(!should_buffer_get_object_in_memory_with_threshold(
&info,
1024 * 1024,
None,
true,
configured_threshold
configured_threshold,
1
));
}
@@ -5566,14 +5713,16 @@ mod tests {
configured_threshold,
None,
false,
configured_threshold
configured_threshold,
1
));
assert!(!should_buffer_get_object_in_memory_with_threshold(
&info,
configured_threshold + 1,
None,
false,
configured_threshold
configured_threshold,
1
));
}
@@ -5587,16 +5736,49 @@ mod tests {
0,
None,
false,
configured_threshold
configured_threshold,
1
));
assert!(!should_buffer_get_object_in_memory_with_threshold(
&info,
-1,
None,
false,
configured_threshold
configured_threshold,
1
));
assert!(!should_buffer_get_object_in_memory_with_threshold(&info, 1024, None, false, 0, 1));
}
#[test]
fn should_buffer_get_object_in_memory_reduces_threshold_under_concurrency() {
let info = ObjectInfo::default();
let configured_threshold = 10_i64 * 1024 * 1024;
assert!(should_buffer_get_object_in_memory_with_threshold(
&info,
configured_threshold,
None,
false,
configured_threshold,
1
));
assert!(!should_buffer_get_object_in_memory_with_threshold(
&info,
configured_threshold,
None,
false,
configured_threshold,
32
));
assert!(should_buffer_get_object_in_memory_with_threshold(
&info,
4_i64 * 1024 * 1024,
None,
false,
configured_threshold,
rustfs_config::DEFAULT_OBJECT_HIGH_CONCURRENCY_THRESHOLD
));
assert!(!should_buffer_get_object_in_memory_with_threshold(&info, 1024, None, false, 0));
}
struct ReadProbeReader {
@@ -5627,6 +5809,7 @@ mod tests {
18_i64 * 1024 * 1024 * 1024,
128 * 1024,
true,
1,
None,
false,
false,
@@ -5659,6 +5842,7 @@ mod tests {
18_i64 * 1024 * 1024 * 1024,
128 * 1024,
true,
1,
None,
false,
true,
@@ -5861,6 +6045,40 @@ mod tests {
assert_eq!(err.code(), &S3ErrorCode::UnexpectedContent);
}
#[tokio::test]
async fn get_object_reader_stream_tracks_remaining_length() {
let mut stream = GetObjectReaderStream::new(std::io::Cursor::new(b"hello".to_vec()), 2, 5);
assert_eq!(stream.remaining_length().exact(), Some(5));
let first = stream
.next()
.await
.expect("reader stream should emit first chunk")
.expect("first chunk should read");
assert_eq!(first.as_ref(), b"he");
assert_eq!(stream.remaining_length().exact(), Some(3));
}
#[tokio::test]
async fn get_object_reader_stream_truncates_to_expected_length() {
let stream = GetObjectReaderStream::new(std::io::Cursor::new(b"hello!".to_vec()), 64, 5);
let chunks = stream
.collect::<Vec<_>>()
.await
.into_iter()
.collect::<Result<Vec<_>, _>>()
.expect("reader stream should read");
let body = chunks.into_iter().fold(Vec::new(), |mut acc, chunk| {
acc.extend_from_slice(&chunk);
acc
});
assert_eq!(body, b"hello");
}
#[tokio::test]
async fn pooled_buffer_reader_keeps_buffer_alive_until_consumed() {
use tokio::io::AsyncReadExt;
+1
View File
@@ -23,6 +23,7 @@ pub(crate) async fn init_observability_runtime(ctx: CancellationToken) {
if startup_runtime_sources::observability_metric_enabled() {
startup_runtime_sources::set_put_stage_metrics_enabled(true);
startup_runtime_sources::set_get_stage_metrics_enabled(true);
startup_runtime_sources::init_metrics_runtime(ctx.clone());
crate::memory_observability::init_memory_observability(ctx.clone());
init_auto_tuner(ctx).await;
+4
View File
@@ -80,6 +80,10 @@ pub(crate) fn set_put_stage_metrics_enabled(enabled: bool) {
rustfs_io_metrics::set_put_stage_metrics_enabled(enabled);
}
pub(crate) fn set_get_stage_metrics_enabled(enabled: bool) {
rustfs_io_metrics::set_get_stage_metrics_enabled(enabled);
}
pub(crate) fn init_tls_metrics() {
rustfs_tls_runtime::init_tls_metrics();
}