库用户指南

自定义表提供者

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

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

自定义 Table Provider

DataFusion 最强大的优点之一就是它的可扩展性。如果你的数据存放在自定义格式中、某个 API 之后,或者 DataFusion 原生不支持的系统里,你可以通过实现自定义 table provider 来教会 DataFusion 读取它们。本文将带你走完设计 table provider 所需理解的三层抽象,并说明规划(planning)与执行(execution)的工作各自应该发生在哪里。

关于主键或唯一约束等表约束如何处理的详细信息,请参阅表约束强制(Table Constraint Enforcement)。

本文大部分内容最初发表于博客Writing Custom Table Providers in Apache DataFusion。

三层抽象

当 DataFusion 针对一张表执行查询时,三个抽象协同工作以产出结果:

  1. TableProvider – 描述该表(schema、能力),并在被查询时生成执行计划。它属于**逻辑计划(Logical Plan)**的一部分。
  2. ExecutionPlan – 描述如何计算结果:分区、排序以及子计划之间的关系。它属于**物理计划(Physical Plan)**的一部分。
  3. SendableRecordBatchStream – 真正干活的异步流,逐个产出 RecordBatch。

可以把它们想象成一个漏斗:TableProvider::scan() 在规划阶段被调用一次以创建 ExecutionPlan,随后 ExecutionPlan::execute() 每个分区各被调用一次以创建流,而这些流正是执行期间真正产生行的地方。

背景:逻辑规划与物理规划

在深入三层抽象之前,先了解一下 DataFusion 如何处理一次查询会很有帮助。从一个 SQL 字符串(或 DataFrame 调用)到流式返回结果之间,要经历几个阶段:

SQL / DataFrame API
  → Logical Plan          (abstract: what to compute)
  → Logical Optimization  (rewrite rules that preserve semantics)
  → Physical Plan         (concrete: how to compute it)
  → Physical Optimization (hardware- and data-aware rewrites)
  → Execution             (streaming RecordBatches)

逻辑计划

逻辑计划描述的是查询要计算什么,而不指定如何计算。它是一棵关系运算符组成的树——TableScan、Filter、Projection、Aggregate、Join、Sort、Limit 等等。逻辑优化器会重写这棵树,在保持查询语义不变的前提下减少计算量。一些逻辑优化包括:

  • 谓词下推(predicate pushdown)——把过滤条件尽可能地向数据源方向移动,从而让流经计划其余部分的行数更少。
  • 投影裁剪(projection pruning)——移除下游从未引用的列,减少内存和 I/O 开销。
  • 表达式简化——把 1 = 1 或 x AND true 这类表达式重写为更简单的形式。
  • 子查询去相关(subquery decorrelation)——把相关子查询形式的 IN / EXISTS 转换为更高效的半连接(semi-join)。
  • Limit 下推——把 LIMIT 在计划中尽量提前,使各个算子产生的数据量更少。

物理计划

物理规划器会把优化后的逻辑计划转换为一棵 ExecutionPlan 树——也就是真正要执行的具体计划。诸如“使用哈希连接还是排序归并连接”“要扫描多少个分区”之类的决策就是在这一步做出的。随后,物理优化器会进一步细化这棵树,采用的重写手段包括:

  • 分布保证(distribution enforcement)——插入 RepartitionExec 节点,使数据按连接和聚合的要求正确分区。
  • 排序保证(sort enforcement)——在需要有序时插入 SortExec 节点,并在数据已经有序时移除它们。
  • 连接策略选择——根据统计信息和表的大小选择最高效的连接策略。
  • 聚合优化——合并部分聚合与最终聚合阶段,并可利用精确统计信息来完全跳过扫描。

这对 Table Provider 为何重要

