servers/batcher/
flow_notifier.rs1use std::num::NonZeroUsize;
16
17use common_batcher::notifier::{Notifier, run_notifier};
18use common_meta::cache::TableFlownodeSetCacheRef;
19use common_meta::node_manager::NodeManagerRef;
20use common_runtime::spawn_global;
21use common_telemetry::{error, warn};
22use prometheus::IntCounterVec;
23use tokio::sync::mpsc::Receiver;
24use tokio::sync::mpsc::error::TrySendError;
25
26use crate::batcher::flow_sender::{FlowNotification, FlowSender};
27
28#[derive(Clone)]
30pub(in crate::batcher) struct FlowNotifier {
31 notifier: Notifier<FlowNotification>,
32 dropped: IntCounterVec,
33}
34
35impl FlowNotifier {
36 pub fn try_new(
38 capacity: usize,
39 dropped: IntCounterVec,
40 ) -> Option<(Self, Receiver<FlowNotification>)> {
41 let (notifier, receiver) = Notifier::try_new(capacity)?;
42 Some((Self { notifier, dropped }, receiver))
43 }
44
45 pub fn try_notify(&self, notification: FlowNotification) -> bool {
47 match self.notifier.try_notify(notification) {
48 Ok(()) => true,
49 Err(TrySendError::Full(notification)) => {
50 self.dropped.with_label_values(&["full"]).inc();
51 warn!(
52 "Dropping flow notification because queue is full, table_id: {}, queue_capacity: {}",
53 notification.table_id,
54 self.notifier.max_capacity()
55 );
56 false
57 }
58 Err(TrySendError::Closed(notification)) => {
59 self.dropped.with_label_values(&["closed"]).inc();
60 error!(
61 "Dropping flow notification because queue is closed, table_id: {}, queue_capacity: {}",
62 notification.table_id,
63 self.notifier.max_capacity()
64 );
65 false
66 }
67 }
68 }
69}
70
71pub(in crate::batcher) fn start_flow_notification_worker(
73 receiver: Receiver<FlowNotification>,
74 cache: TableFlownodeSetCacheRef,
75 node_manager: NodeManagerRef,
76) {
77 let sender = FlowSender::new(cache, node_manager);
78 let concurrency = NonZeroUsize::new(8).unwrap();
79 spawn_global(run_notifier(receiver, concurrency, move |notification| {
80 let sender = sender.clone();
81 async move { sender.send(notification).await }
82 }));
83}
84
85#[cfg(test)]
86mod tests {
87
88 use prometheus::{IntCounterVec, Opts};
89
90 use crate::batcher::flow_notifier::FlowNotifier;
91 use crate::batcher::flow_sender::FlowNotification;
92
93 fn dropped_counter() -> IntCounterVec {
94 IntCounterVec::new(
96 Opts::new("test_flow_dropped", "Dropped notifications"),
97 &["reason"],
98 )
99 .unwrap()
100 }
101
102 fn notification(table_id: u32) -> FlowNotification {
103 FlowNotification {
104 table_id,
105 timestamps: vec![-1, 42, 1_700_000_000_000_000_001, 42],
106 }
107 }
108
109 fn dropped_counts(counter: &IntCounterVec) -> (u64, u64) {
110 (
111 counter.with_label_values(&["full"]).get(),
112 counter.with_label_values(&["closed"]).get(),
113 )
114 }
115
116 #[test]
117 fn test_queues_and_metrics_are_independent() {
118 let first_dropped = dropped_counter();
119 let second_dropped = dropped_counter();
120 let (first, mut first_rx) = FlowNotifier::try_new(1, first_dropped.clone()).unwrap();
121 let (second, mut second_rx) = FlowNotifier::try_new(2, second_dropped.clone()).unwrap();
122
123 assert!(first.try_notify(notification(1)));
124 assert!(!first.try_notify(notification(2)));
125 assert_eq!(dropped_counts(&first_dropped), (1, 0));
126 assert_eq!(dropped_counts(&second_dropped), (0, 0));
127 assert!(second.try_notify(notification(3)));
128 assert!(second.try_notify(notification(4)));
129 assert!(!second.try_notify(notification(5)));
130 assert_eq!(dropped_counts(&second_dropped), (1, 0));
131
132 let payload = first_rx.try_recv().unwrap();
133 assert_eq!(payload.table_id, 1);
134 assert_eq!(payload.timestamps, notification(1).timestamps);
135 assert!(first_rx.try_recv().is_err());
136 for table_id in [3, 4] {
137 let payload = second_rx.try_recv().unwrap();
138 assert_eq!(payload.table_id, table_id);
139 assert_eq!(payload.timestamps, notification(table_id).timestamps);
140 }
141 assert!(second_rx.try_recv().is_err());
142
143 drop(first_rx);
144 assert!(!first.try_notify(notification(6)));
145 assert_eq!(dropped_counts(&first_dropped), (1, 1));
146 assert_eq!(dropped_counts(&second_dropped), (1, 0));
147 assert!(second.try_notify(notification(7)));
148 assert_eq!(second_rx.try_recv().unwrap().table_id, 7);
149 drop(second_rx);
150 assert!(!second.try_notify(notification(8)));
151 assert_eq!(dropped_counts(&second_dropped), (1, 1));
152 }
153
154 #[test]
155 fn test_clones_share_queue_and_drop_metrics() {
156 let dropped = dropped_counter();
157 let (notifier, mut receiver) = FlowNotifier::try_new(1, dropped.clone()).unwrap();
158 let clone = notifier.clone();
159 assert!(notifier.try_notify(notification(1)));
160 assert!(!clone.try_notify(notification(2)));
161 assert_eq!(dropped_counts(&dropped), (1, 0));
162 assert_eq!(receiver.try_recv().unwrap().table_id, 1);
163 drop(notifier);
164 assert!(clone.try_notify(notification(3)));
165 let payload = receiver.try_recv().unwrap();
166 assert_eq!(payload.table_id, 3);
167 assert_eq!(payload.timestamps, notification(3).timestamps);
168 drop(receiver);
169 assert!(!clone.try_notify(notification(4)));
170 assert_eq!(dropped_counts(&dropped), (1, 1));
171 }
172}