Skip to main content

risingwave_stream/executor/monitor/
streaming_stats.rs

1// Copyright 2022 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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    // Executor metrics (disabled by default)
49    pub executor_row_count: RelabeledGuardedIntCounterVec,
50
51    // Profiling Metrics:
52    // Aggregated per operator rather than per actor.
53    // These are purely in-memory, never collected by prometheus.
54    pub mem_stream_node_output_row_count: CountMap<ExecutorId>,
55    pub mem_stream_node_output_blocking_duration_ns: CountMap<ExecutorId>,
56
57    // Streaming actor metrics from tokio (disabled by default)
58    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    // Streaming actor
66    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    // Source
74    pub source_output_row_count: LabelGuardedIntCounterVec,
75    pub source_split_change_count: LabelGuardedIntCounterVec,
76    pub source_backfill_row_count: LabelGuardedIntCounterVec,
77
78    // Sink
79    sink_input_row_count: LabelGuardedIntCounterVec,
80    sink_input_bytes: LabelGuardedIntCounterVec,
81    sink_chunk_buffer_size: LabelGuardedIntGaugeVec,
82
83    // Exchange (see also `compute::ExchangeServiceMetrics`)
84    pub exchange_frag_recv_size: LabelGuardedIntCounterVec,
85
86    // Streaming Merge (We breakout this metric from `barrier_align_duration` because
87    // the alignment happens on different levels)
88    pub merge_barrier_align_duration: RelabeledGuardedIntCounterVec,
89
90    // Backpressure
91    pub actor_output_buffer_blocking_duration_ns: RelabeledGuardedIntCounterVec,
92    actor_input_buffer_blocking_duration_ns: RelabeledGuardedIntCounterVec,
93
94    // Streaming Join
95    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    // Streaming Join, Streaming Dynamic Filter and Streaming Union
104    pub barrier_align_duration: RelabeledGuardedIntCounterVec,
105
106    // Streaming Aggregation
107    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    // Streaming TopN
121    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    // TODO(rc): why not just use the above three?
125    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 executor
130    lookup_cache_miss_count: LabelGuardedIntCounterVec,
131    lookup_total_query_cache_count: LabelGuardedIntCounterVec,
132    lookup_cached_entry_count: LabelGuardedIntGaugeVec,
133
134    // temporal join
135    temporal_join_cache_miss_count: LabelGuardedIntCounterVec,
136    temporal_join_total_query_cache_count: LabelGuardedIntCounterVec,
137    temporal_join_cached_entry_count: LabelGuardedIntGaugeVec,
138
139    // Backfill
140    backfill_snapshot_read_row_count: LabelGuardedIntCounterVec,
141    backfill_upstream_output_row_count: LabelGuardedIntCounterVec,
142
143    // CDC Backfill
144    cdc_backfill_snapshot_read_row_count: LabelGuardedIntCounterVec,
145    cdc_backfill_upstream_output_row_count: LabelGuardedIntCounterVec,
146
147    // Snapshot Backfill
148    pub(crate) snapshot_backfill_consume_row_count: LabelGuardedIntCounterVec,
149
150    // Over Window
151    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
164    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    /// The duration from receipt of barrier to all actors collection.
171    /// The max of all nodes' `barrier_inflight_latency` for a partial graph is the latency for a
172    /// barrier to flow through that partial graph.
173    pub barrier_inflight_latency: LabelGuardedHistogramVec,
174    /// The duration of sync to storage.
175    pub barrier_sync_latency: LabelGuardedHistogramVec,
176    pub barrier_batch_size: Histogram,
177    /// The progress made by the earliest in-flight barriers in the local barrier manager.
178    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    // Memory management
208    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    // Materialized view
221    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    // PostgreSQL CDC LSN monitoring
228    pub pg_cdc_state_table_lsn: LabelGuardedIntGaugeVec,
229    pub pg_cdc_jni_commit_offset_lsn: LabelGuardedIntGaugeVec,
230
231    // MySQL CDC binlog monitoring
232    pub mysql_cdc_state_binlog_file_seq: LabelGuardedIntGaugeVec,
233    pub mysql_cdc_state_binlog_position: LabelGuardedIntGaugeVec,
234
235    // SQL Server CDC LSN monitoring
236    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    // Gap Fill
241    pub gap_fill_generated_rows_count: RelabeledGuardedIntCounterVec,
242
243    // Now (temporal filter clock)
244    pub now_streaming_clock_ms: LabelGuardedIntGaugeVec,
245    pub now_wall_clock_drift_ms: LabelGuardedIntGaugeVec,
246
247    // State Table
248    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            // mask the first label `actor_id` if the level is less verbose than `Debug`
416            .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            // mask the first label `actor_id` if the level is less verbose than `Debug`
427            .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() // max 2^31
612        );
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    /// Create a new `StreamingMetrics` instance used in tests or other places.
1550    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
1976/// Tokio metrics for actors
1977pub 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    /// Rows currently retained in memory across the actor's partitions. Retention is bounded only
2070    /// by match liveness and `WITHIN`, so this gauge is the one signal of a partition set growing
2071    /// toward memory exhaustion (a pattern whose closer never arrives retains its rows forever).
2072    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}