你的 TableProvider 位于逻辑规划与物理规划的交界处。在逻辑优化期间,DataFusion 会确定哪些过滤条件和投影可以下推到数据源。当物理规划阶段调用 scan() 时,这些提示会被传递给你。通过实现 supports_filters_pushdown 这类能力,你能影响优化器可以做哪些优化——而你在 ExecutionPlan 中声明的元数据(分区方式、排序方式)会直接影响哪些物理优化能够适用。

选择正确的起点

并非每个自定义数据源都需要从零开始实现全部三个层次。DataFusion 提供了一系列构建块,让你可以在任何合适的层次上接入:

如果你的数据是……从这里开始你需要实现
已经以内存中 RecordBatch 的形式存在MemTable无需实现——直接构建它即可
批次的异步流StreamTable一个流工厂
其他表的逻辑变换ViewTable,包装一个逻辑计划该逻辑计划
现有文件格式的变体ListingTable,配合自定义的 FileFormat 包装现有格式一个轻量的 FileFormat 包装器
磁盘或对象存储上自定义格式的文件ListingTable,配合自定义的 FileFormat、FileSource 和 FileOpener格式、数据源和打开器
需要完全控制的自定义数据源TableProvider + ExecutionPlan + 流全部三个层面

如果你的数据基于文件,ListingTable 已经负责了文件发现、分区列推断和计划构建——你只需要实现 FileFormat、FileSource 和 FileOpener,用来描述如何读取你的文件。可以参考 custom_file_format 示例了解最简包装方式,或参考 ParquetSource 和 ParquetOpener 这一完整自定义实现作为参考。

本文的其余部分聚焦于完整的 TableProvider + ExecutionPlan + 流式路径,它能让你完全掌控,并且适用于任何数据源。

第一层:TableProvider

TableProvider 表示一个可查询的数据源。要实现一个最简的只读表,你需要三个方法:

impl TableProvider for MyTable {
    fn schema(&self) -> SchemaRef {
        Arc::clone(&self.schema)
    }

    fn table_type(&self) -> TableType {
        TableType::Base
    }

    async fn scan(
        &self,
        state: &dyn Session,
        projection: Option<&[usize]>,
        filters: &[Expr],
        limit: Option<usize>,
    ) -> Result<Arc<dyn ExecutionPlan>> {
        // Build and return an ExecutionPlan -- don't do any execution work here -- keep lightweight!
        Ok(Arc::new(MyExecPlan::new(
            Arc::clone(&self.schema),
            projection,
            limit,
        )))
    }
}

scan 方法是 TableProvider 的核心。它从优化器接收三个下推提示,每一个都能减少数据源需要产出的数据量:

  • projection — 需要哪些列。这会减少输出的宽度。如果数据源支持,就只读取这些列,而不是整个 schema。
  • filters — 引擎希望你在扫描过程中应用的谓词。通过跳过不匹配的数据,这会减少行数。实现 supports_filters_pushdown 来声明你能处理哪些过滤条件。
  • limit — 行数上限。这同样会减少行数 —— 如果你在产出足够行数后可以提前停止读取,就能避免不必要的工作。

你也可以使用 scan_with_args() 变体,它会为其他高级场景提供额外的下推信息。

让 scan() 保持轻量

这是一个关键点:scan() 运行在规划阶段,而非执行阶段。 它应该快速返回。最佳实践是避免在这里执行 I/O、网络调用或繁重的计算。scan 方法的职责是描述数据将如何被产出,而不是产出数据。所有真正的工作都属于流(第 3 层)。

一个常见的误区是在 scan() 中获取数据或建立连接。这会阻塞规划线程,可能导致超时或死锁,尤其是当查询涉及多个表或子查询,而它们都必须在执行开始前完成规划时。

可以参考的现有实现

