Skip to content

Commit d01c70f

Browse files
committed
Merge branch 'main' of https://github.com/apache/datafusion into feat/explode_outer
2 parents 1c7d262 + f7aef23 commit d01c70f

60 files changed

Lines changed: 2249 additions & 1662 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎.github/workflows/codeql.yml‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -45,11 +45,11 @@ jobs:
4545
persist-credentials: false
4646

4747
- name: Initialize CodeQL
48-
uses: github/codeql-action/init@99df26d4f13ea111d4ec1a7dddef6063f76b97e9 # v4
48+
uses: github/codeql-action/init@7188fc363630916deb702c7fdcf4e481b751f97a # v4
4949
with:
5050
languages: actions
5151

5252
- name: Perform CodeQL Analysis
53-
uses: github/codeql-action/analyze@99df26d4f13ea111d4ec1a7dddef6063f76b97e9 # v4
53+
uses: github/codeql-action/analyze@7188fc363630916deb702c7fdcf4e481b751f97a # v4
5454
with:
5555
category: "/language:actions"

‎Cargo.lock‎

Lines changed: 1 addition & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

‎Cargo.toml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -78,7 +78,7 @@ license = "Apache-2.0"
7878
readme = "README.md"
7979
repository = "https://github.com/apache/datafusion"
8080
# Define Minimum Supported Rust Version (MSRV)
81-
rust-version = "1.88.0"
81+
rust-version = "1.94.0"
8282
# Define DataFusion version
8383
version = "54.1.0"
8484

‎datafusion/catalog/src/catalog.rs‎

Lines changed: 5 additions & 192 deletions
Original file line numberDiff line numberDiff line change
@@ -15,195 +15,8 @@
1515
// specific language governing permissions and limitations
1616
// under the License.
1717

