From 7ff18fa8e7177318018fe1450335c420a1d1e925 Mon Sep 17 00:00:00 2001 From: ysq Date: Mon, 10 Nov 2025 23:03:54 +0800 Subject: [PATCH] feat: add gRPC latency monitoring and warning mechanism --- src/streaming/common/constants.rs | 6 +++++ src/streaming/common/event_processor.rs | 25 ++++++++++++++++--- src/streaming/common/metrics.rs | 33 +++++++++++++++++++++++++ 3 files changed, 61 insertions(+), 3 deletions(-) diff --git a/src/streaming/common/constants.rs b/src/streaming/common/constants.rs index e54309c..6be8e37 100644 --- a/src/streaming/common/constants.rs +++ b/src/streaming/common/constants.rs @@ -10,3 +10,9 @@ pub const DEFAULT_MAX_DECODING_MESSAGE_SIZE: usize = 1024 * 1024 * 10; pub const DEFAULT_METRICS_WINDOW_SECONDS: u64 = 5; pub const DEFAULT_METRICS_PRINT_INTERVAL_SECONDS: u64 = 10; pub const SLOW_PROCESSING_THRESHOLD_US: f64 = 3000.0; + +// gRPC 延迟监控 +// Solana 不存储毫秒,所以我们用500ms来校准以获得更好的近似值 +pub const SOLANA_BLOCK_TIME_ADJUSTMENT_MS: i64 = 500; +// 默认最大延迟阈值(毫秒) +pub const MAX_LATENCY_THRESHOLD_MS: i64 = 1000; diff --git a/src/streaming/common/event_processor.rs b/src/streaming/common/event_processor.rs index e415807..87e986d 100644 --- a/src/streaming/common/event_processor.rs +++ b/src/streaming/common/event_processor.rs @@ -18,12 +18,19 @@ fn create_metrics_callback( callback: Arc, ) -> Arc { Arc::new(move |event: DexEvent| { - let processing_time_us = event.metadata().handle_us as f64; + let metadata = event.metadata(); + let processing_time_us = metadata.handle_us as f64; + let recv_us = metadata.recv_us; + let block_time_ms = metadata.block_time_ms; + callback(event); - MetricsManager::global().update_metrics( + + update_metrics_with_latency( MetricsEventType::Transaction, 1, processing_time_us, + recv_us, + block_time_ms, ); }) } @@ -144,8 +151,20 @@ pub async fn process_shred_transaction( Ok(()) } -/// Update metrics for event processing +/// Update metrics for event processing (with optional latency check) #[inline] fn update_metrics(ty: MetricsEventType, count: u64, time_us: f64) { MetricsManager::global().update_metrics(ty, count, time_us); } + +/// Update metrics with latency check +#[inline] +fn update_metrics_with_latency( + ty: MetricsEventType, + count: u64, + time_us: f64, + recv_us: i64, + block_time_ms: i64, +) { + MetricsManager::global().update_metrics_with_latency(ty, count, time_us, recv_us, block_time_ms); +} diff --git a/src/streaming/common/metrics.rs b/src/streaming/common/metrics.rs index dea8bbf..86e9666 100644 --- a/src/streaming/common/metrics.rs +++ b/src/streaming/common/metrics.rs @@ -369,6 +369,25 @@ impl MetricsManager { } } + /// 检查并警告高延迟 (校准后的 gRPC latency) + /// latency = recv_time - (block_time + 500ms) + #[inline] + pub fn check_and_warn_high_latency(&self, recv_us: i64, block_time_ms: i64) { + let recv_ms = recv_us / 1000; + // 校准延迟: recv_time - (block_time + 500ms) + let adjusted_latency_ms = recv_ms - (block_time_ms + SOLANA_BLOCK_TIME_ADJUSTMENT_MS); + + if adjusted_latency_ms > MAX_LATENCY_THRESHOLD_MS { + log::warn!( + "⚠️ High gRPC latency: {}ms (threshold: {}ms, raw: recv={}ms, block={}ms)", + adjusted_latency_ms, + MAX_LATENCY_THRESHOLD_MS, + recv_ms, + block_time_ms + ); + } + } + /// 获取运行时长 pub fn get_uptime(&self) -> std::time::Duration { std::time::Duration::from_secs_f64(GLOBAL_METRICS.get_uptime_seconds()) @@ -481,6 +500,20 @@ impl MetricsManager { self.log_slow_processing(processing_time_us, events_processed as usize); } + /// 更新指标并检查延迟 + #[inline] + pub fn update_metrics_with_latency( + &self, + event_type: MetricsEventType, + events_processed: u64, + processing_time_us: f64, + recv_us: i64, + block_time_ms: i64, + ) { + self.check_and_warn_high_latency(recv_us, block_time_ms); + self.update_metrics(event_type, events_processed, processing_time_us); + } + /// 增加丢弃事件计数 #[inline] pub fn increment_dropped_events(&self) {