DataFusion 自带了几个 TableProvider 实现,它们是极好的参考:

  • MemTable – 以内存中的 Vec<RecordBatch> 形式保存数据。最简单的 provider,非常适合测试和小型数据集。
  • StreamTable – 包装用户提供的流工厂。当你的数据以持续流的形式到达时(例如来自 Kafka 或套接字),它很有用。
  • ListingTable – DataFusion 内置的 Parquet、CSV 和 JSON 支持背后的基于文件的数据源。它展示了复杂的过滤下推、投影下推、文件裁剪以及 schema 推断。
  • ViewTable – 包装一个逻辑计划,代表一个 SQL 视图。如果你的 provider 更适合表达为对其他表的变换,它会很有用。

第 2 层:ExecutionPlan

ExecutionPlan 是物理查询计划树中的一个节点。你的表提供者的 scan() 方法会返回一个这样的节点。需要实现的方法有:

impl ExecutionPlan for MyExecPlan {
    fn name(&self) -> &str { "MyExecPlan" }

    fn properties(&self) -> &Arc<PlanProperties> {
        &self.properties
    }

    fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {
        vec![]  // Leaf node -- no children
    }

    fn replace_children(
        self: Arc<Self>,
        children: Vec<Arc<dyn ExecutionPlan>>,
        _: ReplaceChildrenOptions,
    ) -> Result<Arc<dyn ExecutionPlan>> {
        assert!(children.is_empty());
        Ok(self)
    }

    fn with_new_children(
        self: Arc<Self>,
        children: Vec<Arc<dyn ExecutionPlan>>,
    ) -> Result<Arc<dyn ExecutionPlan>> {
        self.replace_children(children, ReplaceChildrenOptions::new(ChildrenPropertiesMode::Recompute))
    }

    fn execute(
        &self,
        partition: usize,
        context: Arc<TaskContext>,
    ) -> Result<SendableRecordBatchStream> {
        // This is where you build and return your stream
        // ...
    }
}

在 PlanProperties 中需要正确设置的关键属性是输出分区和输出排序。

输出分区告诉引擎你的数据有多少个分区,从而决定并行度。如果你的数据源天然就会对数据分区(例如按文件或按分片),就在这里暴露出来。

输出排序声明你的数据是否天然有序。这使优化器在查询需要有序数据时可以避免插入 SortExec。正确设置这一点可能带来显著的性能提升。

分区策略

由于 execute() 会对每个分区各调用一次,分区方式直接控制表扫描的并行度。每个分区都会产生一个独立的数据流,DataFusion 会将其作为一个任务调度到 tokio 运行时上。区分任务与线程很重要:任务是轻量级的异步工作单元,它们被多路复用到线程池上。任务(分区)的数量可以远多于物理线程数——当它们等待 I/O 或主动让出执行权时,运行时会高效地交替调度它们。

先从简单开始:与数据的自然布局保持一致。 如果你有 4 个文件,就暴露 4 个分区。如果你的数据源有 8 个分片,就暴露 8 个分区。当下游算子需要不同的分布时,DataFusion 会在你的扫描之上插入一个 RepartitionExec。你也可以在 ExecutionPlan 上实现 repartitioned 方法,让 DataFusion 直接向数据源请求不同的分区数,从而完全省去这个额外的算子。

思考一下你的数据源是如何自然地划分数据的:

  • 按文件或对象划分: 如果你从 S3 读取数据,每个文件可以作为一个分区。DataFusion 会并行读取它们。
  • 按分片或区域划分: 如果你的数据源是分片数据库,每个分片自然对应一个分区。
  • 按键范围划分: 如果你的数据带有键(例如时间戳或客户 ID),可以按范围将其切分。

进阶:与 target_partitions 对齐。 当你的实现跑通之后,还可以进一步调优。分区过多并非没有代价:每个分区都会增加调度开销,而且下游算子可能无论如何都需要重新分区数据。会话配置中暴露了一个目标分区数,它反映了优化器预期要处理的分区数量:

async fn scan(
    &self,
    state: &dyn Session,
    projection: Option<&[usize]>,
    filters: &[Expr],
    limit: Option<usize>,
) -> Result<Arc<dyn ExecutionPlan>> {
    let target_partitions = state.config().target_partitions();
    // Optionally coalesce or split partitions to match target_partitions.
    // ...
}

