ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Apache DataFusion 库嵌入指南:在 Rust 项目中以依赖方式使用并扩展查询引擎

Apache DataFusion 库嵌入指南:在 Rust 项目中以依赖方式使用并扩展查询引擎 大数据数据分析后端【免费下载链接】datafusionApache DataFusion SQL Query Engine项目地址https://gitcode.com/gh_mirrors/datafu/datafusion点击查看免费下载DataFusion 不仅仅是一个可独立运行的 SQL 引擎它更是一个设计为可嵌入、可扩展的 Rust 库。本指南面向把 DataFusion 作为依赖集成进自己 Rust 项目的开发者系统讲解如何在Cargo.toml中引入 DataFusion、通过SessionContext执行 SQL 与 DataFrame 查询并围绕其六大扩展点——标量/聚合/窗口/表值用户自定义函数UDF/UDAF/UDWF/UDTF、自定义TableProvider、自定义优化器规则Optimizer Pass、自定义逻辑计划节点LogicalPlan与自定义物理计划节点ExecutionPlan——完成从用库到造库的进阶。读完本文你将掌握一条完整的、可落地的 DataFusion 二次开发路径。快速开始把 DataFusion 引入你的 Rust 工程DataFusion 的库用户指南docs/source/library-user-guide/index.md开篇即说明它讲解的是如何把 DataFusion 作为 Rust 项目中的一个依赖来使用并通过其扩展 API 定制行为。首先在Cargo.toml中加入依赖当前仓库Cargo.toml声明的版本为 55.1.0datafusion 55.1.0 tokio { version 1.0, features [rt-multi-thread] }由于SessionContext的绝大多数 API 是异步的需要 tokio 的多线程运行时配合#[tokio::main]。最小可用示例对应docs/source/user-guide/example-usage.md中的Example Usage如下use datafusion::prelude::*; #[tokio::main] async fn main() - datafusion::error::Result() { // 注册一个 CSV 文件为名为 example 的表 let ctx SessionContext::new(); ctx.register_csv(example, tests/data/example.csv, CsvReadOptions::new()).await?; // 创建并执行 SQL 查询 let df ctx.sql(SELECT a, MIN(b) FROM example WHERE a b GROUP BY a LIMIT 100).await?; df.show().await?; Ok(()) }同样的逻辑也可以完全用 DataFrame API 表达更详细的用法见 使用 DataFrame APIuse datafusion::prelude::*; use datafusion::functions_aggregate::expr_fn::min; #[tokio::main] async fn main() - datafusion::error::Result() { let ctx SessionContext::new(); let df ctx.read_csv(tests/data/example.csv, CsvReadOptions::new()).await?; let df df.filter(col(a).lt_eq(col(b)))? .aggregate(vec![col(a)], vec![min(col(b))])? .limit(0, Some(100))?; df.show().await?; Ok(()) }两段代码输出一致----------- | a | MIN(b) | ----------- | 1 | 2 | -----------需要留意的是SQL 中的所有标识符都会被转成小写。如果 CSV 的列名含大写字母如Name查询时必须给列名加双引号Name或把datafusion.sql_parser.enable_ident_normalization配置设为false来关闭标识符归一化。DataFrame API 侧则可直接用ident(Name)按原样传列名docs/source/user-guide/example-usage.md中有完整对照示例与输出。关于 Arrow 版本匹配的注意事项DataFusion 的大量公开 API 直接使用arrow与parquetcrate 中的类型。如果在你的工程里单独依赖了arrow其版本必须与 DataFusion 所依赖的版本一致否则会报类似mismatched types [E0308] expected Schema, found arrow_schema::Schema的编译错误或出现downcast_ref意外返回None的情况。最稳妥的做法是直接使用 DataFusion 再导出的 Arrow 类型use datafusion::arrow::datatypes::Schema;六大扩展点总览DataFusion 的设计目标是在所有层面都可扩展。库用户指南明确列出了六个官方支持的扩展点对应仓库内 扩展机制文档扩展点作用对象说明仓库内可参考的实现标量 UDFUDF表达式输入一行、返回一个值按 Arrow 批量向量化求值simple_udf.rs、advanced_udf.rs、async_udf.rs聚合 UDFUDAF表达式输入一组行、返回一个值等价于SUM/COUNTsimple_udaf.rs、advanced_udaf.rs、struct_returning_udaf.rs窗口 UDFUDWF表达式输入一行、可访问其周围的行实现如移动平均simple_udwf.rs、advanced_udwf.rs表值 UDFUDTF表接收参数并返回一个TableProvider参与查询计划simple_udtf.rs、table_list_udtf.rs自定义TableProvider表教会 DataFusion 读取任何自定义格式/API/存储系统中的数据custom_datasource.rs、custom_file_format.rs自定义优化器规则计划在逻辑/物理优化阶段做计划重写plan rewriteoptimizer_rule.rs自定义LogicalPlan节点计划在逻辑计划树中引入新算子building-logical-plans.md自定义ExecutionPlan节点计划在物理计划树中引入新的执行算子custom_datasource.rs 中的执行计划部分下文按表达式扩展 → 表扩展 → 计划层扩展的顺序逐一展开深入细节可查阅 adding-udfs.md 与 custom-table-providers.md。扩展点一标量 UDFScalar UDF标量函数接收一行数据、返回一个值。为保证性能DataFusion 对标量 UDF 采用向量化执行函数拿到的是一个或多个 Arrow Array输出一个行数相同的 Array。方式 A实现ScalarUDFImpltrait推荐功能最全两步走先实现ScalarUDFImpl告诉 DataFusion 函数的名称、签名、返回类型与计算逻辑再用ScalarUDF::from(...)包装并通过SessionContext::register_udf注册。以对 Int32 加一为例详见 adding-udfs.mduse std::sync::Arc; use arrow::datatypes::DataType; use datafusion_common::cast::as_int64_array; use datafusion_common::{plan_err, Result}; use datafusion_expr::{ColumnarValue, ScalarFunctionArgs, Signature, Volatility}; use datafusion::arrow::array::{ArrayRef, Int64Array}; use datafusion_expr::{ScalarUDFImpl, ScalarUDF}; #[derive(Debug, PartialEq, Eq, Hash)] struct AddOne { signature: Signature, } impl AddOne { fn new() - Self { Self { signature: Signature::uniform(1, vec![DataType::Int32], Volatility::Immutable), } } } impl ScalarUDFImpl for AddOne { fn name(self) - str { add_one } fn signature(self) - Signature { self.signature } fn return_type(self, args: [DataType]) - ResultDataType { if !matches!(args.get(0), Some(DataType::Int32)) { return plan_err!(add_one only accepts Int32 arguments); } Ok(DataType::Int32) } fn invoke_with_args(self, args: ScalarFunctionArgs) - ResultColumnarValue { let args ColumnarValue::values_to_arrays(args.args)?; let i64s as_int64_array(args[0])?; let new_array i64s .iter() .map(|array_elem| array_elem.map(|value| value 1)) .collect::Int64Array(); Ok(ColumnarValue::from(Arc::new(new_array) as ArrayRef)) } }随后注册并调用use datafusion::execution::context::SessionContext; let add_one ScalarUDF::from(AddOne::new()); let expr add_one.call(vec![col(a)]); // 构造表达式 add_one(col(a)) let mut ctx SessionContext::new(); ctx.register_udf(add_one.clone()); // 注册后即可在 SQL 中使用 add_one(...)实现时值得注意生产级代码还应校验args.args.len()与期望参数个数一致return_type应对每个入参类型做匹配检查。更完整的低层 API 用法含user_doc宏生成函数文档等参见 advanced_udf.rs。方式 Bcreate_udf便捷构造更简短但能力有限先把核心逻辑写成一个接收[ColumnarValue]、返回ResultColumnarValue的纯函数use std::sync::Arc; use datafusion::arrow::array::{ArrayRef, Int64Array}; use datafusion::common::cast::as_int64_array; use datafusion::common::Result; use datafusion::logical_expr::ColumnarValue; pub fn add_one(args: [ColumnarValue]) - ResultColumnarValue { let args ColumnarValue::values_to_arrays(args)?; let i64s as_int64_array(args[0])?; let new_array i64s .iter() .map(|array_elem| array_elem.map(|value| value 1)) .collect::Int64Array(); Ok(ColumnarValue::from(Arc::new(new_array) as ArrayRef)) }然后用create_udf包装成ScalarUDF并注册use datafusion::logical_expr::{Volatility, create_udf}; use datafusion::arrow::datatypes::DataType; use datafusion::execution::context::SessionContext; let udf create_udf( add_one, // 第 1 参SQL 中使用的函数名 vec![DataType::Int64], // 第 2 参入参类型列表 DataType::Int64, // 第 3 参返回类型 Volatility::Immutable, // 第 4 参易变性见下 Arc::new(add_one), // 第 5 参函数实现 ); let mut ctx SessionContext::new(); ctx.register_udf(udf); let df ctx.sql(SELECT add_one(1)).await.unwrap();create_udf五个参数各有讲究第 1 参是 SQL 可见的函数名第 2 参声明可接受的入参类型第 3 参声明返回类型第 4 参Volatility决定优化器能否做常量折叠等优化——Immutable表示同输入必同输出如本函数随机数生成器则应标为Volatile第 5 参是上面写的实现函数。异步标量 UDF在 UDF 里做网络/I/O 调用当函数内部需要执行异步操作如远程调用、文件 I/O时可改用AsyncScalarUDFImpltrait 实现invoke_async_with_args再用AsyncScalarUDF::new(...)包装、into_scalar_udf()转成普通标量 UDF 后注册完整示例见 async_udf.rslet async_upper AsyncUpper::new(); let udf AsyncScalarUDF::new(Arc::new(async_upper)); let mut ctx SessionContext::new(); ctx.register_udf(udf.into_scalar_udf()); // 之后可直接查询SELECT async_upper(datafusion);异步实现里ideal_batch_size()可提示引擎按多大批次切分输入示例中返回Some(10)而普通invoke_with_args通常返回not_impl_err!(... can only be called from async contexts)表示该函数只能走异步路径执行。扩展点二聚合 UDFUDAF聚合函数接收一组行、返回单个值与内置SUM、COUNT同族。核心是Accumulatortrait——它持有跨多行的状态。以几何平均数geo_mean为例见 adding-udfs.md 与 simple_udaf.rsuse datafusion::arrow::array::ArrayRef; use datafusion::scalar::ScalarValue; use datafusion::{error::Result, physical_plan::Accumulator}; #[derive(Debug)] struct GeometricMean { n: u32, prod: f64 } impl GeometricMean { pub fn new() - Self { GeometricMean { n: 0, prod: 1.0 } } } impl Accumulator for GeometricMean { // 把累加器状态序列化为 ScalarValue供跨执行阶段传递 fn state(mut self) - ResultVecScalarValue { Ok(vec![ ScalarValue::from(self.prod), ScalarValue::from(self.n), ]) } // 返回最终聚合值几何平均 prod^(1/n) fn evaluate(mut self) - ResultScalarValue { let value self.prod.powf(1.0 / self.n as f64); Ok(ScalarValue::from(value)) } // 用一批输入行更新状态 fn update_batch(mut self, values: [ArrayRef]) - Result() { if values.is_empty() { return Ok(()); } let arr values[0]; (0..arr.len()).try_for_each(|index| { let v ScalarValue::try_from_array(arr, index)?; if let ScalarValue::Float64(Some(value)) v { self.prod * value; self.n 1; } Ok(()) }) } // 合并其他分区/阶段的中间状态 fn merge_batch(mut self, states: [ArrayRef]) - Result() { if states.is_empty() { return Ok(()); } let arr states[0]; (0..arr.len()).try_for_each(|index| { let v states .iter() .map(|array| ScalarValue::try_from_array(array, index)) .collect::ResultVec_()?; if let (ScalarValue::Float64(Some(prod)), ScalarValue::UInt32(Some(n))) (v[0], v[1]) { self.prod * prod; self.n n; } Ok(()) }) } fn size(self) - usize { std::mem::size_of_val(self) } }create_udaf有六个参数与create_udf相比多出一个状态描述let geometric_mean create_udaf( geo_mean, // 函数名 vec![DataType::Float64], // 入参类型 Arc::new(DataType::Float64), // 返回类型 Volatility::Immutable, // 易变性 Arc::new(|_| Ok(Box::new(GeometricMean::new()))), // 累加器工厂 Arc::new(vec![DataType::Float64, DataType::UInt32]), // 状态描述须与 state() 的类型一致 ); let ctx SessionContext::new(); ctx.register_udaf(geometric_mean);声明 UDAF 对DISTINCT的处理方式默认情况下 DataFusion 假定聚合函数对DISTINCT敏感即累加器需要读取AccumulatorArgs::is_distinct并自行去重。若你的函数不是这样应覆写AggregateUDFImpl::distinct_handling返回DistinctHandling::Insensitive重复值不影响结果合并已见过的值是空操作如min、max、bool_and、bit_or。此时优化器会把f(DISTINCT x)直接规划为f(x)省掉哈希集合与SingleDistinctToGroupBy引入的额外分组阶段。返回DistinctHandling::Unsupported累加器不实现DISTINCT由规划器先做去重或拒绝查询目前仅作为声明。保持默认DistinctHandling::Sensitive累加器自己读is_distinct并去重。这个声明直接影响查询结果只有合并操作真正幂等时才可声明Insensitive。返回多值的聚合 UDF若一个聚合结果需要携带多个值如时间窗口扩展中同时返回窗口起止时间与聚合值可以让聚合返回DataType::Structevaluate返回ScalarValue::Struct调用方再用[...]取字段例如augmented_avg(time, value)[window_start]。仓库中 struct_returning_udaf.rs 提供了完整可运行示例。扩展点三窗口 UDFUDWF窗口函数与标量函数相似但能访问目标行周围的行。实现上需提供PartitionEvaluator每个PARTITION BY分区一份最简单的求值方式是实现evaluate(values, range)range指明当前窗口帧覆盖的索引区间。以下是一个移动平均smooth_it的实现骨架完整代码见 adding-udfs.md 与 simple_udwf.rsuse datafusion::arrow::{array::{ArrayRef, Float64Array, AsArray}, datatypes::Float64Type}; use datafusion::logical_expr::PartitionEvaluator; use datafusion::common::ScalarValue; use datafusion::error::Result; #[derive(Clone, Debug)] struct MyPartitionEvaluator {} impl PartitionEvaluator for MyPartitionEvaluator { // 告知 DataFusion函数结果随窗口帧变化 fn uses_window_frame(self) - bool { true } // 逐行调用range 指明参与计算的 values 索引范围 fn evaluate(mut self, values: [ArrayRef], range: std::ops::Rangeusize) - ResultScalarValue { let arr: Float64Array values[0].as_ref().as_primitive::Float64Type(); let range_len range.end - range.start; let output if range_len 0 { let sum: f64 arr.values().iter().skip(range.start).take(range_len).sum(); Some(sum / range_len as f64) } else { None }; Ok(ScalarValue::Float64(output)) } } fn make_partition_evaluator() - ResultBoxdyn PartitionEvaluator { Ok(Box::new(MyPartitionEvaluator::new())) }注册同样走create_udwf帮助函数注意第二个参数是单个输入DataType而非类型列表use datafusion::logical_expr::{Volatility, create_udwf}; use datafusion::arrow::datatypes::DataType; use datafusion::execution::context::SessionContext; let smooth_it create_udwf( smooth_it, // 函数名 DataType::Float64, // 输入类型单个非列表 Arc::new(DataType::Float64), // 返回类型 Volatility::Immutable, // 易变性 Arc::new(make_partition_evaluator), // 分区求值器工厂 ); let ctx SessionContext::new(); ctx.register_udwf(smooth_it);注册后即可在 SQL 中使用对 cars.csv 按 car 分区、按 time 排序求移动平均SELECT car, speed, smooth_it(speed) OVER (PARTITION BY car ORDER BY time) as smooth_speed, time FROM cars ORDER BY car;evaluate是最通用但也最慢的求值方式PartitionEvaluator还提供了evaluate_all、evaluate_all_with_rank等批量接口在性能敏感场景应优先实现它们。更高级的低层 API 见 advanced_udwf.rs。扩展点四表值 UDFUDTF表值函数接收参数并返回一个TableProvider。实现只需TableFunctionImpltrait 的单一方法call_with_args。下面这个echo函数接收一个Int64字面量返回含单列单行的表见 simple_udtf.rsuse std::sync::Arc; use datafusion::common::{plan_err, ScalarValue, Result}; use datafusion::catalog::{TableFunctionArgs, TableFunctionImpl, TableProvider}; use datafusion::arrow::array::Int64Array; use datafusion::datasource::memory::MemTable; use arrow::record_batch::RecordBatch; use arrow::datatypes::{DataType, Field, Schema}; #[derive(Debug, Default)] pub struct EchoFunction {} impl TableFunctionImpl for EchoFunction { fn call_with_args(self, args: TableFunctionArgs) - ResultArcdyn TableProvider { let exprs args.exprs(); let Some(Expr::Literal(ScalarValue::Int64(Some(value)), _)) exprs.get(0) else { return plan_err!(First argument must be an integer); }; let schema Arc::new(Schema::new(vec![Field::new(a, DataType::Int64, false)])); let batch RecordBatch::try_new(schema.clone(), vec![Arc::new(Int64Array::from(vec![*value]))])?; let provider MemTable::try_new(schema, vec![vec![batch]])?; Ok(Arc::new(provider)) } }注册与使用use datafusion::execution::context::SessionContext; let ctx SessionContext::new(); ctx.register_udtf(echo, Arc::new(EchoFunction::default())); let results ctx.sql(SELECT * FROM echo(1)).await?.collect().await?;UDTF 特别适合读取外部数据源与交互式分析。DataFusion 自带的内置 UDTFparquet_metadata便是一个真实案例——在 CLI 中可直接查询 Parquet 文件的元数据adding-udfs.md 给出了hits.parquet的 row group 统计输出示例。扩展点五自定义TableProvider接入任意数据源这是 DataFusion 可扩展性最强的地方数据在自定义格式、API 背后或 DataFusion 原生不支持的系统中时实现一个custom table provider即可教会引擎读取它。相关完整讲解见 custom-table-providers.md可运行示例见 custom_datasource.rs 与 custom_file_format.rs。三层协作架构查询执行时三个抽象依次协作可以理解为漏斗TableProvider——描述表的 schema 与能力被查询时产出执行计划属于逻辑计划层。ExecutionPlan——描述如何计算结果分区、排序、子计划关系属于物理计划层。SendableRecordBatchStream——真正干活的异步流逐个产出RecordBatch。调用关系物理规划时调用一次TableProvider::scan()生成ExecutionPlan执行阶段对每个分区调用一次ExecutionPlan::execute()生成流行数据在流被轮询poll时产生。关键原则scan()在规划阶段运行必须保持轻量——不要做 I/O、网络调用或重计算它只负责描述数据如何产生所有重活应下沉到流中。如果scan()里取数据/开连接会阻塞规划线程在多表或子查询场景下可能引发超时甚至死锁。scan()接收三个来自优化器的下推提示参数作用减少什么projection指明需要哪些列输出宽度可只读这些列filters希望源在扫描时应用的谓词输出行数跳过不匹配数据limit行数上限输出行数产够即可提前停止想声明自己可以处理哪些谓词覆写supports_filters_pushdown对每个过滤器返回三选一Exact源保证输出中没有使该谓词为假的行DataFusion 不会再叠加FilterExec。Inexact源能减少数据量但可能仍有漏网行如按文件元数据跳文件但不过滤文件内行引擎仍会在扫描之上保留FilterExec。Unsupported源忽略该过滤器由 DataFusion 处理。ExecutionPlan 与并行度ExecutionPlan最关键的两个属性是输出分区partitioning与输出排序ordering前者决定并行度execute()每个分区调用一次每个分区对应 tokio 运行时上的一个 tasktask 是复用在线程池上的轻量异步单元后者声明数据天然有序时可让优化器省掉SortExec。起步建议匹配数据自然布局4 个文件就暴露 4 个分区、8 个分片就暴露 8 个分区下游需要其他分布时 DataFusion 会自动插入RepartitionExec。进阶做法是通过state.config().target_partitions()读取会话的目标分区数并尽量对齐或在源本身就是按某 key 预分区的场景下声明哈希分区如Hash([customer_id], N)这样GROUP BY customer_id的聚合可以免去重分区算子。反之若一律报告UnknownPartitioning引擎只能按最坏情况插入重分区。以上都可以用EXPLAIN验证——EXPLAIN也是排查表提供器问题多余的SortExec、RepartitionExec、FilterExec的第一工具。流的实现与阻塞工作隔离创建SendableRecordBatchStream最简便的方式是RecordBatchStreamAdapter它把任意futures::StreamItem ResultRecordBatch桥接为目标类型。若流里要做长时间阻塞工作无法让出执行权的同步 I/O 或耗时数百毫秒的 CPU 任务务必用tokio::task::spawn_blocking投递到独立线程池再通过tokio::sync::mpsc通道把结果送回流避免阻塞 tokio 异步运行时。线程池隔离的完整示例见 thread_pools.rs。三层职责速查层运行时机应该做不应该做TableProvider::scan()规划期构造带元数据的ExecutionPlanI/O、网络、重计算ExecutionPlan::execute()执行期每分区一次构造流、建立通道阻塞异步操作、读取数据RecordBatchStream轮询执行期全部 I/O、计算、数据产出——总体原则是尽可能把工作推迟到最晚阶段。表提供器还支持行级 DML实现TableProvider::delete_from()与TableProvider::update()即可支持DELETE/UPDATE默认实现返回未实现错误方法返回的执行计划执行后通过count列报告受影响行数MemTable提供了现成的内存实现作为参考。一个完整的最小表提供器CountingTable流式惰性生成数据在 custom-table-providers.md 的Putting It All Together一节注册方式为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?;若数据已在内存RecordBatch、是异步批次流、是其他表的逻辑变换、或是已有文件格式的变体则不必从零实现三层可分别直接使用MemTable、StreamTable、ViewTable、ListingTable配合自定义FileFormat/FileSource/FileOpener作为起点——参考选择表见 custom-table-providers.md。扩展点六自定义优化器规则与计划节点优化器规则Optimizer PassDataFusion 允许注册自定义优化器规则来做计划重写plan rewrite在保持查询语义的前提下减少工作量。仓库内 optimizer_rule.rs 演示了如何实现并注册一条优化器规则。优化器的整体架构、常见规则谓词下推、投影裁剪、表达式简化、子查询去相关、Limit 下推等见 query-optimizer.md。自定义LogicalPlan节点与ExecutionPlan节点逻辑计划节点在LogicalPlan树中引入新的关系算子需要实现相应的 plan node 类型并接入逻辑规划/优化流程指导见 building-logical-plans.md。物理计划节点在ExecutionPlan树中引入新的执行算子前面自定义表提供器里的MyExecPlan/CountingExec就是典型例子需要实现name、properties、children、replace_children、execute等方法并在properties()中正确设置PlanProperties输出分区与排序。扩展 SQL 语法若需要支持这类引擎未内置的运算符或TABLESAMPLE等方言特性可通过实现ExprPlanner、TypePlanner、RelationPlanner并注册到SessionContext来接管默认规划逻辑完整指南见 extending-sql.md。命名参数Named ArgumentsDataFusion 的标量、窗口、聚合 UDF 都支持按参数名传参。只要在Signature上调用.with_parameter_names(vec![...])声明参数名顺序须与签名一致用户即可任意调换命名参数顺序或将命名参数与位置参数混用位置参数必须在前。例如SELECT power(base 2.0, exponent 3.0)与SELECT power(exponent 3.0, base 2.0)等价调用出错时错误消息会显示参数名帮助排查No function matches the given name and argument types substr(Utf8). Candidate functions: substr(str: Any, start_pos: Any) substr(str: Any, start_pos: Any, length: Any)架构入门与进一步学习路径在深入某个扩展点之前建议先了解整体架构。DataFusion 的架构文档位于 datafusion/doc/src/lib.rs对应docs.rs上datafusioncrate 的 Architecture 章节它帮助你建立逻辑计划 → 逻辑优化 → 物理计划 → 物理优化 → 执行的全链路心智模型。查询处理管线大致为SQL / DataFrame API → Logical Plan (抽象计算什么) → Logical Optimization (保持语义的重写规则) → Physical Plan (具体如何计算) → Physical Optimization (面向硬件与数据的重写) → Execution (流式 RecordBatch)按主题深入时可依次阅读想先跑通 SQL 与 DataFrame查看 docs/source/user-guide/example-usage.md本文开头示例即出自该页以及 使用 SQL APISessionContext::register_csv/register_parquet/register_avro注册表、CREATE EXTERNAL TABLE语句、assert_batches_eq!断言宏。想参与社区贡献查阅 contributor-guidePR 流程、./dev/rust_lint.sh检查、Conventional Commits 约定等。想加深表达式知识阅读 working-with-exprs.md 与 extending-operators.md。想管理目录与约束阅读 catalogs.md 与 table-constraints.md。结语从加一行依赖跑通 CSV 查询到实现向量化标量函数、带状态的聚合累加器、按分区求值的窗口函数、返回表提供器的表值函数再到用TableProvider三层架构接入任意数据源、用优化器规则重写计划、用自定义LogicalPlan/ExecutionPlan节点引入新算子——DataFusion 的六大扩展点覆盖了表达式、表与计划三个层面形成了完整而自洽的扩展体系。仓库中 datafusion-examples/examples 目录提供了全部扩展点的可运行参考实现是动手实践时最直接的对照样本。把这些能力组合起来你就能把 DataFusion 真正改造成你自己的查询引擎。赞分享大数据数据分析后端【免费下载链接】datafusionApache DataFusion SQL Query Engine项目地址https://gitcode.com/gh_mirrors/datafu/datafusion点击查看免费下载相关推荐Apache DataFusion 查询引擎全景解析Rust Arrow 打造可嵌入式、可扩展的数据系统基座Apache DataFusion 查询引擎全景解析Rust Arrow 打造可嵌入式、可扩展的数据系统基座 Apache DataFusion 是一个用大数据数据分析后端Apache DataFusion与DuckDB对比嵌入式查询引擎测评Apache DataFusion与DuckDB对比嵌入式查询引擎测评 你是否在为数据处理工具的选择而困扰当需要在应用中嵌入高性能查询能力时Apache大数据数据分析后端终极指南如何在Rust项目中集成Apache DataFusion高性能查询引擎终极指南如何在Rust项目中集成Apache DataFusion高性能查询引擎 Apache DataFusion是一个基于Rust构建的 高性能SQL查询大数据数据分析后端上一篇如何突破平台限制实现高效数据采集MediaCrawler跨平台聚合方案解析下一篇cds-textarea 渲染结构全解析从快照到源码的 Carbon Web Components 多行文本框深度指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表