servers/batcher/logical_table/
flow_notifier.rs1#[cfg(test)]
16use common_meta::cache::TableFlownodeSetCacheRef;
17#[cfg(test)]
18use common_meta::node_manager::NodeManagerRef;
19use common_telemetry::error;
20use datatypes::timestamp::append_timestamps;
21
22use crate::batcher::flow_notifier::FlowNotifier;
23#[cfg(test)]
24use crate::batcher::flow_notifier::start_flow_notification_worker;
25use crate::batcher::flow_sender::FlowNotification;
26use crate::batcher::logical_table::batch_convert::TableBatch;
27#[cfg(test)]
28use crate::metrics::FLOW_NOTIFICATION_DROPPED;
29
30pub(in crate::batcher::logical_table) fn extract_timestamps(table_batch: &TableBatch) -> Vec<i64> {
31 let mut timestamps = Vec::with_capacity(table_batch.row_count);
32 for batch in &table_batch.batches {
33 let timestamp_column = batch.batch.column(batch.timestamp_index);
34 let Some(()) = append_timestamps(timestamp_column, &mut timestamps) else {
35 error!(
36 "Failed to extract timestamps from record batch, table_id: {}, timestamp_index: {}",
37 table_batch.table_id, batch.timestamp_index
38 );
39 continue;
40 };
41 }
42 timestamps
43}
44
45pub(in crate::batcher::logical_table) fn enqueue_flow_notifications(
46 table_batches: Vec<TableBatch>,
47 tx: &FlowNotifier,
48) {
49 for table_batch in table_batches {
50 let timestamps = extract_timestamps(&table_batch);
51 if timestamps.is_empty() {
52 continue;
53 }
54 tx.try_notify(FlowNotification {
55 table_id: table_batch.table_id,
56 timestamps,
57 });
58 }
59}
60
61#[cfg(test)]
62pub(in crate::batcher::logical_table) fn notify_flow_dirty_windows_after_flush(
63 table_batches: Vec<TableBatch>,
64 table_flownode_set_cache: TableFlownodeSetCacheRef,
65 node_manager: NodeManagerRef,
66) {
67 let (tx, rx) = FlowNotifier::try_new(
68 table_batches.len().max(1),
69 FLOW_NOTIFICATION_DROPPED.clone(),
70 )
71 .unwrap();
72 start_flow_notification_worker(rx, table_flownode_set_cache, node_manager);
73 enqueue_flow_notifications(table_batches, &tx);
74}
75
76#[cfg(test)]
77mod tests {
78 use std::any::Any;
79 use std::sync::{Arc, Mutex};
80 use std::time::Duration;
81
82 use api::v1::meta::Peer;
83 use async_trait::async_trait;
84 use common_meta::cache::new_table_flownode_set_cache;
85 use common_meta::error::Result as MetaResult;
86 use common_meta::instruction::{CacheIdent, CreateFlow};
87 use common_meta::kv_backend::{KvBackend, TxnService};
88 use common_meta::node_manager::NodeManagerRef;
89 use common_meta::rpc::store::{
90 BatchDeleteRequest, BatchDeleteResponse, BatchGetRequest, BatchGetResponse,
91 BatchPutRequest, BatchPutResponse, DeleteRangeRequest, DeleteRangeResponse, PutRequest,
92 PutResponse, RangeRequest, RangeResponse,
93 };
94 use moka::future::CacheBuilder;
95 use tokio::sync::{Notify, mpsc, oneshot};
96
97 use crate::batcher::flow_notifier::FlowNotifier;
98 use crate::batcher::flow_sender::FlowNotification;
99 use crate::batcher::logical_table::batch_convert::TableBatch;
100 use crate::batcher::logical_table::flow_notifier::notify_flow_dirty_windows_after_flush;
101 use crate::batcher::logical_table::test_util::{
102 FlowNotificationMockNodeManager, RecordingFlownode, mock_aligned_tag_batch,
103 };
104 use crate::batcher::test_util::mock_table_flownode_cache;
105 use crate::metrics::FLOW_NOTIFICATION_DROPPED;
106
107 #[test]
108 fn test_flow_notification_queue_drops_when_full() {
109 let (tx, mut rx) = FlowNotifier::try_new(1, FLOW_NOTIFICATION_DROPPED.clone()).unwrap();
110 let notification = |table_id| FlowNotification {
111 table_id,
112 timestamps: vec![table_id as i64],
113 };
114 let dropped = FLOW_NOTIFICATION_DROPPED.with_label_values(&["full"]);
115 let dropped_before = dropped.get();
116
117 assert!(tx.try_notify(notification(1)));
118 assert!(!tx.try_notify(notification(2)));
119
120 assert_eq!(1, rx.try_recv().unwrap().table_id);
121 assert_eq!(dropped_before + 1, dropped.get());
122
123 let closed = FLOW_NOTIFICATION_DROPPED.with_label_values(&["closed"]);
124 let closed_before = closed.get();
125 drop(rx);
126 assert!(!tx.try_notify(notification(3)));
127 assert_eq!(closed_before + 1, closed.get());
128 }
129
130 #[tokio::test]
131 async fn test_flow_notifications_do_not_block_on_previous_table_cache_lookup() {
132 let blocked_table_id = 41;
133 let cached_table_id = 42;
134 let peer = Peer {
135 id: 7,
136 addr: "flow-7".to_string(),
137 };
138 let (range_started_tx, range_started_rx) = oneshot::channel();
139 let range_release = Arc::new(Notify::new());
140 let cache = Arc::new(new_table_flownode_set_cache(
141 "test".to_string(),
142 CacheBuilder::new(2).build(),
143 Arc::new(BlockingRangeKvBackend {
144 range_started: Mutex::new(Some(range_started_tx)),
145 range_release: range_release.clone(),
146 }),
147 ));
148 cache
149 .invalidate(&[CacheIdent::CreateFlow(CreateFlow {
150 flow_id: 1,
151 source_table_ids: vec![cached_table_id],
152 partition_to_peer_mapping: vec![(0, peer)],
153 })])
154 .await
155 .unwrap();
156 let (requests_tx, mut requests_rx) = mpsc::unbounded_channel();
157 let _requests_tx = requests_tx.clone();
158 let node_manager: NodeManagerRef = Arc::new(FlowNotificationMockNodeManager {
159 flownode: Arc::new(RecordingFlownode { requests_tx }),
160 });
161 let table_batches = vec![
162 TableBatch {
163 table_name: "blocked".to_string(),
164 table_id: blocked_table_id,
165 batches: vec![mock_aligned_tag_batch("tag1", "host-1", 1000, 1.0)],
166 row_count: 1,
167 },
168 TableBatch {
169 table_name: "cached".to_string(),
170 table_id: cached_table_id,
171 batches: vec![mock_aligned_tag_batch("tag1", "host-2", 2000, 2.0)],
172 row_count: 1,
173 },
174 ];
175
176 notify_flow_dirty_windows_after_flush(table_batches, cache, node_manager);
177
178 tokio::time::timeout(Duration::from_secs(1), range_started_rx)
179 .await
180 .unwrap()
181 .unwrap();
182 let requests = tokio::time::timeout(Duration::from_secs(1), requests_rx.recv())
183 .await
184 .unwrap()
185 .unwrap();
186 assert_eq!(cached_table_id, requests.requests[0].table_id);
187 range_release.notify_one();
188 }
189
190 #[tokio::test]
191 async fn test_successful_flush_notifies_flownode_with_logical_table_timestamps() {
192 let table_id = 42;
193 let peer = Peer {
194 id: 7,
195 addr: "flow-7".to_string(),
196 };
197 let cache = mock_table_flownode_cache(table_id, vec![(0, peer.clone()), (1, peer)]).await;
198 let (requests_tx, mut requests_rx) = mpsc::unbounded_channel();
199 let _requests_tx = requests_tx.clone();
200 let node_manager: NodeManagerRef = Arc::new(FlowNotificationMockNodeManager {
201 flownode: Arc::new(RecordingFlownode { requests_tx }),
202 });
203 let table_batches = vec![TableBatch {
204 table_name: "cpu".to_string(),
205 table_id,
206 batches: vec![mock_aligned_tag_batch("tag1", "host-1", 1000, 1.0)],
207 row_count: 1,
208 }];
209
210 notify_flow_dirty_windows_after_flush(table_batches, cache, node_manager);
211
212 let requests = tokio::time::timeout(Duration::from_secs(1), requests_rx.recv())
213 .await
214 .unwrap()
215 .unwrap();
216 assert_eq!(
217 vec![api::v1::flow::DirtyWindowRequest {
218 table_id,
219 timestamps: vec![1000],
220 time_ranges: Vec::new(),
221 }],
222 requests.requests
223 );
224 assert!(
225 tokio::time::timeout(Duration::from_millis(50), requests_rx.recv())
226 .await
227 .is_err()
228 );
229 }
230
231 #[tokio::test]
232 async fn test_successful_flush_coalesces_logical_batches_per_flownode() {
233 let table_id = 42;
234 let peer = Peer {
235 id: 7,
236 addr: "flow-7".to_string(),
237 };
238 let cache = mock_table_flownode_cache(table_id, vec![(0, peer.clone()), (1, peer)]).await;
239 let (requests_tx, mut requests_rx) = mpsc::unbounded_channel();
240 let _requests_tx = requests_tx.clone();
241 let node_manager: NodeManagerRef = Arc::new(FlowNotificationMockNodeManager {
242 flownode: Arc::new(RecordingFlownode { requests_tx }),
243 });
244 let table_batches = vec![TableBatch {
245 table_name: "cpu".to_string(),
246 table_id,
247 batches: vec![
248 mock_aligned_tag_batch("tag1", "host-1", 1000, 1.0),
249 mock_aligned_tag_batch("tag1", "host-1", 2000, 2.0),
250 ],
251 row_count: 2,
252 }];
253
254 notify_flow_dirty_windows_after_flush(table_batches, cache, node_manager);
255
256 let requests = tokio::time::timeout(Duration::from_secs(1), requests_rx.recv())
257 .await
258 .unwrap()
259 .unwrap();
260 assert_eq!(
261 vec![api::v1::flow::DirtyWindowRequest {
262 table_id,
263 timestamps: vec![1000, 2000],
264 time_ranges: Vec::new(),
265 }],
266 requests.requests
267 );
268 assert!(
269 tokio::time::timeout(Duration::from_millis(50), requests_rx.recv())
270 .await
271 .is_err()
272 );
273 }
274
275 struct BlockingRangeKvBackend {
276 range_started: Mutex<Option<oneshot::Sender<()>>>,
277 range_release: Arc<Notify>,
278 }
279
280 impl TxnService for BlockingRangeKvBackend {
281 type Error = common_meta::error::Error;
282 }
283
284 #[async_trait]
285 impl KvBackend for BlockingRangeKvBackend {
286 fn name(&self) -> &str {
287 "blocking_range"
288 }
289
290 fn as_any(&self) -> &dyn Any {
291 self
292 }
293
294 async fn range(&self, _req: RangeRequest) -> MetaResult<RangeResponse> {
295 let range_started = self.range_started.lock().unwrap().take();
296 if let Some(range_started) = range_started {
297 let _ = range_started.send(());
298 self.range_release.notified().await;
299 }
300 Ok(RangeResponse {
301 kvs: Vec::new(),
302 more: false,
303 })
304 }
305
306 async fn put(&self, _req: PutRequest) -> MetaResult<PutResponse> {
307 unimplemented!()
308 }
309
310 async fn batch_put(&self, _req: BatchPutRequest) -> MetaResult<BatchPutResponse> {
311 unimplemented!()
312 }
313
314 async fn batch_get(&self, _req: BatchGetRequest) -> MetaResult<BatchGetResponse> {
315 unimplemented!()
316 }
317
318 async fn delete_range(&self, _req: DeleteRangeRequest) -> MetaResult<DeleteRangeResponse> {
319 unimplemented!()
320 }
321
322 async fn batch_delete(&self, _req: BatchDeleteRequest) -> MetaResult<BatchDeleteResponse> {
323 unimplemented!()
324 }
325 }
326}