如果数据源产生的数据恰好就是 target_partitions 个分区,优化器就不太可能在你的扫描算子之上插入 RepartitionExec。对于小型数据集,target_partitions 可能被设置为 1,这样就完全避免了任何重分区的开销。

进阶:声明哈希分区。 如果你的数据源按照某个特定键(例如 customer_id)预先对数据进行了分区,你可以在输出分区方式中声明这一点。对于如下查询:

SELECT customer_id, SUM(amount)
FROM my_table
GROUP BY customer_id;

如果你声明输出分区方式为 Hash([customer_id], N),优化器就会识别出数据已经按照聚合所需的分布进行了划分,从而消除原本会出现在执行计划中的 RepartitionExec。你可以用 EXPLAIN 来验证这一点(下文会详细介绍)。

反过来,如果你报告的是 UnknownPartitioning,DataFusion 就必须按最坏情况处理,总会根据需要插入重新分区算子。

让 execute() 保持轻量

与 scan() 一样,execute() 方法应该只负责构造并返回一个流,而不做繁重的工作。实际的数据生产发生在流被轮询时。不要在这里阻塞等待异步操作——构建好流,然后交给运行时去驱动它。

可以参考的现有实现

  • StreamingTableExec – 执行流式表扫描。它接收一个流工厂(用于产生流的闭包)并处理分区。是封装外部流的良好参考。
  • DataSourceExec – DataFusion 内置文件扫描(Parquet、CSV、JSON)背后的执行计划。它展示了精细的分区处理、过滤下推和投影下推。

第 3 层:SendableRecordBatchStream

SendableRecordBatchStream 才是真正干活的地方。它的定义如下:

type SendableRecordBatchStream =
    Pin<Box<dyn RecordBatchStream<Item = Result<RecordBatch>> + Send>>;

这是一个可以跨线程传递的 RecordBatch 异步流。当 DataFusion 运行时轮询这个流时,你的代码就会执行:读取文件、调用 API、转换数据等等。

使用 RecordBatchStreamAdapter

创建 SendableRecordBatchStream 最简单的方式是使用 RecordBatchStreamAdapter。它能将任何 futures::Stream<Item = Result<RecordBatch>> 桥接为 SendableRecordBatchStream 类型:

use datafusion::physical_plan::stream::RecordBatchStreamAdapter;

fn execute(
    &self,
    partition: usize,
    context: Arc<TaskContext>,
) -> Result<SendableRecordBatchStream> {
    let schema = self.schema();
    let config = self.config.clone();

    let stream = futures::stream::once(async move {
        // ALL the heavy work happens here, inside the stream:
        // - Open connections
        // - Read data from external sources
        // - Transform and batch the results
        let batches = fetch_data_from_source(&config).await?;
        Ok(batches)
    })
    .flat_map(|result| match result {
        Ok(batch) => futures::stream::iter(vec![Ok(batch)]),
        Err(e) => futures::stream::iter(vec![Err(e)]),
    });

    Ok(Box::pin(RecordBatchStreamAdapter::new(schema, stream)))
}

阻塞型工作:使用独立线程池

如果你的流执行的是阻塞型工作——例如阻塞式 I/O,或者持续数百毫秒且不让出执行权的 CPU 计算——就必须避免阻塞 tokio 异步运行时。短暂的 CPU 工作(例如用几毫秒解析一个批次)可以内联完成,只要你的代码经常让出执行权回到运行时即可。但对于无法让出执行权的长时间同步工作,请将其卸载到专用线程池,并通过 channel 将结果传回:

