1use std::sync::OnceLock;
16
17use prometheus::{
18 Histogram, IntCounter, IntGauge, Registry, exponential_buckets, histogram_opts,
19 register_histogram_with_registry, register_int_counter_with_registry,
20 register_int_gauge_with_registry,
21};
22use risingwave_common::catalog::TableId;
23use risingwave_common::config::MetricLevel;
24use risingwave_common::metrics::{
25 LabelGuardedHistogramVec, LabelGuardedIntCounter, LabelGuardedIntCounterVec,
26 LabelGuardedIntGauge, LabelGuardedIntGaugeVec, MetricVecRelabelExt,
27 RelabeledGuardedHistogramVec, RelabeledGuardedIntCounterVec, RelabeledGuardedIntGaugeVec,
28};
29use risingwave_common::monitor::GLOBAL_METRICS_REGISTRY;
30use risingwave_common::monitor::in_mem::CountMap;
31use risingwave_common::{
32 register_guarded_histogram_vec_with_registry, register_guarded_int_counter_vec_with_registry,
33 register_guarded_int_gauge_vec_with_registry,
34};
35use risingwave_connector::sink::catalog::SinkId;
36use risingwave_pb::id::ExecutorId;
37
38use crate::common::log_store_impl::kv_log_store::{
39 REWIND_BACKOFF_MULTIPLIER, REWIND_INITIAL_DELAY, REWIND_MAX_DELAY,
40};
41use crate::executor::prelude::ActorId;
42use crate::task::FragmentId;
43
44#[derive(Clone)]
45pub struct StreamingMetrics {
46 pub level: MetricLevel,
47
48 pub executor_row_count: RelabeledGuardedIntCounterVec,
50
51 pub mem_stream_node_output_row_count: CountMap<ExecutorId>,
55 pub mem_stream_node_output_blocking_duration_ns: CountMap<ExecutorId>,
56
57 actor_scheduled_duration: RelabeledGuardedIntCounterVec,
59 actor_scheduled_cnt: RelabeledGuardedIntCounterVec,
60 actor_poll_duration: RelabeledGuardedIntCounterVec,
61 actor_poll_cnt: RelabeledGuardedIntCounterVec,
62 actor_idle_duration: RelabeledGuardedIntCounterVec,
63 actor_idle_cnt: RelabeledGuardedIntCounterVec,
64
65 pub actor_count: LabelGuardedIntGaugeVec,
67 pub actor_in_record_cnt: RelabeledGuardedIntCounterVec,
68 pub actor_out_record_cnt: RelabeledGuardedIntCounterVec,
69 pub fragment_channel_buffered_bytes: LabelGuardedIntGaugeVec,
70 pub actor_current_epoch: RelabeledGuardedIntGaugeVec,
71 pub project_expr_inflight_window_size: LabelGuardedIntGaugeVec,
72
73 pub source_output_row_count: LabelGuardedIntCounterVec,
75 pub source_split_change_count: LabelGuardedIntCounterVec,
76 pub source_backfill_row_count: LabelGuardedIntCounterVec,
77
78 sink_input_row_count: LabelGuardedIntCounterVec,
80 sink_input_bytes: LabelGuardedIntCounterVec,
81 sink_chunk_buffer_size: LabelGuardedIntGaugeVec,
82
83 pub exchange_frag_recv_size: LabelGuardedIntCounterVec,
85
86 pub merge_barrier_align_duration: RelabeledGuardedIntCounterVec,
89
90 pub actor_output_buffer_blocking_duration_ns: RelabeledGuardedIntCounterVec,
92 actor_input_buffer_blocking_duration_ns: RelabeledGuardedIntCounterVec,
93
94 pub join_lookup_miss_count: LabelGuardedIntCounterVec,
96 pub join_lookup_total_count: LabelGuardedIntCounterVec,
97 pub join_insert_cache_miss_count: LabelGuardedIntCounterVec,
98 pub join_actor_input_waiting_duration_ns: LabelGuardedIntCounterVec,
99 pub join_match_duration_ns: LabelGuardedIntCounterVec,
100 pub join_cached_entry_count: LabelGuardedIntGaugeVec,
101 pub join_matched_join_keys: RelabeledGuardedHistogramVec,
102
103 pub barrier_align_duration: RelabeledGuardedIntCounterVec,
105
106 agg_lookup_miss_count: LabelGuardedIntCounterVec,
108 agg_total_lookup_count: LabelGuardedIntCounterVec,
109 agg_cached_entry_count: LabelGuardedIntGaugeVec,
110 agg_chunk_lookup_miss_count: LabelGuardedIntCounterVec,
111 agg_chunk_total_lookup_count: LabelGuardedIntCounterVec,
112 agg_dirty_groups_count: LabelGuardedIntGaugeVec,
113 agg_dirty_groups_heap_size: LabelGuardedIntGaugeVec,
114 agg_distinct_cache_miss_count: LabelGuardedIntCounterVec,
115 agg_distinct_total_cache_count: LabelGuardedIntCounterVec,
116 agg_distinct_cached_entry_count: LabelGuardedIntGaugeVec,
117 agg_state_cache_lookup_count: LabelGuardedIntCounterVec,
118 agg_state_cache_miss_count: LabelGuardedIntCounterVec,
119
120 group_top_n_cache_miss_count: LabelGuardedIntCounterVec,
122 group_top_n_total_query_cache_count: LabelGuardedIntCounterVec,
123 group_top_n_cached_entry_count: LabelGuardedIntGaugeVec,
124 group_top_n_appendonly_cache_miss_count: LabelGuardedIntCounterVec,
126 group_top_n_appendonly_total_query_cache_count: LabelGuardedIntCounterVec,
127 group_top_n_appendonly_cached_entry_count: LabelGuardedIntGaugeVec,
128
129 lookup_cache_miss_count: LabelGuardedIntCounterVec,
131 lookup_total_query_cache_count: LabelGuardedIntCounterVec,
132 lookup_cached_entry_count: LabelGuardedIntGaugeVec,
133
134 temporal_join_cache_miss_count: LabelGuardedIntCounterVec,
136 temporal_join_total_query_cache_count: LabelGuardedIntCounterVec,
137 temporal_join_cached_entry_count: LabelGuardedIntGaugeVec,
138
139 backfill_snapshot_read_row_count: LabelGuardedIntCounterVec,
141 backfill_upstream_output_row_count: LabelGuardedIntCounterVec,
142
143 cdc_backfill_snapshot_read_row_count: LabelGuardedIntCounterVec,
145 cdc_backfill_upstream_output_row_count: LabelGuardedIntCounterVec,
146
147 pub(crate) snapshot_backfill_consume_row_count: LabelGuardedIntCounterVec,
149
150 over_window_cached_entry_count: LabelGuardedIntGaugeVec,
152 over_window_cache_lookup_count: LabelGuardedIntCounterVec,
153 over_window_cache_miss_count: LabelGuardedIntCounterVec,
154 over_window_range_cache_entry_count: LabelGuardedIntGaugeVec,
155 over_window_range_cache_lookup_count: LabelGuardedIntCounterVec,
156 over_window_range_cache_left_miss_count: LabelGuardedIntCounterVec,
157 over_window_range_cache_right_miss_count: LabelGuardedIntCounterVec,
158 over_window_accessed_entry_count: LabelGuardedIntCounterVec,
159 over_window_compute_count: LabelGuardedIntCounterVec,
160 over_window_same_output_count: LabelGuardedIntCounterVec,
161 over_window_state_cleaned_row_count: LabelGuardedIntCounterVec,
162
163 match_recognize_matches_emitted_count: LabelGuardedIntCounterVec,
165 match_recognize_evicted_rows_count: LabelGuardedIntCounterVec,
166 match_recognize_scan_budget_exhausted_count: LabelGuardedIntCounterVec,
167 match_recognize_within_deadline_overflow_count: LabelGuardedIntCounterVec,
168 match_recognize_retained_rows: LabelGuardedIntGaugeVec,
169
170 pub barrier_inflight_latency: LabelGuardedHistogramVec,
174 pub barrier_sync_latency: LabelGuardedHistogramVec,
176 pub barrier_batch_size: Histogram,
177 pub barrier_manager_progress: LabelGuardedIntCounterVec,
179
180 pub kv_log_store_storage_write_count: LabelGuardedIntCounterVec,
181 pub kv_log_store_storage_write_size: LabelGuardedIntCounterVec,
182 pub kv_log_store_rewind_count: LabelGuardedIntCounterVec,
183 pub kv_log_store_rewind_delay: LabelGuardedHistogramVec,
184 pub kv_log_store_storage_read_count: LabelGuardedIntCounterVec,
185 pub kv_log_store_storage_read_size: LabelGuardedIntCounterVec,
186 pub kv_log_store_buffer_unconsumed_item_count: LabelGuardedIntGaugeVec,
187 pub kv_log_store_buffer_unconsumed_row_count: LabelGuardedIntGaugeVec,
188 pub kv_log_store_buffer_unconsumed_epoch_count: LabelGuardedIntGaugeVec,
189 pub kv_log_store_buffer_unconsumed_min_epoch: LabelGuardedIntGaugeVec,
190 pub kv_log_store_buffer_memory_bytes: LabelGuardedIntGaugeVec,
191
192 pub crossdb_last_consumed_min_epoch: LabelGuardedIntGaugeVec,
193
194 pub sync_kv_log_store_read_count: LabelGuardedIntCounterVec,
195 pub sync_kv_log_store_read_size: LabelGuardedIntCounterVec,
196 pub sync_kv_log_store_write_pause_duration_ns: LabelGuardedIntCounterVec,
197 pub sync_kv_log_store_state: LabelGuardedIntCounterVec,
198 pub sync_kv_log_store_wait_next_poll_ns: LabelGuardedIntCounterVec,
199 pub sync_kv_log_store_storage_write_count: LabelGuardedIntCounterVec,
200 pub sync_kv_log_store_storage_write_size: LabelGuardedIntCounterVec,
201 pub sync_kv_log_store_buffer_unconsumed_item_count: LabelGuardedIntGaugeVec,
202 pub sync_kv_log_store_buffer_unconsumed_row_count: LabelGuardedIntGaugeVec,
203 pub sync_kv_log_store_buffer_unconsumed_epoch_count: LabelGuardedIntGaugeVec,
204 pub sync_kv_log_store_buffer_unconsumed_min_epoch: LabelGuardedIntGaugeVec,
205 pub sync_kv_log_store_buffer_memory_bytes: LabelGuardedIntGaugeVec,
206
207 pub lru_runtime_loop_count: IntCounter,
209 pub lru_latest_sequence: IntGauge,
210 pub lru_watermark_sequence: IntGauge,
211 pub lru_eviction_policy: IntGauge,
212 pub jemalloc_allocated_bytes: IntGauge,
213 pub jemalloc_active_bytes: IntGauge,
214 pub jemalloc_resident_bytes: IntGauge,
215 pub jemalloc_metadata_bytes: IntGauge,
216 pub jvm_allocated_bytes: IntGauge,
217 pub jvm_active_bytes: IntGauge,
218 pub stream_memory_usage: RelabeledGuardedIntGaugeVec,
219
220 materialize_cache_hit_count: RelabeledGuardedIntCounterVec,
222 materialize_data_exist_count: RelabeledGuardedIntCounterVec,
223 materialize_cache_total_count: RelabeledGuardedIntCounterVec,
224 materialize_input_row_count: RelabeledGuardedIntCounterVec,
225 pub materialize_current_epoch: RelabeledGuardedIntGaugeVec,
226
227 pub pg_cdc_state_table_lsn: LabelGuardedIntGaugeVec,
229 pub pg_cdc_jni_commit_offset_lsn: LabelGuardedIntGaugeVec,
230
231 pub mysql_cdc_state_binlog_file_seq: LabelGuardedIntGaugeVec,
233 pub mysql_cdc_state_binlog_position: LabelGuardedIntGaugeVec,
234
235 pub sqlserver_cdc_state_change_lsn: LabelGuardedIntGaugeVec,
237 pub sqlserver_cdc_state_commit_lsn: LabelGuardedIntGaugeVec,
238 pub sqlserver_cdc_jni_commit_offset_lsn: LabelGuardedIntGaugeVec,
239
240 pub gap_fill_generated_rows_count: RelabeledGuardedIntCounterVec,
242
243 pub now_streaming_clock_ms: LabelGuardedIntGaugeVec,
245 pub now_wall_clock_drift_ms: LabelGuardedIntGaugeVec,
246
247 pub state_table_iter_count: RelabeledGuardedIntCounterVec,
249 pub state_table_get_count: RelabeledGuardedIntCounterVec,
250 pub state_table_iter_vnode_pruned_count: RelabeledGuardedIntCounterVec,
251 pub state_table_get_vnode_pruned_count: RelabeledGuardedIntCounterVec,
252}
253
254pub static GLOBAL_STREAMING_METRICS: OnceLock<StreamingMetrics> = OnceLock::new();
255
256fn latency_buckets(max: f64, count: usize) -> Vec<f64> {
257 const MIN: f64 = 0.1;
258
259 assert!(count > 1);
260 let factor = (max / MIN).powf(1.0 / (count - 1) as f64);
261 let mut buckets = exponential_buckets(MIN, factor, count).unwrap();
262 *buckets.last_mut().unwrap() = max;
263 buckets
264}
265
266pub fn global_streaming_metrics(metric_level: MetricLevel) -> StreamingMetrics {
267 GLOBAL_STREAMING_METRICS
268 .get_or_init(|| StreamingMetrics::new(&GLOBAL_METRICS_REGISTRY, metric_level))
269 .clone()
270}
271
272impl StreamingMetrics {
273 pub fn new(registry: &Registry, level: MetricLevel) -> Self {
274 let executor_row_count = register_guarded_int_counter_vec_with_registry!(
275 "stream_executor_row_count",
276 "Total number of rows that have been output from each executor",
277 &["actor_id", "fragment_id", "executor_identity"],
278 registry
279 )
280 .unwrap()
281 .relabel_debug_1(level);
282
283 let stream_node_output_row_count = CountMap::new();
284 let stream_node_output_blocking_duration_ns = CountMap::new();
285
286 let source_output_row_count = register_guarded_int_counter_vec_with_registry!(
287 "stream_source_output_rows_counts",
288 "Total number of rows that have been output from source",
289 &["source_id", "source_name", "actor_id", "fragment_id"],
290 registry
291 )
292 .unwrap();
293
294 let source_split_change_count = register_guarded_int_counter_vec_with_registry!(
295 "stream_source_split_change_event_count",
296 "Total number of split change events that have been operated by source",
297 &["source_id", "source_name", "actor_id", "fragment_id"],
298 registry
299 )
300 .unwrap();
301
302 let source_backfill_row_count = register_guarded_int_counter_vec_with_registry!(
303 "stream_source_backfill_rows_counts",
304 "Total number of rows that have been backfilled for source",
305 &["source_id", "source_name", "actor_id", "fragment_id"],
306 registry
307 )
308 .unwrap();
309
310 let sink_input_row_count = register_guarded_int_counter_vec_with_registry!(
311 "stream_sink_input_row_count",
312 "Total number of rows streamed into sink executors",
313 &["sink_id", "actor_id", "fragment_id"],
314 registry
315 )
316 .unwrap();
317
318 let sink_input_bytes = register_guarded_int_counter_vec_with_registry!(
319 "stream_sink_input_bytes",
320 "Total size of chunks streamed into sink executors",
321 &["sink_id", "actor_id", "fragment_id"],
322 registry
323 )
324 .unwrap();
325
326 let materialize_input_row_count = register_guarded_int_counter_vec_with_registry!(
327 "stream_mview_input_row_count",
328 "Total number of rows streamed into materialize executors",
329 &["actor_id", "table_id", "fragment_id"],
330 registry
331 )
332 .unwrap()
333 .relabel_debug_1(level);
334
335 let materialize_current_epoch = register_guarded_int_gauge_vec_with_registry!(
336 "stream_mview_current_epoch",
337 "The current epoch of the materialized executor",
338 &["actor_id", "table_id", "fragment_id"],
339 registry
340 )
341 .unwrap()
342 .relabel_debug_1(level);
343
344 let pg_cdc_state_table_lsn = register_guarded_int_gauge_vec_with_registry!(
345 "stream_pg_cdc_state_table_lsn",
346 "Current LSN value stored in PostgreSQL CDC state table",
347 &["source_id"],
348 registry,
349 )
350 .unwrap();
351
352 let pg_cdc_jni_commit_offset_lsn = register_guarded_int_gauge_vec_with_registry!(
353 "stream_pg_cdc_jni_commit_offset_lsn",
354 "LSN value when JNI commit offset is called for PostgreSQL CDC",
355 &["source_id"],
356 registry,
357 )
358 .unwrap();
359
360 let mysql_cdc_state_binlog_file_seq = register_guarded_int_gauge_vec_with_registry!(
361 "stream_mysql_cdc_state_binlog_file_seq",
362 "Current binlog file sequence number stored in MySQL CDC state table",
363 &["source_id"],
364 registry,
365 )
366 .unwrap();
367
368 let mysql_cdc_state_binlog_position = register_guarded_int_gauge_vec_with_registry!(
369 "stream_mysql_cdc_state_binlog_position",
370 "Current binlog position stored in MySQL CDC state table",
371 &["source_id"],
372 registry,
373 )
374 .unwrap();
375
376 let sqlserver_cdc_state_change_lsn = register_guarded_int_gauge_vec_with_registry!(
377 "stream_sqlserver_cdc_state_change_lsn",
378 "Current change_lsn value stored in SQL Server CDC state table",
379 &["source_id"],
380 registry,
381 )
382 .unwrap();
383
384 let sqlserver_cdc_state_commit_lsn = register_guarded_int_gauge_vec_with_registry!(
385 "stream_sqlserver_cdc_state_commit_lsn",
386 "Current commit_lsn value stored in SQL Server CDC state table",
387 &["source_id"],
388 registry,
389 )
390 .unwrap();
391
392 let sqlserver_cdc_jni_commit_offset_lsn = register_guarded_int_gauge_vec_with_registry!(
393 "stream_sqlserver_cdc_jni_commit_offset_lsn",
394 "LSN value when JNI commit offset is called for SQL Server CDC",
395 &["source_id"],
396 registry,
397 )
398 .unwrap();
399
400 let sink_chunk_buffer_size = register_guarded_int_gauge_vec_with_registry!(
401 "stream_sink_chunk_buffer_size",
402 "Total size of chunks buffered in a barrier",
403 &["sink_id", "actor_id", "fragment_id"],
404 registry
405 )
406 .unwrap();
407 let actor_output_buffer_blocking_duration_ns =
408 register_guarded_int_counter_vec_with_registry!(
409 "stream_actor_output_buffer_blocking_duration_ns",
410 "Total blocking duration (ns) of output buffer",
411 &["actor_id", "fragment_id", "downstream_fragment_id"],
412 registry
413 )
414 .unwrap()
415 .relabel_debug_1(level);
417
418 let actor_input_buffer_blocking_duration_ns =
419 register_guarded_int_counter_vec_with_registry!(
420 "stream_actor_input_buffer_blocking_duration_ns",
421 "Total blocking duration (ns) of input buffer",
422 &["actor_id", "fragment_id", "upstream_fragment_id"],
423 registry
424 )
425 .unwrap()
426 .relabel_debug_1(level);
428
429 let fragment_channel_buffered_bytes = register_guarded_int_gauge_vec_with_registry!(
430 "stream_fragment_channel_buffered_bytes",
431 "Estimated buffered bytes for actor channels by fragment",
432 &["fragment_id"],
433 registry
434 )
435 .unwrap();
436
437 let exchange_frag_recv_size = register_guarded_int_counter_vec_with_registry!(
438 "stream_exchange_frag_recv_size",
439 "Total size of messages that have been received from upstream Fragment",
440 &["up_fragment_id", "down_fragment_id"],
441 registry
442 )
443 .unwrap();
444
445 let actor_poll_duration = register_guarded_int_counter_vec_with_registry!(
446 "stream_actor_poll_duration",
447 "tokio's metrics",
448 &["actor_id", "fragment_id"],
449 registry
450 )
451 .unwrap()
452 .relabel_debug_1(level);
453
454 let actor_poll_cnt = register_guarded_int_counter_vec_with_registry!(
455 "stream_actor_poll_cnt",
456 "tokio's metrics",
457 &["actor_id", "fragment_id"],
458 registry
459 )
460 .unwrap()
461 .relabel_debug_1(level);
462
463 let actor_scheduled_duration = register_guarded_int_counter_vec_with_registry!(
464 "stream_actor_scheduled_duration",
465 "tokio's metrics",
466 &["actor_id", "fragment_id"],
467 registry
468 )
469 .unwrap()
470 .relabel_debug_1(level);
471
472 let actor_scheduled_cnt = register_guarded_int_counter_vec_with_registry!(
473 "stream_actor_scheduled_cnt",
474 "tokio's metrics",
475 &["actor_id", "fragment_id"],
476 registry
477 )
478 .unwrap()
479 .relabel_debug_1(level);
480
481 let actor_idle_duration = register_guarded_int_counter_vec_with_registry!(
482 "stream_actor_idle_duration",
483 "tokio's metrics",
484 &["actor_id", "fragment_id"],
485 registry
486 )
487 .unwrap()
488 .relabel_debug_1(level);
489
490 let actor_idle_cnt = register_guarded_int_counter_vec_with_registry!(
491 "stream_actor_idle_cnt",
492 "tokio's metrics",
493 &["actor_id", "fragment_id"],
494 registry
495 )
496 .unwrap()
497 .relabel_debug_1(level);
498
499 let actor_in_record_cnt = register_guarded_int_counter_vec_with_registry!(
500 "stream_actor_in_record_cnt",
501 "Total number of rows actor received",
502 &["actor_id", "fragment_id", "upstream_fragment_id"],
503 registry
504 )
505 .unwrap()
506 .relabel_debug_1(level);
507
508 let actor_out_record_cnt = register_guarded_int_counter_vec_with_registry!(
509 "stream_actor_out_record_cnt",
510 "Total number of rows actor sent",
511 &["actor_id", "fragment_id"],
512 registry
513 )
514 .unwrap()
515 .relabel_debug_1(level);
516
517 let actor_current_epoch = register_guarded_int_gauge_vec_with_registry!(
518 "stream_actor_current_epoch",
519 "Current epoch of actor",
520 &["actor_id", "fragment_id"],
521 registry
522 )
523 .unwrap()
524 .relabel_debug_1(level);
525
526 let project_expr_inflight_window_size = register_guarded_int_gauge_vec_with_registry!(
527 "stream_project_expr_inflight_window_size",
528 "Number of messages waiting in ProjectExecutor's ordered projection window",
529 &["actor_id", "fragment_id"],
530 registry
531 )
532 .unwrap();
533
534 let actor_count = register_guarded_int_gauge_vec_with_registry!(
535 "stream_actor_count",
536 "Total number of actors (parallelism)",
537 &["fragment_id"],
538 registry
539 )
540 .unwrap();
541
542 let merge_barrier_align_duration = register_guarded_int_counter_vec_with_registry!(
543 "stream_merge_barrier_align_duration_ns",
544 "Total merge barrier alignment duration (ns)",
545 &["actor_id", "fragment_id"],
546 registry
547 )
548 .unwrap()
549 .relabel_debug_1(level);
550
551 let join_lookup_miss_count = register_guarded_int_counter_vec_with_registry!(
552 "stream_join_lookup_miss_count",
553 "Join executor lookup miss duration",
554 &["side", "join_table_id", "actor_id", "fragment_id"],
555 registry
556 )
557 .unwrap();
558
559 let join_lookup_total_count = register_guarded_int_counter_vec_with_registry!(
560 "stream_join_lookup_total_count",
561 "Join executor lookup total operation",
562 &["side", "join_table_id", "actor_id", "fragment_id"],
563 registry
564 )
565 .unwrap();
566
567 let join_insert_cache_miss_count = register_guarded_int_counter_vec_with_registry!(
568 "stream_join_insert_cache_miss_count",
569 "Join executor cache miss when insert operation",
570 &["side", "join_table_id", "actor_id", "fragment_id"],
571 registry
572 )
573 .unwrap();
574
575 let join_actor_input_waiting_duration_ns = register_guarded_int_counter_vec_with_registry!(
576 "stream_join_actor_input_waiting_duration_ns",
577 "Total waiting duration (ns) of input buffer of join actor",
578 &["actor_id", "fragment_id"],
579 registry
580 )
581 .unwrap();
582
583 let join_match_duration_ns = register_guarded_int_counter_vec_with_registry!(
584 "stream_join_match_duration_ns",
585 "Matching duration for each side",
586 &["actor_id", "fragment_id", "side"],
587 registry
588 )
589 .unwrap();
590
591 let barrier_align_duration = register_guarded_int_counter_vec_with_registry!(
592 "stream_barrier_align_duration_ns",
593 "Duration of join align barrier",
594 &["actor_id", "fragment_id", "wait_side", "executor"],
595 registry
596 )
597 .unwrap()
598 .relabel_debug_1(level);
599
600 let join_cached_entry_count = register_guarded_int_gauge_vec_with_registry!(
601 "stream_join_cached_entry_count",
602 "Number of cached entries in streaming join operators",
603 &["actor_id", "fragment_id", "side"],
604 registry
605 )
606 .unwrap();
607
608 let join_matched_join_keys_opts = histogram_opts!(
609 "stream_join_matched_join_keys",
610 "The number of keys matched in the opposite side",
611 exponential_buckets(16.0, 2.0, 28).unwrap() );
613
614 let join_matched_join_keys = register_guarded_histogram_vec_with_registry!(
615 join_matched_join_keys_opts,
616 &["actor_id", "fragment_id", "table_id"],
617 registry
618 )
619 .unwrap()
620 .relabel_debug_1(level);
621
622 let agg_lookup_miss_count = register_guarded_int_counter_vec_with_registry!(
623 "stream_agg_lookup_miss_count",
624 "Aggregation executor lookup miss duration",
625 &["table_id", "actor_id", "fragment_id"],
626 registry
627 )
628 .unwrap();
629
630 let agg_total_lookup_count = register_guarded_int_counter_vec_with_registry!(
631 "stream_agg_lookup_total_count",
632 "Aggregation executor lookup total operation",
633 &["table_id", "actor_id", "fragment_id"],
634 registry
635 )
636 .unwrap();
637
638 let agg_distinct_cache_miss_count = register_guarded_int_counter_vec_with_registry!(
639 "stream_agg_distinct_cache_miss_count",
640 "Aggregation executor dinsinct miss duration",
641 &["table_id", "actor_id", "fragment_id"],
642 registry
643 )
644 .unwrap();
645
646 let agg_distinct_total_cache_count = register_guarded_int_counter_vec_with_registry!(
647 "stream_agg_distinct_total_cache_count",
648 "Aggregation executor distinct total operation",
649 &["table_id", "actor_id", "fragment_id"],
650 registry
651 )
652 .unwrap();
653
654 let agg_distinct_cached_entry_count = register_guarded_int_gauge_vec_with_registry!(
655 "stream_agg_distinct_cached_entry_count",
656 "Total entry counts in distinct aggregation executor cache",
657 &["table_id", "actor_id", "fragment_id"],
658 registry
659 )
660 .unwrap();
661
662 let agg_dirty_groups_count = register_guarded_int_gauge_vec_with_registry!(
663 "stream_agg_dirty_groups_count",
664 "Total dirty group counts in aggregation executor",
665 &["table_id", "actor_id", "fragment_id"],
666 registry
667 )
668 .unwrap();
669
670 let agg_dirty_groups_heap_size = register_guarded_int_gauge_vec_with_registry!(
671 "stream_agg_dirty_groups_heap_size",
672 "Total dirty group heap size in aggregation executor",
673 &["table_id", "actor_id", "fragment_id"],
674 registry
675 )
676 .unwrap();
677
678 let agg_state_cache_lookup_count = register_guarded_int_counter_vec_with_registry!(
679 "stream_agg_state_cache_lookup_count",
680 "Aggregation executor state cache lookup count",
681 &["table_id", "actor_id", "fragment_id"],
682 registry
683 )
684 .unwrap();
685
686 let agg_state_cache_miss_count = register_guarded_int_counter_vec_with_registry!(
687 "stream_agg_state_cache_miss_count",
688 "Aggregation executor state cache miss count",
689 &["table_id", "actor_id", "fragment_id"],
690 registry
691 )
692 .unwrap();
693
694 let group_top_n_cache_miss_count = register_guarded_int_counter_vec_with_registry!(
695 "stream_group_top_n_cache_miss_count",
696 "Group top n executor cache miss count",
697 &["table_id", "actor_id", "fragment_id"],
698 registry
699 )
700 .unwrap();
701
702 let group_top_n_total_query_cache_count = register_guarded_int_counter_vec_with_registry!(
703 "stream_group_top_n_total_query_cache_count",
704 "Group top n executor query cache total count",
705 &["table_id", "actor_id", "fragment_id"],
706 registry
707 )
708 .unwrap();
709
710 let group_top_n_cached_entry_count = register_guarded_int_gauge_vec_with_registry!(
711 "stream_group_top_n_cached_entry_count",
712 "Total entry counts in group top n executor cache",
713 &["table_id", "actor_id", "fragment_id"],
714 registry
715 )
716 .unwrap();
717
718 let group_top_n_appendonly_cache_miss_count =
719 register_guarded_int_counter_vec_with_registry!(
720 "stream_group_top_n_appendonly_cache_miss_count",
721 "Group top n appendonly executor cache miss count",
722 &["table_id", "actor_id", "fragment_id"],
723 registry
724 )
725 .unwrap();
726
727 let group_top_n_appendonly_total_query_cache_count =
728 register_guarded_int_counter_vec_with_registry!(
729 "stream_group_top_n_appendonly_total_query_cache_count",
730 "Group top n appendonly executor total cache count",
731 &["table_id", "actor_id", "fragment_id"],
732 registry
733 )
734 .unwrap();
735
736 let group_top_n_appendonly_cached_entry_count =
737 register_guarded_int_gauge_vec_with_registry!(
738 "stream_group_top_n_appendonly_cached_entry_count",
739 "Total entry counts in group top n appendonly executor cache",
740 &["table_id", "actor_id", "fragment_id"],
741 registry
742 )
743 .unwrap();
744
745 let lookup_cache_miss_count = register_guarded_int_counter_vec_with_registry!(
746 "stream_lookup_cache_miss_count",
747 "Lookup executor cache miss count",
748 &["table_id", "actor_id", "fragment_id"],
749 registry
750 )
751 .unwrap();
752
753 let lookup_total_query_cache_count = register_guarded_int_counter_vec_with_registry!(
754 "stream_lookup_total_query_cache_count",
755 "Lookup executor query cache total count",
756 &["table_id", "actor_id", "fragment_id"],
757 registry
758 )
759 .unwrap();
760
761 let lookup_cached_entry_count = register_guarded_int_gauge_vec_with_registry!(
762 "stream_lookup_cached_entry_count",
763 "Total entry counts in lookup executor cache",
764 &["table_id", "actor_id", "fragment_id"],
765 registry
766 )
767 .unwrap();
768
769 let temporal_join_cache_miss_count = register_guarded_int_counter_vec_with_registry!(
770 "stream_temporal_join_cache_miss_count",
771 "Temporal join executor cache miss count",
772 &["table_id", "actor_id", "fragment_id"],
773 registry
774 )
775 .unwrap();
776
777 let temporal_join_total_query_cache_count =
778 register_guarded_int_counter_vec_with_registry!(
779 "stream_temporal_join_total_query_cache_count",
780 "Temporal join executor query cache total count",
781 &["table_id", "actor_id", "fragment_id"],
782 registry
783 )
784 .unwrap();
785
786 let temporal_join_cached_entry_count = register_guarded_int_gauge_vec_with_registry!(
787 "stream_temporal_join_cached_entry_count",
788 "Total entry count in temporal join executor cache",
789 &["table_id", "actor_id", "fragment_id"],
790 registry
791 )
792 .unwrap();
793
794 let agg_cached_entry_count = register_guarded_int_gauge_vec_with_registry!(
795 "stream_agg_cached_entry_count",
796 "Number of cached keys in streaming aggregation operators",
797 &["table_id", "actor_id", "fragment_id"],
798 registry
799 )
800 .unwrap();
801
802 let agg_chunk_lookup_miss_count = register_guarded_int_counter_vec_with_registry!(
803 "stream_agg_chunk_lookup_miss_count",
804 "Aggregation executor chunk-level lookup miss duration",
805 &["table_id", "actor_id", "fragment_id"],
806 registry
807 )
808 .unwrap();
809
810 let agg_chunk_total_lookup_count = register_guarded_int_counter_vec_with_registry!(
811 "stream_agg_chunk_lookup_total_count",
812 "Aggregation executor chunk-level lookup total operation",
813 &["table_id", "actor_id", "fragment_id"],
814 registry
815 )
816 .unwrap();
817
818 let backfill_snapshot_read_row_count = register_guarded_int_counter_vec_with_registry!(
819 "stream_backfill_snapshot_read_row_count",
820 "Total number of rows that have been read from the backfill snapshot",
821 &["table_id", "actor_id"],
822 registry
823 )
824 .unwrap();
825
826 let backfill_upstream_output_row_count = register_guarded_int_counter_vec_with_registry!(
827 "stream_backfill_upstream_output_row_count",
828 "Total number of rows that have been output from the backfill upstream",
829 &["table_id", "actor_id"],
830 registry
831 )
832 .unwrap();
833
834 let cdc_backfill_snapshot_read_row_count = register_guarded_int_counter_vec_with_registry!(
835 "stream_cdc_backfill_snapshot_read_row_count",
836 "Total number of rows that have been read from the cdc_backfill snapshot",
837 &["table_id", "actor_id"],
838 registry
839 )
840 .unwrap();
841
842 let cdc_backfill_upstream_output_row_count =
843 register_guarded_int_counter_vec_with_registry!(
844 "stream_cdc_backfill_upstream_output_row_count",
845 "Total number of rows that have been output from the cdc_backfill upstream",
846 &["table_id", "actor_id"],
847 registry
848 )
849 .unwrap();
850
851 let snapshot_backfill_consume_row_count = register_guarded_int_counter_vec_with_registry!(
852 "stream_snapshot_backfill_consume_snapshot_row_count",
853 "Total number of rows that have been output from snapshot backfill",
854 &["table_id", "actor_id", "stage"],
855 registry
856 )
857 .unwrap();
858
859 let over_window_cached_entry_count = register_guarded_int_gauge_vec_with_registry!(
860 "stream_over_window_cached_entry_count",
861 "Total entry (partition) count in over window executor cache",
862 &["table_id", "actor_id", "fragment_id"],
863 registry
864 )
865 .unwrap();
866
867 let over_window_cache_lookup_count = register_guarded_int_counter_vec_with_registry!(
868 "stream_over_window_cache_lookup_count",
869 "Over window executor cache lookup count",
870 &["table_id", "actor_id", "fragment_id"],
871 registry
872 )
873 .unwrap();
874
875 let over_window_cache_miss_count = register_guarded_int_counter_vec_with_registry!(
876 "stream_over_window_cache_miss_count",
877 "Over window executor cache miss count",
878 &["table_id", "actor_id", "fragment_id"],
879 registry
880 )
881 .unwrap();
882
883 let over_window_range_cache_entry_count = register_guarded_int_gauge_vec_with_registry!(
884 "stream_over_window_range_cache_entry_count",
885 "Over window partition range cache entry count",
886 &["table_id", "actor_id", "fragment_id"],
887 registry,
888 )
889 .unwrap();
890
891 let over_window_range_cache_lookup_count = register_guarded_int_counter_vec_with_registry!(
892 "stream_over_window_range_cache_lookup_count",
893 "Over window partition range cache lookup count",
894 &["table_id", "actor_id", "fragment_id"],
895 registry
896 )
897 .unwrap();
898
899 let over_window_range_cache_left_miss_count =
900 register_guarded_int_counter_vec_with_registry!(
901 "stream_over_window_range_cache_left_miss_count",
902 "Over window partition range cache left miss count",
903 &["table_id", "actor_id", "fragment_id"],
904 registry
905 )
906 .unwrap();
907
908 let over_window_range_cache_right_miss_count =
909 register_guarded_int_counter_vec_with_registry!(
910 "stream_over_window_range_cache_right_miss_count",
911 "Over window partition range cache right miss count",
912 &["table_id", "actor_id", "fragment_id"],
913 registry
914 )
915 .unwrap();
916
917 let over_window_accessed_entry_count = register_guarded_int_counter_vec_with_registry!(
918 "stream_over_window_accessed_entry_count",
919 "Over window accessed entry count",
920 &["table_id", "actor_id", "fragment_id"],
921 registry
922 )
923 .unwrap();
924
925 let over_window_compute_count = register_guarded_int_counter_vec_with_registry!(
926 "stream_over_window_compute_count",
927 "Over window compute count",
928 &["table_id", "actor_id", "fragment_id"],
929 registry
930 )
931 .unwrap();
932
933 let over_window_same_output_count = register_guarded_int_counter_vec_with_registry!(
934 "stream_over_window_same_output_count",
935 "Over window same output count",
936 &["table_id", "actor_id", "fragment_id"],
937 registry
938 )
939 .unwrap();
940
941 let over_window_state_cleaned_row_count = register_guarded_int_counter_vec_with_registry!(
942 "stream_over_window_state_cleaned_row_count",
943 "Over window number of stale state rows deleted by watermark-driven state cleaning",
944 &["table_id", "actor_id", "fragment_id"],
945 registry
946 )
947 .unwrap();
948
949 let match_recognize_matches_emitted_count =
950 register_guarded_int_counter_vec_with_registry!(
951 "stream_match_recognize_matches_emitted_count",
952 "Matches emitted by the match recognize executor",
953 &["table_id", "actor_id", "fragment_id"],
954 registry
955 )
956 .unwrap();
957
958 let match_recognize_evicted_rows_count = register_guarded_int_counter_vec_with_registry!(
959 "stream_match_recognize_evicted_rows_count",
960 "Buffered rows evicted by the match recognize executor",
961 &["table_id", "actor_id", "fragment_id"],
962 registry
963 )
964 .unwrap();
965
966 let match_recognize_scan_budget_exhausted_count = register_guarded_int_counter_vec_with_registry!(
967 "stream_match_recognize_scan_budget_exhausted_count",
968 "Partition visits that exhausted the match recognize scan budget and degraded conservatively",
969 &["table_id", "actor_id", "fragment_id"],
970 registry
971 )
972 .unwrap();
973
974 let match_recognize_within_deadline_overflow_count = register_guarded_int_counter_vec_with_registry!(
975 "stream_match_recognize_within_deadline_overflow_count",
976 "Buffered rows whose WITHIN deadline lies past the ORDER BY type's range, so their window never closes",
977 &["table_id", "actor_id", "fragment_id"],
978 registry
979 )
980 .unwrap();
981
982 let match_recognize_retained_rows = register_guarded_int_gauge_vec_with_registry!(
983 "stream_match_recognize_retained_rows",
984 "Rows currently retained in memory by the match recognize executor, summed over the actor's partitions",
985 &["table_id", "actor_id", "fragment_id"],
986 registry
987 )
988 .unwrap();
989
990 let barrier_inflight_latency = register_guarded_histogram_vec_with_registry!(
991 "stream_barrier_inflight_duration_seconds",
992 "barrier_inflight_latency",
993 &["partial_graph"],
994 latency_buckets(600.0, 20),
995 registry
996 )
997 .unwrap();
998
999 let barrier_sync_latency = register_guarded_histogram_vec_with_registry!(
1000 "stream_barrier_sync_storage_duration_seconds",
1001 "barrier_sync_latency",
1002 &["partial_graph"],
1003 latency_buckets(600.0, 20),
1004 registry
1005 )
1006 .unwrap();
1007
1008 let opts = histogram_opts!(
1009 "stream_barrier_batch_size",
1010 "barrier_batch_size",
1011 exponential_buckets(1.0, 2.0, 8).unwrap()
1012 );
1013 let barrier_batch_size = register_histogram_with_registry!(opts, registry).unwrap();
1014
1015 let barrier_manager_progress = register_guarded_int_counter_vec_with_registry!(
1016 "stream_barrier_manager_progress",
1017 "The number of actors that have processed the earliest in-flight barriers",
1018 &["partial_graph"],
1019 registry
1020 )
1021 .unwrap();
1022
1023 let sync_kv_log_store_wait_next_poll_ns = register_guarded_int_counter_vec_with_registry!(
1024 "sync_kv_log_store_wait_next_poll_ns",
1025 "Total duration (ns) of waiting for next poll",
1026 &["actor_id", "target", "fragment_id", "relation"],
1027 registry
1028 )
1029 .unwrap();
1030
1031 let sync_kv_log_store_read_count = register_guarded_int_counter_vec_with_registry!(
1032 "sync_kv_log_store_read_count",
1033 "read row count throughput of sync_kv log store",
1034 &["type", "actor_id", "target", "fragment_id", "relation"],
1035 registry
1036 )
1037 .unwrap();
1038
1039 let sync_kv_log_store_read_size = register_guarded_int_counter_vec_with_registry!(
1040 "sync_kv_log_store_read_size",
1041 "read size throughput of sync_kv log store",
1042 &["type", "actor_id", "target", "fragment_id", "relation"],
1043 registry
1044 )
1045 .unwrap();
1046
1047 let sync_kv_log_store_write_pause_duration_ns =
1048 register_guarded_int_counter_vec_with_registry!(
1049 "sync_kv_log_store_write_pause_duration_ns",
1050 "Duration (ns) of sync_kv log store write pause",
1051 &["actor_id", "target", "fragment_id", "relation"],
1052 registry
1053 )
1054 .unwrap();
1055
1056 let sync_kv_log_store_state = register_guarded_int_counter_vec_with_registry!(
1057 "sync_kv_log_store_state",
1058 "clean/unclean state transition for sync_kv log store",
1059 &["state", "actor_id", "target", "fragment_id", "relation"],
1060 registry
1061 )
1062 .unwrap();
1063
1064 let sync_kv_log_store_storage_write_count =
1065 register_guarded_int_counter_vec_with_registry!(
1066 "sync_kv_log_store_storage_write_count",
1067 "Write row count throughput of sync_kv log store",
1068 &["actor_id", "target", "fragment_id", "relation"],
1069 registry
1070 )
1071 .unwrap();
1072
1073 let sync_kv_log_store_storage_write_size = register_guarded_int_counter_vec_with_registry!(
1074 "sync_kv_log_store_storage_write_size",
1075 "Write size throughput of sync_kv log store",
1076 &["actor_id", "target", "fragment_id", "relation"],
1077 registry
1078 )
1079 .unwrap();
1080
1081 let sync_kv_log_store_buffer_unconsumed_item_count =
1082 register_guarded_int_gauge_vec_with_registry!(
1083 "sync_kv_log_store_buffer_unconsumed_item_count",
1084 "Number of Unconsumed Item in buffer",
1085 &["actor_id", "target", "fragment_id", "relation"],
1086 registry
1087 )
1088 .unwrap();
1089
1090 let sync_kv_log_store_buffer_unconsumed_row_count =
1091 register_guarded_int_gauge_vec_with_registry!(
1092 "sync_kv_log_store_buffer_unconsumed_row_count",
1093 "Number of Unconsumed Row in buffer",
1094 &["actor_id", "target", "fragment_id", "relation"],
1095 registry
1096 )
1097 .unwrap();
1098
1099 let sync_kv_log_store_buffer_unconsumed_epoch_count =
1100 register_guarded_int_gauge_vec_with_registry!(
1101 "sync_kv_log_store_buffer_unconsumed_epoch_count",
1102 "Number of Unconsumed Epoch in buffer",
1103 &["actor_id", "target", "fragment_id", "relation"],
1104 registry
1105 )
1106 .unwrap();
1107
1108 let sync_kv_log_store_buffer_unconsumed_min_epoch =
1109 register_guarded_int_gauge_vec_with_registry!(
1110 "sync_kv_log_store_buffer_unconsumed_min_epoch",
1111 "Number of Unconsumed Epoch in buffer",
1112 &["actor_id", "target", "fragment_id", "relation"],
1113 registry
1114 )
1115 .unwrap();
1116 let sync_kv_log_store_buffer_memory_bytes =
1117 register_guarded_int_gauge_vec_with_registry!(
1118 "sync_kv_log_store_buffer_memory_bytes",
1119 "Estimated heap bytes used by synced kv log store buffer (unconsumed + consumed but not truncated)",
1120 &["actor_id", "target", "fragment_id", "relation"],
1121 registry
1122 )
1123 .unwrap();
1124
1125 let kv_log_store_storage_write_count = register_guarded_int_counter_vec_with_registry!(
1126 "kv_log_store_storage_write_count",
1127 "Write row count throughput of kv log store",
1128 &["actor_id", "connector", "sink_id", "sink_name"],
1129 registry
1130 )
1131 .unwrap();
1132
1133 let kv_log_store_storage_write_size = register_guarded_int_counter_vec_with_registry!(
1134 "kv_log_store_storage_write_size",
1135 "Write size throughput of kv log store",
1136 &["actor_id", "connector", "sink_id", "sink_name"],
1137 registry
1138 )
1139 .unwrap();
1140
1141 let kv_log_store_storage_read_count = register_guarded_int_counter_vec_with_registry!(
1142 "kv_log_store_storage_read_count",
1143 "Write row count throughput of kv log store",
1144 &["actor_id", "connector", "sink_id", "sink_name", "read_type"],
1145 registry
1146 )
1147 .unwrap();
1148
1149 let kv_log_store_storage_read_size = register_guarded_int_counter_vec_with_registry!(
1150 "kv_log_store_storage_read_size",
1151 "Write size throughput of kv log store",
1152 &["actor_id", "connector", "sink_id", "sink_name", "read_type"],
1153 registry
1154 )
1155 .unwrap();
1156
1157 let kv_log_store_rewind_count = register_guarded_int_counter_vec_with_registry!(
1158 "kv_log_store_rewind_count",
1159 "Kv log store rewind rate",
1160 &["actor_id", "connector", "sink_id", "sink_name"],
1161 registry
1162 )
1163 .unwrap();
1164
1165 let kv_log_store_rewind_delay_opts = {
1166 let first_delay_secs = REWIND_INITIAL_DELAY.as_secs_f64();
1167 let base = REWIND_BACKOFF_MULTIPLIER as f64;
1168 let bucket_count = (REWIND_MAX_DELAY.as_secs_f64() / first_delay_secs)
1169 .log(base)
1170 .ceil() as usize;
1171 let buckets = exponential_buckets(first_delay_secs, base, bucket_count).unwrap();
1172 histogram_opts!(
1173 "kv_log_store_rewind_delay",
1174 "Kv log store rewind delay",
1175 buckets,
1176 )
1177 };
1178
1179 let kv_log_store_rewind_delay = register_guarded_histogram_vec_with_registry!(
1180 kv_log_store_rewind_delay_opts,
1181 &["actor_id", "connector", "sink_id", "sink_name"],
1182 registry
1183 )
1184 .unwrap();
1185
1186 let kv_log_store_buffer_unconsumed_item_count =
1187 register_guarded_int_gauge_vec_with_registry!(
1188 "kv_log_store_buffer_unconsumed_item_count",
1189 "Number of Unconsumed Item in buffer",
1190 &["actor_id", "connector", "sink_id", "sink_name"],
1191 registry
1192 )
1193 .unwrap();
1194
1195 let kv_log_store_buffer_unconsumed_row_count =
1196 register_guarded_int_gauge_vec_with_registry!(
1197 "kv_log_store_buffer_unconsumed_row_count",
1198 "Number of Unconsumed Row in buffer",
1199 &["actor_id", "connector", "sink_id", "sink_name"],
1200 registry
1201 )
1202 .unwrap();
1203
1204 let kv_log_store_buffer_unconsumed_epoch_count =
1205 register_guarded_int_gauge_vec_with_registry!(
1206 "kv_log_store_buffer_unconsumed_epoch_count",
1207 "Number of Unconsumed Epoch in buffer",
1208 &["actor_id", "connector", "sink_id", "sink_name"],
1209 registry
1210 )
1211 .unwrap();
1212
1213 let kv_log_store_buffer_unconsumed_min_epoch =
1214 register_guarded_int_gauge_vec_with_registry!(
1215 "kv_log_store_buffer_unconsumed_min_epoch",
1216 "Number of Unconsumed Epoch in buffer",
1217 &["actor_id", "connector", "sink_id", "sink_name"],
1218 registry
1219 )
1220 .unwrap();
1221
1222 let crossdb_last_consumed_min_epoch = register_guarded_int_gauge_vec_with_registry!(
1223 "crossdb_last_consumed_min_epoch",
1224 "Last consumed min epoch for cross-database changelog stream scan",
1225 &["table_id", "actor_id", "fragment_id"],
1226 registry
1227 )
1228 .unwrap();
1229
1230 let kv_log_store_buffer_memory_bytes =
1231 register_guarded_int_gauge_vec_with_registry!(
1232 "kv_log_store_buffer_memory_bytes",
1233 "Estimated heap bytes used by kv log store buffer (unconsumed + consumed but not truncated)",
1234 &["actor_id", "connector", "sink_id", "sink_name"],
1235 registry
1236 )
1237 .unwrap();
1238
1239 let lru_runtime_loop_count = register_int_counter_with_registry!(
1240 "lru_runtime_loop_count",
1241 "The counts of the eviction loop in LRU manager per second",
1242 registry
1243 )
1244 .unwrap();
1245
1246 let lru_latest_sequence = register_int_gauge_with_registry!(
1247 "lru_latest_sequence",
1248 "Current LRU global sequence",
1249 registry,
1250 )
1251 .unwrap();
1252
1253 let lru_watermark_sequence = register_int_gauge_with_registry!(
1254 "lru_watermark_sequence",
1255 "Current LRU watermark sequence",
1256 registry,
1257 )
1258 .unwrap();
1259
1260 let lru_eviction_policy = register_int_gauge_with_registry!(
1261 "lru_eviction_policy",
1262 "Current LRU eviction policy",
1263 registry,
1264 )
1265 .unwrap();
1266
1267 let jemalloc_allocated_bytes = register_int_gauge_with_registry!(
1268 "jemalloc_allocated_bytes",
1269 "The allocated memory jemalloc, got from jemalloc_ctl",
1270 registry
1271 )
1272 .unwrap();
1273
1274 let jemalloc_active_bytes = register_int_gauge_with_registry!(
1275 "jemalloc_active_bytes",
1276 "The active memory jemalloc, got from jemalloc_ctl",
1277 registry
1278 )
1279 .unwrap();
1280
1281 let jemalloc_resident_bytes = register_int_gauge_with_registry!(
1282 "jemalloc_resident_bytes",
1283 "The active memory jemalloc, got from jemalloc_ctl",
1284 registry
1285 )
1286 .unwrap();
1287
1288 let jemalloc_metadata_bytes = register_int_gauge_with_registry!(
1289 "jemalloc_metadata_bytes",
1290 "The active memory jemalloc, got from jemalloc_ctl",
1291 registry
1292 )
1293 .unwrap();
1294
1295 let jvm_allocated_bytes = register_int_gauge_with_registry!(
1296 "jvm_allocated_bytes",
1297 "The allocated jvm memory",
1298 registry
1299 )
1300 .unwrap();
1301
1302 let jvm_active_bytes = register_int_gauge_with_registry!(
1303 "jvm_active_bytes",
1304 "The active jvm memory",
1305 registry
1306 )
1307 .unwrap();
1308
1309 let materialize_cache_hit_count = register_guarded_int_counter_vec_with_registry!(
1310 "stream_materialize_cache_hit_count",
1311 "Materialize executor cache hit count",
1312 &["actor_id", "table_id", "fragment_id"],
1313 registry
1314 )
1315 .unwrap()
1316 .relabel_debug_1(level);
1317
1318 let materialize_data_exist_count = register_guarded_int_counter_vec_with_registry!(
1319 "stream_materialize_data_exist_count",
1320 "Materialize executor data exist count",
1321 &["actor_id", "table_id", "fragment_id"],
1322 registry
1323 )
1324 .unwrap()
1325 .relabel_debug_1(level);
1326
1327 let materialize_cache_total_count = register_guarded_int_counter_vec_with_registry!(
1328 "stream_materialize_cache_total_count",
1329 "Materialize executor cache total operation",
1330 &["actor_id", "table_id", "fragment_id"],
1331 registry
1332 )
1333 .unwrap()
1334 .relabel_debug_1(level);
1335
1336 let stream_memory_usage = register_guarded_int_gauge_vec_with_registry!(
1337 "stream_memory_usage",
1338 "Memory usage for stream executors",
1339 &["actor_id", "table_id", "desc"],
1340 registry
1341 )
1342 .unwrap()
1343 .relabel_debug_1(level);
1344
1345 let gap_fill_generated_rows_count = register_guarded_int_counter_vec_with_registry!(
1346 "gap_fill_generated_rows_count",
1347 "Total number of rows generated by gap fill executor",
1348 &["actor_id", "fragment_id"],
1349 registry
1350 )
1351 .unwrap()
1352 .relabel_debug_1(level);
1353
1354 let state_table_iter_count = register_guarded_int_counter_vec_with_registry!(
1355 "state_table_iter_count",
1356 "Total number of state table iter operations",
1357 &["actor_id", "fragment_id", "table_id"],
1358 registry
1359 )
1360 .unwrap()
1361 .relabel_debug_1(level);
1362
1363 let state_table_get_count = register_guarded_int_counter_vec_with_registry!(
1364 "state_table_get_count",
1365 "Total number of state table get operations",
1366 &["actor_id", "fragment_id", "table_id"],
1367 registry
1368 )
1369 .unwrap()
1370 .relabel_debug_1(level);
1371
1372 let state_table_iter_vnode_pruned_count = register_guarded_int_counter_vec_with_registry!(
1373 "state_table_iter_vnode_pruned_count",
1374 "Total number of state table iter operations pruned by vnode statistics",
1375 &["actor_id", "fragment_id", "table_id"],
1376 registry
1377 )
1378 .unwrap()
1379 .relabel_debug_1(level);
1380
1381 let state_table_get_vnode_pruned_count = register_guarded_int_counter_vec_with_registry!(
1382 "state_table_get_vnode_pruned_count",
1383 "Total number of state table get operations pruned by vnode statistics",
1384 &["actor_id", "fragment_id", "table_id"],
1385 registry
1386 )
1387 .unwrap()
1388 .relabel_debug_1(level);
1389
1390 let now_streaming_clock_ms = register_guarded_int_gauge_vec_with_registry!(
1391 "stream_now_streaming_clock_ms",
1392 "Latest streaming NOW() timestamp (milliseconds since Unix epoch) \
1393 emitted as a watermark by this fragment on this compute node.",
1394 &["fragment_id"],
1395 registry,
1396 )
1397 .unwrap();
1398 let now_wall_clock_drift_ms = register_guarded_int_gauge_vec_with_registry!(
1399 "stream_now_wall_clock_drift_ms",
1400 "Latest processed barrier epoch minus the latest streaming NOW() timestamp \
1401 emitted as a watermark by this fragment on this compute node, \
1402 in milliseconds. This is relative to processed barriers, not scrape-time wall clock.",
1403 &["fragment_id"],
1404 registry,
1405 )
1406 .unwrap();
1407
1408 Self {
1409 level,
1410 executor_row_count,
1411 mem_stream_node_output_row_count: stream_node_output_row_count,
1412 mem_stream_node_output_blocking_duration_ns: stream_node_output_blocking_duration_ns,
1413 actor_scheduled_duration,
1414 actor_scheduled_cnt,
1415 actor_poll_duration,
1416 actor_poll_cnt,
1417 actor_idle_duration,
1418 actor_idle_cnt,
1419 actor_count,
1420 actor_in_record_cnt,
1421 actor_out_record_cnt,
1422 fragment_channel_buffered_bytes,
1423 actor_current_epoch,
1424 project_expr_inflight_window_size,
1425 source_output_row_count,
1426 source_split_change_count,
1427 source_backfill_row_count,
1428 sink_input_row_count,
1429 sink_input_bytes,
1430 sink_chunk_buffer_size,
1431 exchange_frag_recv_size,
1432 merge_barrier_align_duration,
1433 actor_output_buffer_blocking_duration_ns,
1434 actor_input_buffer_blocking_duration_ns,
1435 join_lookup_miss_count,
1436 join_lookup_total_count,
1437 join_insert_cache_miss_count,
1438 join_actor_input_waiting_duration_ns,
1439 join_match_duration_ns,
1440 join_cached_entry_count,
1441 join_matched_join_keys,
1442 barrier_align_duration,
1443 agg_lookup_miss_count,
1444 agg_total_lookup_count,
1445 agg_cached_entry_count,
1446 agg_chunk_lookup_miss_count,
1447 agg_chunk_total_lookup_count,
1448 agg_dirty_groups_count,
1449 agg_dirty_groups_heap_size,
1450 agg_distinct_cache_miss_count,
1451 agg_distinct_total_cache_count,
1452 agg_distinct_cached_entry_count,
1453 agg_state_cache_lookup_count,
1454 agg_state_cache_miss_count,
1455 group_top_n_cache_miss_count,
1456 group_top_n_total_query_cache_count,
1457 group_top_n_cached_entry_count,
1458 group_top_n_appendonly_cache_miss_count,
1459 group_top_n_appendonly_total_query_cache_count,
1460 group_top_n_appendonly_cached_entry_count,
1461 lookup_cache_miss_count,
1462 lookup_total_query_cache_count,
1463 lookup_cached_entry_count,
1464 temporal_join_cache_miss_count,
1465 temporal_join_total_query_cache_count,
1466 temporal_join_cached_entry_count,
1467 backfill_snapshot_read_row_count,
1468 backfill_upstream_output_row_count,
1469 cdc_backfill_snapshot_read_row_count,
1470 cdc_backfill_upstream_output_row_count,
1471 snapshot_backfill_consume_row_count,
1472 over_window_cached_entry_count,
1473 over_window_cache_lookup_count,
1474 over_window_cache_miss_count,
1475 over_window_range_cache_entry_count,
1476 over_window_range_cache_lookup_count,
1477 over_window_range_cache_left_miss_count,
1478 over_window_range_cache_right_miss_count,
1479 over_window_accessed_entry_count,
1480 over_window_compute_count,
1481 over_window_same_output_count,
1482 over_window_state_cleaned_row_count,
1483 match_recognize_matches_emitted_count,
1484 match_recognize_evicted_rows_count,
1485 match_recognize_scan_budget_exhausted_count,
1486 match_recognize_within_deadline_overflow_count,
1487 match_recognize_retained_rows,
1488 barrier_inflight_latency,
1489 barrier_sync_latency,
1490 barrier_batch_size,
1491 barrier_manager_progress,
1492 kv_log_store_storage_write_count,
1493 kv_log_store_storage_write_size,
1494 kv_log_store_rewind_count,
1495 kv_log_store_rewind_delay,
1496 kv_log_store_storage_read_count,
1497 kv_log_store_storage_read_size,
1498 kv_log_store_buffer_unconsumed_item_count,
1499 kv_log_store_buffer_unconsumed_row_count,
1500 kv_log_store_buffer_unconsumed_epoch_count,
1501 kv_log_store_buffer_unconsumed_min_epoch,
1502 kv_log_store_buffer_memory_bytes,
1503 crossdb_last_consumed_min_epoch,
1504 sync_kv_log_store_read_count,
1505 sync_kv_log_store_read_size,
1506 sync_kv_log_store_write_pause_duration_ns,
1507 sync_kv_log_store_state,
1508 sync_kv_log_store_wait_next_poll_ns,
1509 sync_kv_log_store_storage_write_count,
1510 sync_kv_log_store_storage_write_size,
1511 sync_kv_log_store_buffer_unconsumed_item_count,
1512 sync_kv_log_store_buffer_unconsumed_row_count,
1513 sync_kv_log_store_buffer_unconsumed_epoch_count,
1514 sync_kv_log_store_buffer_unconsumed_min_epoch,
1515 sync_kv_log_store_buffer_memory_bytes,
1516 lru_runtime_loop_count,
1517 lru_latest_sequence,
1518 lru_watermark_sequence,
1519 lru_eviction_policy,
1520 jemalloc_allocated_bytes,
1521 jemalloc_active_bytes,
1522 jemalloc_resident_bytes,
1523 jemalloc_metadata_bytes,
1524 jvm_allocated_bytes,
1525 jvm_active_bytes,
1526 stream_memory_usage,
1527 materialize_cache_hit_count,
1528 materialize_data_exist_count,
1529 materialize_cache_total_count,
1530 materialize_input_row_count,
1531 materialize_current_epoch,
1532 pg_cdc_state_table_lsn,
1533 pg_cdc_jni_commit_offset_lsn,
1534 mysql_cdc_state_binlog_file_seq,
1535 mysql_cdc_state_binlog_position,
1536 sqlserver_cdc_state_change_lsn,
1537 sqlserver_cdc_state_commit_lsn,
1538 sqlserver_cdc_jni_commit_offset_lsn,
1539 gap_fill_generated_rows_count,
1540 now_streaming_clock_ms,
1541 now_wall_clock_drift_ms,
1542 state_table_iter_count,
1543 state_table_get_count,
1544 state_table_iter_vnode_pruned_count,
1545 state_table_get_vnode_pruned_count,
1546 }
1547 }
1548
1549 pub fn unused() -> Self {
1551 global_streaming_metrics(MetricLevel::Disabled)
1552 }
1553
1554 pub fn new_actor_metrics(&self, actor_id: ActorId, fragment_id: FragmentId) -> ActorMetrics {
1555 let label_list: &[&str; 2] = &[&actor_id.to_string(), &fragment_id.to_string()];
1556 let actor_scheduled_duration = self
1557 .actor_scheduled_duration
1558 .with_guarded_label_values(label_list);
1559 let actor_scheduled_cnt = self
1560 .actor_scheduled_cnt
1561 .with_guarded_label_values(label_list);
1562 let actor_poll_duration = self
1563 .actor_poll_duration
1564 .with_guarded_label_values(label_list);
1565 let actor_poll_cnt = self.actor_poll_cnt.with_guarded_label_values(label_list);
1566 let actor_idle_duration = self
1567 .actor_idle_duration
1568 .with_guarded_label_values(label_list);
1569 let actor_idle_cnt = self.actor_idle_cnt.with_guarded_label_values(label_list);
1570 ActorMetrics {
1571 actor_scheduled_duration,
1572 actor_scheduled_cnt,
1573 actor_poll_duration,
1574 actor_poll_cnt,
1575 actor_idle_duration,
1576 actor_idle_cnt,
1577 }
1578 }
1579
1580 pub(crate) fn new_actor_input_metrics(
1581 &self,
1582 actor_id: ActorId,
1583 fragment_id: FragmentId,
1584 upstream_fragment_id: FragmentId,
1585 ) -> ActorInputMetrics {
1586 let actor_id_str = actor_id.to_string();
1587 let fragment_id_str = fragment_id.to_string();
1588 let upstream_fragment_id_str = upstream_fragment_id.to_string();
1589 ActorInputMetrics {
1590 actor_in_record_cnt: self.actor_in_record_cnt.with_guarded_label_values(&[
1591 &actor_id_str,
1592 &fragment_id_str,
1593 &upstream_fragment_id_str,
1594 ]),
1595 actor_input_buffer_blocking_duration_ns: self
1596 .actor_input_buffer_blocking_duration_ns
1597 .with_guarded_label_values(&[
1598 &actor_id_str,
1599 &fragment_id_str,
1600 &upstream_fragment_id_str,
1601 ]),
1602 }
1603 }
1604
1605 pub fn new_sink_exec_metrics(
1606 &self,
1607 id: SinkId,
1608 actor_id: ActorId,
1609 fragment_id: FragmentId,
1610 ) -> SinkExecutorMetrics {
1611 let label_list: &[&str; 3] = &[
1612 &id.to_string(),
1613 &actor_id.to_string(),
1614 &fragment_id.to_string(),
1615 ];
1616 SinkExecutorMetrics {
1617 sink_input_row_count: self
1618 .sink_input_row_count
1619 .with_guarded_label_values(label_list),
1620 sink_input_bytes: self.sink_input_bytes.with_guarded_label_values(label_list),
1621 sink_chunk_buffer_size: self
1622 .sink_chunk_buffer_size
1623 .with_guarded_label_values(label_list),
1624 }
1625 }
1626
1627 pub fn new_group_top_n_metrics(
1628 &self,
1629 table_id: TableId,
1630 actor_id: ActorId,
1631 fragment_id: FragmentId,
1632 ) -> GroupTopNMetrics {
1633 let label_list: &[&str; 3] = &[
1634 &table_id.to_string(),
1635 &actor_id.to_string(),
1636 &fragment_id.to_string(),
1637 ];
1638
1639 GroupTopNMetrics {
1640 group_top_n_cache_miss_count: self
1641 .group_top_n_cache_miss_count
1642 .with_guarded_label_values(label_list),
1643 group_top_n_total_query_cache_count: self
1644 .group_top_n_total_query_cache_count
1645 .with_guarded_label_values(label_list),
1646 group_top_n_cached_entry_count: self
1647 .group_top_n_cached_entry_count
1648 .with_guarded_label_values(label_list),
1649 }
1650 }
1651
1652 pub fn new_append_only_group_top_n_metrics(
1653 &self,
1654 table_id: TableId,
1655 actor_id: ActorId,
1656 fragment_id: FragmentId,
1657 ) -> GroupTopNMetrics {
1658 let label_list: &[&str; 3] = &[
1659 &table_id.to_string(),
1660 &actor_id.to_string(),
1661 &fragment_id.to_string(),
1662 ];
1663
1664 GroupTopNMetrics {
1665 group_top_n_cache_miss_count: self
1666 .group_top_n_appendonly_cache_miss_count
1667 .with_guarded_label_values(label_list),
1668 group_top_n_total_query_cache_count: self
1669 .group_top_n_appendonly_total_query_cache_count
1670 .with_guarded_label_values(label_list),
1671 group_top_n_cached_entry_count: self
1672 .group_top_n_appendonly_cached_entry_count
1673 .with_guarded_label_values(label_list),
1674 }
1675 }
1676
1677 pub fn new_lookup_executor_metrics(
1678 &self,
1679 table_id: TableId,
1680 actor_id: ActorId,
1681 fragment_id: FragmentId,
1682 ) -> LookupExecutorMetrics {
1683 let label_list: &[&str; 3] = &[
1684 &table_id.to_string(),
1685 &actor_id.to_string(),
1686 &fragment_id.to_string(),
1687 ];
1688
1689 LookupExecutorMetrics {
1690 lookup_cache_miss_count: self
1691 .lookup_cache_miss_count
1692 .with_guarded_label_values(label_list),
1693 lookup_total_query_cache_count: self
1694 .lookup_total_query_cache_count
1695 .with_guarded_label_values(label_list),
1696 lookup_cached_entry_count: self
1697 .lookup_cached_entry_count
1698 .with_guarded_label_values(label_list),
1699 }
1700 }
1701
1702 pub fn new_hash_agg_metrics(
1703 &self,
1704 table_id: TableId,
1705 actor_id: ActorId,
1706 fragment_id: FragmentId,
1707 ) -> HashAggMetrics {
1708 let label_list: &[&str; 3] = &[
1709 &table_id.to_string(),
1710 &actor_id.to_string(),
1711 &fragment_id.to_string(),
1712 ];
1713 HashAggMetrics {
1714 agg_lookup_miss_count: self
1715 .agg_lookup_miss_count
1716 .with_guarded_label_values(label_list),
1717 agg_total_lookup_count: self
1718 .agg_total_lookup_count
1719 .with_guarded_label_values(label_list),
1720 agg_cached_entry_count: self
1721 .agg_cached_entry_count
1722 .with_guarded_label_values(label_list),
1723 agg_chunk_lookup_miss_count: self
1724 .agg_chunk_lookup_miss_count
1725 .with_guarded_label_values(label_list),
1726 agg_chunk_total_lookup_count: self
1727 .agg_chunk_total_lookup_count
1728 .with_guarded_label_values(label_list),
1729 agg_dirty_groups_count: self
1730 .agg_dirty_groups_count
1731 .with_guarded_label_values(label_list),
1732 agg_dirty_groups_heap_size: self
1733 .agg_dirty_groups_heap_size
1734 .with_guarded_label_values(label_list),
1735 agg_state_cache_lookup_count: self
1736 .agg_state_cache_lookup_count
1737 .with_guarded_label_values(label_list),
1738 agg_state_cache_miss_count: self
1739 .agg_state_cache_miss_count
1740 .with_guarded_label_values(label_list),
1741 }
1742 }
1743
1744 pub fn new_agg_distinct_dedup_metrics(
1745 &self,
1746 table_id: TableId,
1747 actor_id: ActorId,
1748 fragment_id: FragmentId,
1749 ) -> AggDistinctDedupMetrics {
1750 let label_list: &[&str; 3] = &[
1751 &table_id.to_string(),
1752 &actor_id.to_string(),
1753 &fragment_id.to_string(),
1754 ];
1755 AggDistinctDedupMetrics {
1756 agg_distinct_cache_miss_count: self
1757 .agg_distinct_cache_miss_count
1758 .with_guarded_label_values(label_list),
1759 agg_distinct_total_cache_count: self
1760 .agg_distinct_total_cache_count
1761 .with_guarded_label_values(label_list),
1762 agg_distinct_cached_entry_count: self
1763 .agg_distinct_cached_entry_count
1764 .with_guarded_label_values(label_list),
1765 }
1766 }
1767
1768 pub fn new_temporal_join_metrics(
1769 &self,
1770 table_id: TableId,
1771 actor_id: ActorId,
1772 fragment_id: FragmentId,
1773 ) -> TemporalJoinMetrics {
1774 let label_list: &[&str; 3] = &[
1775 &table_id.to_string(),
1776 &actor_id.to_string(),
1777 &fragment_id.to_string(),
1778 ];
1779 TemporalJoinMetrics {
1780 temporal_join_cache_miss_count: self
1781 .temporal_join_cache_miss_count
1782 .with_guarded_label_values(label_list),
1783 temporal_join_total_query_cache_count: self
1784 .temporal_join_total_query_cache_count
1785 .with_guarded_label_values(label_list),
1786 temporal_join_cached_entry_count: self
1787 .temporal_join_cached_entry_count
1788 .with_guarded_label_values(label_list),
1789 }
1790 }
1791
1792 pub fn new_backfill_metrics(&self, table_id: TableId, actor_id: ActorId) -> BackfillMetrics {
1793 let label_list: &[&str; 2] = &[&table_id.to_string(), &actor_id.to_string()];
1794 BackfillMetrics {
1795 backfill_snapshot_read_row_count: self
1796 .backfill_snapshot_read_row_count
1797 .with_guarded_label_values(label_list),
1798 backfill_upstream_output_row_count: self
1799 .backfill_upstream_output_row_count
1800 .with_guarded_label_values(label_list),
1801 }
1802 }
1803
1804 pub fn new_cdc_backfill_metrics(
1805 &self,
1806 table_id: TableId,
1807 actor_id: ActorId,
1808 ) -> CdcBackfillMetrics {
1809 let label_list: &[&str; 2] = &[&table_id.to_string(), &actor_id.to_string()];
1810 CdcBackfillMetrics {
1811 cdc_backfill_snapshot_read_row_count: self
1812 .cdc_backfill_snapshot_read_row_count
1813 .with_guarded_label_values(label_list),
1814 cdc_backfill_upstream_output_row_count: self
1815 .cdc_backfill_upstream_output_row_count
1816 .with_guarded_label_values(label_list),
1817 }
1818 }
1819
1820 pub fn new_over_window_metrics(
1821 &self,
1822 table_id: TableId,
1823 actor_id: ActorId,
1824 fragment_id: FragmentId,
1825 ) -> OverWindowMetrics {
1826 let label_list: &[&str; 3] = &[
1827 &table_id.to_string(),
1828 &actor_id.to_string(),
1829 &fragment_id.to_string(),
1830 ];
1831 OverWindowMetrics {
1832 over_window_cached_entry_count: self
1833 .over_window_cached_entry_count
1834 .with_guarded_label_values(label_list),
1835 over_window_cache_lookup_count: self
1836 .over_window_cache_lookup_count
1837 .with_guarded_label_values(label_list),
1838 over_window_cache_miss_count: self
1839 .over_window_cache_miss_count
1840 .with_guarded_label_values(label_list),
1841 over_window_range_cache_entry_count: self
1842 .over_window_range_cache_entry_count
1843 .with_guarded_label_values(label_list),
1844 over_window_range_cache_lookup_count: self
1845 .over_window_range_cache_lookup_count
1846 .with_guarded_label_values(label_list),
1847 over_window_range_cache_left_miss_count: self
1848 .over_window_range_cache_left_miss_count
1849 .with_guarded_label_values(label_list),
1850 over_window_range_cache_right_miss_count: self
1851 .over_window_range_cache_right_miss_count
1852 .with_guarded_label_values(label_list),
1853 over_window_accessed_entry_count: self
1854 .over_window_accessed_entry_count
1855 .with_guarded_label_values(label_list),
1856 over_window_compute_count: self
1857 .over_window_compute_count
1858 .with_guarded_label_values(label_list),
1859 over_window_same_output_count: self
1860 .over_window_same_output_count
1861 .with_guarded_label_values(label_list),
1862 over_window_state_cleaned_row_count: self
1863 .over_window_state_cleaned_row_count
1864 .with_guarded_label_values(label_list),
1865 }
1866 }
1867
1868 pub fn new_match_recognize_metrics(
1869 &self,
1870 table_id: TableId,
1871 actor_id: ActorId,
1872 fragment_id: FragmentId,
1873 ) -> MatchRecognizeMetrics {
1874 let label_list: &[&str; 3] = &[
1875 &table_id.to_string(),
1876 &actor_id.to_string(),
1877 &fragment_id.to_string(),
1878 ];
1879 MatchRecognizeMetrics {
1880 match_recognize_matches_emitted_count: self
1881 .match_recognize_matches_emitted_count
1882 .with_guarded_label_values(label_list),
1883 match_recognize_evicted_rows_count: self
1884 .match_recognize_evicted_rows_count
1885 .with_guarded_label_values(label_list),
1886 match_recognize_scan_budget_exhausted_count: self
1887 .match_recognize_scan_budget_exhausted_count
1888 .with_guarded_label_values(label_list),
1889 match_recognize_within_deadline_overflow_count: self
1890 .match_recognize_within_deadline_overflow_count
1891 .with_guarded_label_values(label_list),
1892 match_recognize_retained_rows: self
1893 .match_recognize_retained_rows
1894 .with_guarded_label_values(label_list),
1895 }
1896 }
1897
1898 pub fn new_materialize_cache_metrics(
1899 &self,
1900 table_id: TableId,
1901 actor_id: ActorId,
1902 fragment_id: FragmentId,
1903 ) -> MaterializeCacheMetrics {
1904 let label_list: &[&str; 3] = &[
1905 &actor_id.to_string(),
1906 &table_id.to_string(),
1907 &fragment_id.to_string(),
1908 ];
1909 MaterializeCacheMetrics {
1910 materialize_cache_hit_count: self
1911 .materialize_cache_hit_count
1912 .with_guarded_label_values(label_list),
1913 materialize_data_exist_count: self
1914 .materialize_data_exist_count
1915 .with_guarded_label_values(label_list),
1916 materialize_cache_total_count: self
1917 .materialize_cache_total_count
1918 .with_guarded_label_values(label_list),
1919 }
1920 }
1921
1922 pub fn new_materialize_metrics(
1923 &self,
1924 table_id: TableId,
1925 actor_id: ActorId,
1926 fragment_id: FragmentId,
1927 ) -> MaterializeMetrics {
1928 let label_list: &[&str; 3] = &[
1929 &actor_id.to_string(),
1930 &table_id.to_string(),
1931 &fragment_id.to_string(),
1932 ];
1933 MaterializeMetrics {
1934 materialize_input_row_count: self
1935 .materialize_input_row_count
1936 .with_guarded_label_values(label_list),
1937 materialize_current_epoch: self
1938 .materialize_current_epoch
1939 .with_guarded_label_values(label_list),
1940 }
1941 }
1942
1943 pub fn new_state_table_metrics(
1944 &self,
1945 table_id: TableId,
1946 actor_id: ActorId,
1947 fragment_id: FragmentId,
1948 ) -> StateTableMetrics {
1949 let label_list: &[&str; 3] = &[
1950 &actor_id.to_string(),
1951 &fragment_id.to_string(),
1952 &table_id.to_string(),
1953 ];
1954 StateTableMetrics {
1955 iter_count: self
1956 .state_table_iter_count
1957 .with_guarded_label_values(label_list),
1958 get_count: self
1959 .state_table_get_count
1960 .with_guarded_label_values(label_list),
1961 iter_vnode_pruned_count: self
1962 .state_table_iter_vnode_pruned_count
1963 .with_guarded_label_values(label_list),
1964 get_vnode_pruned_count: self
1965 .state_table_get_vnode_pruned_count
1966 .with_guarded_label_values(label_list),
1967 }
1968 }
1969}
1970
1971pub(crate) struct ActorInputMetrics {
1972 pub(crate) actor_in_record_cnt: LabelGuardedIntCounter,
1973 pub(crate) actor_input_buffer_blocking_duration_ns: LabelGuardedIntCounter,
1974}
1975
1976pub struct ActorMetrics {
1978 pub actor_scheduled_duration: LabelGuardedIntCounter,
1979 pub actor_scheduled_cnt: LabelGuardedIntCounter,
1980 pub actor_poll_duration: LabelGuardedIntCounter,
1981 pub actor_poll_cnt: LabelGuardedIntCounter,
1982 pub actor_idle_duration: LabelGuardedIntCounter,
1983 pub actor_idle_cnt: LabelGuardedIntCounter,
1984}
1985
1986pub struct SinkExecutorMetrics {
1987 pub sink_input_row_count: LabelGuardedIntCounter,
1988 pub sink_input_bytes: LabelGuardedIntCounter,
1989 pub sink_chunk_buffer_size: LabelGuardedIntGauge,
1990}
1991
1992pub struct MaterializeCacheMetrics {
1993 pub materialize_cache_hit_count: LabelGuardedIntCounter,
1994 pub materialize_data_exist_count: LabelGuardedIntCounter,
1995 pub materialize_cache_total_count: LabelGuardedIntCounter,
1996}
1997
1998pub struct MaterializeMetrics {
1999 pub materialize_input_row_count: LabelGuardedIntCounter,
2000 pub materialize_current_epoch: LabelGuardedIntGauge,
2001}
2002
2003pub struct GroupTopNMetrics {
2004 pub group_top_n_cache_miss_count: LabelGuardedIntCounter,
2005 pub group_top_n_total_query_cache_count: LabelGuardedIntCounter,
2006 pub group_top_n_cached_entry_count: LabelGuardedIntGauge,
2007}
2008
2009pub struct LookupExecutorMetrics {
2010 pub lookup_cache_miss_count: LabelGuardedIntCounter,
2011 pub lookup_total_query_cache_count: LabelGuardedIntCounter,
2012 pub lookup_cached_entry_count: LabelGuardedIntGauge,
2013}
2014
2015pub struct HashAggMetrics {
2016 pub agg_lookup_miss_count: LabelGuardedIntCounter,
2017 pub agg_total_lookup_count: LabelGuardedIntCounter,
2018 pub agg_cached_entry_count: LabelGuardedIntGauge,
2019 pub agg_chunk_lookup_miss_count: LabelGuardedIntCounter,
2020 pub agg_chunk_total_lookup_count: LabelGuardedIntCounter,
2021 pub agg_dirty_groups_count: LabelGuardedIntGauge,
2022 pub agg_dirty_groups_heap_size: LabelGuardedIntGauge,
2023 pub agg_state_cache_lookup_count: LabelGuardedIntCounter,
2024 pub agg_state_cache_miss_count: LabelGuardedIntCounter,
2025}
2026
2027pub struct AggDistinctDedupMetrics {
2028 pub agg_distinct_cache_miss_count: LabelGuardedIntCounter,
2029 pub agg_distinct_total_cache_count: LabelGuardedIntCounter,
2030 pub agg_distinct_cached_entry_count: LabelGuardedIntGauge,
2031}
2032
2033pub struct TemporalJoinMetrics {
2034 pub temporal_join_cache_miss_count: LabelGuardedIntCounter,
2035 pub temporal_join_total_query_cache_count: LabelGuardedIntCounter,
2036 pub temporal_join_cached_entry_count: LabelGuardedIntGauge,
2037}
2038
2039pub struct BackfillMetrics {
2040 pub backfill_snapshot_read_row_count: LabelGuardedIntCounter,
2041 pub backfill_upstream_output_row_count: LabelGuardedIntCounter,
2042}
2043
2044#[derive(Clone)]
2045pub struct CdcBackfillMetrics {
2046 pub cdc_backfill_snapshot_read_row_count: LabelGuardedIntCounter,
2047 pub cdc_backfill_upstream_output_row_count: LabelGuardedIntCounter,
2048}
2049
2050pub struct OverWindowMetrics {
2051 pub over_window_cached_entry_count: LabelGuardedIntGauge,
2052 pub over_window_cache_lookup_count: LabelGuardedIntCounter,
2053 pub over_window_cache_miss_count: LabelGuardedIntCounter,
2054 pub over_window_range_cache_entry_count: LabelGuardedIntGauge,
2055 pub over_window_range_cache_lookup_count: LabelGuardedIntCounter,
2056 pub over_window_range_cache_left_miss_count: LabelGuardedIntCounter,
2057 pub over_window_range_cache_right_miss_count: LabelGuardedIntCounter,
2058 pub over_window_accessed_entry_count: LabelGuardedIntCounter,
2059 pub over_window_compute_count: LabelGuardedIntCounter,
2060 pub over_window_same_output_count: LabelGuardedIntCounter,
2061 pub over_window_state_cleaned_row_count: LabelGuardedIntCounter,
2062}
2063
2064pub struct MatchRecognizeMetrics {
2065 pub match_recognize_matches_emitted_count: LabelGuardedIntCounter,
2066 pub match_recognize_evicted_rows_count: LabelGuardedIntCounter,
2067 pub match_recognize_scan_budget_exhausted_count: LabelGuardedIntCounter,
2068 pub match_recognize_within_deadline_overflow_count: LabelGuardedIntCounter,
2069 pub match_recognize_retained_rows: LabelGuardedIntGauge,
2073}
2074
2075#[derive(Clone)]
2076pub struct StateTableMetrics {
2077 pub iter_count: LabelGuardedIntCounter,
2078 pub get_count: LabelGuardedIntCounter,
2079 pub iter_vnode_pruned_count: LabelGuardedIntCounter,
2080 pub get_vnode_pruned_count: LabelGuardedIntCounter,
2081}