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_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 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 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; let result = compute_persist_region_stat(®ion_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); 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; let persisted_stat = PersistedRegionStat {
335 region_id: region_stat.id,
336 written_bytes: 1000, };
338 let result = compute_persist_region_stat(
339 ®ion_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); 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; let persisted_stat = PersistedRegionStat {
363 region_id: region_stat.id,
364 written_bytes: 1200, };
366 let result = compute_persist_region_stat(
367 ®ion_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); 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; let persisted_stat = PersistedRegionStat {
391 region_id: region_stat.id,
392 written_bytes: 2000, };
394 let result = compute_persist_region_stat(
395 ®ion_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); 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; let persisted_stat = PersistedRegionStat {
419 region_id: region_stat.id,
420 written_bytes: 1000, };
422 let result = compute_persist_region_stat(
423 ®ion_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); 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 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 ®ion_stat,
522 &handler.last_persisted_region_stats,
523 datanode_id,
524 timestamp_millis,
525 )
526 .unwrap();
527 let before_insert_time = Instant::now();
528 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 assert!(
546 handler
547 .last_persisted_time
548 .get(&datanode_id)
549 .unwrap()
550 .gt(&before_insert_time)
551 );
552
553 assert_eq!(
555 handler
556 .last_persisted_region_stats
557 .get(®ion_id)
558 .unwrap()
559 .value(),
560 &expected_persisted_region_stat
561 );
562 }
563}