18-
use std::any::Any;
19-
use std::fmt::Debug;
20-
use std::sync::Arc;
21-
22-
pub use crate::schema::SchemaProvider;
23-
use datafusion_common::Result;
24-
use datafusion_common::not_impl_err;
25-
26-
/// Represents a catalog, comprising a number of named schemas.
27-
///
28-
/// # Catalog Overview
29-
///
30-
/// To plan and execute queries, DataFusion needs a "Catalog" that provides
31-
/// metadata such as which schemas and tables exist, their columns and data
32-
/// types, and how to access the data.
33-
///
34-
/// The Catalog API consists:
35-
/// * [`CatalogProviderList`]: a collection of `CatalogProvider`s
36-
/// * [`CatalogProvider`]: a collection of `SchemaProvider`s (sometimes called a "database" in other systems)
37-
/// * [`SchemaProvider`]: a collection of `TableProvider`s (often called a "schema" in other systems)
38-
/// * [`TableProvider`]: individual tables
39-
///
40-
/// # Implementing Catalogs
41-
///
42-
/// To implement a catalog, you implement at least one of the [`CatalogProviderList`],
43-
/// [`CatalogProvider`] and [`SchemaProvider`] traits and register them
44-
/// appropriately in the `SessionContext`.
45-
///
46-
/// DataFusion comes with a simple in-memory catalog implementation,
47-
/// `MemoryCatalogProvider`, that is used by default and has no persistence.
48-
/// DataFusion does not include more complex Catalog implementations because
49-
/// catalog management is a key design choice for most data systems, and thus
50-
/// it is unlikely that any general-purpose catalog implementation will work
51-
/// well across many use cases.
52-
///
53-
/// # Implementing "Remote" catalogs
54-
///
55-
/// See [`remote_catalog`] for an end to end example of how to implement a
56-
/// remote catalog.
57-
///
58-
/// Sometimes catalog information is stored remotely and requires a network call
59-
/// to retrieve. For example, the [Delta Lake] table format stores table
60-
/// metadata in files on S3 that must be first downloaded to discover what
61-
/// schemas and tables exist.
62-
///
63-
/// [Delta Lake]: https://delta.io/
64-
/// [`remote_catalog`]: https://github.com/apache/datafusion/blob/main/datafusion-examples/examples/data_io/remote_catalog.rs
65-
///
66-
/// The [`CatalogProvider`] can support this use case, but it takes some care.
67-
/// The planning APIs in DataFusion are not `async` and thus network IO can not
68-
/// be performed "lazily" / "on demand" during query planning. The rationale for
69-
/// this design is that using remote procedure calls for all catalog accesses
70-
/// required for query planning would likely result in multiple network calls
71-
/// per plan, resulting in very poor planning performance.
72-
///
73-
/// To implement [`CatalogProvider`] and [`SchemaProvider`] for remote catalogs,
74-
/// you need to provide an in memory snapshot of the required metadata. Most
75-
/// systems typically either already have this information cached locally or can
76-
/// batch access to the remote catalog to retrieve multiple schemas and tables
77-
/// in a single network call.
78-
///
79-
/// Note that [`SchemaProvider::table`] **is** an `async` function in order to
80-
/// simplify implementing simple [`SchemaProvider`]s. For many table formats it
81-
/// is easy to list all available tables but there is additional non trivial
82-
/// access required to read table details (e.g. statistics).
83-
///
84-
/// The pattern that DataFusion itself uses to plan SQL queries is to walk over
85-
/// the query to find all table references, performing required remote catalog
86-
/// lookups in parallel, storing the results in a cached snapshot, and then plans
87-
/// the query using that snapshot.
88-
///
89-
/// # Example Catalog Implementations
90-
///
91-
/// Here are some examples of how to implement custom catalogs:
92-
///
93-
/// * [`datafusion-cli`]: [`DynamicFileCatalogProvider`] catalog provider
94-
/// that treats files and directories on a filesystem as tables.
95-
///
96-
/// * The [`catalog.rs`]: a simple directory based catalog.
97-
///
98-
/// * [delta-rs]: [`UnityCatalogProvider`] implementation that can
99-
/// read from Delta Lake tables
100-
///
101-
/// [`datafusion-cli`]: https://datafusion.apache.org/user-guide/cli/index.html
102-
/// [`DynamicFileCatalogProvider`]: https://github.com/apache/datafusion/blob/31b9b48b08592b7d293f46e75707aad7dadd7cbc/datafusion-cli/src/catalog.rs#L75
103-
/// [`catalog.rs`]: https://github.com/apache/datafusion/blob/main/datafusion-examples/examples/data_io/catalog.rs
104-
/// [delta-rs]: https://github.com/delta-io/delta-rs
105-
/// [`UnityCatalogProvider`]: https://github.com/delta-io/delta-rs/blob/951436ecec476ce65b5ed3b58b50fb0846ca7b91/crates/deltalake-core/src/data_catalog/unity/datafusion.rs#L111-L123
106-
///
107-
/// [`TableProvider`]: crate::TableProvider
108-
pub trait CatalogProvider: Any + Debug + Sync + Send {
109-
/// Retrieves the list of available schema names in this catalog.
110-
fn schema_names(&self) -> Vec<String>;
111-
112-
/// Retrieves a specific schema from the catalog by name, provided it exists.
113-
fn schema(&self, name: &str) -> Option<Arc<dyn SchemaProvider>>;
114-
115-
/// Adds a new schema to this catalog.
116-
///
117-
/// If a schema of the same name existed before, it is replaced in
118-
/// the catalog and returned.
119-
///
120-
/// By default returns a "Not Implemented" error
121-
fn register_schema(
122-
&self,
123-
name: &str,
124-
schema: Arc<dyn SchemaProvider>,
125-
) -> Result<Option<Arc<dyn SchemaProvider>>> {
126-
// use variables to avoid unused variable warnings
127-
let _ = name;
128-
let _ = schema;
129-
not_impl_err!("Registering new schemas is not supported")
130-
}
131-
132-
/// Removes a schema from this catalog. Implementations of this method should return
133-
/// errors if the schema exists but cannot be dropped. For example, in DataFusion's
134-
/// default in-memory catalog, `MemoryCatalogProvider`, a non-empty schema
135-
/// will only be successfully dropped when `cascade` is true.
136-
/// This is equivalent to how DROP SCHEMA works in PostgreSQL.
137-
///
138-
/// Implementations of this method should return None if schema with `name`
139-
/// does not exist.
140-
///
141-
/// By default returns a "Not Implemented" error
142-
fn deregister_schema(
143-
&self,
144-
_name: &str,
145-
_cascade: bool,
146-
) -> Result<Option<Arc<dyn SchemaProvider>>> {
147-
not_impl_err!("Deregistering new schemas is not supported")
148-
}
149-
}
150-
151-
impl dyn CatalogProvider {
152-
/// Returns `true` if the catalog provider is of type `T`.
153-
///
154-
/// Prefer this over `downcast_ref::<T>().is_some()`. Works correctly when
155-
/// called on `Arc<dyn CatalogProvider>` via auto-deref.
156-
pub fn is<T: CatalogProvider>(&self) -> bool {
157-
(self as &dyn Any).is::<T>()
158-
}
159-
160-
/// Attempts to downcast this catalog provider to a concrete type `T`,
161-
/// returning `None` if the provider is not of that type.
162-
///
163-
/// Works correctly when called on `Arc<dyn CatalogProvider>` via auto-deref,
164-
/// unlike `(&arc as &dyn Any).downcast_ref::<T>()` which would attempt to
165-
/// downcast the `Arc` itself.
166-
pub fn downcast_ref<T: CatalogProvider>(&self) -> Option<&T> {
167-
(self as &dyn Any).downcast_ref()
168-
}
169-
}
170-
171-
/// Represent a list of named [`CatalogProvider`]s.
172-
///
173-
/// Please see the documentation on [`CatalogProvider`] for details of
174-
/// implementing a custom catalog.
175-
pub trait CatalogProviderList: Any + Debug + Sync + Send {
176-
/// Adds a new catalog to this catalog list
177-
/// If a catalog of the same name existed before, it is replaced in the list and returned.
178-
fn register_catalog(
179-
&self,
180-
name: String,
181-
catalog: Arc<dyn CatalogProvider>,
182-
) -> Option<Arc<dyn CatalogProvider>>;
183-
184-
/// Retrieves the list of available catalog names
185-
fn catalog_names(&self) -> Vec<String>;
186-
187-
/// Retrieves a specific catalog by name, provided it exists.
188-
fn catalog(&self, name: &str) -> Option<Arc<dyn CatalogProvider>>;
189-
}
190-
191-
impl dyn CatalogProviderList {
192-
/// Returns `true` if the catalog provider list is of type `T`.
193-
///
194-
/// Prefer this over `downcast_ref::<T>().is_some()`. Works correctly when
195-
/// called on `Arc<dyn CatalogProviderList>` via auto-deref.
196-
pub fn is<T: CatalogProviderList>(&self) -> bool {
197-
(self as &dyn Any).is::<T>()
198-
}
199-
200-
/// Attempts to downcast this catalog provider list to a concrete type `T`,
201-
/// returning `None` if the provider list is not of that type.
202-
///
203-
/// Works correctly when called on `Arc<dyn CatalogProviderList>` via
204-
/// auto-deref, unlike `(&arc as &dyn Any).downcast_ref::<T>()` which would
205-
/// attempt to downcast the `Arc` itself.
206-
pub fn downcast_ref<T: CatalogProviderList>(&self) -> Option<&T> {
207-
(self as &dyn Any).downcast_ref()
208-
}
209-
}
18+
// Re-export from this module for backwards compatibility.
19+
pub use datafusion_session::{CatalogProvider, CatalogProviderList};
20+
// Re-export so users can access this type through `datafusion_catalog` and
21+
// `datafusion::catalog` without depending directly on `datafusion_session`.
22+
pub use datafusion_session::EmptyCatalogProviderList;

