Skip to main content

servers/batcher/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
15use std::num::NonZeroUsize;
16
17use arrow::record_batch::RecordBatch;
18use common_meta::cache::TableFlownodeSetCacheRef;
19use common_meta::node_manager::NodeManagerRef;
20use common_telemetry::error;
21use lazy_static::lazy_static;
22use operator::req_convert::insert::extract_timestamps;
23use prometheus::{IntCounterVec, register_int_counter_vec};
24use table::metadata::TableInfoRef;
25
26use crate::batcher::flow_notifier::{
27    FlowNotifier as NotificationQueue, start_flow_notification_worker,
28};
29use crate::batcher::flow_sender::FlowNotification;
30
31lazy_static! {
32    static ref FLOW_NOTIFICATION_DROPPED: IntCounterVec = register_int_counter_vec!(
33        "greptime_table_batcher_flow_notification_dropped_total",
34        "Ordinary batch flow notifications dropped before delivery",
35        &["reason"]
36    )
37    .unwrap();
38}
39
40/// Enqueues best-effort flow notifications independently of write completion.
41#[derive(Clone)]
42pub(in crate::batcher::table) struct FlowNotifier {
43    notifier: NotificationQueue,
44}
45
46impl FlowNotifier {
47    /// Starts bounded delivery on the shared runtime, or rejects unsupported capacity.
48    pub fn new(
49        cache: TableFlownodeSetCacheRef,
50        node_manager: NodeManagerRef,
51        capacity: NonZeroUsize,
52    ) -> Option<Self> {
53        let (notifier, receiver) =
54            NotificationQueue::try_new(capacity.get(), FLOW_NOTIFICATION_DROPPED.clone())?;
55        start_flow_notification_worker(receiver, cache, node_manager);
56        Some(Self { notifier })
57    }
58
59    /// Never waits for queue capacity or delivery and never changes the write result.
60    pub fn notify(&self, table_info: TableInfoRef, batch: &RecordBatch) {
61        let Some(timestamp_column) = table_info.meta.schema.timestamp_column() else {
62            return;
63        };
64        let timestamps = match extract_timestamps(batch, &timestamp_column.name) {
65            Ok(timestamps) => timestamps,
66            Err(error) => {
67                error!(error; "Failed to extract flow notification timestamps, table_id: {}", table_info.table_id());
68                return;
69            }
70        };
71        if timestamps.is_empty() {
72            return;
73        }
74        self.enqueue(FlowNotification {
75            table_id: table_info.table_id(),
76            timestamps,
77        });
78    }
79
80    fn enqueue(&self, notification: FlowNotification) {
81        self.notifier.try_notify(notification);
82    }
83}
84
85#[cfg(test)]
86mod tests {
87    use std::sync::Arc;
88
89    use api::helper::ColumnDataTypeWrapper;
90    use api::v1::ColumnDataType;
91    use arrow::array::{
92        ArrayRef, TimestampMicrosecondArray, TimestampMillisecondArray, TimestampNanosecondArray,
93        TimestampSecondArray,
94    };
95    use arrow::record_batch::RecordBatch;
96    use datatypes::data_type::ConcreteDataType;
97    use datatypes::schema::{ColumnSchema, Schema};
98    use operator::test_util::new_test_table_info;
99
100    use crate::batcher::flow_notifier::FlowNotifier as NotificationQueue;
101    use crate::batcher::table::flow_notifier::{
102        FLOW_NOTIFICATION_DROPPED, FlowNotification, FlowNotifier,
103    };
104
105    #[test]
106    fn test_full_and_closed_queues_do_not_block() {
107        let (notifier, receiver) =
108            NotificationQueue::try_new(1, FLOW_NOTIFICATION_DROPPED.clone()).unwrap();
109        let notifier = FlowNotifier { notifier };
110        let full = FLOW_NOTIFICATION_DROPPED.with_label_values(&["full"]);
111        let closed = FLOW_NOTIFICATION_DROPPED.with_label_values(&["closed"]);
112        let before_full = full.get();
113        let before_closed = closed.get();
114        for _ in 0..2 {
115            notifier.enqueue(FlowNotification {
116                table_id: 1,
117                timestamps: vec![42],
118            });
119        }
120        assert!(full.get() > before_full);
121        drop(receiver);
122        notifier.enqueue(FlowNotification {
123            table_id: 1,
124            timestamps: vec![42],
125        });
126        assert!(closed.get() > before_closed);
127    }
128    #[test]
129    fn test_notification_preserves_native_timestamp_units() {
130        let cases: Vec<(ColumnDataType, ArrayRef)> = vec![
131            (
132                ColumnDataType::TimestampSecond,
133                Arc::new(TimestampSecondArray::from(vec![Some(42), None])),
134            ),
135            (
136                ColumnDataType::TimestampMillisecond,
137                Arc::new(TimestampMillisecondArray::from(vec![Some(42), None])),
138            ),
139            (
140                ColumnDataType::TimestampMicrosecond,
141                Arc::new(TimestampMicrosecondArray::from(vec![Some(42), None])),
142            ),
143            (
144                ColumnDataType::TimestampNanosecond,
145                Arc::new(TimestampNanosecondArray::from(vec![Some(42), None])),
146            ),
147        ];
148        for (datatype, array) in cases {
149            let mut table = new_test_table_info(1, "test", [0].into_iter());
150            let column = ColumnSchema::new(
151                "ts",
152                ConcreteDataType::from(ColumnDataTypeWrapper::new(datatype, None)),
153                true,
154            )
155            .with_time_index(true);
156            table.meta.schema = Arc::new(Schema::new(vec![column]));
157            let batch = RecordBatch::try_new(table.meta.schema.arrow_schema().clone(), vec![array])
158                .unwrap();
159            let (notifier, mut receiver) =
160                NotificationQueue::try_new(1, FLOW_NOTIFICATION_DROPPED.clone()).unwrap();
161            FlowNotifier { notifier }.notify(Arc::new(table), &batch);
162            let notification = receiver.try_recv().unwrap();
163            assert_eq!(notification.table_id, 1);
164            assert_eq!(notification.timestamps, vec![42]);
165        }
166    }
167}