ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Arrow Acero 执行引擎:基于 ExecPlan、ExecNode 与 ExecBatch 的流式数据处理核心概念详解

Apache Arrow Acero 执行引擎:基于 ExecPlan、ExecNode 与 ExecBatch 的流式数据处理核心概念详解 Apache Arrow Acero 执行引擎基于 ExecPlan、ExecNode 与 ExecBatch 的流式数据处理核心概念详解【免费下载链接】arrowApache Arrow is the universal columnar format and multi-language toolbox for fast data interchange and in-memory analytics项目地址: https://gitcode.com/GitHub_Trending/arrow3/arrowAcero 是 Apache Arrow C 实现中的流式执行引擎streaming execution engine它让你把大规模甚至无限数据的计算表达为一张由节点组成的执行计划图。本文以官方概览文档docs/source/cpp/acero/overview.rst为主线完整还原 Acero 的定位边界它不是什么、它与 Arrow Compute / Datasets / Substrait 的关系并深入源码剖析ExecNode、ExecBatch、ExecPlan、Declaration四大核心概念的落地形态帮助你建立前端产出 Substrait 计划、Acero 精确执行的完整心智模型。一、什么是 Acero以执行计划为单位的流式计算库Acero 是一个 C 库用于分析大型甚至潜在无限的数据流。它允许把计算表达为一个执行计划ExecPlan计划接收零个或多个输入数据流产出一个输出数据流并描述数据在流经各节点时的变换方式。典型的计划可以是用公共列合并join两个数据流对已有列求值表达式从而创建新列把流式数据写盘落成分区partitioned布局。Acero 的节点实现全部位于 cpp/src/arrow/acero 目录下从源码结构看核心算子文件与概念一一对应filter_node.cc、hash_join_node.cc、asof_join_node.cc、order_by_node.cc、fetch_node.cc、aggregate_internal.cc、groupby_aggregate_node.cc等头部统一由 cpp/src/arrow/acero/exec_plan.h 与 cpp/src/arrow/acero/options.h 组织。Acero 不是什么四条清晰的定位边界官方概览文档花了相当篇幅说明 Acero 的非定位这四条边界对选型至关重要1. 不是给数据科学家直接用的库。Acero 不预期被终端用户直接调用——通常用户使用的是某种前端如 Pandas、Ibis 或 SQL。Acero 的 API 聚焦于可用能力与算法本身。不过了解 Acero 内部机制有助于前端用户理解其库的后端处理过程。2. 不是数据库DBMS。数据库通常是更庞大的独立服务。Acero 可以成为数据库的组件几乎所有数据库都有某种执行引擎也可以成为与数据库毫无相似之处的数据处理应用中的组件。Acero 不管用户管理、外部通信、隔离、持久性或一致性而且 Acero 主要聚焦读路径其写工具不具备任何事务支持。3. 不是优化器。Acero 没有 SQL 解析器、没有查询规划器、没有任何优化器。它期待被给出如何操纵数据的非常细致、低层的指令然后严格按照描述执行。原文档明确指出创建最优执行计划非常困难小的细节会显著影响性能作者团队认为优化器很重要但应独立于 Acero 实现并希望通过 Substrait 这类标准以可组合的方式存在让任意后端都能受益。4. 不是分布式引擎。Acero 不提供分布式执行但目标是可被分布式查询执行引擎使用——Acero 不会配置和协调 worker但它预期作为 worker 被使用。两者边界有时模糊例如某个 Acero source 可能是一台能执行过滤等高级分析的智能存储设备可视为分布式计划的一部分。关键区分在于Acero 没有把逻辑计划转换成分布式执行计划的能力这一步必须在别处完成。二、Acero 与 Arrow 其他模块及生态的对比Arrow Compute流式 vs 全内存核心区别在于Acero 处理数据流streams而 Arrow Compute 处理全量内存中的数据。这一点在源码上体现得很直接——Compute 的函数签名接收的是数组/批次/表整体而 Acero 的ExecNode::InputReceived(ExecNode* input, ExecBatch batch)接口见 cpp/src/arrow/acero/exec_plan.h则是一批一批地喂入数据。Arrow Datasets文件格式复杂性的隔离层Arrow Datasets 库提供发现、扫描、写入文件集合的基础例程并且datasets 模块依赖 Acero扫描和写入 datasets 都使用 Acero其中 scan 节点与 write 节点属于 datasets 模块本身。这样做的好处是把文件格式与文件系统的复杂性隔离在 Acero 核心逻辑之外。这也解释了为什么scan节点的 options 类型是arrow::dataset::ScanNodeOptions而非 Acero 自身类型——datasets 模块把自身挂在 Acero 的节点工厂注册表上对应ARROW_REGISTER_EXEC_NODE_FACTORY机制。Substrait标准查询计划语言的消费者Substrait 是一个为查询计划制定标准的项目。Acero 执行查询计划并生成数据因此Acero 是 Substrait 的消费者consumer。仓库中的 Substrait 协议副本见 format/substrait/substrait.yamlAcero 的 Substrait 消费端示例见 cpp/examples/arrow/engine_substrait_consumption.cc专题文档见 docs/source/cpp/acero/substrait.rst。与 DataFusion / DuckDB / Velox 等列式引擎的关系列式数据引擎不断涌现原文档对此持开放态度鼓励 Substrait 类标准让使用者按需在不同引擎间切换并通常不鼓励对比基准测试——因为 benchmark 几乎必然由工作负载驱动很难做到同条件apples-to-apples比较。三、Acero 与 Arrow C 的分层关系Acero 是 Arrow C 实现的一部分它作为独立模块构建但依赖核心 Arrow 模块不能独立存在。官方文档用三层结构描述它与 C 库的关系见 docs/source/cpp/acero/overview.rst 中的分层图第一层核心 Arrow 库。提供按 Arrow 列式布局组织的 buffer、array 容器。除少数例外核心库不检查也不修改 buffer 的内容——例如把字符串数组从小写转大写不属于核心库因为那需要检查数组内容。第二层Compute 模块。在核心库之上提供分析与变换数据的功能能力全部通过函数注册表FunctionRegistry暴露。一个 Arrow 函数接收零个或多个数组、批次或表产出数组、批次或表函数调用还可以与字段引用、字面量组合成表达式一棵函数调用树由 Compute 模块求值。例如给定含列x、y的表计算x (y * 3)。第三层Acero。在前两层之上为数据流添加计算操作。例如一个 project 节点可以对一批一批的 batch 流应用 compute 表达式产出把表达式结果作为新列追加的新批次流这些节点可以组合成图形成更复杂的执行计划——这与函数组合成树形成复杂表达式的思路完全同构。一个值得注意的设计决策官方 note 明确指出Acero不使用核心库的arrow::Table或arrow::ChunkedArray容器。因为 Acero 操作的是批次流无需多批次容器这简化了 Acero 的复杂度也避免了表中各列 chunk 大小不一致这类棘手问题。Acero 中会大量使用arrow::Datum——一个可持有多种类型的变体在 Acero 内datum 总是且仅持有arrow::Array或arrow::Scalar之一。四、核心概念一ExecNode——节点是计划图的基本单元Acero 最基本的概念是ExecNode。一个 ExecNode 有零个或多个输入、零个或一个输出没有输入的节点叫source数据源没有输出的节点叫sink汇点。节点种类繁多各自以不同方式变换输入例如scan 节点从文件读取数据的 source 节点属于 datasets 模块aggregate 节点累积批次以计算汇总统计filter 节点按过滤表达式删除行table sink 节点把数据累积成一张表。从源码看cpp/src/arrow/acero/exec_plan.h 中的ExecNode是一个抽象基类自定义算子需实现两个纯虚函数InputReceived(ExecNode* input, ExecBatch batch)上游节点把批次交给本节点节点通常对批次做某种操作然后调用自己输出的InputReceived传递结果。需要累积若干输入才能产出输出的节点会把批次加入内存累积队列对应accumulation_queue.hInputFinished(ExecNode* input, int total_batches)标记某个输入流结束即使调用时未必已收齐全部批次——它固定了该输入的最终批次数让节点知道何时收齐所有输入。此外还有一个值得注意的生命周期钩子节点在ExecPlan创建之后、StartProducing之前有一个初始化钩子This hook performs any actions in between creation of ExecPlan and the call to StartProducing典型用途如Bloom filter 下推见 cpp/src/arrow/acero/bloom_filter.h。另一个设计细节是节点间传递的排序契约Ordering。源码中ordering()的注释cpp/src/arrow/acero/exec_plan.h说明了排序保证的精确语义它不保证批次按序发出而是保证批次的ExecBatch::index属性尊重该排序filter/project不改变排序order-by节点引入全新排序而 hash-join、聚合可能破坏排序。依赖排序的节点如fetch、asofjoin会因此有输入约束——这是低层指令严格照办哲学的具体体现。完整算子清单factory name、options 类型与说明见用户指南 docs/source/cpp/acero/user_guide.rst 的 Available ExecNode Implementations 表格包括 source 类source、table_source、record_batch_source、scan等、compute 类filter、project、aggregate、pivot_longer、arrangement 类hash_join、asofjoin、union、order_by、fetch以及 sink 类sink、write、consuming_sink、table_sink、order_by_sink。五、核心概念二ExecBatch——无 Schema、可含标量的二维批次数据批次由ExecBatch类表示。它是一个二维结构非常类似RecordBatch可以有零个或多个列且所有列长度必须相同。它有三个关键差异没有 schema。因为ExecBatch被视为批次流中的一员而流被认为具有一致的 schema——所以 schema 通常存放在 ExecNode 里对应ExecNode::output_schema()接口。列可以是Array或Scalar。当某列是Scalar时表示该列对批次中每一行取同一个值ExecBatch还有一个length属性描述批次的行数。因此从另一个角度看Scalar就是一个含length个元素的常量数组。携带执行计划所需的附加信息。例如index可以描述批次在有序流中的位置官方预期ExecBatch未来还会演化出 selection vector选择向量等字段。零拷贝转换的边界RecordBatch 转 ExecBatch 永远是零拷贝两者引用完全相同的底层数组反过来ExecBatch 转 RecordBatch仅当批次中没有标量时才是零拷贝——标量列需要物化成等长数组。轻量数组容器。Acero 与 Compute 模块各自有批次的轻量版Compute 模块里是BatchSpan、ArraySpan、BufferSpanAcero 里对应概念叫KeyColumnArray定义见 cpp/src/arrow/compute/light_array_internal.h。两者同期开发、目的一致提供一种可完全栈分配的数组容器前提是非嵌套类型避免堆分配开销。官方 note 直言这两个概念理想情况下终有一天会合并。一个值得留意的工程约束藏在ExecPlan的头几行cpp/src/arrow/acero/exec_plan.h// This allows operators to rely on signed 16-bit indices static const uint32_t kMaxBatchSize 1 15;即单个批次的上限为 32768 行以便算子内部可以安全使用有符号 16 位行索引。而 cpp/src/arrow/acero/options.h 中TableSourceNodeOptions的默认切片大小是kDefaultMaxBatchSize 1 201048576 行——超过该尺寸的输入会被切成kMaxBatchSize的小块再流入计划。这就是为什么用DeclarationToTable收集结果时输出的 chunk 划分往往与输入不同如输入是一个 200 万行的单 chunk输出可能是 64 个 32K 行的 chunk。六、核心概念三ExecPlan——一次执行的生命周期载体ExecPlan表示一张 ExecNode 对象组成的图。有效的 ExecPlan必须至少有一个 source 节点但技术上不一定要有 sink 节点。计划包含所有节点共享的资源并提供控制节点启停的工具函数。两个关键事实ExecPlan 和 ExecNode 都绑定单次执行的生命周期它们携带状态不预期可重启实验性警告Acero 内部结构包括ExecBatch仍是实验性的ExecBatch不应在 Acero 之外使用应转换为RecordBatch等标准结构同样ExecPlan 是内部概念——用户构建计划应使用 Declaration 对象消费/执行计划的 API 应抽象掉底层计划细节、不直接暴露计划对象。从源码接口看cpp/src/arrow/acero/exec_plan.h 给出了执行的生命周期方法ExecPlan::Make()创建空计划可传入QueryOptions与ExecContextAddNode()/EmplaceNode()添加节点Validate()校验StartProducing()按逆拓扑序启动所有节点保证任何节点在其输入之前启动StopProducing()触发所有 source 停止产出新数据已在执行的任务仍会跑完最后等待finished()返回的 Future 完成。与执行相关的另一个实用配置是未对齐 buffer 的处理策略cpp/src/arrow/acero/exec_plan.h可通过环境变量ACERO_ALIGNMENT_HANDLING设为warn默认、ignore、reallocate或error。七、核心概念四Declaration——可序列化、可转换的蓝图如果说 ExecPlan 是已实例化的一次执行那么Declaration 就是 ExecNode 的蓝图。声明可以组合成图形成 ExecPlan 的蓝图。Declaration 描述需要做什么计算但不负责实际执行——这让它类似于表达式expression。官方预期 Declaration 需要与各种查询表示例如 Substrait互相转换。Declaration 对象连同DeclarationToXyz系列方法就是 Acero 当前的公共 API。源码中Declaration结构定义于 cpp/src/arrow/acero/exec_plan.h只有四个成员factory_name节点工厂名须已注册进执行节点注册表inputs输入其他声明或已实例化的ExecNode*options控制节点行为的 options 对象label用于区分同类节点的可选标签。它还提供一个便捷工厂Declaration::Sequence({...})把一组线性声明依次串联免去手工嵌套输入的可读性灾难——源码注释里直接给出了不用 Sequence 时的丑陋嵌套写法作为对照。一个可复制的最小示例Scan → Project → Table以下代码节选自官方文档配套示例 cpp/examples/arrow/execution_plan_documentation_examples.cc演示了构建声明图 DeclarationToTable收集结果的完整闭环// 构建数据源datasets 模块的 Dataset ScanOptions auto options std::make_sharedarrow::dataset::ScanOptions(); // 表达式a 列乘以 2 cp::Expression a_times_2 cp::call(multiply, {cp::field_ref(a), cp::literal(2)}); auto scan_node_options arrow::dataset::ScanNodeOptions{dataset, options}; // 方式一显式声明图——scan 是 source无输入project 以 scan 为输入 ac::Declaration scan{scan, std::move(scan_node_options)}; ac::Declaration project{ project, {std::move(scan)}, ac::ProjectNodeOptions({a_times_2})}; // 方式二线性序列用 Sequence 简写project 无需再传输入 ac::Declaration plan ac::Declaration::Sequence({{scan, std::move(scan_node_options)}, {project, ac::ProjectNodeOptions({a_times_2})}});执行侧DeclarationToXyz方法族定义见 cpp/src/arrow/acero/exec_plan.h提供了不同的结果形态各有异步版本DeclarationToTableAsync等DeclarationToTable把所有结果累积成一张arrow::Table最简但需要全量驻留内存注意输出 chunk 划分由执行引擎决定可能与输入不同前述kMaxBatchSize机制DeclarationToReader返回arrow::RecordBatchReader让你迭代消费读得不够快会触发背压使计划暂停关闭 reader 会取消计划析构时会等待计划完成当前工作可能阻塞DeclarationToStatus只跑计划不消费结果适合基准测试或计划有副作用如 dataset write 节点的场景产出的任何结果会被立即丢弃。如果DeclarationToXyz都不满足例如自建了自定义 sink 节点、或需要多个输出也可以直接操作ExecPlan创建ExecPlan→ 为 sink 节点建声明并加入图 →Declaration::AddToPlan多输出场景不可用→ExecPlan::Validate()→ExecPlan::StartProducing()→ 等待finished()Future。需要说明的是Acero 的设计并不硬性禁止多个 sink 节点尽管学术文献和多数系统默认计划至多一个输出。八、对扩展者的启示如何理解 Acero 的架构哲学把以上概念串起来Acero 的架构可以概括为一条两层世界的链路Declaration可序列化蓝图层面向外部——可与 Substrait 等标准表示互转、可跨语言传递ExecPlan/ExecNode/ExecBatch运行时层面向内部——按批次流推进、带背压控制backpressure_handler.h、逆拓扑序启动、不可重启。这个切分正是文档反复强调用户应停留在 Declaration 侧、把 ExecPlan 当内部概念的原因。对想在 Acero 之上做研究或商业用途的团队源码给出的扩展点也很明确实现ExecNode的子类遵循InputReceived/InputFinished/初始化钩子的契约注册工厂名后即可像scan、write那样被Declaration引用。相关示例可直接参考 cpp/examples/arrow/acero_register_example.cc自定义节点注册与 cpp/examples/arrow/engine_substrait_consumption.ccSubstrait 计划消费端。最后重申适用前提与限制Acero 是实验性演进中的 C 模块依赖 Arrow 核心与 Compute 模块构建不提供 SQL 解析、计划优化、事务或分布式协调写路径无事务支持面向作为分布式引擎的 worker而非编排分布式引擎。理解这些边界才能把它放到正确的架构位置——前端产出 Substrait 计划Acero 忠实执行优化与分布式编排留在体系之外的组件中完成。【免费下载链接】arrowApache Arrow is the universal columnar format and multi-language toolbox for fast data interchange and in-memory analytics项目地址: https://gitcode.com/GitHub_Trending/arrow3/arrow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表