Reload table information from ObjectStore
(&self, state: &dyn Session)
| 88 | |
| 89 | /// Reload table information from ObjectStore |
| 90 | pub async fn refresh(&self, state: &dyn Session) -> datafusion_common::Result<()> { |
| 91 | let entries: Vec<_> = self.store.list(Some(&self.path)).try_collect().await?; |
| 92 | let base = Path::new(self.path.as_ref()); |
| 93 | let mut tables = HashSet::new(); |
| 94 | for file in entries.iter() { |
| 95 | // The listing will initially be a file. However if we've recursed up to match our base, we know our path is a directory. |
| 96 | let mut is_dir = false; |
| 97 | let mut parent = Path::new(file.location.as_ref()); |
| 98 | while let Some(p) = parent.parent() { |
| 99 | if p == base { |
| 100 | tables.insert(TablePath { |
| 101 | is_dir, |
| 102 | path: parent, |
| 103 | }); |
| 104 | } |
| 105 | parent = p; |
| 106 | is_dir = true; |
| 107 | } |
| 108 | } |
| 109 | for table in tables.iter() { |
| 110 | let file_name = table |
| 111 | .path |
| 112 | .file_name() |
| 113 | .ok_or_else(|| internal_datafusion_err!("Cannot parse file name!"))? |
| 114 | .to_str() |
| 115 | .ok_or_else(|| internal_datafusion_err!("Cannot parse file name!"))?; |
| 116 | let table_name = file_name.split('.').collect_vec()[0]; |
| 117 | let table_path = table |
| 118 | .to_string() |
| 119 | .ok_or_else(|| internal_datafusion_err!("Cannot parse file name!"))?; |
| 120 | |
| 121 | if !self.table_exist(table_name) { |
| 122 | let table_url = format!("{}/{}", self.authority, table_path); |
| 123 | |
| 124 | let name = TableReference::bare(table_name); |
| 125 | let provider = self |
| 126 | .factory |
| 127 | .create( |
| 128 | state, |
| 129 | &CreateExternalTable::builder( |
| 130 | name, |
| 131 | table_url, |
| 132 | self.format.clone(), |
| 133 | Arc::new(DFSchema::empty()), |
| 134 | ) |
| 135 | .build(), |
| 136 | ) |
| 137 | .await?; |
| 138 | let _ = |
| 139 | self.register_table(table_name.to_string(), Arc::clone(&provider))?; |
| 140 | } |
| 141 | } |
| 142 | Ok(()) |
| 143 | } |
| 144 | } |
| 145 | |
| 146 | #[async_trait] |
no test coverage detected