fn execute(
    &self,
    partition: usize,
    context: Arc<TaskContext>,
) -> Result<SendableRecordBatchStream> {
    let schema = self.schema();
    let config = self.config.clone();

    let (tx, rx) = tokio::sync::mpsc::channel(2);

    // Spawn blocking work on a dedicated thread pool
    tokio::task::spawn_blocking(move || {
        let batches = generate_data(&config);
        for batch in batches {
            if tx.blocking_send(Ok(batch)).is_err() {
                break; // Receiver dropped, query was cancelled
            }
        }
    });

    let stream = tokio_stream::wrappers::ReceiverStream::new(rx);
    Ok(Box::pin(RecordBatchStreamAdapter::new(schema, stream)))
}

这种模式让异步运行时保持响应能力,而长时间运行的同步工作则在各自的线程上执行。若要查看一个展示如何为 I/O 工作和 CPU 工作分别配置线程池的可运行示例,请参阅 DataFusion 仓库中的 thread_pools 示例。

工作应该在哪里发生?

下表总结了每个层次各自应负责的内容:

层次运行时机应该做什么不应该做什么
TableProvider::scan()规划阶段构建带有元数据的 ExecutionPlanI/O、网络调用、繁重计算
ExecutionPlan::execute()执行阶段(每个分区一次)构造数据流、建立通道阻塞在异步工作上、读取数据
RecordBatchStream(轮询)执行阶段所有 I/O、计算、数据生产–

指导原则:尽可能把工作推迟。 规划应当快速完成,这样优化器才能发挥作用。执行的初始化也应当快速完成,以便所有分区都能及时启动。数据流才是你花费时间生产数据的地方。

为什么这很重要

当 scan() 承担繁重工作时,会出现以下问题:

  1. 规划变得缓慢。 如果一个查询涉及 10 个表,每个 scan() 耗时 500 毫秒,那么在任何数据开始流动之前,仅规划就需要 5 秒。
  2. 执行是单线程的。 scan() 在规划期间于单个线程上运行,因此在那里完成的任何工作都无法受益于 DataFusion 跨分区提供的并行执行。
  3. 优化器无从帮助。 优化器运行在规划与执行之间。如果你在规划阶段就已取回数据,谓词下推或分区裁剪等优化就无法减少工作量。
  4. 资源管理失效。 DataFusion 在执行期间管理并发和内存,而在规划期间完成的工作会绕过这些控制。

过滤下推:减少工作量

你可以为自定义表提供者添加的最有价值的优化之一就是过滤下推——让数据源跳过查询不需要的数据,而不是读取全部数据之后再进行过滤。

过滤下推的工作原理

当 DataFusion 规划带有 WHERE 子句的查询时,它会将过滤谓词作为 filters 参数传递给你的 scan() 方法。默认情况下,DataFusion 假定你的 provider 无法处理任何过滤条件,因此会在你的扫描之上插入一个 FilterExec 节点来执行这些过滤。但如果你的数据源能够在扫描过程中评估部分谓词——例如,通过跳过不可能匹配的文件、分区或行组——就能消除大量不必要的 I/O。

要启用这一能力,请实现 supports_filters_pushdown:

fn supports_filters_pushdown(
    &self,
    filters: &[&Expr],
) -> Result<Vec<TableProviderFilterPushDown>> {
    Ok(filters.iter().map(|f| {
        match f {
            // We can fully evaluate equality filters on
            // the partition column at the source
            Expr::BinaryExpr(BinaryExpr {
                left, op: Operator::Eq, right
            }) if is_partition_column(left) || is_partition_column(right) => {
                TableProviderFilterPushDown::Exact
            }
            // All other filters: let DataFusion handle them
            _ => TableProviderFilterPushDown::Unsupported,
        }
    }).collect())
}

对于每个过滤条件,有三种可能的响应:

  • Exact – 你的数据源保证不会有输出行对该谓词求值为 false。由于过滤器完全在数据源处求值,DataFusion 将不会为它添加 FilterExec。
  • Inexact – 你的数据源能够减少产生的数据量,但输出中仍可能包含不满足该谓词的行。例如,你可能基于元数据统计信息跳过整个文件,但不会过滤文件内部的单个行。DataFusion 仍然会在你的扫描之上添加一个 FilterExec,以清除任何遗漏的行。
  • Unsupported – 你的数据源完全忽略该过滤条件,由 DataFusion 处理它。

