1use 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#[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 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(®ion_stat.engine));
284 self.region_roles.push(Some(®ion_stat.role.to_string()));
285 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}