1use std::collections::{HashMap, HashSet};
16use std::str::FromStr;
17
18use api::v1::meta::{DatanodeWorkloads, HeartbeatRequest, RequestHeader};
19use common_time::{Timestamp, util as time_util};
20use lazy_static::lazy_static;
21use regex::Regex;
22use serde::{Deserialize, Serialize};
23use snafu::{OptionExt, ResultExt, ensure};
24use store_api::region_engine::{RegionRole, RegionStatistic};
25use store_api::storage::RegionId;
26use table::metadata::TableId;
27
28use crate::error::{self, DeserializeFromJsonSnafu, Result};
29use crate::heartbeat::utils::get_datanode_workloads;
30
31const DATANODE_STAT_PREFIX: &str = "__meta_datanode_stat";
32
33pub const REGION_STATISTIC_KEY: &str = "__region_statistic";
34pub const REGION_STATS_HISTORY_TABLE_NAME: &str = "region_statistics_history";
36
37lazy_static! {
38 pub(crate) static ref DATANODE_LEASE_KEY_PATTERN: Regex =
39 Regex::new("^__meta_datanode_lease-([0-9]+)-([0-9]+)$").unwrap();
40 static ref DATANODE_STAT_KEY_PATTERN: Regex =
41 Regex::new(&format!("^{DATANODE_STAT_PREFIX}-([0-9]+)-([0-9]+)$")).unwrap();
42 static ref INACTIVE_REGION_KEY_PATTERN: Regex =
43 Regex::new("^__meta_inactive_region-([0-9]+)-([0-9]+)-([0-9]+)$").unwrap();
44}
45
46#[derive(Debug, Clone, Default, Serialize, Deserialize)]
50pub struct Stat {
51 pub timestamp_millis: i64,
52 pub id: u64,
54 pub addr: String,
56 pub rcus: i64,
58 pub wcus: i64,
60 pub region_num: u64,
62 pub region_stats: Vec<RegionStat>,
64 pub topic_stats: Vec<TopicStat>,
66 pub node_epoch: u64,
68 pub datanode_workloads: DatanodeWorkloads,
70 pub gc_stat: Option<GcStat>,
72}
73
74#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
76pub struct RegionStat {
77 pub id: RegionId,
79 pub rcus: i64,
81 pub wcus: i64,
83 pub approximate_bytes: u64,
85 pub engine: String,
87 pub role: RegionRole,
89 pub num_rows: u64,
91 pub memtable_size: u64,
93 pub manifest_size: u64,
95 pub sst_size: u64,
97 pub sst_num: u64,
99 pub index_size: u64,
101 pub region_manifest: RegionManifestInfo,
103 pub written_bytes: u64,
105 #[serde(default)]
109 pub query_cpu_time: u64,
110 #[serde(default)]
112 pub query_scanned_bytes: u64,
113 pub data_topic_latest_entry_id: u64,
116 pub metadata_topic_latest_entry_id: u64,
120 #[serde(default)]
122 pub min_timestamp: Option<Timestamp>,
123 #[serde(default)]
125 pub max_timestamp: Option<Timestamp>,
126}
127
128#[derive(Debug, Clone, Serialize, Deserialize)]
129pub struct TopicStat {
130 pub topic: String,
132 pub latest_entry_id: u64,
134 pub record_size: u64,
136 pub record_num: u64,
138}
139
140pub trait TopicStatsReporter: Send + Sync {
142 fn reportable_topics(&mut self) -> Vec<TopicStat>;
144}
145
146#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Hash)]
147pub enum RegionManifestInfo {
148 Mito {
149 manifest_version: u64,
150 flushed_entry_id: u64,
151 file_removed_cnt: u64,
153 },
154 Metric {
155 data_manifest_version: u64,
156 data_flushed_entry_id: u64,
157 metadata_manifest_version: u64,
158 metadata_flushed_entry_id: u64,
159 },
160}
161
162impl Stat {
163 #[inline]
164 pub fn is_empty(&self) -> bool {
165 self.region_stats.is_empty()
166 }
167
168 pub fn stat_key(&self) -> DatanodeStatKey {
169 DatanodeStatKey { node_id: self.id }
170 }
171
172 pub fn regions(&self) -> Vec<(RegionId, RegionRole)> {
174 self.region_stats.iter().map(|s| (s.id, s.role)).collect()
175 }
176
177 pub fn table_ids(&self) -> HashSet<TableId> {
179 self.region_stats.iter().map(|s| s.id.table_id()).collect()
180 }
181
182 pub fn retain_active_region_stats(&mut self, inactive_region_ids: &HashSet<RegionId>) {
184 if inactive_region_ids.is_empty() {
185 return;
186 }
187
188 self.region_stats
189 .retain(|r| !inactive_region_ids.contains(&r.id));
190 self.rcus = self.region_stats.iter().map(|s| s.rcus).sum();
191 self.wcus = self.region_stats.iter().map(|s| s.wcus).sum();
192 self.region_num = self.region_stats.len() as u64;
193 }
194
195 pub fn memory_size(&self) -> usize {
196 std::mem::size_of::<i64>() * 3 +
198 std::mem::size_of::<u64>() * 3 +
200 std::mem::size_of::<String>() + self.addr.capacity() +
202 self.region_stats.iter().map(|s| s.memory_size()).sum::<usize>()
204 }
205}
206
207impl RegionStat {
208 pub fn memory_size(&self) -> usize {
209 std::mem::size_of::<RegionRole>() +
211 std::mem::size_of::<RegionId>() +
213 std::mem::size_of::<i64>() * 4 +
215 std::mem::size_of::<u64>() * 5 +
217 std::mem::size_of::<String>() + self.engine.capacity() +
219 self.region_manifest.memory_size()
221 }
222}
223
224impl RegionManifestInfo {
225 pub fn memory_size(&self) -> usize {
226 match self {
227 RegionManifestInfo::Mito { .. } => std::mem::size_of::<u64>() * 2,
228 RegionManifestInfo::Metric { .. } => std::mem::size_of::<u64>() * 4,
229 }
230 }
231}
232
233impl TryFrom<&HeartbeatRequest> for Stat {
234 type Error = Option<RequestHeader>;
235
236 fn try_from(value: &HeartbeatRequest) -> std::result::Result<Self, Self::Error> {
237 let HeartbeatRequest {
238 header,
239 peer,
240 region_stats,
241 node_epoch,
242 node_workloads,
243 topic_stats,
244 extensions,
245 ..
246 } = value;
247
248 match (header, peer) {
249 (Some(header), Some(peer)) => {
250 let region_stats = region_stats
251 .iter()
252 .map(RegionStat::from)
253 .collect::<Vec<_>>();
254 let topic_stats = topic_stats.iter().map(TopicStat::from).collect::<Vec<_>>();
255
256 let datanode_workloads = get_datanode_workloads(node_workloads.as_ref());
257
258 let gc_stat = GcStat::from_extensions(extensions).map_err(|err| {
259 common_telemetry::error!(
260 "Failed to deserialize GcStat from extensions: {}",
261 err
262 );
263 header.clone()
264 })?;
265 Ok(Self {
266 timestamp_millis: time_util::current_time_millis(),
267 id: peer.id,
269 addr: peer.addr.clone(),
271 rcus: region_stats.iter().map(|s| s.rcus).sum(),
272 wcus: region_stats.iter().map(|s| s.wcus).sum(),
273 region_num: region_stats.len() as u64,
274 region_stats,
275 topic_stats,
276 node_epoch: *node_epoch,
277 datanode_workloads,
278 gc_stat,
279 })
280 }
281 (header, _) => Err(header.clone()),
282 }
283 }
284}
285
286impl From<store_api::region_engine::RegionManifestInfo> for RegionManifestInfo {
287 fn from(value: store_api::region_engine::RegionManifestInfo) -> Self {
288 match value {
289 store_api::region_engine::RegionManifestInfo::Mito {
290 manifest_version,
291 flushed_entry_id,
292 file_removed_cnt,
293 } => RegionManifestInfo::Mito {
294 manifest_version,
295 flushed_entry_id,
296 file_removed_cnt,
297 },
298 store_api::region_engine::RegionManifestInfo::Metric {
299 data_manifest_version,
300 data_flushed_entry_id,
301 metadata_manifest_version,
302 metadata_flushed_entry_id,
303 } => RegionManifestInfo::Metric {
304 data_manifest_version,
305 data_flushed_entry_id,
306 metadata_manifest_version,
307 metadata_flushed_entry_id,
308 },
309 }
310 }
311}
312
313impl From<&api::v1::meta::RegionStat> for RegionStat {
314 fn from(value: &api::v1::meta::RegionStat) -> Self {
315 let region_stat = value
316 .extensions
317 .get(REGION_STATISTIC_KEY)
318 .and_then(|value| RegionStatistic::deserialize_from_slice(value))
319 .unwrap_or_default();
320
321 Self {
322 id: RegionId::from_u64(value.region_id),
323 rcus: value.rcus,
324 wcus: value.wcus,
325 approximate_bytes: value.approximate_bytes as u64,
326 engine: value.engine.clone(),
327 role: RegionRole::from(value.role()),
328 num_rows: region_stat.num_rows,
329 memtable_size: region_stat.memtable_size,
330 manifest_size: region_stat.manifest_size,
331 sst_size: region_stat.sst_size,
332 sst_num: region_stat.sst_num,
333 index_size: region_stat.index_size,
334 region_manifest: region_stat.manifest.into(),
335 written_bytes: region_stat.written_bytes,
336 query_cpu_time: region_stat.query_cpu_time,
337 query_scanned_bytes: region_stat.query_scanned_bytes,
338 data_topic_latest_entry_id: region_stat.data_topic_latest_entry_id,
339 metadata_topic_latest_entry_id: region_stat.metadata_topic_latest_entry_id,
340 min_timestamp: region_stat.min_timestamp,
341 max_timestamp: region_stat.max_timestamp,
342 }
343 }
344}
345
346impl From<&api::v1::meta::TopicStat> for TopicStat {
347 fn from(value: &api::v1::meta::TopicStat) -> Self {
348 Self {
349 topic: value.topic_name.clone(),
350 latest_entry_id: value.latest_entry_id,
351 record_size: value.record_size,
352 record_num: value.record_num,
353 }
354 }
355}
356
357#[derive(Debug, Clone, Serialize, Deserialize, Default)]
358pub struct GcStat {
359 pub running_gc_tasks: u32,
361 pub gc_concurrency: u32,
363}
364
365impl GcStat {
366 pub const GC_STAT_KEY: &str = "__gc_stat";
367
368 pub fn new(running_gc_tasks: u32, gc_concurrency: u32) -> Self {
369 Self {
370 running_gc_tasks,
371 gc_concurrency,
372 }
373 }
374
375 pub fn into_extensions(&self, extensions: &mut std::collections::HashMap<String, Vec<u8>>) {
376 let bytes = serde_json::to_vec(self).unwrap_or_default();
377 extensions.insert(Self::GC_STAT_KEY.to_string(), bytes);
378 }
379
380 pub fn from_extensions(
381 extensions: &std::collections::HashMap<String, Vec<u8>>,
382 ) -> Result<Option<Self>> {
383 extensions
384 .get(Self::GC_STAT_KEY)
385 .map(|bytes| {
386 serde_json::from_slice(bytes).with_context(|_| DeserializeFromJsonSnafu {
387 input: String::from_utf8_lossy(bytes).to_string(),
388 })
389 })
390 .transpose()
391 }
392}
393
394#[derive(Debug, Clone, Serialize, Deserialize, Default)]
396pub struct EnvVars {
397 pub vars: HashMap<String, String>,
398}
399
400impl EnvVars {
401 pub const ENV_VARS_KEY: &str = "__env_vars";
402
403 pub fn new(vars: HashMap<String, String>) -> Self {
404 Self { vars }
405 }
406
407 pub fn from_config(keys: &[String]) -> Self {
409 let vars = keys
410 .iter()
411 .filter_map(|key| std::env::var(key).ok().map(|value| (key.clone(), value)))
412 .collect();
413 Self { vars }
414 }
415
416 pub fn into_extensions(&self, extensions: &mut HashMap<String, Vec<u8>>) {
417 if self.vars.is_empty() {
418 return;
419 }
420 let bytes = serde_json::to_vec(self).unwrap_or_default();
421 extensions.insert(Self::ENV_VARS_KEY.to_string(), bytes);
422 }
423
424 pub fn from_extensions(extensions: &HashMap<String, Vec<u8>>) -> Result<Option<Self>> {
425 extensions
426 .get(Self::ENV_VARS_KEY)
427 .map(|bytes| {
428 serde_json::from_slice(bytes).with_context(|_| DeserializeFromJsonSnafu {
429 input: String::from_utf8_lossy(bytes).to_string(),
430 })
431 })
432 .transpose()
433 }
434}
435
436#[derive(Debug, Clone, Copy, Eq, PartialEq, Hash)]
440pub struct DatanodeStatKey {
441 pub node_id: u64,
442}
443
444impl DatanodeStatKey {
445 pub fn prefix_key() -> Vec<u8> {
447 format!("{DATANODE_STAT_PREFIX}-0-").into_bytes()
449 }
450}
451
452impl From<DatanodeStatKey> for Vec<u8> {
453 fn from(value: DatanodeStatKey) -> Self {
454 format!("{}-0-{}", DATANODE_STAT_PREFIX, value.node_id).into_bytes()
456 }
457}
458
459impl FromStr for DatanodeStatKey {
460 type Err = error::Error;
461
462 fn from_str(key: &str) -> Result<Self> {
463 let caps = DATANODE_STAT_KEY_PATTERN
464 .captures(key)
465 .context(error::InvalidStatKeySnafu { key })?;
466
467 ensure!(caps.len() == 3, error::InvalidStatKeySnafu { key });
468 let node_id = caps[2].to_string();
469 let node_id: u64 = node_id.parse().context(error::ParseNumSnafu {
470 err_msg: format!("invalid node_id: {node_id}"),
471 })?;
472
473 Ok(Self { node_id })
474 }
475}
476
477impl TryFrom<Vec<u8>> for DatanodeStatKey {
478 type Error = error::Error;
479
480 fn try_from(bytes: Vec<u8>) -> Result<Self> {
481 String::from_utf8(bytes)
482 .context(error::FromUtf8Snafu {
483 name: "DatanodeStatKey",
484 })
485 .map(|x| x.parse())?
486 }
487}
488
489#[derive(Debug, Clone, Serialize, Deserialize)]
491#[serde(transparent)]
492pub struct DatanodeStatValue {
493 pub stats: Vec<Stat>,
494}
495
496impl DatanodeStatValue {
497 pub fn region_num(&self) -> Option<u64> {
499 self.stats.last().map(|x| x.region_num)
500 }
501
502 pub fn node_addr(&self) -> Option<String> {
504 self.stats.last().map(|x| x.addr.clone())
505 }
506}
507
508impl TryFrom<DatanodeStatValue> for Vec<u8> {
509 type Error = error::Error;
510
511 fn try_from(stats: DatanodeStatValue) -> Result<Self> {
512 Ok(serde_json::to_string(&stats)
513 .context(error::SerializeToJsonSnafu {
514 input: format!("{stats:?}"),
515 })?
516 .into_bytes())
517 }
518}
519
520impl FromStr for DatanodeStatValue {
521 type Err = error::Error;
522
523 fn from_str(value: &str) -> Result<Self> {
524 serde_json::from_str(value).context(error::DeserializeFromJsonSnafu { input: value })
525 }
526}
527
528impl TryFrom<Vec<u8>> for DatanodeStatValue {
529 type Error = error::Error;
530
531 fn try_from(value: Vec<u8>) -> Result<Self> {
532 String::from_utf8(value)
533 .context(error::FromUtf8Snafu {
534 name: "DatanodeStatValue",
535 })
536 .map(|x| x.parse())?
537 }
538}
539
540#[cfg(test)]
541mod tests {
542 use super::*;
543
544 #[test]
545 fn test_stat_key() {
546 let stat = Stat {
547 id: 101,
548 region_num: 10,
549 ..Default::default()
550 };
551
552 let stat_key = stat.stat_key();
553
554 assert_eq!(101, stat_key.node_id);
555 }
556
557 #[test]
558 fn test_stat_val_round_trip() {
559 let stat = Stat {
560 id: 101,
561 region_num: 100,
562 ..Default::default()
563 };
564
565 let stat_val = DatanodeStatValue { stats: vec![stat] };
566
567 let bytes: Vec<u8> = stat_val.try_into().unwrap();
568 let stat_val: DatanodeStatValue = bytes.try_into().unwrap();
569 let stats = stat_val.stats;
570
571 assert_eq!(1, stats.len());
572
573 let stat = stats.first().unwrap();
574 assert_eq!(101, stat.id);
575 assert_eq!(100, stat.region_num);
576 }
577
578 #[test]
579 fn test_stat_val_deserializes_without_query_stats() {
580 let stat = Stat {
581 region_stats: vec![RegionStat {
582 id: RegionId::new(1024, 1),
583 rcus: 0,
584 wcus: 0,
585 approximate_bytes: 0,
586 engine: "mito".to_string(),
587 role: RegionRole::Leader,
588 num_rows: 0,
589 memtable_size: 0,
590 manifest_size: 0,
591 sst_size: 0,
592 sst_num: 0,
593 index_size: 0,
594 region_manifest: RegionManifestInfo::Mito {
595 manifest_version: 0,
596 flushed_entry_id: 0,
597 file_removed_cnt: 0,
598 },
599 written_bytes: 0,
600 query_cpu_time: 10,
601 query_scanned_bytes: 20,
602 data_topic_latest_entry_id: 0,
603 metadata_topic_latest_entry_id: 0,
604 min_timestamp: None,
605 max_timestamp: None,
606 }],
607 ..Default::default()
608 };
609 let stat_val = DatanodeStatValue { stats: vec![stat] };
610 let mut value = serde_json::to_value(stat_val).unwrap();
611 let region_stat = value[0]["region_stats"][0].as_object_mut().unwrap();
612 region_stat.remove("query_cpu_time");
613 region_stat.remove("query_scanned_bytes");
614
615 let stat_val: DatanodeStatValue = serde_json::from_value(value).unwrap();
616 let region_stat = &stat_val.stats[0].region_stats[0];
617
618 assert_eq!(region_stat.query_cpu_time, 0);
619 assert_eq!(region_stat.query_scanned_bytes, 0);
620 }
621
622 #[test]
623 fn test_get_addr_from_stat_val() {
624 let empty = DatanodeStatValue { stats: vec![] };
625 let addr = empty.node_addr();
626 assert!(addr.is_none());
627
628 let stat_val = DatanodeStatValue {
629 stats: vec![
630 Stat {
631 addr: "1".to_string(),
632 ..Default::default()
633 },
634 Stat {
635 addr: "2".to_string(),
636 ..Default::default()
637 },
638 Stat {
639 addr: "3".to_string(),
640 ..Default::default()
641 },
642 ],
643 };
644 let addr = stat_val.node_addr().unwrap();
645 assert_eq!("3", addr);
646 }
647
648 #[test]
649 fn test_get_region_num_from_stat_val() {
650 let empty = DatanodeStatValue { stats: vec![] };
651 let region_num = empty.region_num();
652 assert!(region_num.is_none());
653
654 let wrong = DatanodeStatValue {
655 stats: vec![Stat {
656 region_num: 0,
657 ..Default::default()
658 }],
659 };
660 let right = wrong.region_num();
661 assert_eq!(Some(0), right);
662
663 let stat_val = DatanodeStatValue {
664 stats: vec![
665 Stat {
666 region_num: 1,
667 ..Default::default()
668 },
669 Stat {
670 region_num: 0,
671 ..Default::default()
672 },
673 Stat {
674 region_num: 2,
675 ..Default::default()
676 },
677 ],
678 };
679 let region_num = stat_val.region_num().unwrap();
680 assert_eq!(2, region_num);
681 }
682
683 #[test]
684 fn test_region_stat_from_heartbeat_preserves_staging_leader_role() {
685 let request = HeartbeatRequest {
686 header: Some(RequestHeader::default()),
687 peer: Some(api::v1::meta::Peer {
688 id: 1,
689 addr: "127.0.0.1:3001".to_string(),
690 }),
691 region_stats: vec![api::v1::meta::RegionStat {
692 region_id: RegionId::new(1024, 1).as_u64(),
693 engine: "mito".to_string(),
694 role: api::v1::meta::RegionRole::StagingLeader.into(),
695 ..Default::default()
696 }],
697 ..Default::default()
698 };
699
700 let stat = Stat::try_from(&request).unwrap();
701
702 assert_eq!(stat.region_stats.len(), 1);
703 assert_eq!(stat.region_stats[0].role, RegionRole::StagingLeader);
704 }
705
706 #[test]
707 fn test_env_vars_round_trip() {
708 let mut vars = HashMap::new();
709 vars.insert("AZ".to_string(), "us-east-1a".to_string());
710 vars.insert("REGION".to_string(), "us-east-1".to_string());
711 let env_vars = EnvVars::new(vars);
712
713 let mut extensions = HashMap::new();
714 env_vars.into_extensions(&mut extensions);
715
716 let extracted = EnvVars::from_extensions(&extensions).unwrap().unwrap();
717 assert_eq!(extracted.vars.get("AZ").unwrap(), "us-east-1a");
718 assert_eq!(extracted.vars.get("REGION").unwrap(), "us-east-1");
719 }
720
721 #[test]
722 fn test_env_vars_empty_not_written() {
723 let env_vars = EnvVars::default();
724 let mut extensions = HashMap::new();
725 env_vars.into_extensions(&mut extensions);
726 assert!(extensions.is_empty());
727 }
728
729 #[test]
730 fn test_env_vars_from_extensions_missing() {
731 let extensions = HashMap::new();
732 let result = EnvVars::from_extensions(&extensions).unwrap();
733 assert!(result.is_none());
734 }
735}