Skip to main content

catalog/system_schema/information_schema/
region_statistics.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::sync::{Arc, Weak};
16
17use arrow_schema::SchemaRef as ArrowSchemaRef;
18use common_catalog::consts::INFORMATION_SCHEMA_REGION_STATISTICS_TABLE_ID;
19use common_error::ext::BoxedError;
20use common_meta::datanode::RegionStat;
21use common_recordbatch::adapter::RecordBatchStreamAdapter;
22use common_recordbatch::{DfSendableRecordBatchStream, RecordBatch, SendableRecordBatchStream};
23use common_time::timestamp::TimeUnit;
24use datafusion::execution::TaskContext;
25use datafusion::physical_plan::stream::RecordBatchStreamAdapter as DfRecordBatchStreamAdapter;
26use datafusion::physical_plan::streaming::PartitionStream as DfPartitionStream;
27use datatypes::prelude::{ConcreteDataType, ScalarVectorBuilder, VectorRef};
28use datatypes::schema::{ColumnSchema, Schema, SchemaRef};
29use datatypes::timestamp::TimestampMillisecond;
30use datatypes::value::Value;
31use datatypes::vectors::{
32    StringVectorBuilder, TimestampMillisecondVectorBuilder, UInt32VectorBuilder,
33    UInt64VectorBuilder,
34};
35use snafu::ResultExt;
36use store_api::storage::{ScanRequest, TableId};
37
38use crate::CatalogManager;
39use crate::error::{CreateRecordBatchSnafu, InternalSnafu, Result};
40use crate::information_schema::Predicates;
41use crate::system_schema::information_schema::{InformationTable, REGION_STATISTICS};
42use crate::system_schema::utils;
43
44const REGION_ID: &str = "region_id";
45const TABLE_ID: &str = "table_id";
46const REGION_NUMBER: &str = "region_number";
47const REGION_ROWS: &str = "region_rows";
48const WRITTEN_BYTES: &str = "written_bytes_since_open";
49const QUERY_CPU_TIME_MILLIS: &str = "query_cpu_time_millis";
50const QUERY_SCANNED_BYTES: &str = "query_scanned_bytes";
51const DISK_SIZE: &str = "disk_size";
52const MEMTABLE_SIZE: &str = "memtable_size";
53const MANIFEST_SIZE: &str = "manifest_size";
54const SST_SIZE: &str = "sst_size";
55const SST_NUM: &str = "sst_num";
56const INDEX_SIZE: &str = "index_size";
57const ENGINE: &str = "engine";
58const REGION_ROLE: &str = "region_role";
59const MIN_TIMESTAMP: &str = "min_timestamp";
60const MAX_TIMESTAMP: &str = "max_timestamp";
61
62const INIT_CAPACITY: usize = 42;
63
64/// The `REGION_STATISTICS` table provides information about the region statistics. Including fields:
65///
66/// - `region_id`: The region id.
67/// - `table_id`: The table id.
68/// - `region_number`: The region number.
69/// - `region_rows`: The number of rows in region.
70/// - `written_bytes_since_open`: The total bytes written of the region since region opened.
71/// - `query_cpu_time_millis`: The total query CPU time of the region since region opened, in milliseconds.
72/// - `query_scanned_bytes`: The total bytes scanned by queries since region opened.
73/// - `memtable_size`: The memtable size in bytes.
74/// - `disk_size`: The approximate disk size in bytes.
75/// - `manifest_size`: The manifest size in bytes.
76/// - `sst_size`: The sst data files size in bytes.
77/// - `index_size`: The sst index files size in bytes.
78/// - `engine`: The engine type.
79/// - `region_role`: The region role.
80#[derive(Debug)]
81pub(super) struct InformationSchemaRegionStatistics {
82    schema: SchemaRef,
83    catalog_manager: Weak<dyn CatalogManager>,
84}
85
86impl InformationSchemaRegionStatistics {
87    pub(super) fn new(catalog_manager: Weak<dyn CatalogManager>) -> Self {
88        Self {
89            schema: Self::schema(),
90            catalog_manager,
91        }
92    }
93
94    pub(crate) fn schema() -> SchemaRef {
95        Arc::new(Schema::new(vec![
96            ColumnSchema::new(REGION_ID, ConcreteDataType::uint64_datatype(), false),
97            ColumnSchema::new(TABLE_ID, ConcreteDataType::uint32_datatype(), false),
98            ColumnSchema::new(REGION_NUMBER, ConcreteDataType::uint32_datatype(), false),
99            ColumnSchema::new(REGION_ROWS, ConcreteDataType::uint64_datatype(), true),
100            ColumnSchema::new(WRITTEN_BYTES, ConcreteDataType::uint64_datatype(), true),
101            ColumnSchema::new(
102                QUERY_CPU_TIME_MILLIS,
103                ConcreteDataType::uint64_datatype(),
104                true,
105            ),
106            ColumnSchema::new(
107                QUERY_SCANNED_BYTES,
108                ConcreteDataType::uint64_datatype(),
109                true,
110            ),
111            ColumnSchema::new(DISK_SIZE, ConcreteDataType::uint64_datatype(), true),
112            ColumnSchema::new(MEMTABLE_SIZE, ConcreteDataType::uint64_datatype(), true),
113            ColumnSchema::new(MANIFEST_SIZE, ConcreteDataType::uint64_datatype(), true),
114            ColumnSchema::new(SST_SIZE, ConcreteDataType::uint64_datatype(), true),
115            ColumnSchema::new(SST_NUM, ConcreteDataType::uint64_datatype(), true),
116            ColumnSchema::new(INDEX_SIZE, ConcreteDataType::uint64_datatype(), true),
117            ColumnSchema::new(ENGINE, ConcreteDataType::string_datatype(), true),
118            ColumnSchema::new(REGION_ROLE, ConcreteDataType::string_datatype(), true),
119            ColumnSchema::new(
120                MIN_TIMESTAMP,
121                ConcreteDataType::timestamp_millisecond_datatype(),
122                true,
123            ),
124            ColumnSchema::new(
125                MAX_TIMESTAMP,
126                ConcreteDataType::timestamp_millisecond_datatype(),
127                true,
128            ),
129        ]))
130    }
131
132    fn builder(&self) -> InformationSchemaRegionStatisticsBuilder {
133        InformationSchemaRegionStatisticsBuilder::new(
134            self.schema.clone(),
135            self.catalog_manager.clone(),
136        )
137    }
138}
139
140impl InformationTable for InformationSchemaRegionStatistics {
141    fn table_id(&self) -> TableId {
142        INFORMATION_SCHEMA_REGION_STATISTICS_TABLE_ID
143    }
144
145    fn table_name(&self) -> &'static str {
146        REGION_STATISTICS
147    }
148
149    fn schema(&self) -> SchemaRef {
150        self.schema.clone()
151    }
152
153    fn to_stream(&self, request: ScanRequest) -> Result<SendableRecordBatchStream> {
154        let schema = self.schema.arrow_schema().clone();
155        let mut builder = self.builder();
156
157        let stream = Box::pin(DfRecordBatchStreamAdapter::new(
158            schema,
159            futures::stream::once(async move {
160                builder
161                    .make_region_statistics(Some(request))
162                    .await
163                    .map(|x| x.into_df_record_batch())
164                    .map_err(Into::into)
165            }),
166        ));
167
168        Ok(Box::pin(
169            RecordBatchStreamAdapter::try_new(stream)
170                .map_err(BoxedError::new)
171                .context(InternalSnafu)?,
172        ))
173    }
174}
175
176struct InformationSchemaRegionStatisticsBuilder {
177    schema: SchemaRef,
178    catalog_manager: Weak<dyn CatalogManager>,
179
180    region_ids: UInt64VectorBuilder,
181    table_ids: UInt32VectorBuilder,
182    region_numbers: UInt32VectorBuilder,
183    region_rows: UInt64VectorBuilder,
184    written_bytes: UInt64VectorBuilder,
185    query_cpu_time_millis: UInt64VectorBuilder,
186    query_scanned_bytes: UInt64VectorBuilder,
187    disk_sizes: UInt64VectorBuilder,
188    memtable_sizes: UInt64VectorBuilder,
189    manifest_sizes: UInt64VectorBuilder,
190    sst_sizes: UInt64VectorBuilder,
191    sst_nums: UInt64VectorBuilder,
192    index_sizes: UInt64VectorBuilder,
193    engines: StringVectorBuilder,
194    region_roles: StringVectorBuilder,
195    min_timestamps: TimestampMillisecondVectorBuilder,
196    max_timestamps: TimestampMillisecondVectorBuilder,
197}
198
199impl InformationSchemaRegionStatisticsBuilder {
200    fn new(schema: SchemaRef, catalog_manager: Weak<dyn CatalogManager>) -> Self {
201        Self {
202            schema,
203            catalog_manager,
204            region_ids: UInt64VectorBuilder::with_capacity(INIT_CAPACITY),
205            table_ids: UInt32VectorBuilder::with_capacity(INIT_CAPACITY),
206            region_numbers: UInt32VectorBuilder::with_capacity(INIT_CAPACITY),
207            region_rows: UInt64VectorBuilder::with_capacity(INIT_CAPACITY),
208            written_bytes: UInt64VectorBuilder::with_capacity(INIT_CAPACITY),
209            query_cpu_time_millis: UInt64VectorBuilder::with_capacity(INIT_CAPACITY),
210            query_scanned_bytes: UInt64VectorBuilder::with_capacity(INIT_CAPACITY),
211            disk_sizes: UInt64VectorBuilder::with_capacity(INIT_CAPACITY),
212            memtable_sizes: UInt64VectorBuilder::with_capacity(INIT_CAPACITY),
213            manifest_sizes: UInt64VectorBuilder::with_capacity(INIT_CAPACITY),
214            sst_sizes: UInt64VectorBuilder::with_capacity(INIT_CAPACITY),
215            sst_nums: UInt64VectorBuilder::with_capacity(INIT_CAPACITY),
216            index_sizes: UInt64VectorBuilder::with_capacity(INIT_CAPACITY),
217            engines: StringVectorBuilder::with_capacity(INIT_CAPACITY),
218            region_roles: StringVectorBuilder::with_capacity(INIT_CAPACITY),
219            min_timestamps: TimestampMillisecondVectorBuilder::with_capacity(INIT_CAPACITY),
220            max_timestamps: TimestampMillisecondVectorBuilder::with_capacity(INIT_CAPACITY),
221        }
222    }
223
224    /// Construct a new `InformationSchemaRegionStatistics` from the collected data.
225    async fn make_region_statistics(
226        &mut self,
227        request: Option<ScanRequest>,
228    ) -> Result<RecordBatch> {
229        let predicates = Predicates::from_scan_request(&request);
230        let information_extension = utils::information_extension(&self.catalog_manager)?;
231        let region_stats = information_extension.region_stats().await?;
232        for region_stat in region_stats {
233            self.add_region_statistic(&predicates, region_stat);
234        }
235        self.finish()
236    }
237
238    fn add_region_statistic(&mut self, predicate: &Predicates, region_stat: RegionStat) {
239        let row = [
240            (REGION_ID, &Value::from(region_stat.id.as_u64())),
241            (TABLE_ID, &Value::from(region_stat.id.table_id())),
242            (REGION_NUMBER, &Value::from(region_stat.id.region_number())),
243            (REGION_ROWS, &Value::from(region_stat.num_rows)),
244            (WRITTEN_BYTES, &Value::from(region_stat.written_bytes)),
245            (
246                QUERY_CPU_TIME_MILLIS,
247                &Value::from(region_stat.query_cpu_time / 1_000_000),
248            ),
249            (
250                QUERY_SCANNED_BYTES,
251                &Value::from(region_stat.query_scanned_bytes),
252            ),
253            (DISK_SIZE, &Value::from(region_stat.approximate_bytes)),
254            (MEMTABLE_SIZE, &Value::from(region_stat.memtable_size)),
255            (MANIFEST_SIZE, &Value::from(region_stat.manifest_size)),
256            (SST_SIZE, &Value::from(region_stat.sst_size)),
257            (SST_NUM, &Value::from(region_stat.sst_num)),
258            (INDEX_SIZE, &Value::from(region_stat.index_size)),
259            (ENGINE, &Value::from(region_stat.engine.as_str())),
260            (REGION_ROLE, &Value::from(region_stat.role.to_string())),
261        ];
262
263        if !predicate.eval(&row) {
264            return;
265        }
266
267        self.region_ids.push(Some(region_stat.id.as_u64()));
268        self.table_ids.push(Some(region_stat.id.table_id()));
269        self.region_numbers
270            .push(Some(region_stat.id.region_number()));
271        self.region_rows.push(Some(region_stat.num_rows));
272        self.written_bytes.push(Some(region_stat.written_bytes));
273        self.query_cpu_time_millis
274            .push(Some(region_stat.query_cpu_time / 1_000_000));
275        self.query_scanned_bytes
276            .push(Some(region_stat.query_scanned_bytes));
277        self.disk_sizes.push(Some(region_stat.approximate_bytes));
278        self.memtable_sizes.push(Some(region_stat.memtable_size));
279        self.manifest_sizes.push(Some(region_stat.manifest_size));
280        self.sst_sizes.push(Some(region_stat.sst_size));
281        self.sst_nums.push(Some(region_stat.sst_num));
282        self.index_sizes.push(Some(region_stat.index_size));
283        self.engines.push(Some(&region_stat.engine));
284        self.region_roles.push(Some(&region_stat.role.to_string()));
285        // Floor the min and ceil the max so the window stays a superset: a narrower
286        // one would hide regions from a time-bounded lookup.
287        self.min_timestamps.push(
288            region_stat
289                .min_timestamp
290                .and_then(|ts| ts.convert_to(TimeUnit::Millisecond))
291                .map(TimestampMillisecond),
292        );
293        self.max_timestamps.push(
294            region_stat
295                .max_timestamp
296                .and_then(|ts| ts.convert_to_ceil(TimeUnit::Millisecond))
297                .map(TimestampMillisecond),
298        );
299    }
300
301    fn finish(&mut self) -> Result<RecordBatch> {
302        let columns: Vec<VectorRef> = vec![
303            Arc::new(self.region_ids.finish()),
304            Arc::new(self.table_ids.finish()),
305            Arc::new(self.region_numbers.finish()),
306            Arc::new(self.region_rows.finish()),
307            Arc::new(self.written_bytes.finish()),
308            Arc::new(self.query_cpu_time_millis.finish()),
309            Arc::new(self.query_scanned_bytes.finish()),
310            Arc::new(self.disk_sizes.finish()),
311            Arc::new(self.memtable_sizes.finish()),
312            Arc::new(self.manifest_sizes.finish()),
313            Arc::new(self.sst_sizes.finish()),
314            Arc::new(self.sst_nums.finish()),
315            Arc::new(self.index_sizes.finish()),
316            Arc::new(self.engines.finish()),
317            Arc::new(self.region_roles.finish()),
318            Arc::new(self.min_timestamps.finish()),
319            Arc::new(self.max_timestamps.finish()),
320        ];
321
322        RecordBatch::new(self.schema.clone(), columns).context(CreateRecordBatchSnafu)
323    }
324}
325
326impl DfPartitionStream for InformationSchemaRegionStatistics {
327    fn schema(&self) -> &ArrowSchemaRef {
328        self.schema.arrow_schema()
329    }
330
331    fn execute(&self, _: Arc<TaskContext>) -> DfSendableRecordBatchStream {
332        let schema = self.schema.arrow_schema().clone();
333        let mut builder = self.builder();
334        Box::pin(DfRecordBatchStreamAdapter::new(
335            schema,
336            futures::stream::once(async move {
337                builder
338                    .make_region_statistics(None)
339                    .await
340                    .map(|x| x.into_df_record_batch())
341                    .map_err(Into::into)
342            }),
343        ))
344    }
345}