mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-04 03:05:39 +00:00
perf(io-metrics): make internode peer-health read-mostly with an RwLock (#4746)
`cluster_peer_should_bypass` is called before every internode RPC when the offline-bypass feature is enabled (remote disk `get_client`, remote locker), and it took a single process-global `Mutex<HashMap>` even for the overwhelmingly common case of an online or unknown peer, which is read-only. That serialized all concurrent internode client acquisitions on one lock. Switch `CLUSTER_PEER_HEALTH` to a `RwLock`: - the hot check takes a shared read lock and returns immediately for unknown or online peers, so concurrent RPCs no longer serialize; - only an offline peer (rare) drops to the write lock to record a re-probe, with a re-fetch/re-check because the state can flip back online between releasing the read lock and taking the write lock; - the write paths (dial reachable/unreachable) and the read-only `cluster_peer_is_offline` move to `write()`/`read()` accordingly. Behavior is unchanged on every path (online/unknown -> not bypassed; offline -> bypass with one re-probe per interval); `Instant::now()` now runs only on the offline path. Poison recovery is preserved. All internode unit tests pass. Addresses rustfs/backlog#1185 (P2). Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
@@ -15,7 +15,7 @@
|
|||||||
use metrics::{counter, gauge};
|
use metrics::{counter, gauge};
|
||||||
use std::collections::HashMap;
|
use std::collections::HashMap;
|
||||||
use std::sync::{
|
use std::sync::{
|
||||||
Arc, LazyLock, Mutex,
|
Arc, LazyLock, RwLock,
|
||||||
atomic::{AtomicU64, Ordering},
|
atomic::{AtomicU64, Ordering},
|
||||||
};
|
};
|
||||||
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
|
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
|
||||||
@@ -460,7 +460,11 @@ impl Default for PeerHealthState {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
static CLUSTER_PEER_HEALTH: LazyLock<Mutex<HashMap<String, PeerHealthState>>> = LazyLock::new(|| Mutex::new(HashMap::new()));
|
/// Read-mostly: the hot `cluster_peer_should_bypass` check per internode RPC only
|
||||||
|
/// reads the map for the common (unknown/online) peer; a `RwLock` lets those run
|
||||||
|
/// concurrently instead of serializing every internode RPC on one mutex. Writes
|
||||||
|
/// (dial reachable/unreachable, and recording an offline peer's re-probe) are rare.
|
||||||
|
static CLUSTER_PEER_HEALTH: LazyLock<RwLock<HashMap<String, PeerHealthState>>> = LazyLock::new(|| RwLock::new(HashMap::new()));
|
||||||
|
|
||||||
fn publish_offline_gauge(peers: &HashMap<String, PeerHealthState>) {
|
fn publish_offline_gauge(peers: &HashMap<String, PeerHealthState>) {
|
||||||
let offline = peers.values().filter(|peer| !peer.online).count();
|
let offline = peers.values().filter(|peer| !peer.online).count();
|
||||||
@@ -483,7 +487,7 @@ fn normalize_peer_key(addr: &str) -> &str {
|
|||||||
/// counter. Called on a successful dial to `addr`.
|
/// counter. Called on a successful dial to `addr`.
|
||||||
pub fn record_peer_reachable(addr: &str) {
|
pub fn record_peer_reachable(addr: &str) {
|
||||||
// Recover from a poisoned lock so peer-health tracking and the offline gauge never stall permanently.
|
// Recover from a poisoned lock so peer-health tracking and the offline gauge never stall permanently.
|
||||||
let mut peers = CLUSTER_PEER_HEALTH.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
|
let mut peers = CLUSTER_PEER_HEALTH.write().unwrap_or_else(|poisoned| poisoned.into_inner());
|
||||||
let entry = peers.entry(normalize_peer_key(addr).to_string()).or_default();
|
let entry = peers.entry(normalize_peer_key(addr).to_string()).or_default();
|
||||||
entry.online = true;
|
entry.online = true;
|
||||||
entry.consecutive_failures = 0;
|
entry.consecutive_failures = 0;
|
||||||
@@ -495,7 +499,7 @@ pub fn record_peer_reachable(addr: &str) {
|
|||||||
/// `failure_threshold` (>= 1) consecutive failures the peer flips offline.
|
/// `failure_threshold` (>= 1) consecutive failures the peer flips offline.
|
||||||
pub fn record_peer_unreachable(addr: &str, failure_threshold: u32) {
|
pub fn record_peer_unreachable(addr: &str, failure_threshold: u32) {
|
||||||
// Recover from a poisoned lock so peer-health tracking and the offline gauge never stall permanently.
|
// Recover from a poisoned lock so peer-health tracking and the offline gauge never stall permanently.
|
||||||
let mut peers = CLUSTER_PEER_HEALTH.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
|
let mut peers = CLUSTER_PEER_HEALTH.write().unwrap_or_else(|poisoned| poisoned.into_inner());
|
||||||
let entry = peers.entry(normalize_peer_key(addr).to_string()).or_default();
|
let entry = peers.entry(normalize_peer_key(addr).to_string()).or_default();
|
||||||
entry.consecutive_failures = entry.consecutive_failures.saturating_add(1);
|
entry.consecutive_failures = entry.consecutive_failures.saturating_add(1);
|
||||||
if entry.consecutive_failures >= failure_threshold.max(1) {
|
if entry.consecutive_failures >= failure_threshold.max(1) {
|
||||||
@@ -506,7 +510,7 @@ pub fn record_peer_unreachable(addr: &str, failure_threshold: u32) {
|
|||||||
|
|
||||||
/// Whether a cluster peer is currently considered offline (known and marked offline).
|
/// Whether a cluster peer is currently considered offline (known and marked offline).
|
||||||
pub fn cluster_peer_is_offline(addr: &str) -> bool {
|
pub fn cluster_peer_is_offline(addr: &str) -> bool {
|
||||||
let peers = CLUSTER_PEER_HEALTH.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
|
let peers = CLUSTER_PEER_HEALTH.read().unwrap_or_else(|poisoned| poisoned.into_inner());
|
||||||
peers.get(normalize_peer_key(addr)).map(|peer| !peer.online).unwrap_or(false)
|
peers.get(normalize_peer_key(addr)).map(|peer| !peer.online).unwrap_or(false)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -518,8 +522,26 @@ pub fn cluster_peer_is_offline(addr: &str) -> bool {
|
|||||||
/// can recover via a normal dial even if no background monitor is running. Online peers are never
|
/// can recover via a normal dial even if no background monitor is running. Online peers are never
|
||||||
/// bypassed.
|
/// bypassed.
|
||||||
pub fn cluster_peer_should_bypass(addr: &str, reprobe_interval: Duration) -> bool {
|
pub fn cluster_peer_should_bypass(addr: &str, reprobe_interval: Duration) -> bool {
|
||||||
let mut peers = CLUSTER_PEER_HEALTH.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
|
let key = normalize_peer_key(addr);
|
||||||
let Some(entry) = peers.get_mut(normalize_peer_key(addr)) else {
|
|
||||||
|
// Fast path: the overwhelmingly common cases — an unknown peer or an online
|
||||||
|
// one — are read-only, so take a shared read lock and let concurrent internode
|
||||||
|
// RPCs check peer health without serializing on a single lock.
|
||||||
|
{
|
||||||
|
let peers = CLUSTER_PEER_HEALTH.read().unwrap_or_else(|poisoned| poisoned.into_inner());
|
||||||
|
match peers.get(key) {
|
||||||
|
None => return false,
|
||||||
|
Some(entry) if entry.online => return false,
|
||||||
|
Some(_) => {} // offline: fall through to the write path below
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Slow path: the peer is offline, which may require recording a re-probe, so
|
||||||
|
// take the write lock. Re-fetch and re-check because the state can change
|
||||||
|
// between releasing the read lock and acquiring the write lock (e.g. a
|
||||||
|
// successful dial flipped it back online).
|
||||||
|
let mut peers = CLUSTER_PEER_HEALTH.write().unwrap_or_else(|poisoned| poisoned.into_inner());
|
||||||
|
let Some(entry) = peers.get_mut(key) else {
|
||||||
return false;
|
return false;
|
||||||
};
|
};
|
||||||
if entry.online {
|
if entry.online {
|
||||||
@@ -542,7 +564,7 @@ pub fn cluster_peer_should_bypass(addr: &str, reprobe_interval: Duration) -> boo
|
|||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
fn cluster_peer_online(addr: &str) -> Option<bool> {
|
fn cluster_peer_online(addr: &str) -> Option<bool> {
|
||||||
CLUSTER_PEER_HEALTH
|
CLUSTER_PEER_HEALTH
|
||||||
.lock()
|
.read()
|
||||||
.ok()?
|
.ok()?
|
||||||
.get(normalize_peer_key(addr))
|
.get(normalize_peer_key(addr))
|
||||||
.map(|peer| peer.online)
|
.map(|peer| peer.online)
|
||||||
@@ -551,7 +573,7 @@ fn cluster_peer_online(addr: &str) -> Option<bool> {
|
|||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
fn cluster_peer_health_keys() -> Vec<String> {
|
fn cluster_peer_health_keys() -> Vec<String> {
|
||||||
CLUSTER_PEER_HEALTH
|
CLUSTER_PEER_HEALTH
|
||||||
.lock()
|
.read()
|
||||||
.map(|peers| peers.keys().cloned().collect())
|
.map(|peers| peers.keys().cloned().collect())
|
||||||
.unwrap_or_default()
|
.unwrap_or_default()
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user