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/arrow
Acero 是 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.yaml,Acero 的 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.h);InputFinished(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、BufferSpan,Acero 里对应概念叫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 << 20(1048576 行)——超过该尺寸的输入会被切成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与ExecContext),AddNode()/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_shared<arrow::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.cc(Substrait 计划消费端)。
最后重申适用前提与限制: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),仅供参考