servers/batcher/table/
flow_notifier.rs1use 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#[derive(Clone)]
42pub(in crate::batcher::table) struct FlowNotifier {
43 notifier: NotificationQueue,
44}
45
46impl FlowNotifier {
47 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 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, ×tamp_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}