为什么过滤下推很重要

考虑一张按 region 分区、包含 10 亿行的表,以及如下查询:

SELECT * FROM events WHERE region = 'us-east-1' AND event_type = 'click';

**不进行过滤下推时:**你的表提供者需要读取所有地区全部 10 亿行数据,然后 DataFusion 再应用这两个过滤条件,丢弃绝大多数数据。

**对 region 进行过滤下推时:**你的 scan() 方法会看到 region = 'us-east-1' 这个过滤条件,并构建出只读取 us-east-1 分区的执行计划。如果该分区包含 1 亿行数据,你就直接消除了 90% 的 I/O。如果你将 event_type 报告为 Unsupported,DataFusion 仍会通过 FilterExec 应用该过滤条件。

仅在数据源能够做得更好时才下推过滤条件

DataFusion 已经会尽可能将过滤条件下推到靠近数据源的位置,通常是直接放在扫描之上。FilterExec 也经过了高度优化,具备向量化求值和针对类型的专用内核,能够快速计算谓词。

正因如此,只有当你的数据源确实能够做得更好时,才应该实现过滤下推——例如,通过基于元数据跳过文件或分区来完全避免 I/O。如果你的数据源无法以这种方式消除 I/O,通常更好的做法是让 DataFusion 处理过滤条件,因为它的内存内执行已经非常高效。

使用 EXPLAIN 调试你的表提供者

EXPLAIN 语句是了解 DataFusion 究竟对你的表提供者做了什么的最佳工具。它会展示 DataFusion 将要执行的物理计划,其中包含它插入的所有算子:

EXPLAIN SELECT * FROM events WHERE region = 'us-east-1' AND event_type = 'click';

如果使用 DataFrame,可以调用 .explain(false, false) 查看逻辑计划,或调用 .explain(false, true) 查看物理计划。也可以使用 .explain(true, true) 以详细模式打印计划。

在过滤条件下推之前,计划可能如下所示:

FilterExec: region@0 = us-east-1 AND event_type@1 = click
  MyExecPlan: partitions=50

在这里,DataFusion 会读取全部 50 个分区,然后再进行过滤。位于扫描之上的 FilterExec 承担了所有的谓词计算工作。

在为 region 实现下推之后(该列被报告为 Exact):

FilterExec: event_type@1 = click
  MyExecPlan: partitions=5, filter=[region = us-east-1]

现在你的执行计划只读取 us-east-1 的 5 个分区,剩余的 FilterExec 只需处理 event_type 谓词。region 过滤条件已被你的扫描完全吸收。

为两个过滤条件都实现下推之后(两者均为 Exact):

MyExecPlan: partitions=5, filter=[region = us-east-1 AND event_type = click]

完全没有 FilterExec —— 你的数据源处理了所有工作。

同样地,EXPLAIN 会显示 DataFusion 是否插入了不必要的 SortExec 或 RepartitionExec 节点,而通过声明更好的输出属性,你可以消除这些节点。每当查询看起来比预期慢时,EXPLAIN 都是首先要查看的地方。

一个完整的过滤下推示例

为了让过滤下推更加具体,这里给出一个说明性的例子。设想一个表提供者,它从磁盘上一组按日期分区的目录读取数据(例如 data/2026-03-01/、data/2026-03-02/ 等)。每个目录包含该日期对应的一个或多个 Parquet 文件。通过下推针对 date 列的过滤条件,该提供者可以跳过整个目录——避免列出和读取那些绝不可能匹配查询的文件所带来的 I/O 开销。

/// A table provider backed by date-partitioned directories.
/// Each date directory contains data files; by filtering on the
/// `date` column we can skip entire directories of I/O.
struct DatePartitionedTable {
    schema: SchemaRef,
    /// Maps date strings ("2026-03-01") to directory paths
    partitions: HashMap<String, String>,
}

