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),仅供参考

相关新闻

Mage AI 集成指南:使用 Snowflake Source 连接器读取云数据仓库数据

Mage AI 集成指南:使用 Snowflake Source 连接器读取云数据仓库数据

数据工程数据编排ETL任务调度批处理流处理数据集成后端 【免费下载链接】mage-ai 🧙 Build, run, and manage data pipelines for integrating and transforming data. 项目地址: https://gitcode.com/gh_mirrors/ma/mage-ai 点击查看 免费下载 本指南基…

2026/9/25 7:57:13 阅读更多 →
PHP图书管理系统老代码改造:从部署到借还书事务与乱码修复

PHP图书管理系统老代码改造:从部署到借还书事务与乱码修复

简介:一套面向PHP学习者和毕业设计的图书管理系统源代码,采用PHPMySQL实现,覆盖图书录入、分类管理、模糊搜索、在线借阅、归还处理及用户权限控制等完整业务闭环,适合用于课程实践、毕业设计或作为企业级Web开发的入门范本。压缩…

2026/9/25 7:57:13 阅读更多 →
JUnit 4.5 版本解析:BlockJUnit4ClassRunner 架构演进与 Theories 数据驱动增强

JUnit 4.5 版本解析:BlockJUnit4ClassRunner 架构演进与 Theories 数据驱动增强

测试开发工具 【免费下载链接】junit4 A programmer-oriented testing framework for Java — :warning: maintenance mode 项目地址: https://gitcode.com/gh_mirrors/ju/junit4 点击查看 免费下载 JUnit 4.5(Release Notes 见 doc/ReleaseNotes4.5.md…

2026/9/25 7:57:13 阅读更多 →

最新新闻

网络安全人才缺口480万背后:从入门到入行的实战路线全解析

网络安全人才缺口480万背后:从入门到入行的实战路线全解析

这几年,只要聊到网络安全行业,有一句话几乎必被提到:全球网络安全人才缺口480万。这话我听过不下100遍,也经常被准备入行的朋友追问:既然缺口这么大,为什么我投出去的简历还是没人看?这个问题其…

2026/9/25 8:42:05 阅读更多 →
SQL注入靶场实战:Pikachu平台从搭建到通关的完整渗透笔记

SQL注入靶场实战:Pikachu平台从搭建到通关的完整渗透笔记

1. 靶场环境准备:从零搭建一套可复现的注入实验环境先说句实话。SQL注入的教程在网上一抓一大把,但很多人看完依然不会做题、不会挖洞,核心问题就一个:缺少一套能随手复现、反复折腾的靶场环境。纸上谈兵永远练不出手感&#xff0…

2026/9/25 8:42:05 阅读更多 →
GK7205V300与HI3516EV300对比:H.265编码实测与IPC选型指南

GK7205V300与HI3516EV300对比:H.265编码实测与IPC选型指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/25 8:42:04 阅读更多 →
Windows Server下530/930阵列卡驱动安装与避坑指南

Windows Server下530/930阵列卡驱动安装与避坑指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/25 8:42:04 阅读更多 →
Atlas 300V 24G推理加速卡部署YOLO:从硬件到模型转换全流程

Atlas 300V 24G推理加速卡部署YOLO:从硬件到模型转换全流程

最近总有人拿同一个词条来问我——"atlas 300v 24g 是运算加速卡吗"。会这么问的朋友,手里大概率已经有一块 Atlas 300V,或者正想从某渠道收一块回来跑 YOLO 模型。我每次的答复都一样:它确实是加速卡,但准确说是 AI 推…

2026/9/25 8:42:04 阅读更多 →
IronClaw Reborn CLI 的 Docker 与 Railway 部署实战指南

IronClaw Reborn CLI 的 Docker 与 Railway 部署实战指南

