1use 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
35pub 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
43const 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 name = "greptime_timestamp",
80 semantic = "Timestamp",
81 datatype = "TimestampMillisecond"
82 )]
83 timestamp_millis: i64,
84}
85
86fn 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(®ion_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
142fn 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 pub fn new(inserter: Box<dyn Inserter>, mut persist_interval: Duration) -> Self {
157 if persist_interval < Duration::from_mins(10) {
158 warn!("persist_interval is less than 10 minutes, set to 10 minutes");
159 persist_interval = Duration::from_mins(10);
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 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 ¤t_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 }
306 }
307
308 #[test]
309 fn test_compute_persist_region_stat_with_no_persisted_stat() {
310 let region_stat = create_test_region_stat(1, 1, 1000, "mito");
311 let datanode_id = 123;
312 let timestamp_millis = 1640995200000; let result = compute_persist_region_stat(®ion_stat, datanode_id, timestamp_millis, None);
314 assert_eq!(result.table_id, 1);
315 assert_eq!(result.region_id, region_stat.id.as_u64());
316 assert_eq!(result.region_number, 1);
317 assert_eq!(result.manifest_size, 256);
318 assert_eq!(result.datanode_id, datanode_id);
319 assert_eq!(result.engine, "mito");
320 assert_eq!(result.num_rows, 1000);
321 assert_eq!(result.sst_num, 5);
322 assert_eq!(result.sst_size, 2048);
323 assert_eq!(result.write_bytes_delta, 0); assert_eq!(result.timestamp_millis, timestamp_millis);
325 }
326
327 #[test]
328 fn test_compute_persist_region_stat_with_persisted_stat_increase() {
329 let region_stat = create_test_region_stat(2, 3, 1500, "mito");
330 let datanode_id = 456;
331 let timestamp_millis = 1640995260000; let persisted_stat = PersistedRegionStat {
333 region_id: region_stat.id,
334 written_bytes: 1000, };
336 let result = compute_persist_region_stat(
337 ®ion_stat,
338 datanode_id,
339 timestamp_millis,
340 Some(persisted_stat),
341 );
342 assert_eq!(result.table_id, 2);
343 assert_eq!(result.region_id, region_stat.id.as_u64());
344 assert_eq!(result.region_number, 3);
345 assert_eq!(result.manifest_size, 256);
346 assert_eq!(result.datanode_id, datanode_id);
347 assert_eq!(result.engine, "mito");
348 assert_eq!(result.num_rows, 1000);
349 assert_eq!(result.sst_num, 5);
350 assert_eq!(result.sst_size, 2048);
351 assert_eq!(result.write_bytes_delta, 500); assert_eq!(result.timestamp_millis, timestamp_millis);
353 }
354
355 #[test]
356 fn test_compute_persist_region_stat_with_persisted_stat_decrease() {
357 let region_stat = create_test_region_stat(3, 5, 800, "mito");
358 let datanode_id = 789;
359 let timestamp_millis = 1640995320000; let persisted_stat = PersistedRegionStat {
361 region_id: region_stat.id,
362 written_bytes: 1200, };
364 let result = compute_persist_region_stat(
365 ®ion_stat,
366 datanode_id,
367 timestamp_millis,
368 Some(persisted_stat),
369 );
370 assert_eq!(result.table_id, 3);
371 assert_eq!(result.region_id, region_stat.id.as_u64());
372 assert_eq!(result.region_number, 5);
373 assert_eq!(result.manifest_size, 256);
374 assert_eq!(result.datanode_id, datanode_id);
375 assert_eq!(result.engine, "mito");
376 assert_eq!(result.num_rows, 1000);
377 assert_eq!(result.sst_num, 5);
378 assert_eq!(result.sst_size, 2048);
379 assert_eq!(result.write_bytes_delta, 0); assert_eq!(result.timestamp_millis, timestamp_millis);
381 }
382
383 #[test]
384 fn test_compute_persist_region_stat_with_persisted_stat_equal() {
385 let region_stat = create_test_region_stat(4, 7, 2000, "mito");
386 let datanode_id = 101;
387 let timestamp_millis = 1640995380000; let persisted_stat = PersistedRegionStat {
389 region_id: region_stat.id,
390 written_bytes: 2000, };
392 let result = compute_persist_region_stat(
393 ®ion_stat,
394 datanode_id,
395 timestamp_millis,
396 Some(persisted_stat),
397 );
398 assert_eq!(result.table_id, 4);
399 assert_eq!(result.region_id, region_stat.id.as_u64());
400 assert_eq!(result.region_number, 7);
401 assert_eq!(result.manifest_size, 256);
402 assert_eq!(result.datanode_id, datanode_id);
403 assert_eq!(result.engine, "mito");
404 assert_eq!(result.num_rows, 1000);
405 assert_eq!(result.sst_num, 5);
406 assert_eq!(result.sst_size, 2048);
407 assert_eq!(result.write_bytes_delta, 0); assert_eq!(result.timestamp_millis, timestamp_millis);
409 }
410
411 #[test]
412 fn test_compute_persist_region_stat_with_overflow_protection() {
413 let region_stat = create_test_region_stat(8, 15, 500, "mito");
414 let datanode_id = 505;
415 let timestamp_millis = 1640995620000; let persisted_stat = PersistedRegionStat {
417 region_id: region_stat.id,
418 written_bytes: 1000, };
420 let result = compute_persist_region_stat(
421 ®ion_stat,
422 datanode_id,
423 timestamp_millis,
424 Some(persisted_stat),
425 );
426 assert_eq!(result.table_id, 8);
427 assert_eq!(result.region_id, region_stat.id.as_u64());
428 assert_eq!(result.region_number, 15);
429 assert_eq!(result.manifest_size, 256);
430 assert_eq!(result.datanode_id, datanode_id);
431 assert_eq!(result.engine, "mito");
432 assert_eq!(result.num_rows, 1000);
433 assert_eq!(result.sst_num, 5);
434 assert_eq!(result.sst_size, 2048);
435 assert_eq!(result.write_bytes_delta, 0); assert_eq!(result.timestamp_millis, timestamp_millis);
437 }
438
439 struct MockInserter {
440 requests: Arc<Mutex<Vec<api::v1::RowInsertRequest>>>,
441 }
442
443 #[async_trait::async_trait]
444 impl Inserter for MockInserter {
445 async fn insert_rows(
446 &self,
447 _context: &InserterContext<'_>,
448 requests: api::v1::RowInsertRequests,
449 ) -> client::error::Result<()> {
450 self.requests.lock().unwrap().extend(requests.inserts);
451
452 Ok(())
453 }
454
455 fn set_options(&mut self, _options: &InsertOptions) {}
456 }
457
458 #[tokio::test]
459 async fn test_not_persist_region_stats() {
460 let env = TestEnv::new();
461 let mut ctx = env.ctx();
462
463 let requests = Arc::new(Mutex::new(vec![]));
464 let inserter = MockInserter {
465 requests: requests.clone(),
466 };
467 let handler = PersistStatsHandler::new(Box::new(inserter), Duration::from_secs(10));
468 let mut acc = HeartbeatAccumulator {
469 stat: Some(Stat {
470 id: 1,
471 timestamp_millis: 1640995200000,
472 region_stats: vec![create_test_region_stat(1, 1, 1000, "mito")],
473 ..Default::default()
474 }),
475 ..Default::default()
476 };
477 handler.last_persisted_time.insert(1, Instant::now());
478 handler
480 .handle(&HeartbeatRequest::default(), &mut ctx, &mut acc)
481 .await
482 .unwrap();
483 assert!(requests.lock().unwrap().is_empty());
484 }
485
486 #[tokio::test]
487 async fn test_persist_region_stats() {
488 let env = TestEnv::new();
489 let mut ctx = env.ctx();
490 let requests = Arc::new(Mutex::new(vec![]));
491 let inserter = MockInserter {
492 requests: requests.clone(),
493 };
494
495 let handler = PersistStatsHandler::new(Box::new(inserter), Duration::from_secs(10));
496
497 let region_stat = create_test_region_stat(1, 1, 1000, "mito");
498 let timestamp_millis = 1640995200000;
499 let datanode_id = 1;
500 let region_id = RegionId::new(1, 1);
501 let mut acc = HeartbeatAccumulator {
502 stat: Some(Stat {
503 id: datanode_id,
504 timestamp_millis,
505 region_stats: vec![region_stat.clone()],
506 ..Default::default()
507 }),
508 ..Default::default()
509 };
510
511 handler.last_persisted_region_stats.insert(
512 region_id,
513 PersistedRegionStat {
514 region_id,
515 written_bytes: 500,
516 },
517 );
518 let (expected_row, expected_persisted_region_stat) = to_persisted_if_leader(
519 ®ion_stat,
520 &handler.last_persisted_region_stats,
521 datanode_id,
522 timestamp_millis,
523 )
524 .unwrap();
525 let before_insert_time = Instant::now();
526 handler
528 .handle(&HeartbeatRequest::default(), &mut ctx, &mut acc)
529 .await
530 .unwrap();
531 let request = {
532 let mut requests = requests.lock().unwrap();
533 assert_eq!(requests.len(), 1);
534 requests.pop().unwrap()
535 };
536 assert_eq!(
537 request.table_name,
538 REGION_STATS_HISTORY_TABLE_NAME.to_string()
539 );
540 assert_eq!(request.rows.unwrap().rows, vec![expected_row]);
541
542 assert!(
544 handler
545 .last_persisted_time
546 .get(&datanode_id)
547 .unwrap()
548 .gt(&before_insert_time)
549 );
550
551 assert_eq!(
553 handler
554 .last_persisted_region_stats
555 .get(®ion_id)
556 .unwrap()
557 .value(),
558 &expected_persisted_region_stat
559 );
560 }
561}