#[async_trait::async_trait]
impl TableProvider for DatePartitionedTable {
    fn schema(&self) -> SchemaRef { Arc::clone(&self.schema) }
    fn table_type(&self) -> TableType { TableType::Base }

    fn supports_filters_pushdown(
        &self,
        filters: &[&Expr],
    ) -> Result<Vec<TableProviderFilterPushDown>> {
        Ok(filters.iter().map(|f| {
            if Self::is_date_equality_filter(f) {
                // We can fully evaluate this: we will only read
                // directories matching the date, so no rows with
                // a different date will appear in the output.
                TableProviderFilterPushDown::Exact
            } else {
                TableProviderFilterPushDown::Unsupported
            }
        }).collect())
    }

    async fn scan(
        &self,
        _state: &dyn Session,
        projection: Option<&[usize]>,
        filters: &[Expr],
        limit: Option<usize>,
    ) -> Result<Arc<dyn ExecutionPlan>> {
        // Determine which date partitions to read by inspecting
        // the pushed-down filters. This is the key optimization:
        // we decide *during planning* which directories to scan,
        // so that execution never touches irrelevant data.
        let dates_to_read: Vec<String> = self
            .extract_date_values(filters)
            .unwrap_or_else(||
                self.partitions.keys().cloned().collect()
            );

        let dirs: Vec<String> = dates_to_read
            .iter()
            .filter_map(|d| self.partitions.get(d).cloned())
            .collect();
        let num_dirs = dirs.len();

        Ok(Arc::new(DatePartitionedExec {
            schema: Arc::clone(&self.schema),
            directories: dirs,
            properties: Arc::new(PlanProperties::new(
                EquivalenceProperties::new(
                    Arc::clone(&self.schema),
                ),
                // One partition per date directory -- these
                // will be read in parallel.
                Partitioning::UnknownPartitioning(num_dirs),
                EmissionType::Incremental,
                Boundedness::Bounded,
            )),
        }))
    }
}

impl DatePartitionedTable {
    /// Check if a filter is an equality comparison on the `date` column.
    fn is_date_equality_filter(expr: &Expr) -> bool {
        // In practice, match on BinaryExpr { left, op: Eq, right }
        // and check if either side references the "date" column.
        // Simplified here for clarity.
        todo!("match on date equality expressions")
    }

    /// Extract date literal values from pushed-down equality filters.
    fn extract_date_values(&self, filters: &[Expr]) -> Option<Vec<String>> {
        // Parse filters like `date = '2026-03-01'` and return
        // the literal date strings. Returns None if no date
        // filters are present (meaning: read all partitions).
        todo!("extract date literals from filter expressions")
    }
}

其核心要点在于:过滤下推决策(supports_filters_pushdown)与分区裁剪(scan())相互配合——前者告知 DataFusion 针对 date 谓词无需 FilterExec,后者则确保只扫描相关的目录。实际的文件读取稍后才会发生,即在 execute() 产生的流中进行。

整合起来

下面是一个最小但完整的自定义表提供者示例,它在流式处理期间按需延迟生成数据:

use std::any::Any;
use std::sync::Arc;

use arrow::array::Int64Array;
use arrow::datatypes::{DataType, Field, Schema, SchemaRef};
use arrow::record_batch::RecordBatch;
use datafusion::catalog::TableProvider;
use datafusion::common::Result;
use datafusion::datasource::TableType;
use datafusion::catalog::Session;
use datafusion::execution::SendableRecordBatchStream;
use datafusion::logical_expr::Expr;
use datafusion::physical_expr::EquivalenceProperties;
use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType};
use datafusion::physical_plan::stream::RecordBatchStreamAdapter;
use datafusion::physical_plan::{
    ChildrenPropertiesMode, ReplaceChildrenOptions, ExecutionPlan, Partitioning, PlanProperties,
};
use futures::stream;

