大数据数据分析后端【免费下载链接】datafusionApache DataFusion SQL Query Engine项目地址https://gitcode.com/gh_mirrors/datafu/datafusion点击查看免费下载导读本文以 Apache DataFusion本仓库gh_mirrors/datafu/datafusion的官方架构文档为核心系统讲解这套用 Rust 编写、以 Apache Arrow 为内存格式的列式查询引擎的代码组织与扩展机制为什么社区推荐用扩展 API 而非维护 Fork、如何在规划/优化/执行各阶段插入自定义逻辑、SQL 与 DataFrame 两条 API 如何汇聚到同一套 LogicalPlan 管线以及各子 crate 的职责边界。读完本文你将掌握 DataFusion 的架构分层、六大扩展点的源码级位置以及从提 issue 到合入 PR 的完整扩展 API 创建流程。架构文档的权威来源与阅读指引DataFusion 的代码结构与组织方式由官方文档给出为保持与源码尽可能一致最权威的架构说明并不在docs/source/contributor-guide/architecture.md这一篇中而是内嵌于主 crate 的 rustdoc 注释里——即 datafusion/core/src/lib.rs 的# Architecture章节。该章节与源码同文件、同提交永远不会出现文档滞后于代码的问题。此外工作区各 crate 之间的内部依赖关系见 docs/source/contributor-guide/architecture/dependency-graph.md其中的交互式 SVG 依赖图由 docs/scripts/generate_dependency_graph.sh 通过cargo depgraph自动生成并在docs/build.sh构建文档时自动更新保证图与代码同步。社区维护的扩展清单见 docs/source/library-user-guide/extensions.md。设计目标三个原则决定架构走向从 datafusion/core/src/lib.rs 的文档可知DataFusion 的架构设计遵循三条目标开箱即用Work out of the box以最小配置提供一个非常快、世界级的查询引擎用户不需要做大量设置就能跑起来。一切皆可定制Customizable everything所有行为都应能通过实现 trait 来定制。这条原则直接决定了后续的扩展 API 设计。架构上保守Architecturally boring遵循工业界成熟最佳实践而不是追逐未经大规模验证的激进技术。在这三条原则下用户从一台基础但高性能的引擎出发随着需求与工程能力的增长逐步把它特化到自己的业务场景——这正是默认引擎 扩展点架构存在的理由。Forks 与扩展 API为什么社区强烈建议后者架构文档docs/source/contributor-guide/architecture.md用专门一节讨论了 Fork 与扩展 API 的取舍这是理解 DataFusion 社区协作模式的关键。DataFusion 是一个快速演进的活跃项目内部结构频繁调整。这对 DataFusion 本身是好事能快速响应需求、持续进化但对带着大改动维护一个 Fork的团队而言意味着持续的合并成本每次上游重构都需要手工移植。相比之下通过 crates.io 发布使用的公共 API通常稳定得多虽然每个版本之间也会变动。因此官方文档的明确建议是与其维护 Fork不如使用丰富的扩展 API来定制 DataFusion例如TableProvider、OptimizerRule、ExecutionPlan等。如果现有 API 无法满足需求社区欢迎你通过下文所述流程共同设计新 API而不是分叉代码。注datafusion/core/src/lib.rs中列出了完整的扩展点清单详见下文六大扩展点一节各扩展的社区实现与贡献入口见 docs/source/library-user-guide/extensions.md。创建新的扩展 API从 issue 到 PR 的标准流程架构文档给出了新增扩展 API 的典型流程docs/source/contributor-guide/architecture.md查找或创建 issue先在 issue 列表里搜索是否已有描述你需求的议题没有就新建一个。讨论 API 形态通过提及相关贡献者可以从最近改动过的 PR 和 issue 中找到这些人获取反馈把 API 设计讨论清楚。用示例原型化新 API典型做法是在 datafusion-examples/examples 中添加一个示例或重构现有代码展示新 API 如何工作。提交 PR 并合入带着新 API 创建 PR与社区协作直到合并。采用示例驱动方式有额外收益未来 API 变更时示例会迫使你同步更新保证功能不回归同时示例也是你代码需要改动时的蓝本只需看示例里改了什么。文档还以SQL Extension Planning API的创建过程作为该流程的实例参考。什么 API 适合进入 DataFusion 核心扩展 API 是否适合被纳入 DataFusion 核心取决于其默认行为是否安全提供安全默认行为的扩展 API 更可能被接受需要大幅改动内置算子实现的 API 则较难被接受。例如若某流处理特性会导致内置算子性能下降那它就不太适合直接进入核心——这种场景更合理的做法是核心提供扩展 API而把具体算子实现留给下游项目。六大扩展点在哪些位置插入自定义逻辑datafusion/core/src/lib.rs的 Customization and Extension 一节datafusion/core/src/lib.rs明确了 DataFusion 支持在几乎所有环节定制从源码结构看可归纳为以下六大扩展点扩展点对应 trait / 类型用途数据源TableProvider从任意数据源读取数据目录与 SchemaCatalogProvider、SchemaProvider自定义 catalog、schema 与表列表查询语言与计划构建LogicalPlanBuilder自建查询语言直接构建LogicalPlan而非走内置 SQL 规划器用户自定义函数ScalarUDF、AggregateUDF、WindowUDF声明并使用标量/聚合/窗口函数计划重写AnalyzerRule、OptimizerRule、PhysicalOptimizerRule添加自定义计划改写与优化规则规划器扩展QueryPlanner让规划器支持用户自定义的逻辑/物理节点各扩展点的 trait 定义可分别在datafusion_optimizer::analyzer::AnalyzerRule、datafusion_optimizer::optimizer::OptimizerRule、datafusion_physical_optimizer::PhysicalOptimizerRule、datafusion_catalog::CatalogProvider、datafusion_expr::logical_plan::builder::LogicalPlanBuilder等模块中找到。每个扩展点都有对应示例位于 datafusion-examples/examples 目录下按custom_data_source、udf、extension_types、query_planning等子目录组织是学习各扩展 API 用法的最佳起点。查询规划与执行概览从 SQL 到结果的两阶段管线DataFusion 的查询处理分为逻辑规划与物理执行两个阶段完整描述见 datafusion/core/src/lib.rs。阶段一SQL 字符串 → LogicalPlanParsed with SqlToRel creates sqlparser initial plan ┌───────────────┐ ┌─────────┐ ┌─────────────┐ │ SELECT * │ │Query { │ │Project │ │ FROM ... │──────────▶│.. │────────────▶│ TableScan │ │ │ │} │ │ ... │ └───────────────┘ └─────────┘ └─────────────┘ SQL String sqlparser LogicalPlan AST nodes查询字符串先由 [sqlparser] 解析为抽象语法树ASTStatement再由SqlToRel位于datafusion_sqlcrate把 AST 转换为LogicalPlan与逻辑表达式Expr该阶段同时完成名称解析与类型解析即 binding。而使用 DataFrame API 时流程与 SQL 完全一致唯一区别是 DataFrame API 直接通过LogicalPlanBuilder构建LogicalPlan。自带自定义查询语言的系统通常也直接构建LogicalPlan。阶段二LogicalPlan → 优化 → ExecutionPlanAnalyzerRules and PhysicalPlanner PhysicalOptimizerRules OptimizerRules creates ExecutionPlan improve performance rewrite plan ┌─────────────┐ ┌─────────────┐ ┌─────────────────┐ ┌─────────────────┐ │Project │ │Project(x, y)│ │ProjectExec │ │ProjectExec │ │ TableScan │──...──▶│ TableScan │─────▶│ ... │──...──▶│ ... │ │ ... │ │ ... │ │ DataSourceExec│ │ DataSourceExec│ └─────────────┘ └─────────────┘ └─────────────────┘ └─────────────────┘ LogicalPlan LogicalPlan ExecutionPlan ExecutionPlan为了尽可能高效地处理海量行数据DataFusion 在规划与优化阶段投入了大量工作按以下顺序进行AnalyzerRule检查并重写LogicalPlan强制执行类型转换等语义规则OptimizerRule重写LogicalPlan例如投影下推、过滤下推以提升效率PhysicalPlanner把LogicalPlan转换为可执行的ExecutionPlanPhysicalOptimizerRule重写ExecutionPlan例如选择排序与连接算法进一步提升效率。这一流程意味着无论你的定制发生在哪个阶段分析期、逻辑优化期、物理优化期都有对应的 trait 扩展点可以插入。数据源抽象TableProvider 连接规划与执行DataFusion 的TableProvider是数据接入层的核心抽象datafusion/core/src/lib.rs规划阶段向它请求 schema 等元信息执行阶段由它的scan方法创建ExecutionPlan。DataFusion 内置了三个开箱即用的TableProviderListingTable从一个或多个本地/远程目录读取 Parquet、JSON、CSV、Avro 文件支持 Hive 风格分区、可选压缩、直接从远程对象存储读取、文件元数据缓存等MemTable从内存中的RecordBatch读取数据StreamingTable从潜在无界输入读取数据。要实现任何其他数据源或文件格式只需实现TableProvidertrait 即可接入。计划的两种表示LogicalPlan 与 ExecutionPlan逻辑计划LogicalPlan由LogicalPlan节点与Expr表达式构成是 schema 感知的描述要算什么与物理执行方式无关。LogicalPlan本质上是其他LogicalPlan构成的有向无环图DAG每个节点可能内嵌若干Expr。逻辑计划可用TreeNodeAPI 重写见 datafusion/core/src/lib.rsExpr还可通过ExprSimplifier简化相关示例位于 datafusion-examples/examples/query_planning/expr_api.rs。物理计划ExecutionPlan是可以对着数据执行的计划同样是ExecutionPlan的 DAG每个节点包含实现PhysicalExprtrait 的表达式。相比逻辑计划物理计划携带了具体的计算方式如 hash join 还是 merge join与数据流信息如分区方式、有序性。此外cp_solver对PhysicalExpr做区间传播分析PruningPredicate可借助统计信息证明某些过滤表达式永远不可能为true从而跳过数据读取——这是 DataFusion 行级裁剪性能的重要来源。执行模型基于 Pull 的流式执行与批量处理ExecutionPlan 的执行协议DataFusion 的执行协议datafusion/core/src/lib.rs基于 Apache Arrow 内存格式ExecutionPlan::execute Calling next() on the produces a stream stream produces the data ┌────────────────┐ ┌─────────────────────────┐ ┌────────────┐ │ProjectExec │ │impl │ ┌───▶│RecordBatch │ │ ... │─────▶│SendableRecordBatchStream│────┤ └────────────┘ │ DataSourceExec│ │ │ │ ┌────────────┐ └────────────────┘ └─────────────────────────┘ ├───▶│RecordBatch │调用ExecutionPlan::execute会为每个分区产出一个SendableRecordBatchStream这是一个pull 式执行 API消费者反复调用next().await流会增量计算并返回下一个RecordBatch。值以ColumnarValue表示要么是单个常量ScalarValue要么是 Arrow 数组ArrayRef。跨线程的均衡并行度通过经典的 Volcano 风格 Exchange 算子实现即RepartitionExec。流式执行与 pipeline breakerDataFusion 是流式查询引擎datafusion/core/src/lib.rsExecutionPlan从一个RecordBatch开始增量读取输入、计算输出每个输出/中间RecordBatch约含batch_size行从而摊薄逐批执行开销。以如下 SQL 为例SELECT name FROM data.parquet WHERE id 10其简化的执行管线为Parquet File → DataSource → FilterExec (id 10) → ProjectionExec (keeps name) → Results每一步一次只处理一个RecordBatch多分区时多个批次可在不同 CPU 核上并发处理。需要注意pipeline breaker如全量排序、hash 聚合本质上是非流式的必须读完整输入才能产生任何输出除此之外其他算子尽量做到读一个 batch、出一个 batch。这种 pull 式控制流还带来良好的缓存局部性生产数据的 CPU 核往往立即消费它。线程调度与 Tokio RuntimeCPU/IO 资源管理DataFusion 使用 TokioRuntime作为线程池自动用多核执行每个计划datafusion/core/src/lib.rs使用的核心数由配置项target_partitions决定默认等于 CPU 核心数执行准备阶段会为每个ExecutionPlan创建相应数量的独立异步Stream某些算子的Stream如RepartitionExec、CoalescePartitionsExec会 spawn Tokio task运行在Runtime管理的线程上DataFusion 采用协作式调度每个Stream在完成一定工作量后主动把控制权交还给Runtime详见datafusion_physical_plan::coop模块。高并发负载下的 CPU/IO 分离async设计让TableProvider在执行中能方便地用标准 Rustasync做网络 I/O但也容易把 CPU 密集与延迟敏感的 I/O 混在同一个线程池里。处理本地文件或初始开发时通常没问题但当负载升高或从 AWS S3 等网络源读取时可能出现 CPU 或网络带宽利用率不足、尾延迟如 p99显著升高——此时很可能需要为 DataFusion 计划使用独立的Runtime参考 datafusion-examples/examples/query_planning/thread_pools.rs 中的示例。此外需注意DataFusion不使用tokio::task::spawn_blocking处理 CPU 密集任务因为spawn_blocking是为阻塞 I/O 设计的spawn 的阻塞任务无法在等待输入时让出不能await因而既不能限制并发 CPU 任务数也无法让处理管线停留在同一核上。状态管理与资源控制SessionContext / TaskContext / ExecutionProps执行查询所需的状态由三个结构分层管理datafusion/core/src/lib.rsSessionContext创建LogicalPlan所需的状态如表定义与函数注册表TaskContext执行所需的状态如MemoryPool、DiskManager、ObjectStoreRegistryExecutionProps单次执行相关的属性与数据如起始时间戳等。资源管理方面计划运行时的内存与临时磁盘用量分别由MemoryPool与DiskManager控制其他运行时选项见RuntimeEnv。所有执行选项统一由ConfigOptionsdatafusion_common::config承载可通过会话配置调整。Crate 组织模块化拆分与依赖关系大多数用户通过主 cratedatafusion本仓库 datafusion/core交互它 re-export 了构建与执行查询所需的全部功能模块 re-export 见 datafusion/core/src/lib.rs。另有三个需要直接使用的附加 cratedatafusion-proto计划序列化/反序列化、datafusion-substraitSubstrait 计划序列化格式支持、datafusion-sqllogictestSQL 逻辑测试运行器。DataFusion 内部按多个子 crate 拆分以强制模块化并缩短编译时间主要子 crate 及职责如下子 crate职责datafusion_common公共 trait 与类型datafusion_catalogCatalog APISchemaProvider、CatalogProviderdatafusion_datasource文件与数据 I/OFileSource、DataSinkdatafusion_sessionSession及关联结构datafusion_execution执行所需的状态与结构datafusion_exprLogicalPlan、Expr及逻辑规划结构datafusion_functions标量函数包datafusion_functions_aggregate聚合函数MIN、MAX、SUM等datafusion_functions_nestedARRAY、MAP、STRUCT的标量函数包datafusion_functions_table表函数如GENERATE_SERIESdatafusion_functions_window窗口函数ROW_NUMBER、RANK等datafusion_optimizerOptimizerRule与AnalyzerRuledatafusion_physical_exprPhysicalExpr及关联表达式datafusion_physical_planExecutionPlan及关联表达式datafusion_physical_optimizer物理计划优化规则datafusion_sqlSQL 规划器SqlToRel工作区内各 crate 之间的完整依赖关系图见 docs/source/contributor-guide/architecture/dependency-graph.md黑色连线为普通依赖、蓝色为 dev-dependency、绿色为 build-dependency、虚线为可通过禁用 cargo feature 移除的可选依赖为保持可读性图中刻意忽略了传递依赖。总结如何参与 DataFusion 的架构演进回到架构文档的核心主张DataFusion 是一个快速迭代、高度可扩展的通用查询引擎。对集成方而言正确姿势是基于稳定公共 API 与扩展点定制而不是 Fork对希望新增通用能力的贡献者而言流程是issue → 讨论 API 形态 → 示例原型 → PR 合入且安全默认行为的 API 更易被接受。理解本文介绍的规划/优化/执行两阶段管线、TableProvider数据源抽象、pull 式流式执行与 Tokio 线程模型是深入阅读源码入口 datafusion/core/src/lib.rs与上手编写自定义算子、优化规则、UDF 和表提供者的前提。赞分享大数据数据分析后端【免费下载链接】datafusionApache DataFusion SQL Query Engine项目地址https://gitcode.com/gh_mirrors/datafu/datafusion点击查看免费下载相关推荐Apache DataFusion-Ballista分布式查询引擎架构解析Apache DataFusion Ballista分布式查询引擎架构解析 项目概述 Apache DataFusion Ballista简称Ballista5 分钟定位 LiteParse 解析异常缺字、乱码、整页空白怎么办5 分钟定位 LiteParse 解析异常缺字、乱码、整页空白怎么办 输出缺了半页内容、报错信息却不说原因别慌。LiteParse 是一款在本地高速运行的文大数据数据分析后端AutoTrader高级技巧自定义指标开发与策略信号优化AutoTrader高级技巧自定义指标开发与策略信号优化 AutoTrader是一个基于Python的自动化交易系统开发平台支持从回测到优化再到实盘交易的全上一篇Compound Engineering 插件实战33个技能串起会自我积累的AI开发流程下一篇Godot Tools场景预览功能深度解析可视化编辑.tscn文件的终极指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
