feat(ecstore): expose fdatasync group wait metrics (#6299)

Add PUT-stage diagnostics for file fdatasync group commit wait time, per-group outstanding depth, and rename disk completion position. These metrics keep the existing default-off PUT stage gate and do not change group commit scheduling or quorum behavior.

Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
houseme
2026-08-20 16:59:37 +08:00
committed by GitHub
parent 76eb9c72e4
commit 621fcb93c7
3 changed files with 198 additions and 9 deletions
+44 -4
View File
@@ -23,6 +23,7 @@ use std::{
io,
path::{Component, Path, PathBuf},
sync::{Arc, LazyLock, Weak},
time::Instant,
};
use tokio::fs;
use tokio::sync::{
@@ -766,6 +767,8 @@ type FileFdatasyncGroupKey = usize;
struct FileFdatasyncWaiter {
files: Vec<PathBuf>,
enqueued_at: Option<Instant>,
wait_role: &'static str,
result_tx: oneshot::Sender<SharedFileFdatasyncResult>,
}
@@ -799,6 +802,7 @@ struct FileFdatasyncGroup {
#[derive(Default)]
struct FileFdatasyncGroupInner {
worker_running: bool,
pending_files: usize,
pending: VecDeque<FileFdatasyncWaiter>,
}
@@ -865,11 +869,30 @@ impl FileFdatasyncGroupCommit {
};
let file_count = files.len();
let mut group_state = group.inner.lock();
group_state.pending.push_back(FileFdatasyncWaiter { files, result_tx });
let start_worker = !group_state.worker_running;
let wait_role = if start_worker {
rustfs_io_metrics::PUT_RENAME_FDATASYNC_GROUP_WAIT_ROLE_LEADER
} else {
rustfs_io_metrics::PUT_RENAME_FDATASYNC_GROUP_WAIT_ROLE_FOLLOWER
};
group_state.pending.push_back(FileFdatasyncWaiter {
files,
enqueued_at: rustfs_io_metrics::put_stage_timer(),
wait_role,
result_tx,
});
group_state.pending_files += file_count;
if start_worker {
group_state.worker_running = true;
}
rustfs_io_metrics::record_put_rename_fdatasync_group_outstanding(
rustfs_io_metrics::PUT_RENAME_FDATASYNC_GROUP_OUTSTANDING_STATE_ENQUEUE_WAITERS,
group_state.pending.len(),
);
rustfs_io_metrics::record_put_rename_fdatasync_group_outstanding(
rustfs_io_metrics::PUT_RENAME_FDATASYNC_GROUP_OUTSTANDING_STATE_ENQUEUE_FILES,
group_state.pending_files,
);
registry.total_waiters += 1;
registry.total_files += file_count;
drop(group_state);
@@ -911,9 +934,11 @@ async fn run_file_fdatasync_group_worker(group: Arc<FileFdatasyncGroup>) {
#[cfg(test)]
file_sync_probe::run_before_group_batch();
tokio::task::yield_now().await;
let batch: Vec<FileFdatasyncWaiter> = {
let (batch, batch_file_count): (Vec<FileFdatasyncWaiter>, usize) = {
let mut group_state = group.inner.lock();
group_state.pending.drain(..).collect()
let batch_file_count = group_state.pending_files;
group_state.pending_files = 0;
(group_state.pending.drain(..).collect(), batch_file_count)
};
if batch.is_empty() {
let mut group_state = group.inner.lock();
@@ -923,8 +948,23 @@ async fn run_file_fdatasync_group_worker(group: Arc<FileFdatasyncGroup>) {
return;
}
rustfs_io_metrics::record_put_rename_fdatasync_group_outstanding(
rustfs_io_metrics::PUT_RENAME_FDATASYNC_GROUP_OUTSTANDING_STATE_BATCH_WAITERS,
batch.len(),
);
rustfs_io_metrics::record_put_rename_fdatasync_group_outstanding(
rustfs_io_metrics::PUT_RENAME_FDATASYNC_GROUP_OUTSTANDING_STATE_BATCH_FILES,
batch_file_count,
);
for waiter in &batch {
if let Some(enqueued_at) = waiter.enqueued_at {
rustfs_io_metrics::record_put_rename_fdatasync_group_wait(
waiter.wait_role,
enqueued_at.elapsed().as_secs_f64() * 1000.0,
);
}
}
let batch_files: Vec<PathBuf> = batch.iter().flat_map(|waiter| waiter.files.iter().cloned()).collect();
let batch_file_count = batch_files.len();
#[cfg(test)]
file_sync_probe::record_group_batch(batch_file_count);
rustfs_io_metrics::record_put_rename_fdatasync_batch(
@@ -74,7 +74,10 @@ use std::{
collections::{HashMap, HashSet, VecDeque},
future::Future,
pin::Pin,
sync::OnceLock,
sync::{
OnceLock,
atomic::{AtomicUsize, Ordering},
},
task::{Context, Poll},
time::{Duration, Instant},
};
@@ -3368,6 +3371,8 @@ impl SetDisks {
// preserving slot-indexed quorum and convergence accounting without a
// scheduler task for every disk.
let fanout = tokio::spawn(async move {
let successful_rename_completion_rank =
rustfs_io_metrics::put_stage_metrics_enabled().then(|| Arc::new(AtomicUsize::new(0)));
let futures = fanout_disks
.into_iter()
.zip(fanout_file_infos.iter())
@@ -3377,6 +3382,7 @@ impl SetDisks {
let src_object = fanout_src_object.clone();
let dst_object = fanout_dst_object.clone();
let dst_bucket = fanout_dst_bucket.clone();
let successful_rename_completion_rank = successful_rename_completion_rank.clone();
std::panic::AssertUnwindSafe(async move {
// Test-only introspection guard: counts this operation as
@@ -3409,10 +3415,27 @@ impl SetDisks {
let result = disk
.rename_data_borrowed(&src_bucket, &src_object, file_info, &dst_bucket, &dst_object)
.await;
rustfs_io_metrics::record_put_object_stage_duration_from(
rustfs_io_metrics::PUT_STAGE_SET_DISK_RENAME_DISK_WAIT,
disk_wait_started,
);
if let Some(disk_wait_started) = disk_wait_started {
let duration_ms = disk_wait_started.elapsed().as_secs_f64() * 1000.0;
rustfs_io_metrics::record_put_object_stage_duration(
rustfs_io_metrics::PUT_STAGE_SET_DISK_RENAME_DISK_WAIT,
duration_ms,
);
let position = if result.is_ok() {
let rank = successful_rename_completion_rank
.as_ref()
.map(|rank| rank.fetch_add(1, Ordering::Relaxed) + 1)
.unwrap_or(1);
if rank <= write_quorum {
rustfs_io_metrics::PUT_RENAME_DISK_WAIT_COMPLETION_POSITION_QUORUM_FIRST
} else {
rustfs_io_metrics::PUT_RENAME_DISK_WAIT_COMPLETION_POSITION_QUORUM_TAIL
}
} else {
rustfs_io_metrics::PUT_RENAME_DISK_WAIT_COMPLETION_POSITION_ERROR
};
rustfs_io_metrics::record_put_rename_disk_wait_completion(position, duration_ms);
}
result
})
.catch_unwind()