/// A table provider that generates sequential numbers on demand.
struct CountingTable {
    schema: SchemaRef,
    num_partitions: usize,
    rows_per_partition: usize,
}

impl CountingTable {
    fn new(num_partitions: usize, rows_per_partition: usize) -> Self {
        let schema = Arc::new(Schema::new(vec![
            Field::new("partition", DataType::Int64, false),
            Field::new("value", DataType::Int64, false),
        ]));
        Self { schema, num_partitions, rows_per_partition }
    }
}

#[async_trait::async_trait]
impl TableProvider for CountingTable {
    fn schema(&self) -> SchemaRef { Arc::clone(&self.schema) }
    fn table_type(&self) -> TableType { TableType::Base }

    async fn scan(
        &self,
        _state: &dyn Session,
        projection: Option<&[usize]>,
        _filters: &[Expr],
        limit: Option<usize>,
    ) -> Result<Arc<dyn ExecutionPlan>> {
        // Light work only: build the plan with metadata
        Ok(Arc::new(CountingExec {
            schema: Arc::clone(&self.schema),
            num_partitions: self.num_partitions,
            rows_per_partition: limit
                .unwrap_or(self.rows_per_partition)
                .min(self.rows_per_partition),
            properties: Arc::new(PlanProperties::new(
                EquivalenceProperties::new(Arc::clone(&self.schema)),
                Partitioning::UnknownPartitioning(self.num_partitions),
                EmissionType::Incremental,
                Boundedness::Bounded,
            )),
        }))
    }
}

struct CountingExec {
    schema: SchemaRef,
    num_partitions: usize,
    rows_per_partition: usize,
    properties: Arc<PlanProperties>,
}

impl ExecutionPlan for CountingExec {
    fn name(&self) -> &str { "CountingExec" }
    fn properties(&self) -> &Arc<PlanProperties> { &self.properties }
    fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> { vec![] }

    fn replace_children(
        self: Arc<Self>,
        _: Vec<Arc<dyn ExecutionPlan>>,
        _: ReplaceChildrenOptions,
    ) -> Result<Arc<dyn ExecutionPlan>> {
        Ok(self)
    }

    fn with_new_children(
        self: Arc<Self>,
        children: Vec<Arc<dyn ExecutionPlan>>,
    ) -> Result<Arc<dyn ExecutionPlan>> {
        self.replace_children(children, ReplaceChildrenOptions::new(ChildrenPropertiesMode::Recompute))
    }

    fn execute(
        &self,
        partition: usize,
        _context: Arc<TaskContext>,
    ) -> Result<SendableRecordBatchStream> {
        let schema = Arc::clone(&self.schema);
        let rows = self.rows_per_partition;

        // The heavy work (data generation) happens inside the stream,
        // not here in execute().
        let batch_stream = stream::once(async move {
            let partitions = Int64Array::from(
                vec![partition as i64; rows],
            );
            let values = Int64Array::from(
                (0..rows as i64).collect::<Vec<_>>(),
            );
            let batch = RecordBatch::try_new(
                Arc::clone(&schema),
                vec![Arc::new(partitions), Arc::new(values)],
            )?;
            Ok(batch)
        });

        Ok(Box::pin(RecordBatchStreamAdapter::new(
            Arc::clone(&self.schema),
            batch_stream,
        )))
    }

}

使用你的表提供者

实现 TableProvider 之后,将其注册到 SessionContext 中,即可用于查询:

use datafusion::execution::context::SessionContext;

#[tokio::main]
async fn main() -> Result<()> {
    let ctx = SessionContext::new();

    let provider = CountingTable::new(4, 1000);
    ctx.register_table("counting", Arc::new(provider))?;

    let df = ctx.sql("SELECT * FROM counting LIMIT 10").await?;
    df.show().await?;

    Ok(())
}

延伸阅读

评论

登录后参与评论

正在加载评论…