库用户指南

目录、模式与表

师成师成· 更新于 2026-09-28· 阅读 16 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

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 的映射。

回顾

回顾一下,你需要:

  1. 实现 TableProvider trait 以创建表提供者,或者使用现有的表提供者。
  2. 实现 SchemaProvider trait 以创建 schema 提供者,或者使用现有的 schema 提供者。
  3. 实现 CatalogProvider trait 以创建 catalog 提供者,或者使用现有的 catalog 提供者。
  4. 实现 CatalogProviderList trait 以创建 CatalogProviderList,或者使用现有的 CatalogProviderList。

评论

登录后参与评论

正在加载评论…