Skip to main content

datafusion_cli/
functions.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18//! Functions that are query-able and searchable via the `\h` command
19
20use 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/// PARQUET_META table function
225#[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
304/// Convert to a string if it has utf8 encoding, otherwise print bytes directly
305fn 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
313/// Convert to a string if it has utf8 encoding, otherwise print bytes directly
314fn 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, // single quote: parquet_metadata('x.parquet')
330            Some(Expr::Column(Column { name, .. })) => name, // double quote: parquet_metadata("x.parquet")
331            _ => {
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        // construct record batch from metadata
369        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                // need to collect into Vec to format
426                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/// METADATA_CACHE table function
471#[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        // construct record batch from metadata
536        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/// STATISTICS_CACHE table function
588#[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        // construct record batch from metadata
656        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/// Implementation of the `list_files_cache` table function in datafusion-cli.
718///
719/// This function returns the cached results of running a LIST command on a
720/// particular object store path for a table. The object metadata is returned as
721/// a List of Structs, with one Struct for each object. DataFusion uses these
722/// cached results to plan queries against external tables.
723///
724/// # Schema
725/// ```sql
726/// > describe select * from list_files_cache();
727/// +---------------------+--------------------------------------------------------------------------------------------------------------------------------------------------------------------------+-------------+
728/// | column_name         | data_type                                                                                                                                                                | is_nullable |
729/// +---------------------+--------------------------------------------------------------------------------------------------------------------------------------------------------------------------+-------------+
730/// | table               | Utf8                                                                                                                                                                     | NO          |
731/// | path                | Utf8                                                                                                                                                                     | NO          |
732/// | metadata_size_bytes | UInt64                                                                                                                                                                   | NO          |
733/// | expires_in          | Duration(ms)                                                                                                                                                             | YES         |
734/// | metadata_list       | List(Struct("file_path": non-null Utf8, "file_modified": non-null Timestamp(ms), "file_size_bytes": non-null UInt64, "e_tag": Utf8, "version": Utf8), field: 'metadata') | YES         |
735/// +---------------------+--------------------------------------------------------------------------------------------------------------------------------------------------------------------------+-------------+
736/// ```
737#[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            // expires field in ListFilesEntry has type Instant when set, from which we cannot get "the number of seconds", hence using Duration instead of Timestamp as data type.
806            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                // calculates time left before entry expires
841                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}