自定义表提供者
自定义 Table Provider
DataFusion 最强大的优点之一就是它的可扩展性。如果你的数据存放在自定义格式中、某个 API 之后,或者 DataFusion 原生不支持的系统里,你可以通过实现自定义 table provider 来教会 DataFusion 读取它们。本文将带你走完设计 table provider 所需理解的三层抽象,并说明规划(planning)与执行(execution)的工作各自应该发生在哪里。
关于主键或唯一约束等表约束如何处理的详细信息,请参阅表约束强制(Table Constraint Enforcement)。
本文大部分内容最初发表于博客Writing Custom Table Providers in Apache DataFusion。
三层抽象
当 DataFusion 针对一张表执行查询时,三个抽象协同工作以产出结果:
- TableProvider – 描述该表(schema、能力),并在被查询时生成执行计划。它属于**逻辑计划(Logical Plan)**的一部分。
- ExecutionPlan – 描述如何计算结果:分区、排序以及子计划之间的关系。它属于**物理计划(Physical Plan)**的一部分。
- 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() | 规划阶段 | 构建带有元数据的 ExecutionPlan | I/O、网络调用、繁重计算 |
ExecutionPlan::execute() | 执行阶段(每个分区一次) | 构造数据流、建立通道 | 阻塞在异步工作上、读取数据 |
RecordBatchStream(轮询) | 执行阶段 | 所有 I/O、计算、数据生产 | – |
指导原则:尽可能把工作推迟。 规划应当快速完成,这样优化器才能发挥作用。执行的初始化也应当快速完成,以便所有分区都能及时启动。数据流才是你花费时间生产数据的地方。
为什么这很重要
当 scan() 承担繁重工作时,会出现以下问题:
- 规划变得缓慢。 如果一个查询涉及 10 个表,每个
scan()耗时 500 毫秒,那么在任何数据开始流动之前,仅规划就需要 5 秒。 - 执行是单线程的。
scan()在规划期间于单个线程上运行,因此在那里完成的任何工作都无法受益于 DataFusion 跨分区提供的并行执行。 - 优化器无从帮助。 优化器运行在规划与执行之间。如果你在规划阶段就已取回数据,谓词下推或分区裁剪等优化就无法减少工作量。
- 资源管理失效。 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(())
}延伸阅读
TableProviderAPI 文档ExecutionPlanAPI 文档SendableRecordBatchStreamAPI 文档- DataFusion 示例目录 – 包含可运行的示例,其中有自定义表提供者(table provider)的示例
评论
登录后参与评论
KnowForge