1use std::sync::{Arc, OnceLock};
16
17use prometheus::core::{AtomicU64, Collector, Desc, GenericCounter};
18use prometheus::{
19 Gauge, Histogram, HistogramVec, IntGauge, IntGaugeVec, Opts, Registry, exponential_buckets,
20 histogram_opts, proto, register_histogram_vec_with_registry, register_histogram_with_registry,
21 register_int_counter_vec_with_registry, register_int_gauge_vec_with_registry,
22 register_int_gauge_with_registry,
23};
24use risingwave_common::config::MetricLevel;
25use risingwave_common::metrics::{
26 RelabeledCounterVec, RelabeledGuardedHistogramVec, RelabeledGuardedIntCounterVec,
27 RelabeledGuardedIntGaugeVec, RelabeledHistogramVec, RelabeledMetricVec, UintGauge,
28};
29use risingwave_common::monitor::GLOBAL_METRICS_REGISTRY;
30use risingwave_common::{
31 register_guarded_histogram_vec_with_registry, register_guarded_int_counter_vec_with_registry,
32 register_guarded_int_gauge_vec_with_registry,
33};
34use thiserror_ext::AsReport;
35use tracing::warn;
36
37#[derive(Debug, Clone)]
43pub struct HummockStateStoreMetrics {
44 pub bloom_filter_true_negative_counts: RelabeledGuardedIntCounterVec,
45 pub bloom_filter_check_counts: RelabeledGuardedIntCounterVec,
46 pub iter_merge_sstable_counts: RelabeledHistogramVec,
47 pub vnode_pruning_counts: RelabeledGuardedIntCounterVec,
48 pub sst_store_block_request_counts: RelabeledGuardedIntCounterVec,
49 pub iter_scan_key_counts: RelabeledGuardedIntCounterVec,
50 pub get_shared_buffer_hit_counts: RelabeledCounterVec,
51 pub remote_read_time: RelabeledHistogramVec,
52 pub iter_fetch_meta_duration: RelabeledGuardedHistogramVec,
53 pub iter_fetch_meta_cache_unhits: IntGauge,
54 pub iter_slow_fetch_meta_cache_unhits: IntGauge,
55
56 pub vector_object_request_counts: RelabeledGuardedIntCounterVec,
57 pub vector_request_stats: RelabeledGuardedHistogramVec,
58 pub vector_hnsw_graph_level_node_count: RelabeledGuardedIntGaugeVec,
59 pub vector_index_file_count: RelabeledGuardedIntGaugeVec,
60 pub vector_index_file_size: RelabeledGuardedIntGaugeVec,
61
62 pub read_req_bloom_filter_positive_counts: RelabeledGuardedIntCounterVec,
63 pub read_req_positive_but_non_exist_counts: RelabeledGuardedIntCounterVec,
64 pub read_req_check_bloom_filter_counts: RelabeledGuardedIntCounterVec,
65
66 pub write_batch_tuple_counts: RelabeledGuardedIntCounterVec,
67 pub write_batch_duration: RelabeledGuardedHistogramVec,
68 pub write_batch_size: RelabeledGuardedHistogramVec,
69
70 pub spill_task_counts_from_unsealed: GenericCounter<AtomicU64>,
72 pub spill_task_size_from_unsealed: GenericCounter<AtomicU64>,
74
75 pub uploader_uploading_task_size: UintGauge,
77 pub uploader_uploading_task_count: IntGauge,
78 pub uploader_imm_size: UintGauge,
79 pub uploader_upload_task_latency: Histogram,
80 pub uploader_syncing_epoch_count: IntGauge,
81 pub uploader_wait_poll_latency: Histogram,
82 pub uploader_per_table_imm_size: RelabeledGuardedIntGaugeVec,
83 pub uploader_per_table_imm_count: RelabeledGuardedIntGaugeVec,
84
85 pub replicated_imm_size: UintGauge,
87 pub per_table_imm_size: RelabeledGuardedIntGaugeVec,
88 pub per_table_imm_count: RelabeledGuardedIntGaugeVec,
89 pub mem_table_spill_counts: RelabeledGuardedIntCounterVec,
90 pub old_value_size: RelabeledGuardedIntGaugeVec,
91
92 pub block_efficiency_histogram: Histogram,
94
95 pub event_handler_pending_event: IntGaugeVec,
96 pub event_handler_latency: HistogramVec,
97
98 pub safe_version_hit: GenericCounter<AtomicU64>,
99 pub safe_version_miss: GenericCounter<AtomicU64>,
100 pub table_change_log_fetch_latency: Histogram,
101 pub table_change_log_cache_hit: GenericCounter<AtomicU64>,
102 pub table_change_log_cache_miss: GenericCounter<AtomicU64>,
103}
104
105pub static GLOBAL_HUMMOCK_STATE_STORE_METRICS: OnceLock<HummockStateStoreMetrics> = OnceLock::new();
106
107pub fn global_hummock_state_store_metrics(metric_level: MetricLevel) -> HummockStateStoreMetrics {
108 GLOBAL_HUMMOCK_STATE_STORE_METRICS
109 .get_or_init(|| HummockStateStoreMetrics::new(&GLOBAL_METRICS_REGISTRY, metric_level))
110 .clone()
111}
112
113impl HummockStateStoreMetrics {
114 pub fn new(registry: &Registry, metric_level: MetricLevel) -> Self {
115 let time_buckets = exponential_buckets(0.01, 10.0, 7).unwrap();
117
118 let state_store_read_time_buckets = exponential_buckets(0.001, 10.0, 5).unwrap();
120
121 let bloom_filter_true_negative_counts = register_guarded_int_counter_vec_with_registry!(
122 "state_store_bloom_filter_true_negative_counts",
123 "Total number of sstables that have been considered true negative by SST filters",
124 &["table_id", "type"],
125 registry
126 )
127 .unwrap();
128 let bloom_filter_true_negative_counts = RelabeledMetricVec::with_metric_level(
129 MetricLevel::Debug,
130 bloom_filter_true_negative_counts,
131 metric_level,
132 );
133
134 let bloom_filter_check_counts = register_guarded_int_counter_vec_with_registry!(
135 "state_store_bloom_filter_check_counts",
136 "Total number of read request to check SST filters",
137 &["table_id", "type"],
138 registry
139 )
140 .unwrap();
141 let bloom_filter_check_counts = RelabeledMetricVec::with_metric_level(
142 MetricLevel::Debug,
143 bloom_filter_check_counts,
144 metric_level,
145 );
146
147 let opts = histogram_opts!(
149 "state_store_iter_merge_sstable_counts",
150 "Number of child iterators merged into one MergeIterator",
151 vec![1.0, 10.0, 100.0, 1000.0, 10000.0]
152 );
153 let iter_merge_sstable_counts =
154 register_histogram_vec_with_registry!(opts, &["table_id", "type"], registry).unwrap();
155 let iter_merge_sstable_counts = RelabeledHistogramVec::with_metric_level(
156 MetricLevel::Debug,
157 iter_merge_sstable_counts,
158 metric_level,
159 );
160
161 let vnode_pruning_counts = register_guarded_int_counter_vec_with_registry!(
162 "state_store_vnode_pruning_counts",
163 "Total number of SST pruning operations by vnode key range hints",
164 &["table_id", "operation", "result"],
165 registry
166 )
167 .unwrap();
168
169 let vnode_pruning_counts = RelabeledGuardedIntCounterVec::with_metric_level(
170 MetricLevel::Debug,
171 vnode_pruning_counts,
172 metric_level,
173 );
174
175 let sst_store_block_request_counts = register_guarded_int_counter_vec_with_registry!(
177 "state_store_sst_store_block_request_counts",
178 "Total number of sst block requests that have been issued to sst store",
179 &["table_id", "type"],
180 registry
181 )
182 .unwrap();
183 let sst_store_block_request_counts = RelabeledGuardedIntCounterVec::with_metric_level(
184 MetricLevel::Info,
185 sst_store_block_request_counts,
186 metric_level,
187 );
188
189 let iter_scan_key_counts = register_guarded_int_counter_vec_with_registry!(
190 "state_store_iter_scan_key_counts",
191 "Total number of keys read by iterator",
192 &["table_id", "type"],
193 registry
194 )
195 .unwrap();
196 let iter_scan_key_counts = RelabeledGuardedIntCounterVec::with_metric_level(
197 MetricLevel::Info,
198 iter_scan_key_counts,
199 metric_level,
200 );
201
202 let get_shared_buffer_hit_counts = register_int_counter_vec_with_registry!(
203 "state_store_get_shared_buffer_hit_counts",
204 "Total number of get requests that have been fulfilled by shared buffer",
205 &["table_id"],
206 registry
207 )
208 .unwrap();
209 let get_shared_buffer_hit_counts = RelabeledCounterVec::with_metric_level(
210 MetricLevel::Debug,
211 get_shared_buffer_hit_counts,
212 metric_level,
213 );
214
215 let opts = histogram_opts!(
216 "state_store_remote_read_time_per_task",
217 "Total time of operations which read from remote storage when enable prefetch",
218 time_buckets.clone(),
219 );
220 let remote_read_time =
221 register_histogram_vec_with_registry!(opts, &["table_id"], registry).unwrap();
222 let remote_read_time = RelabeledHistogramVec::with_metric_level(
223 MetricLevel::Debug,
224 remote_read_time,
225 metric_level,
226 );
227
228 let opts = histogram_opts!(
229 "state_store_iter_fetch_meta_duration",
230 "Histogram of iterator fetch SST meta time that have been issued to state store",
231 state_store_read_time_buckets.clone(),
232 );
233 let iter_fetch_meta_duration =
234 register_guarded_histogram_vec_with_registry!(opts, &["table_id"], registry).unwrap();
235 let iter_fetch_meta_duration = RelabeledGuardedHistogramVec::with_metric_level(
236 MetricLevel::Info,
237 iter_fetch_meta_duration,
238 metric_level,
239 );
240
241 let iter_fetch_meta_cache_unhits = register_int_gauge_with_registry!(
242 "state_store_iter_fetch_meta_cache_unhits",
243 "Number of SST meta cache unhit during one iterator meta fetch",
244 registry
245 )
246 .unwrap();
247
248 let iter_slow_fetch_meta_cache_unhits = register_int_gauge_with_registry!(
249 "state_store_iter_slow_fetch_meta_cache_unhits",
250 "Number of SST meta cache unhit during a iterator meta fetch which is slow (costs >5 seconds)",
251 registry
252 )
253 .unwrap();
254
255 let opts = histogram_opts!(
256 "state_store_table_change_log_fetch_latency",
257 "Latency of fetching table change logs",
258 state_store_read_time_buckets,
259 );
260 let table_change_log_fetch_latency =
261 register_histogram_with_registry!(opts, registry).unwrap();
262
263 let table_change_log_cache_counts = register_int_counter_vec_with_registry!(
264 "state_store_table_change_log_cache_counts",
265 "Total number of table change log cache lookups",
266 &["result"],
267 registry
268 )
269 .unwrap();
270 let table_change_log_cache_hit = table_change_log_cache_counts.with_label_values(&["hit"]);
271 let table_change_log_cache_miss =
272 table_change_log_cache_counts.with_label_values(&["miss"]);
273
274 let vector_object_request_counts = register_guarded_int_counter_vec_with_registry!(
276 "state_store_vector_object_request_counts",
277 "Metrics about vector object requests that have been issued",
278 &["table_id", "type", "mode"],
279 registry
280 )
281 .unwrap();
282 let vector_object_request_counts = RelabeledGuardedIntCounterVec::with_metric_level(
283 MetricLevel::Critical,
284 vector_object_request_counts,
285 metric_level,
286 );
287
288 let opts = histogram_opts!(
289 "state_store_vector_request_stats",
290 "Metrics about vector requests",
291 exponential_buckets(100.0, 10.0, 5).unwrap(),
292 );
293
294 let vector_request_stats = register_guarded_histogram_vec_with_registry!(
295 opts,
296 &["table_id", "type", "mode", "top_n", "ef"],
297 registry
298 )
299 .unwrap();
300 let vector_request_stats = RelabeledGuardedHistogramVec::with_metric_level(
301 MetricLevel::Critical,
302 vector_request_stats,
303 metric_level,
304 );
305
306 let vector_hnsw_graph_level_node_count = register_guarded_int_gauge_vec_with_registry!(
307 "state_store_vector_hnsw_graph_level_node_count",
308 "Number of nodes in each level of hnsw graph",
309 &["table_id", "level"],
310 registry
311 )
312 .unwrap();
313 let vector_hnsw_graph_level_node_count = RelabeledGuardedIntGaugeVec::with_metric_level(
314 MetricLevel::Critical,
315 vector_hnsw_graph_level_node_count,
316 metric_level,
317 );
318
319 let vector_index_file_count = register_guarded_int_gauge_vec_with_registry!(
320 "state_store_vector_index_file_count",
321 "Number of vector file",
322 &["table_id"],
323 registry
324 )
325 .unwrap();
326 let vector_index_file_count = RelabeledGuardedIntGaugeVec::with_metric_level(
327 MetricLevel::Critical,
328 vector_index_file_count,
329 metric_level,
330 );
331
332 let vector_index_file_size = register_guarded_int_gauge_vec_with_registry!(
333 "state_store_vector_index_file_size",
334 "total size of vector index file",
335 &["table_id", "type"],
336 registry
337 )
338 .unwrap();
339 let vector_index_file_size = RelabeledGuardedIntGaugeVec::with_metric_level(
340 MetricLevel::Critical,
341 vector_index_file_size,
342 metric_level,
343 );
344
345 let write_batch_tuple_counts = register_guarded_int_counter_vec_with_registry!(
347 "state_store_write_batch_tuple_counts",
348 "Total number of batched write kv pairs requests that have been issued to state store",
349 &["table_id"],
350 registry
351 )
352 .unwrap();
353 let write_batch_tuple_counts = RelabeledGuardedIntCounterVec::with_metric_level(
354 MetricLevel::Debug,
355 write_batch_tuple_counts,
356 metric_level,
357 );
358
359 let opts = histogram_opts!(
360 "state_store_write_batch_duration",
361 "Total time of batched write that have been issued to state store. With shared buffer on, this is the latency writing to the shared buffer",
362 time_buckets.clone()
363 );
364 let write_batch_duration =
365 register_guarded_histogram_vec_with_registry!(opts, &["table_id"], registry).unwrap();
366 let write_batch_duration = RelabeledGuardedHistogramVec::with_metric_level(
367 MetricLevel::Debug,
368 write_batch_duration,
369 metric_level,
370 );
371
372 let opts = histogram_opts!(
373 "state_store_write_batch_size",
374 "Total size of batched write that have been issued to state store",
375 exponential_buckets(256.0, 16.0, 7).unwrap() );
377 let write_batch_size =
378 register_guarded_histogram_vec_with_registry!(opts, &["table_id"], registry).unwrap();
379 let write_batch_size = RelabeledGuardedHistogramVec::with_metric_level(
380 MetricLevel::Debug,
381 write_batch_size,
382 metric_level,
383 );
384
385 let spill_task_counts = register_int_counter_vec_with_registry!(
386 "state_store_spill_task_counts",
387 "Total number of started spill tasks",
388 &["uploader_stage"],
389 registry
390 )
391 .unwrap();
392
393 let spill_task_size = register_int_counter_vec_with_registry!(
394 "state_store_spill_task_size",
395 "Total task of started spill tasks",
396 &["uploader_stage"],
397 registry
398 )
399 .unwrap();
400
401 let uploader_uploading_task_size = UintGauge::new(
402 "state_store_uploader_uploading_task_size",
403 "Total size of uploader uploading tasks",
404 )
405 .unwrap();
406 registry
407 .register(Box::new(uploader_uploading_task_size.clone()))
408 .unwrap();
409
410 let uploader_uploading_task_count = register_int_gauge_with_registry!(
411 "state_store_uploader_uploading_task_count",
412 "Total number of uploader uploading tasks",
413 registry
414 )
415 .unwrap();
416
417 let uploader_imm_size = UintGauge::new(
418 "state_store_uploader_imm_size",
419 "Total size of imms tracked by uploader",
420 )
421 .unwrap();
422 registry
423 .register(Box::new(uploader_imm_size.clone()))
424 .unwrap();
425
426 let opts = histogram_opts!(
427 "state_store_uploader_upload_task_latency",
428 "Latency of uploader uploading tasks",
429 time_buckets
430 );
431
432 let uploader_upload_task_latency =
433 register_histogram_with_registry!(opts, registry).unwrap();
434
435 let opts = histogram_opts!(
436 "state_store_uploader_wait_poll_latency",
437 "Latency of upload uploading task being polled after finish",
438 exponential_buckets(0.001, 5.0, 7).unwrap(), );
440
441 let uploader_wait_poll_latency = register_histogram_with_registry!(opts, registry).unwrap();
442
443 let uploader_syncing_epoch_count = register_int_gauge_with_registry!(
444 "state_store_uploader_syncing_epoch_count",
445 "Total number of syncing epoch",
446 registry
447 )
448 .unwrap();
449
450 let uploader_per_table_imm_size = register_guarded_int_gauge_vec_with_registry!(
451 "state_store_uploader_per_table_imm_size",
452 "Total uploader-tracked imm size per table",
453 &["table_id"],
454 registry
455 )
456 .unwrap();
457
458 let uploader_per_table_imm_size = RelabeledGuardedIntGaugeVec::with_metric_level(
459 MetricLevel::Debug,
460 uploader_per_table_imm_size,
461 metric_level,
462 );
463
464 let uploader_per_table_imm_count = register_guarded_int_gauge_vec_with_registry!(
465 "state_store_uploader_per_table_imm_count",
466 "Total uploader-tracked imm count per table",
467 &["table_id"],
468 registry
469 )
470 .unwrap();
471
472 let uploader_per_table_imm_count = RelabeledGuardedIntGaugeVec::with_metric_level(
473 MetricLevel::Debug,
474 uploader_per_table_imm_count,
475 metric_level,
476 );
477
478 let replicated_imm_size = UintGauge::new(
479 "state_store_replicated_imm_size",
480 "Total retained replicated immutable memtable size",
481 )
482 .unwrap();
483 registry
484 .register(Box::new(replicated_imm_size.clone()))
485 .unwrap();
486
487 let per_table_imm_size = register_guarded_int_gauge_vec_with_registry!(
488 "state_store_per_table_imm_size",
489 "Total imm size per table",
490 &["table_id", "fragment_id"],
491 registry
492 )
493 .unwrap();
494
495 let per_table_imm_size = RelabeledGuardedIntGaugeVec::with_metric_level_relabel_n(
496 MetricLevel::Debug,
497 per_table_imm_size,
498 metric_level,
499 1,
500 );
501
502 let per_table_imm_count = register_guarded_int_gauge_vec_with_registry!(
503 "state_store_per_table_imm_count",
504 "Total imm count per table",
505 &["table_id"],
506 registry
507 )
508 .unwrap();
509
510 let per_table_imm_count = RelabeledGuardedIntGaugeVec::with_metric_level(
511 MetricLevel::Debug,
512 per_table_imm_count,
513 metric_level,
514 );
515
516 let read_req_bloom_filter_positive_counts =
517 register_guarded_int_counter_vec_with_registry!(
518 "state_store_read_req_bloom_filter_positive_counts",
519 "Total number of read request with at least one SST filter check returns positive",
520 &["table_id", "type"],
521 registry
522 )
523 .unwrap();
524 let read_req_bloom_filter_positive_counts =
525 RelabeledGuardedIntCounterVec::with_metric_level_relabel_n(
526 MetricLevel::Info,
527 read_req_bloom_filter_positive_counts,
528 metric_level,
529 1,
530 );
531
532 let read_req_positive_but_non_exist_counts = register_guarded_int_counter_vec_with_registry!(
533 "state_store_read_req_positive_but_non_exist_counts",
534 "Total number of read request on non-existent key/prefix with at least one SST filter check returns positive",
535 &["table_id", "type"],
536 registry
537 )
538 .unwrap();
539 let read_req_positive_but_non_exist_counts =
540 RelabeledGuardedIntCounterVec::with_metric_level(
541 MetricLevel::Info,
542 read_req_positive_but_non_exist_counts,
543 metric_level,
544 );
545
546 let read_req_check_bloom_filter_counts = register_guarded_int_counter_vec_with_registry!(
547 "state_store_read_req_check_bloom_filter_counts",
548 "Total number of read request that checks an SST filter with a prefix hint",
549 &["table_id", "type"],
550 registry
551 )
552 .unwrap();
553
554 let read_req_check_bloom_filter_counts = RelabeledGuardedIntCounterVec::with_metric_level(
555 MetricLevel::Info,
556 read_req_check_bloom_filter_counts,
557 metric_level,
558 );
559
560 let mem_table_spill_counts = register_guarded_int_counter_vec_with_registry!(
561 "state_store_mem_table_spill_counts",
562 "Total number of mem table spill occurs for one table",
563 &["table_id"],
564 registry
565 )
566 .unwrap();
567
568 let mem_table_spill_counts = RelabeledGuardedIntCounterVec::with_metric_level(
569 MetricLevel::Info,
570 mem_table_spill_counts,
571 metric_level,
572 );
573
574 let old_value_size = register_guarded_int_gauge_vec_with_registry!(
575 "state_store_old_value_size",
576 "The size of old value",
577 &["table_id"],
578 registry
579 )
580 .unwrap();
581
582 let old_value_size = RelabeledGuardedIntGaugeVec::with_metric_level(
583 MetricLevel::Info,
584 old_value_size,
585 metric_level,
586 );
587
588 let opts = histogram_opts!(
589 "block_efficiency_histogram",
590 "Access ratio of in-memory block.",
591 exponential_buckets(0.001, 2.0, 11).unwrap(),
592 );
593 let block_efficiency_histogram = register_histogram_with_registry!(opts, registry).unwrap();
594
595 let event_handler_pending_event = register_int_gauge_vec_with_registry!(
596 "state_store_event_handler_pending_event",
597 "The number of sent but unhandled events",
598 &["event_type"],
599 registry,
600 )
601 .unwrap();
602
603 let opts = histogram_opts!(
604 "state_store_event_handler_latency",
605 "Latency to handle event",
606 exponential_buckets(0.001, 5.0, 7).unwrap(), );
608
609 let event_handler_latency =
610 register_histogram_vec_with_registry!(opts, &["event_type"], registry).unwrap();
611
612 let safe_version_hit = GenericCounter::new(
613 "state_store_safe_version_hit",
614 "The total count of a safe version that can be retrieved successfully",
615 )
616 .unwrap();
617 registry
618 .register(Box::new(safe_version_hit.clone()))
619 .unwrap();
620
621 let safe_version_miss = GenericCounter::new(
622 "state_store_safe_version_miss",
623 "The total count of a safe version that cannot be retrieved",
624 )
625 .unwrap();
626 registry
627 .register(Box::new(safe_version_miss.clone()))
628 .unwrap();
629
630 Self {
631 bloom_filter_true_negative_counts,
632 bloom_filter_check_counts,
633 iter_merge_sstable_counts,
634 vnode_pruning_counts,
635 sst_store_block_request_counts,
636 iter_scan_key_counts,
637 get_shared_buffer_hit_counts,
638 remote_read_time,
639 iter_fetch_meta_duration,
640 iter_fetch_meta_cache_unhits,
641 iter_slow_fetch_meta_cache_unhits,
642 vector_object_request_counts,
643 vector_request_stats,
644 vector_hnsw_graph_level_node_count,
645 vector_index_file_count,
646 vector_index_file_size,
647 read_req_bloom_filter_positive_counts,
648 read_req_positive_but_non_exist_counts,
649 read_req_check_bloom_filter_counts,
650 write_batch_tuple_counts,
651 write_batch_duration,
652 write_batch_size,
653 spill_task_counts_from_unsealed: spill_task_counts.with_label_values(&["unsealed"]),
654 spill_task_size_from_unsealed: spill_task_size.with_label_values(&["unsealed"]),
655 uploader_uploading_task_size,
656 uploader_uploading_task_count,
657 uploader_imm_size,
658 uploader_upload_task_latency,
659 uploader_syncing_epoch_count,
660 uploader_wait_poll_latency,
661 uploader_per_table_imm_size,
662 uploader_per_table_imm_count,
663 replicated_imm_size,
664 per_table_imm_size,
665 per_table_imm_count,
666 mem_table_spill_counts,
667 old_value_size,
668
669 block_efficiency_histogram,
670 event_handler_pending_event,
671 event_handler_latency,
672 safe_version_hit,
673 safe_version_miss,
674 table_change_log_fetch_latency,
675 table_change_log_cache_hit,
676 table_change_log_cache_miss,
677 }
678 }
679
680 pub fn unused() -> Self {
681 global_hummock_state_store_metrics(MetricLevel::Disabled)
682 }
683}
684
685pub trait MemoryCollector: Sync + Send {
686 fn get_meta_memory_usage(&self) -> u64;
687 fn get_data_memory_usage(&self) -> u64;
688 fn get_vector_meta_memory_usage(&self) -> u64;
689 fn get_vector_data_memory_usage(&self) -> u64;
690 fn get_uploading_memory_usage(&self) -> u64;
691 fn get_prefetch_memory_usage(&self) -> usize;
692 fn get_meta_cache_memory_usage_ratio(&self) -> f64;
693 fn get_block_cache_memory_usage_ratio(&self) -> f64;
694 fn get_vector_meta_cache_memory_usage_ratio(&self) -> f64;
695 fn get_vector_data_cache_memory_usage_ratio(&self) -> f64;
696 fn get_shared_buffer_usage_ratio(&self) -> f64;
697}
698
699#[derive(Clone)]
700struct StateStoreCollector {
701 memory_collector: Arc<dyn MemoryCollector>,
702 collectors: Vec<Arc<dyn Collector>>,
703 block_cache_size: IntGauge,
704 meta_cache_size: IntGauge,
705 vector_data_cache_size: IntGauge,
706 vector_meta_cache_size: IntGauge,
707 uploading_memory_size: IntGauge,
708 prefetch_memory_size: IntGauge,
709 meta_cache_usage_ratio: Gauge,
710 block_cache_usage_ratio: Gauge,
711 vector_data_cache_usage_ratio: Gauge,
712 vector_meta_cache_usage_ratio: Gauge,
713 uploading_memory_usage_ratio: Gauge,
714}
715
716impl StateStoreCollector {
717 pub fn new(memory_collector: Arc<dyn MemoryCollector>) -> Self {
718 let mut collectors = Vec::new();
719
720 let block_cache_size = IntGauge::with_opts(Opts::new(
721 "state_store_block_cache_size",
722 "the size of cache for data block cache",
723 ))
724 .unwrap();
725 collectors.push(Arc::new(block_cache_size.clone()) as _);
726
727 let block_cache_usage_ratio = Gauge::with_opts(Opts::new(
728 "state_store_block_cache_usage_ratio",
729 "the ratio of block cache to it's pre-allocated memory",
730 ))
731 .unwrap();
732 collectors.push(Arc::new(block_cache_usage_ratio.clone()) as _);
733
734 let meta_cache_size = IntGauge::with_opts(Opts::new(
735 "state_store_meta_cache_size",
736 "the size of cache for meta file cache",
737 ))
738 .unwrap();
739 collectors.push(Arc::new(meta_cache_size.clone()) as _);
740
741 let meta_cache_usage_ratio = Gauge::with_opts(Opts::new(
742 "state_store_meta_cache_usage_ratio",
743 "the ratio of meta cache to it's pre-allocated memory",
744 ))
745 .unwrap();
746 collectors.push(Arc::new(meta_cache_usage_ratio.clone()) as _);
747
748 let vector_data_cache_size = IntGauge::with_opts(Opts::new(
749 "state_store_vector_data_cache_size",
750 "the size of cache for vector data file cache",
751 ))
752 .unwrap();
753 collectors.push(Arc::new(vector_data_cache_size.clone()) as _);
754
755 let vector_data_cache_usage_ratio = Gauge::with_opts(Opts::new(
756 "state_store_vector_data_cache_usage_ratio",
757 "the ratio of vector data cache to it's pre-allocated memory",
758 ))
759 .unwrap();
760 collectors.push(Arc::new(vector_data_cache_usage_ratio.clone()) as _);
761
762 let vector_meta_cache_size = IntGauge::with_opts(Opts::new(
763 "state_store_vector_meta_cache_size",
764 "the size of cache for vector meta file cache",
765 ))
766 .unwrap();
767 collectors.push(Arc::new(vector_meta_cache_size.clone()) as _);
768
769 let vector_meta_cache_usage_ratio = Gauge::with_opts(Opts::new(
770 "state_store_vector_meta_cache_usage_ratio",
771 "the ratio of vector meta cache to it's pre-allocated memory",
772 ))
773 .unwrap();
774 collectors.push(Arc::new(vector_meta_cache_usage_ratio.clone()) as _);
775
776 let uploading_memory_size = IntGauge::with_opts(Opts::new(
777 "uploading_memory_size",
778 "the size of uploading SSTs memory usage",
779 ))
780 .unwrap();
781 collectors.push(Arc::new(uploading_memory_size.clone()) as _);
782
783 let uploading_memory_usage_ratio = Gauge::with_opts(Opts::new(
784 "state_store_uploading_memory_usage_ratio",
785 "the ratio of uploading SSTs memory usage to it's pre-allocated memory",
786 ))
787 .unwrap();
788 collectors.push(Arc::new(uploading_memory_usage_ratio.clone()) as _);
789
790 let prefetch_memory_size = IntGauge::with_opts(Opts::new(
791 "state_store_prefetch_memory_size",
792 "the size of prefetch memory usage",
793 ))
794 .unwrap();
795 collectors.push(Arc::new(prefetch_memory_size.clone()) as _);
796
797 Self {
798 memory_collector,
799 collectors,
800 block_cache_size,
801 meta_cache_size,
802 vector_data_cache_size,
803 vector_meta_cache_size,
804 uploading_memory_size,
805 prefetch_memory_size,
806 meta_cache_usage_ratio,
807 block_cache_usage_ratio,
808
809 vector_data_cache_usage_ratio,
810 vector_meta_cache_usage_ratio,
811 uploading_memory_usage_ratio,
812 }
813 }
814}
815
816impl Collector for StateStoreCollector {
817 fn desc(&self) -> Vec<&Desc> {
818 self.collectors.iter().flat_map(|c| c.desc()).collect()
819 }
820
821 fn collect(&self) -> Vec<proto::MetricFamily> {
822 self.block_cache_size
823 .set(self.memory_collector.get_data_memory_usage() as i64);
824 self.meta_cache_size
825 .set(self.memory_collector.get_meta_memory_usage() as i64);
826 self.vector_data_cache_size
827 .set(self.memory_collector.get_vector_data_memory_usage() as _);
828 self.vector_meta_cache_size
829 .set(self.memory_collector.get_vector_meta_memory_usage() as _);
830 self.uploading_memory_size
831 .set(self.memory_collector.get_uploading_memory_usage() as i64);
832 self.prefetch_memory_size
833 .set(self.memory_collector.get_prefetch_memory_usage() as i64);
834 self.meta_cache_usage_ratio
835 .set(self.memory_collector.get_meta_cache_memory_usage_ratio());
836 self.block_cache_usage_ratio
837 .set(self.memory_collector.get_block_cache_memory_usage_ratio());
838 self.vector_meta_cache_usage_ratio.set(
839 self.memory_collector
840 .get_vector_meta_cache_memory_usage_ratio(),
841 );
842 self.vector_data_cache_usage_ratio.set(
843 self.memory_collector
844 .get_vector_data_cache_memory_usage_ratio(),
845 );
846 self.uploading_memory_usage_ratio
847 .set(self.memory_collector.get_shared_buffer_usage_ratio());
848 self.collectors.iter().flat_map(|c| c.collect()).collect()
850 }
851}
852
853pub fn monitor_cache(memory_collector: Arc<dyn MemoryCollector>) {
854 let collector = Box::new(StateStoreCollector::new(memory_collector));
855 if let Err(e) = GLOBAL_METRICS_REGISTRY.register(collector) {
856 warn!(
857 "unable to monitor cache. May have been registered if in all-in-one deployment: {}",
858 e.as_report()
859 );
860 }
861}