目录、模式与表
Catalog、Schema 与表
本节介绍如何在 DataFusion 中创建和管理 Catalog、Schema 与表。想要快速深入代码的读者,请参阅示例。
基本概念
Catalog Provider、Catalog、Schema 与表按层次结构组织。CatalogProviderList 包含多个 CatalogProvider,CatalogProvider 包含多个 SchemaProvider,而 SchemaProvider 包含多个 TableProvider。
DataFusion 在 catalog 模块中提供了基本的内存 Catalog 功能。你可以直接使用这些内存实现,也可以用自定义的 Catalog 实现来扩展 DataFusion,例如基于本地文件或远程对象存储上的文件。
DataFusion 支持使用本节所述 Catalog API 的 DDL 查询(例如 CREATE TABLE)。有关 DML 查询(例如 INSERT INTO)的信息,请参见 TableProvider 一节。
与 DataFusion 中的其他概念类似,你需要实现各种 trait 来创建自己的 catalog、schema 和表。以下各节将介绍需要实现的 trait。
CatalogProviderList trait 提供了注册新 catalog、按名称获取 catalog 以及列出所有 catalog 的方法。CatalogProvider trait 提供了按名称设置 schema、按名称获取 schema 以及列出所有 schema 的方法。可以注册到 CatalogProvider 中的 SchemaProvider 提供了按名称设置表、按名称获取表、列出所有表、注销表以及检查表是否存在的方法。TableProvider trait 提供了扫描底层数据并将其用于 DataFusion 的方法。关于 TableProvider trait 的更详细介绍请见此处。
在下面的示例中,我们将实现一个内存 catalog,首先从 SchemaProvider trait 开始,因为需要一个这样的实现来注册到 CatalogProvider 中。最后我们将实现 CatalogProviderList 来注册这个 CatalogProvider。
实现 MemorySchemaProvider
MemorySchemaProvider 是 SchemaProvider trait 的一个简单实现。它将状态(即表)存储在 DashMap 中,由后者支撑 SchemaProvider trait。
use std::sync::Arc;
use dashmap::DashMap;
use datafusion::catalog::{TableProvider, SchemaProvider};
#[derive(Debug)]
pub struct MemorySchemaProvider {
tables: DashMap<String, Arc<dyn TableProvider>>,
}tables 就是上述键值对。其底层状态也可以是另一种数据结构,或文件、事务型数据库等其他存储机制。
接下来,我们为 MemorySchemaProvider 实现 SchemaProvider trait。
use std::any::Any;
use datafusion::catalog::SchemaProvider;
use async_trait::async_trait;
use datafusion::common::{Result, exec_err};
#[async_trait]
impl SchemaProvider for MemorySchemaProvider {
fn table_names(&self) -> Vec<String> {
self.tables
.iter()
.map(|table| table.key().clone())
.collect()
}
async fn table(&self, name: &str) -> Result<Option<Arc<dyn TableProvider>>> {
Ok(self.tables.get(name).map(|table| table.value().clone()))
}
fn register_table(
&self,
name: String,
table: Arc<dyn TableProvider>,
) -> Result<Option<Arc<dyn TableProvider>>> {
if self.table_exist(name.as_str()) {
return exec_err!(
"The table {name} already exists"
);
}
Ok(self.tables.insert(name, table))
}
fn deregister_table(&self, name: &str) -> Result<Option<Arc<dyn TableProvider>>> {
Ok(self.tables.remove(name).map(|(_, table)| table))
}
fn table_exist(&self, name: &str) -> bool {
self.tables.contains_key(name)
}
}暂不深入 CatalogProvider 的实现,我们可以创建一个 MemorySchemaProvider,并向其中注册 TableProvider。
use arrow::datatypes::{DataType, Field, Schema, SchemaRef};
use arrow::record_batch::RecordBatch;
use datafusion::datasource::MemTable;
use arrow::array::{self, Array, ArrayRef, Int32Array};
impl MemorySchemaProvider {
/// Instantiates a new MemorySchemaProvider with an empty collection of tables.
pub fn new() -> Self {
Self {
tables: DashMap::new(),
}
}
}
let schema_provider = Arc::new(MemorySchemaProvider::new());
let table_provider = {
let schema = Arc::new(Schema::new(vec![Field::new("i", DataType::Int32, true)]));
let arr = Arc::new(Int32Array::from((1..=1).collect::<Vec<_>>()));
let partitions = vec![vec![RecordBatch::try_new(schema.clone(), vec![arr as ArrayRef]).unwrap()]];
Arc::new(MemTable::try_new(schema, partitions).unwrap())
};
schema_provider.register_table("users".to_string(), table_provider);
let table = schema_provider.table("users");异步 SchemaProvider
从远程数据源获取某个 schema 中包含哪些表的元数据信息,往往非常有用。例如,一个 schema provider 可以从远程数据库中获取元数据。为支持这一点,SchemaProvider trait 提供了一个异步的 table 方法。
该 trait 与前面基本相同,只是 table 方法不同,并增加了 #[async_trait] 属性。
#[async_trait]
impl SchemaProvider for Schema {
async fn table(&self, name: &str) -> Result<Option<Arc<dyn TableProvider>>> {
}
}实现 MemoryCatalogProvider
如前所述,CatalogProvider 可以管理 catalog 中的 schema,而 MemoryCatalogProvider 是 CatalogProvider trait 的一个简单实现。它使用 DashMap 来存储 schema。借助这一点,就可以实现 CatalogProvider trait 了。
use std::any::Any;
use std::sync::Arc;
use dashmap::DashMap;
use datafusion::catalog::{CatalogProvider, SchemaProvider};
use datafusion::common::Result;
#[derive(Debug)]
pub struct MemoryCatalogProvider {
schemas: DashMap<String, Arc<dyn SchemaProvider>>,
}
impl CatalogProvider for MemoryCatalogProvider {
fn schema_names(&self) -> Vec<String> {
self.schemas.iter().map(|s| s.key().clone()).collect()
}
fn schema(&self, name: &str) -> Option<Arc<dyn SchemaProvider>> {
self.schemas.get(name).map(|s| s.value().clone())
}
fn register_schema(
&self,
name: &str,
schema: Arc<dyn SchemaProvider>,
) -> Result<Option<Arc<dyn SchemaProvider>>> {
Ok(self.schemas.insert(name.into(), schema))
}
fn deregister_schema(
&self,
name: &str,
cascade: bool,
) -> Result<Option<Arc<dyn SchemaProvider>>> {
/// `cascade` is not used here, but can be used to control whether
/// to delete all tables in the schema or not.
if let Some(schema) = self.schema(name) {
let (_, removed) = self.schemas.remove(name).unwrap();
Ok(Some(removed))
} else {
Ok(None)
}
}
}同样,这相当直接,因为存在一个底层数据结构通过键值对来存储状态。借助它,就可以实现 CatalogProviderList trait。
实现 MemoryCatalogProviderList
use std::any::Any;
use std::sync::Arc;
use dashmap::DashMap;
use datafusion::catalog::{CatalogProviderList, CatalogProvider};
use datafusion::common::Result;
#[derive(Debug)]
pub struct MemoryCatalogProviderList {
/// Collection of catalogs containing schemas and ultimately TableProviders
pub catalogs: DashMap<String, Arc<dyn CatalogProvider>>,
}
impl CatalogProviderList for MemoryCatalogProviderList {
fn register_catalog(
&self,
name: String,
catalog: Arc<dyn CatalogProvider>,
) -> Option<Arc<dyn CatalogProvider>> {
self.catalogs.insert(name, catalog)
}
fn catalog_names(&self) -> Vec<String> {
self.catalogs.iter().map(|c| c.key().clone()).collect()
}
fn catalog(&self, name: &str) -> Option<Arc<dyn CatalogProvider>> {
self.catalogs.get(name).map(|c| c.value().clone())
}
}与其他 trait 一样,它也会维护 Catalog 名称到 CatalogProvider 的映射。
回顾
回顾一下,你需要:
- 实现
TableProvidertrait 以创建表提供者,或者使用现有的表提供者。 - 实现
SchemaProvidertrait 以创建 schema 提供者,或者使用现有的 schema 提供者。 - 实现
CatalogProvidertrait 以创建 catalog 提供者,或者使用现有的 catalog 提供者。 - 实现
CatalogProviderListtrait 以创建 CatalogProviderList,或者使用现有的 CatalogProviderList。
评论
登录后参与评论
KnowForge