Skip to main content

datafusion_federation/
table_provider.rs

1use std::{borrow::Cow, sync::Arc};
2
3use async_trait::async_trait;
4use datafusion::{
5    arrow::datatypes::SchemaRef,
6    catalog::Session,
7    common::Constraints,
8    datasource::TableProvider,
9    error::{DataFusionError, Result},
10    logical_expr::{
11        dml::InsertOp, Expr, LogicalPlan, TableProviderFilterPushDown, TableSource, TableType,
12    },
13    physical_plan::ExecutionPlan,
14};
15
16use crate::FederationProvider;
17
18// FederatedTableSourceWrapper helps to recover the FederatedTableSource
19// from a TableScan. This wrapper may be avoidable.
20#[derive(Debug)]
21pub struct FederatedTableProviderAdaptor {
22    pub source: Arc<dyn FederatedTableSource>,
23    pub table_provider: Option<Arc<dyn TableProvider>>,
24}
25
26impl FederatedTableProviderAdaptor {
27    pub fn new(source: Arc<dyn FederatedTableSource>) -> Self {
28        Self {
29            source,
30            table_provider: None,
31        }
32    }
33
34    /// Creates a new FederatedTableProviderAdaptor that falls back to the
35    /// provided TableProvider. This is useful if used within a DataFusion
36    /// context without the federation optimizer.
37    pub fn new_with_provider(
38        source: Arc<dyn FederatedTableSource>,
39        table_provider: Arc<dyn TableProvider>,
40    ) -> Self {
41        Self {
42            source,
43            table_provider: Some(table_provider),
44        }
45    }
46}
47
48#[async_trait]
49impl TableProvider for FederatedTableProviderAdaptor {
50    fn schema(&self) -> SchemaRef {
51        if let Some(table_provider) = &self.table_provider {
52            return table_provider.schema();
53        }
54
55        self.source.schema()
56    }
57    fn constraints(&self) -> Option<&Constraints> {
58        if let Some(table_provider) = &self.table_provider {
59            return table_provider
60                .constraints()
61                .or_else(|| self.source.constraints());
62        }
63
64        self.source.constraints()
65    }
66    fn table_type(&self) -> TableType {
67        if let Some(table_provider) = &self.table_provider {
68            return table_provider.table_type();
69        }
70
71        self.source.table_type()
72    }
73    fn get_logical_plan(&self) -> Option<Cow<'_, LogicalPlan>> {
74        if let Some(table_provider) = &self.table_provider {
75            return table_provider
76                .get_logical_plan()
77                .or_else(|| self.source.get_logical_plan());
78        }
79
80        self.source.get_logical_plan()
81    }
82    fn get_column_default(&self, column: &str) -> Option<&Expr> {
83        if let Some(table_provider) = &self.table_provider {
84            return table_provider
85                .get_column_default(column)
86                .or_else(|| self.source.get_column_default(column));
87        }
88
89        self.source.get_column_default(column)
90    }
91    fn supports_filters_pushdown(
92        &self,
93        filters: &[&Expr],
94    ) -> Result<Vec<TableProviderFilterPushDown>> {
95        if let Some(table_provider) = &self.table_provider {
96            return table_provider.supports_filters_pushdown(filters);
97        }
98
99        Ok(vec![
100            TableProviderFilterPushDown::Unsupported;
101            filters.len()
102        ])
103    }
104
105    // Scan is not supported; the adaptor should be replaced
106    // with a virtual TableProvider that provides federation for a sub-plan.
107    async fn scan(
108        &self,
109        state: &dyn Session,
110        projection: Option<&Vec<usize>>,
111        filters: &[Expr],
112        limit: Option<usize>,
113    ) -> Result<Arc<dyn ExecutionPlan>> {
114        if let Some(table_provider) = &self.table_provider {
115            return table_provider.scan(state, projection, filters, limit).await;
116        }
117
118        Err(DataFusionError::NotImplemented(
119            "FederatedTableProviderAdaptor cannot scan".to_string(),
120        ))
121    }
122
123    async fn insert_into(
124        &self,
125        _state: &dyn Session,
126        input: Arc<dyn ExecutionPlan>,
127        insert_op: InsertOp,
128    ) -> Result<Arc<dyn ExecutionPlan>> {
129        if let Some(table_provider) = &self.table_provider {
130            return table_provider.insert_into(_state, input, insert_op).await;
131        }
132
133        Err(DataFusionError::NotImplemented(
134            "FederatedTableProviderAdaptor cannot insert_into".to_string(),
135        ))
136    }
137}
138
139// FederatedTableProvider extends DataFusion's TableProvider trait
140// to allow grouping of TableScans of the same FederationProvider.
141#[async_trait]
142pub trait FederatedTableSource: TableSource {
143    // Return the FederationProvider associated with this Table
144    fn federation_provider(&self) -> Arc<dyn FederationProvider>;
145}
146
147impl std::fmt::Debug for dyn FederatedTableSource {
148    fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
149        write!(
150            f,
151            "FederatedTableSource: {:?}",
152            self.federation_provider().name()
153        )
154    }
155}