1use datafusion_common::instant::Instant;
21use std::fmt;
22use std::fs::File;
23use std::str::FromStr;
24use std::sync::Arc;
25
26use arrow::array::{
27 DurationMillisecondArray, GenericListArray, Int64Array, StringArray, StructArray,
28 TimestampMillisecondArray, UInt64Array,
29};
30use arrow::buffer::{Buffer, OffsetBuffer, ScalarBuffer};
31use arrow::datatypes::{DataType, Field, Fields, Schema, SchemaRef, TimeUnit};
32use arrow::record_batch::RecordBatch;
33use arrow::util::pretty::pretty_format_batches;
34use datafusion::catalog::{Session, TableFunctionArgs, TableFunctionImpl};
35use datafusion::common::{Column, plan_err};
36use datafusion::datasource::TableProvider;
37use datafusion::datasource::memory::MemorySourceConfig;
38use datafusion::error::Result;
39use datafusion::execution::cache::cache_manager::CacheManager;
40use datafusion::logical_expr::Expr;
41use datafusion::physical_plan::ExecutionPlan;
42use datafusion::scalar::ScalarValue;
43
44use async_trait::async_trait;
45use datafusion_common::heap_size::{DFHeapSize, DFHeapSizeCtx};
46use parquet::basic::ConvertedType;
47use parquet::data_type::{ByteArray, FixedLenByteArray};
48use parquet::file::reader::FileReader;
49use parquet::file::serialized_reader::SerializedFileReader;
50use parquet::file::statistics::Statistics;
51
52#[derive(Debug)]
53pub enum Function {
54 Select,
55 Explain,
56 Show,
57 CreateTable,
58 CreateTableAs,
59 Insert,
60 DropTable,
61}
62
63const ALL_FUNCTIONS: [Function; 7] = [
64 Function::CreateTable,
65 Function::CreateTableAs,
66 Function::DropTable,
67 Function::Explain,
68 Function::Insert,
69 Function::Select,
70 Function::Show,
71];
72
73impl Function {
74 pub fn function_details(&self) -> Result<&str> {
75 let details = match self {
76 Function::Select => {
77 r#"
78Command: SELECT
79Description: retrieve rows from a table or view
80Syntax:
81SELECT [ ALL | DISTINCT [ ON ( expression [, ...] ) ] ]
82 [ * | expression [ [ AS ] output_name ] [, ...] ]
83 [ FROM from_item [, ...] ]
84 [ WHERE condition ]
85 [ GROUP BY [ ALL | DISTINCT ] grouping_element [, ...] ]
86 [ HAVING condition ]
87 [ WINDOW window_name AS ( window_definition ) [, ...] ]
88 [ { UNION | INTERSECT | EXCEPT } [ ALL | DISTINCT ] select ]
89 [ ORDER BY expression [ ASC | DESC | USING operator ] [ NULLS { FIRST | LAST } ] [, ...] ]
90 [ LIMIT { count | ALL } ]
91 [ OFFSET start [ ROW | ROWS ] ]
92
93where from_item can be one of:
94
95 [ ONLY ] table_name [ * ] [ [ AS ] alias [ ( column_alias [, ...] ) ] ]
96 [ TABLESAMPLE sampling_method ( argument [, ...] ) [ REPEATABLE ( seed ) ] ]
97 [ LATERAL ] ( select ) [ AS ] alias [ ( column_alias [, ...] ) ]
98 with_query_name [ [ AS ] alias [ ( column_alias [, ...] ) ] ]
99 [ LATERAL ] function_name ( [ argument [, ...] ] )
100 [ WITH ORDINALITY ] [ [ AS ] alias [ ( column_alias [, ...] ) ] ]
101 [ LATERAL ] function_name ( [ argument [, ...] ] ) [ AS ] alias ( column_definition [, ...] )
102 [ LATERAL ] function_name ( [ argument [, ...] ] ) AS ( column_definition [, ...] )
103 [ LATERAL ] ROWS FROM( function_name ( [ argument [, ...] ] ) [ AS ( column_definition [, ...] ) ] [, ...] )
104 [ WITH ORDINALITY ] [ [ AS ] alias [ ( column_alias [, ...] ) ] ]
105 from_item [ NATURAL ] join_type from_item [ ON join_condition | USING ( join_column [, ...] ) [ AS join_using_alias ] ]
106
107and grouping_element can be one of:
108
109 ( )
110 expression
111 ( expression [, ...] )
112
113and with_query is:
114
115 with_query_name [ ( column_name [, ...] ) ] AS [ [ NOT ] MATERIALIZED ] ( select | values | insert | update | delete )
116
117TABLE [ ONLY ] table_name [ * ]"#
118 }
119 Function::Explain => {
120 r#"
121Command: EXPLAIN
122Description: show the execution plan of a statement
123Syntax:
124EXPLAIN [ ANALYZE ] statement
125"#
126 }
127 Function::Show => {
128 r#"
129Command: SHOW
130Description: show the value of a run-time parameter
131Syntax:
132SHOW name
133"#
134 }
135 Function::CreateTable => {
136 r#"
137Command: CREATE TABLE
138Description: define a new table
139Syntax:
140CREATE [ EXTERNAL ] TABLE table_name ( [
141 { column_name data_type }
142 [, ... ]
143] )
144"#
145 }
146 Function::CreateTableAs => {
147 r#"
148Command: CREATE TABLE AS
149Description: define a new table from the results of a query
150Syntax:
151CREATE TABLE table_name
152 [ (column_name [, ...] ) ]
153 AS query
154 [ WITH [ NO ] DATA ]
155"#
156 }
157 Function::Insert => {
158 r#"
159Command: INSERT
160Description: create new rows in a table
161Syntax:
162INSERT INTO table_name [ ( column_name [, ...] ) ]
163 { VALUES ( { expression } [, ...] ) [, ...] }
164"#
165 }
166 Function::DropTable => {
167 r#"
168Command: DROP TABLE
169Description: remove a table
170Syntax:
171DROP TABLE [ IF EXISTS ] name [, ...]
172"#
173 }
174 };
175 Ok(details)
176 }
177}
178
179impl FromStr for Function {
180 type Err = ();
181
182 fn from_str(s: &str) -> Result<Self, Self::Err> {
183 Ok(match s.trim().to_uppercase().as_str() {
184 "SELECT" => Self::Select,
185 "EXPLAIN" => Self::Explain,
186 "SHOW" => Self::Show,
187 "CREATE TABLE" => Self::CreateTable,
188 "CREATE TABLE AS" => Self::CreateTableAs,
189 "INSERT" => Self::Insert,
190 "DROP TABLE" => Self::DropTable,
191 _ => return Err(()),
192 })
193 }
194}
195
196impl fmt::Display for Function {
197 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
198 match *self {
199 Function::Select => write!(f, "SELECT"),
200 Function::Explain => write!(f, "EXPLAIN"),
201 Function::Show => write!(f, "SHOW"),
202 Function::CreateTable => write!(f, "CREATE TABLE"),
203 Function::CreateTableAs => write!(f, "CREATE TABLE AS"),
204 Function::Insert => write!(f, "INSERT"),
205 Function::DropTable => write!(f, "DROP TABLE"),
206 }
207 }
208}
209
210pub fn display_all_functions() -> Result<()> {
211 println!("Available help:");
212 let array = StringArray::from(
213 ALL_FUNCTIONS
214 .iter()
215 .map(|f| format!("{f}"))
216 .collect::<Vec<String>>(),
217 );
218 let schema = Schema::new(vec![Field::new("Function", DataType::Utf8, false)]);
219 let batch = RecordBatch::try_new(Arc::new(schema), vec![Arc::new(array)])?;
220 println!("{}", pretty_format_batches(&[batch]).unwrap());
221 Ok(())
222}
223
224#[derive(Debug)]
226struct ParquetMetadataTable {
227 schema: SchemaRef,
228 batch: RecordBatch,
229}
230
231#[async_trait]
232impl TableProvider for ParquetMetadataTable {
233 fn schema(&self) -> SchemaRef {
234 self.schema.clone()
235 }
236
237 fn table_type(&self) -> datafusion::logical_expr::TableType {
238 datafusion::logical_expr::TableType::Base
239 }
240
241 async fn scan(
242 &self,
243 _state: &dyn Session,
244 projection: Option<&Vec<usize>>,
245 _filters: &[Expr],
246 _limit: Option<usize>,
247 ) -> Result<Arc<dyn ExecutionPlan>> {
248 Ok(MemorySourceConfig::try_new_exec(
249 &[vec![self.batch.clone()]],
250 TableProvider::schema(self),
251 projection.cloned(),
252 )?)
253 }
254}
255
256fn convert_parquet_statistics(
257 value: &Statistics,
258 converted_type: ConvertedType,
259) -> (Option<String>, Option<String>) {
260 match (value, converted_type) {
261 (Statistics::Boolean(val), _) => (
262 val.min_opt().map(|v| v.to_string()),
263 val.max_opt().map(|v| v.to_string()),
264 ),
265 (Statistics::Int32(val), _) => (
266 val.min_opt().map(|v| v.to_string()),
267 val.max_opt().map(|v| v.to_string()),
268 ),
269 (Statistics::Int64(val), _) => (
270 val.min_opt().map(|v| v.to_string()),
271 val.max_opt().map(|v| v.to_string()),
272 ),
273 (Statistics::Int96(val), _) => (
274 val.min_opt().map(|v| v.to_string()),
275 val.max_opt().map(|v| v.to_string()),
276 ),
277 (Statistics::Float(val), _) => (
278 val.min_opt().map(|v| v.to_string()),
279 val.max_opt().map(|v| v.to_string()),
280 ),
281 (Statistics::Double(val), _) => (
282 val.min_opt().map(|v| v.to_string()),
283 val.max_opt().map(|v| v.to_string()),
284 ),
285 (Statistics::ByteArray(val), ConvertedType::UTF8) => (
286 byte_array_to_string(val.min_opt()),
287 byte_array_to_string(val.max_opt()),
288 ),
289 (Statistics::ByteArray(val), _) => (
290 val.min_opt().map(|v| v.to_string()),
291 val.max_opt().map(|v| v.to_string()),
292 ),
293 (Statistics::FixedLenByteArray(val), ConvertedType::UTF8) => (
294 fixed_len_byte_array_to_string(val.min_opt()),
295 fixed_len_byte_array_to_string(val.max_opt()),
296 ),
297 (Statistics::FixedLenByteArray(val), _) => (
298 val.min_opt().map(|v| v.to_string()),
299 val.max_opt().map(|v| v.to_string()),
300 ),
301 }
302}
303
304fn byte_array_to_string(val: Option<&ByteArray>) -> Option<String> {
306 val.map(|v| {
307 v.as_utf8()
308 .map(|s| s.to_string())
309 .unwrap_or_else(|_e| v.to_string())
310 })
311}
312
313fn fixed_len_byte_array_to_string(val: Option<&FixedLenByteArray>) -> Option<String> {
315 val.map(|v| {
316 v.as_utf8()
317 .map(|s| s.to_string())
318 .unwrap_or_else(|_e| v.to_string())
319 })
320}
321
322#[derive(Debug)]
323pub struct ParquetMetadataFunc {}
324
325impl TableFunctionImpl for ParquetMetadataFunc {
326 fn call_with_args(&self, args: TableFunctionArgs) -> Result<Arc<dyn TableProvider>> {
327 let exprs = args.exprs();
328 let filename = match exprs.first() {
329 Some(Expr::Literal(ScalarValue::Utf8(Some(s)), _)) => s, Some(Expr::Column(Column { name, .. })) => name, _ => {
332 return plan_err!(
333 "parquet_metadata requires string argument as its input"
334 );
335 }
336 };
337
338 let file = File::open(filename.clone())?;
339 let reader = SerializedFileReader::new(file)?;
340 let metadata = reader.metadata();
341
342 let schema = Arc::new(Schema::new(vec![
343 Field::new("filename", DataType::Utf8, true),
344 Field::new("row_group_id", DataType::Int64, true),
345 Field::new("row_group_num_rows", DataType::Int64, true),
346 Field::new("row_group_num_columns", DataType::Int64, true),
347 Field::new("row_group_bytes", DataType::Int64, true),
348 Field::new("column_id", DataType::Int64, true),
349 Field::new("file_offset", DataType::Int64, true),
350 Field::new("num_values", DataType::Int64, true),
351 Field::new("path_in_schema", DataType::Utf8, true),
352 Field::new("type", DataType::Utf8, true),
353 Field::new("stats_min", DataType::Utf8, true),
354 Field::new("stats_max", DataType::Utf8, true),
355 Field::new("stats_null_count", DataType::Int64, true),
356 Field::new("stats_distinct_count", DataType::Int64, true),
357 Field::new("stats_min_value", DataType::Utf8, true),
358 Field::new("stats_max_value", DataType::Utf8, true),
359 Field::new("compression", DataType::Utf8, true),
360 Field::new("encodings", DataType::Utf8, true),
361 Field::new("index_page_offset", DataType::Int64, true),
362 Field::new("dictionary_page_offset", DataType::Int64, true),
363 Field::new("data_page_offset", DataType::Int64, true),
364 Field::new("total_compressed_size", DataType::Int64, true),
365 Field::new("total_uncompressed_size", DataType::Int64, true),
366 ]));
367
368 let mut filename_arr = vec![];
370 let mut row_group_id_arr = vec![];
371 let mut row_group_num_rows_arr = vec![];
372 let mut row_group_num_columns_arr = vec![];
373 let mut row_group_bytes_arr = vec![];
374 let mut column_id_arr = vec![];
375 let mut file_offset_arr = vec![];
376 let mut num_values_arr = vec![];
377 let mut path_in_schema_arr = vec![];
378 let mut type_arr = vec![];
379 let mut stats_min_arr = vec![];
380 let mut stats_max_arr = vec![];
381 let mut stats_null_count_arr = vec![];
382 let mut stats_distinct_count_arr = vec![];
383 let mut stats_min_value_arr = vec![];
384 let mut stats_max_value_arr = vec![];
385 let mut compression_arr = vec![];
386 let mut encodings_arr = vec![];
387 let mut index_page_offset_arr = vec![];
388 let mut dictionary_page_offset_arr = vec![];
389 let mut data_page_offset_arr = vec![];
390 let mut total_compressed_size_arr = vec![];
391 let mut total_uncompressed_size_arr = vec![];
392 for (rg_idx, row_group) in metadata.row_groups().iter().enumerate() {
393 for (col_idx, column) in row_group.columns().iter().enumerate() {
394 filename_arr.push(filename.clone());
395 row_group_id_arr.push(rg_idx as i64);
396 row_group_num_rows_arr.push(row_group.num_rows());
397 row_group_num_columns_arr.push(row_group.num_columns() as i64);
398 row_group_bytes_arr.push(row_group.total_byte_size());
399 column_id_arr.push(col_idx as i64);
400 file_offset_arr.push(column.file_offset());
401 num_values_arr.push(column.num_values());
402 path_in_schema_arr.push(column.column_path().to_string());
403 type_arr.push(column.column_type().to_string());
404 let converted_type = column.column_descr().converted_type();
405
406 if let Some(s) = column.statistics() {
407 let (min_val, max_val) =
408 convert_parquet_statistics(s, converted_type);
409 stats_min_arr.push(min_val.clone());
410 stats_max_arr.push(max_val.clone());
411 stats_null_count_arr.push(s.null_count_opt().map(|c| c as i64));
412 stats_distinct_count_arr
413 .push(s.distinct_count_opt().map(|c| c as i64));
414 stats_min_value_arr.push(min_val);
415 stats_max_value_arr.push(max_val);
416 } else {
417 stats_min_arr.push(None);
418 stats_max_arr.push(None);
419 stats_null_count_arr.push(None);
420 stats_distinct_count_arr.push(None);
421 stats_min_value_arr.push(None);
422 stats_max_value_arr.push(None);
423 };
424 compression_arr.push(format!("{:?}", column.compression()));
425 let encodings: Vec<_> = column.encodings().collect();
427 encodings_arr.push(format!("{encodings:?}"));
428 index_page_offset_arr.push(column.index_page_offset());
429 dictionary_page_offset_arr.push(column.dictionary_page_offset());
430 data_page_offset_arr.push(column.data_page_offset());
431 total_compressed_size_arr.push(column.compressed_size());
432 total_uncompressed_size_arr.push(column.uncompressed_size());
433 }
434 }
435
436 let rb = RecordBatch::try_new(
437 schema.clone(),
438 vec![
439 Arc::new(StringArray::from(filename_arr)),
440 Arc::new(Int64Array::from(row_group_id_arr)),
441 Arc::new(Int64Array::from(row_group_num_rows_arr)),
442 Arc::new(Int64Array::from(row_group_num_columns_arr)),
443 Arc::new(Int64Array::from(row_group_bytes_arr)),
444 Arc::new(Int64Array::from(column_id_arr)),
445 Arc::new(Int64Array::from(file_offset_arr)),
446 Arc::new(Int64Array::from(num_values_arr)),
447 Arc::new(StringArray::from(path_in_schema_arr)),
448 Arc::new(StringArray::from(type_arr)),
449 Arc::new(StringArray::from(stats_min_arr)),
450 Arc::new(StringArray::from(stats_max_arr)),
451 Arc::new(Int64Array::from(stats_null_count_arr)),
452 Arc::new(Int64Array::from(stats_distinct_count_arr)),
453 Arc::new(StringArray::from(stats_min_value_arr)),
454 Arc::new(StringArray::from(stats_max_value_arr)),
455 Arc::new(StringArray::from(compression_arr)),
456 Arc::new(StringArray::from(encodings_arr)),
457 Arc::new(Int64Array::from(index_page_offset_arr)),
458 Arc::new(Int64Array::from(dictionary_page_offset_arr)),
459 Arc::new(Int64Array::from(data_page_offset_arr)),
460 Arc::new(Int64Array::from(total_compressed_size_arr)),
461 Arc::new(Int64Array::from(total_uncompressed_size_arr)),
462 ],
463 )?;
464
465 let parquet_metadata = ParquetMetadataTable { schema, batch: rb };
466 Ok(Arc::new(parquet_metadata))
467 }
468}
469
470#[derive(Debug)]
472struct MetadataCacheTable {
473 schema: SchemaRef,
474 batch: RecordBatch,
475}
476
477#[async_trait]
478impl TableProvider for MetadataCacheTable {
479 fn schema(&self) -> SchemaRef {
480 self.schema.clone()
481 }
482
483 fn table_type(&self) -> datafusion::logical_expr::TableType {
484 datafusion::logical_expr::TableType::Base
485 }
486
487 async fn scan(
488 &self,
489 _state: &dyn Session,
490 projection: Option<&Vec<usize>>,
491 _filters: &[Expr],
492 _limit: Option<usize>,
493 ) -> Result<Arc<dyn ExecutionPlan>> {
494 Ok(MemorySourceConfig::try_new_exec(
495 &[vec![self.batch.clone()]],
496 TableProvider::schema(self),
497 projection.cloned(),
498 )?)
499 }
500}
501
502#[derive(Debug)]
503pub struct MetadataCacheFunc {
504 cache_manager: Arc<CacheManager>,
505}
506
507impl MetadataCacheFunc {
508 pub fn new(cache_manager: Arc<CacheManager>) -> Self {
509 Self { cache_manager }
510 }
511}
512
513impl TableFunctionImpl for MetadataCacheFunc {
514 fn call_with_args(&self, args: TableFunctionArgs) -> Result<Arc<dyn TableProvider>> {
515 let exprs = args.exprs();
516 if !exprs.is_empty() {
517 return plan_err!("metadata_cache should have no arguments");
518 }
519
520 let schema = Arc::new(Schema::new(vec![
521 Field::new("path", DataType::Utf8, false),
522 Field::new(
523 "file_modified",
524 DataType::Timestamp(TimeUnit::Millisecond, None),
525 false,
526 ),
527 Field::new("file_size_bytes", DataType::UInt64, false),
528 Field::new("e_tag", DataType::Utf8, true),
529 Field::new("version", DataType::Utf8, true),
530 Field::new("metadata_size_bytes", DataType::UInt64, false),
531 Field::new("hits", DataType::UInt64, false),
532 Field::new("extra", DataType::Utf8, true),
533 ]));
534
535 let mut path_arr = vec![];
537 let mut file_modified_arr = vec![];
538 let mut file_size_bytes_arr = vec![];
539 let mut e_tag_arr = vec![];
540 let mut version_arr = vec![];
541 let mut metadata_size_bytes = vec![];
542 let mut hits_arr = vec![];
543 let mut extra_arr = vec![];
544
545 let cached_entries = self.cache_manager.get_file_metadata_cache().list_entries();
546
547 for (path, entry) in cached_entries {
548 path_arr.push(path.to_string());
549 file_modified_arr
550 .push(Some(entry.value.meta.last_modified.timestamp_millis()));
551 file_size_bytes_arr.push(entry.value.meta.size);
552 e_tag_arr.push(entry.value.meta.e_tag);
553 version_arr.push(entry.value.meta.version);
554 metadata_size_bytes.push(entry.size_bytes as u64);
555 hits_arr.push(entry.hits as u64);
556
557 let mut extra = entry
558 .value
559 .file_metadata
560 .extra_info()
561 .iter()
562 .map(|(k, v)| format!("{k}={v}"))
563 .collect::<Vec<_>>();
564 extra.sort();
565 extra_arr.push(extra.join(" "));
566 }
567
568 let batch = RecordBatch::try_new(
569 schema.clone(),
570 vec![
571 Arc::new(StringArray::from(path_arr)),
572 Arc::new(TimestampMillisecondArray::from(file_modified_arr)),
573 Arc::new(UInt64Array::from(file_size_bytes_arr)),
574 Arc::new(StringArray::from(e_tag_arr)),
575 Arc::new(StringArray::from(version_arr)),
576 Arc::new(UInt64Array::from(metadata_size_bytes)),
577 Arc::new(UInt64Array::from(hits_arr)),
578 Arc::new(StringArray::from(extra_arr)),
579 ],
580 )?;
581
582 let metadata_cache = MetadataCacheTable { schema, batch };
583 Ok(Arc::new(metadata_cache))
584 }
585}
586
587#[derive(Debug)]
589struct StatisticsCacheTable {
590 schema: SchemaRef,
591 batch: RecordBatch,
592}
593
594#[async_trait]
595impl TableProvider for StatisticsCacheTable {
596 fn schema(&self) -> SchemaRef {
597 self.schema.clone()
598 }
599
600 fn table_type(&self) -> datafusion::logical_expr::TableType {
601 datafusion::logical_expr::TableType::Base
602 }
603
604 async fn scan(
605 &self,
606 _state: &dyn Session,
607 projection: Option<&Vec<usize>>,
608 _filters: &[Expr],
609 _limit: Option<usize>,
610 ) -> Result<Arc<dyn ExecutionPlan>> {
611 Ok(MemorySourceConfig::try_new_exec(
612 &[vec![self.batch.clone()]],
613 TableProvider::schema(self),
614 projection.cloned(),
615 )?)
616 }
617}
618
619#[derive(Debug)]
620pub struct StatisticsCacheFunc {
621 cache_manager: Arc<CacheManager>,
622}
623
624impl StatisticsCacheFunc {
625 pub fn new(cache_manager: Arc<CacheManager>) -> Self {
626 Self { cache_manager }
627 }
628}
629
630impl TableFunctionImpl for StatisticsCacheFunc {
631 fn call_with_args(&self, args: TableFunctionArgs) -> Result<Arc<dyn TableProvider>> {
632 let exprs = args.exprs();
633 if !exprs.is_empty() {
634 return plan_err!("statistics_cache should have no arguments");
635 }
636
637 let schema = Arc::new(Schema::new(vec![
638 Field::new("path", DataType::Utf8, false),
639 Field::new("table", DataType::Utf8, false),
640 Field::new(
641 "file_modified",
642 DataType::Timestamp(TimeUnit::Millisecond, None),
643 false,
644 ),
645 Field::new("file_size_bytes", DataType::UInt64, false),
646 Field::new("e_tag", DataType::Utf8, true),
647 Field::new("version", DataType::Utf8, true),
648 Field::new("num_rows", DataType::Utf8, false),
649 Field::new("num_columns", DataType::UInt64, false),
650 Field::new("table_size_bytes", DataType::Utf8, false),
651 Field::new("statistics_size_bytes", DataType::UInt64, false),
652 Field::new("hits", DataType::UInt64, false),
653 ]));
654
655 let mut path_arr = vec![];
657 let mut table_arr = vec![];
658 let mut file_modified_arr = vec![];
659 let mut file_size_bytes_arr = vec![];
660 let mut e_tag_arr = vec![];
661 let mut version_arr = vec![];
662 let mut num_rows_arr = vec![];
663 let mut num_columns_arr = vec![];
664 let mut table_size_bytes_arr = vec![];
665 let mut statistics_size_bytes_arr = vec![];
666 let mut hits_arr = vec![];
667
668 if let Some(file_statistics_cache) = self.cache_manager.get_file_statistic_cache()
669 {
670 for (path, entry) in file_statistics_cache.list_entries() {
671 path_arr.push(path.path.to_string());
672 table_arr
673 .push(path.table.map_or_else(|| "".to_string(), |t| t.to_string()));
674 file_modified_arr
675 .push(Some(entry.value.meta.last_modified.timestamp_millis()));
676 file_size_bytes_arr.push(entry.value.meta.size);
677 e_tag_arr.push(entry.value.meta.e_tag);
678 version_arr.push(entry.value.meta.version);
679 num_rows_arr.push(entry.value.statistics.num_rows.to_string());
680 num_columns_arr
681 .push(entry.value.statistics.column_statistics.len() as u64);
682 table_size_bytes_arr
683 .push(entry.value.statistics.total_byte_size.to_string());
684 statistics_size_bytes_arr.push(
685 entry
686 .value
687 .statistics
688 .heap_size(&mut DFHeapSizeCtx::default())
689 as u64,
690 );
691 hits_arr.push(entry.hits as u64);
692 }
693 }
694
695 let batch = RecordBatch::try_new(
696 schema.clone(),
697 vec![
698 Arc::new(StringArray::from(path_arr)),
699 Arc::new(StringArray::from(table_arr)),
700 Arc::new(TimestampMillisecondArray::from(file_modified_arr)),
701 Arc::new(UInt64Array::from(file_size_bytes_arr)),
702 Arc::new(StringArray::from(e_tag_arr)),
703 Arc::new(StringArray::from(version_arr)),
704 Arc::new(StringArray::from(num_rows_arr)),
705 Arc::new(UInt64Array::from(num_columns_arr)),
706 Arc::new(StringArray::from(table_size_bytes_arr)),
707 Arc::new(UInt64Array::from(statistics_size_bytes_arr)),
708 Arc::new(UInt64Array::from(hits_arr)),
709 ],
710 )?;
711
712 let statistics_cache = StatisticsCacheTable { schema, batch };
713 Ok(Arc::new(statistics_cache))
714 }
715}
716
717#[derive(Debug)]
738struct ListFilesCacheTable {
739 schema: SchemaRef,
740 batch: RecordBatch,
741}
742
743#[async_trait]
744impl TableProvider for ListFilesCacheTable {
745 fn schema(&self) -> SchemaRef {
746 self.schema.clone()
747 }
748
749 fn table_type(&self) -> datafusion::logical_expr::TableType {
750 datafusion::logical_expr::TableType::Base
751 }
752
753 async fn scan(
754 &self,
755 _state: &dyn Session,
756 projection: Option<&Vec<usize>>,
757 _filters: &[Expr],
758 _limit: Option<usize>,
759 ) -> Result<Arc<dyn ExecutionPlan>> {
760 Ok(MemorySourceConfig::try_new_exec(
761 &[vec![self.batch.clone()]],
762 TableProvider::schema(self),
763 projection.cloned(),
764 )?)
765 }
766}
767
768#[derive(Debug)]
769pub struct ListFilesCacheFunc {
770 cache_manager: Arc<CacheManager>,
771}
772
773impl ListFilesCacheFunc {
774 pub fn new(cache_manager: Arc<CacheManager>) -> Self {
775 Self { cache_manager }
776 }
777}
778
779impl TableFunctionImpl for ListFilesCacheFunc {
780 fn call_with_args(&self, args: TableFunctionArgs) -> Result<Arc<dyn TableProvider>> {
781 let exprs = args.exprs();
782 if !exprs.is_empty() {
783 return plan_err!("list_files_cache should have no arguments");
784 }
785
786 let nested_fields = Fields::from(vec![
787 Field::new("file_path", DataType::Utf8, false),
788 Field::new(
789 "file_modified",
790 DataType::Timestamp(TimeUnit::Millisecond, None),
791 false,
792 ),
793 Field::new("file_size_bytes", DataType::UInt64, false),
794 Field::new("e_tag", DataType::Utf8, true),
795 Field::new("version", DataType::Utf8, true),
796 ]);
797
798 let metadata_field =
799 Field::new("metadata", DataType::Struct(nested_fields.clone()), true);
800
801 let schema = Arc::new(Schema::new(vec![
802 Field::new("table", DataType::Utf8, true),
803 Field::new("path", DataType::Utf8, false),
804 Field::new("metadata_size_bytes", DataType::UInt64, false),
805 Field::new(
807 "expires_in",
808 DataType::Duration(TimeUnit::Millisecond),
809 true,
810 ),
811 Field::new(
812 "metadata_list",
813 DataType::List(Arc::new(metadata_field.clone())),
814 true,
815 ),
816 Field::new("hits", DataType::UInt64, false),
817 ]));
818
819 let mut table_arr = vec![];
820 let mut path_arr = vec![];
821 let mut metadata_size_bytes_arr = vec![];
822 let mut expires_arr = vec![];
823
824 let mut file_path_arr = vec![];
825 let mut file_modified_arr = vec![];
826 let mut file_size_bytes_arr = vec![];
827 let mut etag_arr = vec![];
828 let mut version_arr = vec![];
829 let mut offsets: Vec<i32> = vec![0];
830 let mut hits_arr = vec![];
831
832 if let Some(list_files_cache) = self.cache_manager.get_list_files_cache() {
833 let now = Instant::now();
834 let mut current_offset: i32 = 0;
835
836 for (path, entry) in list_files_cache.list_entries() {
837 table_arr.push(path.table.map(|t| t.to_string()));
838 path_arr.push(path.path.to_string());
839 metadata_size_bytes_arr.push(entry.size_bytes as u64);
840 expires_arr.push(
842 entry
843 .expires
844 .map(|t| t.duration_since(now).as_millis() as i64),
845 );
846
847 for meta in entry.value.files.iter() {
848 file_path_arr.push(meta.location.to_string());
849 file_modified_arr.push(meta.last_modified.timestamp_millis());
850 file_size_bytes_arr.push(meta.size);
851 etag_arr.push(meta.e_tag.clone());
852 version_arr.push(meta.version.clone());
853 }
854 current_offset += entry.value.files.len() as i32;
855 offsets.push(current_offset);
856 hits_arr.push(entry.hits as u64);
857 }
858 }
859
860 let struct_arr = StructArray::new(
861 nested_fields,
862 vec![
863 Arc::new(StringArray::from(file_path_arr)),
864 Arc::new(TimestampMillisecondArray::from(file_modified_arr)),
865 Arc::new(UInt64Array::from(file_size_bytes_arr)),
866 Arc::new(StringArray::from(etag_arr)),
867 Arc::new(StringArray::from(version_arr)),
868 ],
869 None,
870 );
871
872 let offsets_buffer: OffsetBuffer<i32> =
873 OffsetBuffer::new(ScalarBuffer::from(Buffer::from_vec(offsets)));
874
875 let batch = RecordBatch::try_new(
876 schema.clone(),
877 vec![
878 Arc::new(StringArray::from(table_arr)),
879 Arc::new(StringArray::from(path_arr)),
880 Arc::new(UInt64Array::from(metadata_size_bytes_arr)),
881 Arc::new(DurationMillisecondArray::from(expires_arr)),
882 Arc::new(GenericListArray::new(
883 Arc::new(metadata_field),
884 offsets_buffer,
885 Arc::new(struct_arr),
886 None,
887 )),
888 Arc::new(UInt64Array::from(hits_arr)),
889 ],
890 )?;
891
892 let list_files_cache = ListFilesCacheTable { schema, batch };
893 Ok(Arc::new(list_files_cache))
894 }
895}