人工智能AI 应用交互助手AI Agent 【免费下载链接】ironclaw IronClaw is an Agent OS focused on privacy, security and extensibility 项目地址: https://gitcode.com/gh_mirrors/iro/ironclaw 点击查看 免费下载 导读 本文以仓库根目录 Dockerfile 与 docker/…

2026/9/25 8:41:04 阅读更多 →

日新闻

AI元人文:从工具使用到思维重构的深度探索

AI元人文:从工具使用到思维重构的深度探索

最近半年我一直在琢磨一件事:AI元人文到底是什么?说白了,就是“用元视角重新审视人与AI的关系”,也在“探索AI如何反向逼着我们发现自己的思考边界”。标题里的“元探索”,在我看就是一层套一层的追问——当你用AI解决…

2026/9/25 0:00:41 阅读更多 →
Python+CNN车牌识别实战:从数据预处理到模型训练与部署

Python+CNN车牌识别实战:从数据预处理到模型训练与部署

简介:基于Python与卷积神经网络的车牌识别项目,面向计算机视觉初学者及智能交通开发者,目标是帮助用户掌握从数据预处理、模型构建到实际部署的完整流程。压缩包共25个文件,包含jpg/png图像样本、py训练脚本、md说明文档、dat数据…

2026/9/25 0:00:41 阅读更多 →
Vim基础操作全攻略:保存退出、模式切换与高频命令实战

Vim基础操作全攻略:保存退出、模式切换与高频命令实战

1. 项目概述1.1 核心需求解析今天聊聊Vim。写这个题目的原因是:几乎每个后端开发者、运维人员、数据工程师某天都会遇到一个场景——深夜加班,服务器登录界面只有黑底白字,编辑器只有vi/vim,你必须在五分钟内完成一次配置修改并保…

2026/9/25 0:00:41 阅读更多 →

周新闻

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

直接铺开项目本身吧。这几个月我一直在折腾一件事:用Flutter给OpenHarmony做一款游戏集合类的App,说白了就是把若干小游戏塞进一个壳里,用统一入口分发。这个方向本身不算新鲜,真正让我花了不少心思的,是首页那堆游戏卡…

2026/9/24 14:34:13 阅读更多 →
Word表格编号全攻略:从列表编号到题注交叉引用

Word表格编号全攻略:从列表编号到题注交叉引用

写Word文档,最让人头疼的往往是那些“看起来不起眼”的小问题。比如表格编号这事:今天在表后面多加了两个空白行,明天给客户交稿前发现整个章节的编号全部错位,光是挨个改序号就能耗掉大半个下午。我前阵子帮人整理一份上百页的技…

2026/9/24 9:10:42 阅读更多 →
从第一个站到第二个站:独立开发者的静态网站选型与落地实践

从第一个站到第二个站:独立开发者的静态网站选型与落地实践

1. 项目概述1.1 核心需求解析做独立开发者这几年,说实话,第一个网站上线的那天晚上我兴奋得没睡着。但等它跑了半年,流量惨淡、功能臃肿、代码自己都懒得看第二遍之后,我才慢慢琢磨明白一个道理:第一个网站是练手&…

2026/9/24 14:33:56 阅读更多 →

月新闻

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能分类:[AI/大模型]细分主题:AI 增强型 CI/CD 流水线自动化与 GitOps 实践:Agent 工作流、工具调用与任务拆解:从原型到生产的验收清单很多团队在尝试用大…

2026/9/24 12:50:34 阅读更多 →
容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场分类:[工程技术]细分主题:Kubernetes 生产环境运维与排障实战:可复制的项目复盘模板与决策记录大部分团队的事故复盘报告,最后都变成了躺在 Confluence 或钉…

2026/9/24 14:33:48 阅读更多 →
容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步分类:[工程技术]细分主题:Docker 容器化技术与镜像安全管理:核心链路的逐步实现与关键代码取舍面对一个积累了五六年历史包袱的单体架构应用(包含 Web 接口、后台…

2026/9/24 12:49:17 阅读更多 →