Skip to main content

meta_srv/handler/
persist_stats_handler.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::time::{Duration, Instant};
16
17use api::v1::meta::{HeartbeatRequest, Role};
18use api::v1::value::ValueData;
19use api::v1::{ColumnSchema, Row, RowInsertRequest, RowInsertRequests, Rows, Value};
20use client::DEFAULT_CATALOG_NAME;
21use client::inserter::{Context as InserterContext, Inserter};
22use common_catalog::consts::DEFAULT_PRIVATE_SCHEMA_NAME;
23use common_macro::{Schema, ToRow};
24use common_meta::DatanodeId;
25use common_meta::datanode::{REGION_STATS_HISTORY_TABLE_NAME, RegionStat};
26use common_telemetry::warn;
27use dashmap::DashMap;
28use store_api::region_engine::RegionRole;
29use store_api::storage::RegionId;
30
31use crate::error::Result;
32use crate::handler::{HandleControl, HeartbeatAccumulator, HeartbeatHandler};
33use crate::metasrv::Context;
34
35/// The handler to persist stats.
36pub struct PersistStatsHandler {
37    inserter: Box<dyn Inserter>,
38    last_persisted_region_stats: DashMap<RegionId, PersistedRegionStat>,
39    last_persisted_time: DashMap<DatanodeId, Instant>,
40    persist_interval: Duration,
41}
42
43/// The default context to persist region stats.
44const DEFAULT_CONTEXT: InserterContext = InserterContext {
45    catalog: DEFAULT_CATALOG_NAME,
46    schema: DEFAULT_PRIVATE_SCHEMA_NAME,
47};
48
49#[derive(Debug, Clone, Copy, PartialEq, Eq)]
50struct PersistedRegionStat {
51    region_id: RegionId,
52    written_bytes: u64,
53}
54
55impl From<&RegionStat> for PersistedRegionStat {
56    fn from(stat: &RegionStat) -> Self {
57        Self {
58            region_id: stat.id,
59            written_bytes: stat.written_bytes,
60        }
61    }
62}
63
64#[derive(ToRow, Schema)]
65struct PersistRegionStat<'a> {
66    table_id: u32,
67    region_id: u64,
68    region_number: u32,
69    manifest_size: u64,
70    datanode_id: u64,
71    #[col(datatype = "string")]
72    engine: &'a str,
73    num_rows: u64,
74    sst_num: u64,
75    sst_size: u64,
76    write_bytes_delta: u64,
77    #[col(
78        // This col name is for the information schema table, so we don't touch it
79        name = "greptime_timestamp",
80        semantic = "Timestamp",
81        datatype = "TimestampMillisecond"
82    )]
83    timestamp_millis: i64,
84}
85
86/// Compute the region stat to persist.
87fn compute_persist_region_stat(
88    region_stat: &RegionStat,
89    datanode_id: DatanodeId,
90    timestamp_millis: i64,
91    persisted_region_stat: Option<PersistedRegionStat>,
92) -> PersistRegionStat<'_> {
93    let write_bytes_delta = persisted_region_stat
94        .and_then(|persisted_region_stat| {
95            region_stat
96                .written_bytes
97                .checked_sub(persisted_region_stat.written_bytes)
98        })
99        .unwrap_or_default();
100
101    PersistRegionStat {
102        table_id: region_stat.id.table_id(),
103        region_id: region_stat.id.as_u64(),
104        region_number: region_stat.id.region_number(),
105        manifest_size: region_stat.manifest_size,
106        datanode_id,
107        engine: region_stat.engine.as_str(),
108        num_rows: region_stat.num_rows,
109        sst_num: region_stat.sst_num,
110        sst_size: region_stat.sst_size,
111        write_bytes_delta,
112        timestamp_millis,
113    }
114}
115
116fn to_persisted_if_leader(
117    region_stat: &RegionStat,
118    last_persisted_region_stats: &DashMap<RegionId, PersistedRegionStat>,
119    datanode_id: DatanodeId,
120    timestamp_millis: i64,
121) -> Option<(Row, PersistedRegionStat)> {
122    if matches!(
123        region_stat.role,
124        RegionRole::Leader | RegionRole::StagingLeader
125    ) {
126        let persisted_region_stat = last_persisted_region_stats.get(&region_stat.id).map(|s| *s);
127        Some((
128            compute_persist_region_stat(
129                region_stat,
130                datanode_id,
131                timestamp_millis,
132                persisted_region_stat,
133            )
134            .to_row(),
135            PersistedRegionStat::from(region_stat),
136        ))
137    } else {
138        None
139    }
140}
141
142/// Align the timestamp to the nearest interval.
143///
144/// # Panics
145/// Panics if `interval` as milliseconds is zero.
146fn align_ts(ts: i64, interval: Duration) -> i64 {
147    assert!(
148        interval.as_millis() != 0,
149        "interval must be greater than zero"
150    );
151    ts / interval.as_millis() as i64 * interval.as_millis() as i64
152}
153
154impl PersistStatsHandler {
155    /// Creates a new [`PersistStatsHandler`].
156    pub fn new(inserter: Box<dyn Inserter>, mut persist_interval: Duration) -> Self {
157        if persist_interval < Duration::from_secs(10 * 60) {
158            warn!("persist_interval is less than 10 minutes, set to 10 minutes");
159            persist_interval = Duration::from_secs(10 * 60);
160        }
161
162        Self {
163            inserter,
164            last_persisted_region_stats: DashMap::new(),
165            last_persisted_time: DashMap::new(),
166            persist_interval,
167        }
168    }
169
170    fn should_persist(&self, datanode_id: DatanodeId) -> bool {
171        let Some(last_persisted_time) = self.last_persisted_time.get(&datanode_id) else {
172            return true;
173        };
174
175        last_persisted_time.elapsed() >= self.persist_interval
176    }
177
178    async fn persist(
179        &self,
180        timestamp_millis: i64,
181        datanode_id: DatanodeId,
182        region_stats: &[RegionStat],
183    ) {
184        // Safety: persist_interval is greater than zero.
185        let aligned_ts = align_ts(timestamp_millis, self.persist_interval);
186        let (rows, incoming_region_stats): (Vec<_>, Vec<_>) = region_stats
187            .iter()
188            .flat_map(|region_stat| {
189                to_persisted_if_leader(
190                    region_stat,
191                    &self.last_persisted_region_stats,
192                    datanode_id,
193                    aligned_ts,
194                )
195            })
196            .unzip();
197
198        if rows.is_empty() {
199            return;
200        }
201
202        if let Err(err) = self
203            .inserter
204            .insert_rows(
205                &DEFAULT_CONTEXT,
206                RowInsertRequests {
207                    inserts: vec![RowInsertRequest {
208                        table_name: REGION_STATS_HISTORY_TABLE_NAME.to_string(),
209                        rows: Some(Rows {
210                            schema: PersistRegionStat::schema(),
211                            rows,
212                        }),
213                    }],
214                },
215            )
216            .await
217        {
218            warn!(
219                "Failed to persist region stats, datanode_id: {}, error: {:?}",
220                datanode_id, err
221            );
222            return;
223        }
224
225        self.last_persisted_time.insert(datanode_id, Instant::now());
226        for s in incoming_region_stats {
227            self.last_persisted_region_stats.insert(s.region_id, s);
228        }
229    }
230}
231
232#[async_trait::async_trait]
233impl HeartbeatHandler for PersistStatsHandler {
234    fn is_acceptable(&self, role: Role) -> bool {
235        role == Role::Datanode
236    }
237
238    async fn handle(
239        &self,
240        _req: &HeartbeatRequest,
241        _: &mut Context,
242        acc: &mut HeartbeatAccumulator,
243    ) -> Result<HandleControl> {
244        let Some(current_stat) = acc.stat.as_ref() else {
245            return Ok(HandleControl::Continue);
246        };
247
248        if !self.should_persist(current_stat.id) {
249            return Ok(HandleControl::Continue);
250        }
251
252        self.persist(
253            current_stat.timestamp_millis,
254            current_stat.id,
255            &current_stat.region_stats,
256        )
257        .await;
258
259        Ok(HandleControl::Continue)
260    }
261}
262
263#[cfg(test)]
264mod tests {
265    use std::sync::{Arc, Mutex};
266
267    use client::inserter::{Context as InserterContext, InsertOptions};
268    use common_meta::datanode::{RegionManifestInfo, RegionStat, Stat};
269    use store_api::region_engine::RegionRole;
270    use store_api::storage::RegionId;
271
272    use super::*;
273    use crate::handler::test_utils::TestEnv;
274
275    fn create_test_region_stat(
276        table_id: u32,
277        region_number: u32,
278        written_bytes: u64,
279        engine: &str,
280    ) -> RegionStat {
281        let region_id = RegionId::new(table_id, region_number);
282        RegionStat {
283            id: region_id,
284            rcus: 100,
285            wcus: 200,
286            approximate_bytes: 1024,
287            engine: engine.to_string(),
288            role: RegionRole::Leader,
289            num_rows: 1000,
290            memtable_size: 512,
291            manifest_size: 256,
292            sst_size: 2048,
293            sst_num: 5,
294            index_size: 128,
295            region_manifest: RegionManifestInfo::Mito {
296                manifest_version: 1,
297                flushed_entry_id: 100,
298                file_removed_cnt: 0,
299            },
300            written_bytes,
301            query_cpu_time: 0,
302            query_scanned_bytes: 0,
303            data_topic_latest_entry_id: 200,
304            metadata_topic_latest_entry_id: 200,
305            min_timestamp: None,
306            max_timestamp: None,
307        }
308    }
309
310    #[test]
311    fn test_compute_persist_region_stat_with_no_persisted_stat() {
312        let region_stat = create_test_region_stat(1, 1, 1000, "mito");
313        let datanode_id = 123;
314        let timestamp_millis = 1640995200000; // 2022-01-01 00:00:00 UTC
315        let result = compute_persist_region_stat(&region_stat, datanode_id, timestamp_millis, None);
316        assert_eq!(result.table_id, 1);
317        assert_eq!(result.region_id, region_stat.id.as_u64());
318        assert_eq!(result.region_number, 1);
319        assert_eq!(result.manifest_size, 256);
320        assert_eq!(result.datanode_id, datanode_id);
321        assert_eq!(result.engine, "mito");
322        assert_eq!(result.num_rows, 1000);
323        assert_eq!(result.sst_num, 5);
324        assert_eq!(result.sst_size, 2048);
325        assert_eq!(result.write_bytes_delta, 0); // No previous stat, so delta is 0
326        assert_eq!(result.timestamp_millis, timestamp_millis);
327    }
328
329    #[test]
330    fn test_compute_persist_region_stat_with_persisted_stat_increase() {
331        let region_stat = create_test_region_stat(2, 3, 1500, "mito");
332        let datanode_id = 456;
333        let timestamp_millis = 1640995260000; // 2022-01-01 00:01:00 UTC
334        let persisted_stat = PersistedRegionStat {
335            region_id: region_stat.id,
336            written_bytes: 1000, // Previous write bytes
337        };
338        let result = compute_persist_region_stat(
339            &region_stat,
340            datanode_id,
341            timestamp_millis,
342            Some(persisted_stat),
343        );
344        assert_eq!(result.table_id, 2);
345        assert_eq!(result.region_id, region_stat.id.as_u64());
346        assert_eq!(result.region_number, 3);
347        assert_eq!(result.manifest_size, 256);
348        assert_eq!(result.datanode_id, datanode_id);
349        assert_eq!(result.engine, "mito");
350        assert_eq!(result.num_rows, 1000);
351        assert_eq!(result.sst_num, 5);
352        assert_eq!(result.sst_size, 2048);
353        assert_eq!(result.write_bytes_delta, 500); // 1500 - 1000 = 500
354        assert_eq!(result.timestamp_millis, timestamp_millis);
355    }
356
357    #[test]
358    fn test_compute_persist_region_stat_with_persisted_stat_decrease() {
359        let region_stat = create_test_region_stat(3, 5, 800, "mito");
360        let datanode_id = 789;
361        let timestamp_millis = 1640995320000; // 2022-01-01 00:02:00 UTC
362        let persisted_stat = PersistedRegionStat {
363            region_id: region_stat.id,
364            written_bytes: 1200, // Previous write bytes (higher than current)
365        };
366        let result = compute_persist_region_stat(
367            &region_stat,
368            datanode_id,
369            timestamp_millis,
370            Some(persisted_stat),
371        );
372        assert_eq!(result.table_id, 3);
373        assert_eq!(result.region_id, region_stat.id.as_u64());
374        assert_eq!(result.region_number, 5);
375        assert_eq!(result.manifest_size, 256);
376        assert_eq!(result.datanode_id, datanode_id);
377        assert_eq!(result.engine, "mito");
378        assert_eq!(result.num_rows, 1000);
379        assert_eq!(result.sst_num, 5);
380        assert_eq!(result.sst_size, 2048);
381        assert_eq!(result.write_bytes_delta, 0); // 800 - 1200 would be negative, so 0 due to checked_sub
382        assert_eq!(result.timestamp_millis, timestamp_millis);
383    }
384
385    #[test]
386    fn test_compute_persist_region_stat_with_persisted_stat_equal() {
387        let region_stat = create_test_region_stat(4, 7, 2000, "mito");
388        let datanode_id = 101;
389        let timestamp_millis = 1640995380000; // 2022-01-01 00:03:00 UTC
390        let persisted_stat = PersistedRegionStat {
391            region_id: region_stat.id,
392            written_bytes: 2000, // Same as current write bytes
393        };
394        let result = compute_persist_region_stat(
395            &region_stat,
396            datanode_id,
397            timestamp_millis,
398            Some(persisted_stat),
399        );
400        assert_eq!(result.table_id, 4);
401        assert_eq!(result.region_id, region_stat.id.as_u64());
402        assert_eq!(result.region_number, 7);
403        assert_eq!(result.manifest_size, 256);
404        assert_eq!(result.datanode_id, datanode_id);
405        assert_eq!(result.engine, "mito");
406        assert_eq!(result.num_rows, 1000);
407        assert_eq!(result.sst_num, 5);
408        assert_eq!(result.sst_size, 2048);
409        assert_eq!(result.write_bytes_delta, 0); // 2000 - 2000 = 0
410        assert_eq!(result.timestamp_millis, timestamp_millis);
411    }
412
413    #[test]
414    fn test_compute_persist_region_stat_with_overflow_protection() {
415        let region_stat = create_test_region_stat(8, 15, 500, "mito");
416        let datanode_id = 505;
417        let timestamp_millis = 1640995620000; // 2022-01-01 00:07:00 UTC
418        let persisted_stat = PersistedRegionStat {
419            region_id: region_stat.id,
420            written_bytes: 1000, // Higher than current, would cause underflow
421        };
422        let result = compute_persist_region_stat(
423            &region_stat,
424            datanode_id,
425            timestamp_millis,
426            Some(persisted_stat),
427        );
428        assert_eq!(result.table_id, 8);
429        assert_eq!(result.region_id, region_stat.id.as_u64());
430        assert_eq!(result.region_number, 15);
431        assert_eq!(result.manifest_size, 256);
432        assert_eq!(result.datanode_id, datanode_id);
433        assert_eq!(result.engine, "mito");
434        assert_eq!(result.num_rows, 1000);
435        assert_eq!(result.sst_num, 5);
436        assert_eq!(result.sst_size, 2048);
437        assert_eq!(result.write_bytes_delta, 0); // checked_sub returns None, so default to 0
438        assert_eq!(result.timestamp_millis, timestamp_millis);
439    }
440
441    struct MockInserter {
442        requests: Arc<Mutex<Vec<api::v1::RowInsertRequest>>>,
443    }
444
445    #[async_trait::async_trait]
446    impl Inserter for MockInserter {
447        async fn insert_rows(
448            &self,
449            _context: &InserterContext<'_>,
450            requests: api::v1::RowInsertRequests,
451        ) -> client::error::Result<()> {
452            self.requests.lock().unwrap().extend(requests.inserts);
453
454            Ok(())
455        }
456
457        fn set_options(&mut self, _options: &InsertOptions) {}
458    }
459
460    #[tokio::test]
461    async fn test_not_persist_region_stats() {
462        let env = TestEnv::new();
463        let mut ctx = env.ctx();
464
465        let requests = Arc::new(Mutex::new(vec![]));
466        let inserter = MockInserter {
467            requests: requests.clone(),
468        };
469        let handler = PersistStatsHandler::new(Box::new(inserter), Duration::from_secs(10));
470        let mut acc = HeartbeatAccumulator {
471            stat: Some(Stat {
472                id: 1,
473                timestamp_millis: 1640995200000,
474                region_stats: vec![create_test_region_stat(1, 1, 1000, "mito")],
475                ..Default::default()
476            }),
477            ..Default::default()
478        };
479        handler.last_persisted_time.insert(1, Instant::now());
480        // Do not persist
481        handler
482            .handle(&HeartbeatRequest::default(), &mut ctx, &mut acc)
483            .await
484            .unwrap();
485        assert!(requests.lock().unwrap().is_empty());
486    }
487
488    #[tokio::test]
489    async fn test_persist_region_stats() {
490        let env = TestEnv::new();
491        let mut ctx = env.ctx();
492        let requests = Arc::new(Mutex::new(vec![]));
493        let inserter = MockInserter {
494            requests: requests.clone(),
495        };
496
497        let handler = PersistStatsHandler::new(Box::new(inserter), Duration::from_secs(10));
498
499        let region_stat = create_test_region_stat(1, 1, 1000, "mito");
500        let timestamp_millis = 1640995200000;
501        let datanode_id = 1;
502        let region_id = RegionId::new(1, 1);
503        let mut acc = HeartbeatAccumulator {
504            stat: Some(Stat {
505                id: datanode_id,
506                timestamp_millis,
507                region_stats: vec![region_stat.clone()],
508                ..Default::default()
509            }),
510            ..Default::default()
511        };
512
513        handler.last_persisted_region_stats.insert(
514            region_id,
515            PersistedRegionStat {
516                region_id,
517                written_bytes: 500,
518            },
519        );
520        let (expected_row, expected_persisted_region_stat) = to_persisted_if_leader(
521            &region_stat,
522            &handler.last_persisted_region_stats,
523            datanode_id,
524            timestamp_millis,
525        )
526        .unwrap();
527        let before_insert_time = Instant::now();
528        // Persist
529        handler
530            .handle(&HeartbeatRequest::default(), &mut ctx, &mut acc)
531            .await
532            .unwrap();
533        let request = {
534            let mut requests = requests.lock().unwrap();
535            assert_eq!(requests.len(), 1);
536            requests.pop().unwrap()
537        };
538        assert_eq!(
539            request.table_name,
540            REGION_STATS_HISTORY_TABLE_NAME.to_string()
541        );
542        assert_eq!(request.rows.unwrap().rows, vec![expected_row]);
543
544        // Check last persisted time
545        assert!(
546            handler
547                .last_persisted_time
548                .get(&datanode_id)
549                .unwrap()
550                .gt(&before_insert_time)
551        );
552
553        // Check last persisted region stats
554        assert_eq!(
555            handler
556                .last_persisted_region_stats
557                .get(&region_id)
558                .unwrap()
559                .value(),
560            &expected_persisted_region_stat
561        );
562    }
563}