datafusion_federation/
table_provider.rs1use 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#[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 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 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#[async_trait]
142pub trait FederatedTableSource: TableSource {
143 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}