mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-11 07:36:53 +00:00
feat(get): tune output response handoff (#3956)
This commit is contained in:
@@ -405,6 +405,8 @@ where
|
||||
E: ErasureDecodeEngine + Clone + Send + Sync + 'static,
|
||||
{
|
||||
fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll<io::Result<()>> {
|
||||
let filled_before_poll = buf.filled().len();
|
||||
|
||||
loop {
|
||||
if self.output_pos < self.output_buf.len() {
|
||||
if self.prefetched_bufs.len() < self.fill_policy.max_inflight()
|
||||
@@ -431,7 +433,10 @@ where
|
||||
copy_start.elapsed().as_secs_f64(),
|
||||
);
|
||||
}
|
||||
return Poll::Ready(Ok(()));
|
||||
if buf.remaining() == 0 {
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
if let Some(next_buf) = self.prefetched_bufs.pop_front() {
|
||||
@@ -441,6 +446,9 @@ where
|
||||
continue;
|
||||
}
|
||||
|
||||
if self.prefetch_error.is_some() && buf.filled().len() > filled_before_poll {
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
if let Some(err) = self.prefetch_error.take() {
|
||||
return Poll::Ready(Err(err));
|
||||
}
|
||||
@@ -452,7 +460,11 @@ where
|
||||
if self.output_wait_started_at.is_none() {
|
||||
self.output_wait_started_at = Some(Instant::now());
|
||||
}
|
||||
let prefetch = ready!(self.poll_prefetch(cx));
|
||||
let prefetch = match self.poll_prefetch(cx) {
|
||||
Poll::Ready(result) => result,
|
||||
Poll::Pending if buf.filled().len() > filled_before_poll => return Poll::Ready(Ok(())),
|
||||
Poll::Pending => return Poll::Pending,
|
||||
};
|
||||
if let Some(started_at) = self.output_wait_started_at.take() {
|
||||
rustfs_io_metrics::record_get_object_fill_waited_by_output(
|
||||
self.metrics_path,
|
||||
|
||||
@@ -339,6 +339,71 @@ pub fn record_get_object_response_handoff(
|
||||
record_get_object_response_handoff_duration("s3_handler", duration_secs);
|
||||
}
|
||||
|
||||
/// Record ReaderStream capacity chosen for GetObject handoff.
|
||||
#[inline(always)]
|
||||
pub fn record_get_object_reader_stream_buffer_size(strategy: &str, buffer_source: &str, buffer_size_bytes: usize) {
|
||||
if !get_stage_metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
histogram!(
|
||||
"rustfs_io_get_object_reader_stream_buffer_size_bytes",
|
||||
"strategy" => strategy.to_string(),
|
||||
"buffer_source" => buffer_source.to_string()
|
||||
)
|
||||
.record(usize_to_f64(buffer_size_bytes));
|
||||
}
|
||||
|
||||
/// Record ReaderStream poll outcomes for GetObject handoff attribution.
|
||||
#[inline(always)]
|
||||
pub fn record_get_object_reader_stream_poll(
|
||||
strategy: &str,
|
||||
buffer_source: &str,
|
||||
outcome: &'static str,
|
||||
remaining_before: usize,
|
||||
bytes: usize,
|
||||
duration_secs: f64,
|
||||
) {
|
||||
if !get_stage_metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
let bytes = u64::try_from(bytes).unwrap_or(u64::MAX);
|
||||
counter!(
|
||||
"rustfs_io_get_object_reader_stream_poll_total",
|
||||
"strategy" => strategy.to_string(),
|
||||
"buffer_source" => buffer_source.to_string(),
|
||||
"outcome" => outcome
|
||||
)
|
||||
.increment(1);
|
||||
counter!(
|
||||
"rustfs_io_get_object_reader_stream_poll_bytes_total",
|
||||
"strategy" => strategy.to_string(),
|
||||
"buffer_source" => buffer_source.to_string(),
|
||||
"outcome" => outcome
|
||||
)
|
||||
.increment(bytes);
|
||||
histogram!(
|
||||
"rustfs_io_get_object_reader_stream_poll_remaining_bytes",
|
||||
"strategy" => strategy.to_string(),
|
||||
"buffer_source" => buffer_source.to_string(),
|
||||
"outcome" => outcome
|
||||
)
|
||||
.record(usize_to_f64(remaining_before));
|
||||
histogram!(
|
||||
"rustfs_io_get_object_reader_stream_poll_bytes",
|
||||
"strategy" => strategy.to_string(),
|
||||
"buffer_source" => buffer_source.to_string(),
|
||||
"outcome" => outcome
|
||||
)
|
||||
.record(usize_to_f64(bytes as usize));
|
||||
histogram!(
|
||||
"rustfs_io_get_object_reader_stream_poll_duration_seconds",
|
||||
"strategy" => strategy.to_string(),
|
||||
"buffer_source" => buffer_source.to_string(),
|
||||
"outcome" => outcome
|
||||
)
|
||||
.record(duration_secs);
|
||||
}
|
||||
|
||||
/// Record I/O queue congestion observation.
|
||||
#[inline(always)]
|
||||
pub fn record_io_queue_congestion() {
|
||||
@@ -1548,6 +1613,8 @@ mod tests {
|
||||
record_get_object_fill_waited_by_output("codec_streaming", "single_inflight", 0.0003);
|
||||
record_get_object_fill_cancelled_on_drop("codec_streaming", "single_inflight");
|
||||
record_get_object_reader_prefetch_bytes("codec_streaming", "single_inflight", 4096);
|
||||
record_get_object_reader_stream_buffer_size("standard", "selected", 131072);
|
||||
record_get_object_reader_stream_poll("standard", "selected", "ready_data", 8192, 4096, 0.0002);
|
||||
|
||||
assert!(0.0003_f64.is_sign_positive());
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user