‎datafusion/catalog/src/lib.rs‎

Lines changed: 4 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -25,7 +25,10 @@
2525
#![cfg_attr(not(test), deny(clippy::clone_on_ref_ptr))]
2626
#![cfg_attr(test, allow(clippy::needless_pass_by_value))]
2727

28-
//! Interfaces and default implementations of catalogs and schemas.
28+
//! Default implementations of catalogs and schemas.
29+
//!
30+
//! The catalog interfaces are defined in [`datafusion_session`] and re-exported
31+
//! by this crate.
2932
//!
3033
//! Implementations
3134
//! * Information schema: [`information_schema`]
@@ -57,8 +60,3 @@ pub use memory::{
5760
};
5861
pub use schema::*;
5962
pub use table::*;
60-
61-
// For backwards compatibility,
62-
mod session {
63-
pub use datafusion_session::Session;
64-
}

‎datafusion/catalog/src/schema.rs‎

Lines changed: 2 additions & 90 deletions
Original file line numberDiff line numberDiff line change
@@ -15,93 +15,5 @@
1515
// specific language governing permissions and limitations
1616
// under the License.
1717

18-
//! Describes the interface and built-in implementations of schemas,
19-
//! representing collections of named tables.
20-
21-
use async_trait::async_trait;
22-
use datafusion_common::{DataFusionError, exec_err};
23-
use std::any::Any;
24-
use std::fmt::Debug;
25-
use std::sync::Arc;
26-
27-
use crate::table::TableProvider;
28-
use datafusion_common::Result;
29-
use datafusion_expr::TableType;
30-
31-
/// Represents a schema, comprising a number of named tables.
32-
///
33-
/// Please see [`CatalogProvider`] for details of implementing a custom catalog.
34-
///
35-
/// [`CatalogProvider`]: super::CatalogProvider
36-
#[async_trait]
37-
pub trait SchemaProvider: Any + Debug + Sync + Send {
38-
/// Returns the owner of the Schema, default is None. This value is reported
39-
/// as part of `information_tables.schemata
40-
fn owner_name(&self) -> Option<&str> {
41-
None
42-
}
43-
44-
/// Retrieves the list of available table names in this schema.
45-
fn table_names(&self) -> Vec<String>;
46-
47-
/// Retrieves a specific table from the schema by name, if it exists,
48-
/// otherwise returns `None`.
49-
async fn table(
50-
&self,
51-
name: &str,
52-
) -> Result<Option<Arc<dyn TableProvider>>, DataFusionError>;
53-
54-
/// Retrieves the type of a specific table from the schema by name, if it exists, otherwise
55-
/// returns `None`. Implementations for which this operation is cheap but [Self::table] is
56-
/// expensive can override this to improve operations that only need the type, e.g.
57-
/// `SELECT * FROM information_schema.tables`.
58-
async fn table_type(&self, name: &str) -> Result<Option<TableType>> {
59-
self.table(name).await.map(|o| o.map(|t| t.table_type()))
60-
}
61-
62-
/// If supported by the implementation, adds a new table named `name` to
63-
/// this schema.
64-
///
65-
/// If a table of the same name was already registered, returns "Table
66-
/// already exists" error.
67-
#[expect(unused_variables)]
68-
fn register_table(
69-
&self,
70-
name: String,
71-
table: Arc<dyn TableProvider>,
72-
) -> Result<Option<Arc<dyn TableProvider>>> {
73-
exec_err!("schema provider does not support registering tables")
74-
}
75-
76-
/// If supported by the implementation, removes the `name` table from this
77-
/// schema and returns the previously registered [`TableProvider`], if any.
78-
///
79-
/// If no `name` table exists, returns Ok(None).
80-
#[expect(unused_variables)]
81-
fn deregister_table(&self, name: &str) -> Result<Option<Arc<dyn TableProvider>>> {
82-
exec_err!("schema provider does not support deregistering tables")
83-
}
84-
85-
/// Returns true if table exist in the schema provider, false otherwise.
86-
fn table_exist(&self, name: &str) -> bool;
87-
}
88-
89-
impl dyn SchemaProvider {
90-
/// Returns `true` if the schema provider is of type `T`.
91-
///
92-
/// Prefer this over `downcast_ref::<T>().is_some()`. Works correctly when
93-
/// called on `Arc<dyn SchemaProvider>` via auto-deref.
94-
pub fn is<T: SchemaProvider>(&self) -> bool {
95-
(self as &dyn Any).is::<T>()
96-
}
97-
98-
/// Attempts to downcast this schema provider to a concrete type `T`,
99-
/// returning `None` if the provider is not of that type.
100-
///
101-
/// Works correctly when called on `Arc<dyn SchemaProvider>` via auto-deref,
102-
/// unlike `(&arc as &dyn Any).downcast_ref::<T>()` which would attempt to
103-
/// downcast the `Arc` itself.
104-
pub fn downcast_ref<T: SchemaProvider>(&self) -> Option<&T> {
105-
(self as &dyn Any).downcast_ref()
106-
}
107-
}
18+
// Re-export from this module for backwards compatibility.
19+
pub use datafusion_session::SchemaProvider;

0 commit comments

Comments
 (0)