Skip to main content

servers/batcher/logical_table/
flow_notifier.rs

1// Copyright 2023 Greptime Team
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
15#[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}