ARTICLE DETAIL

资讯详情

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

Apache Arrow C++ Acero 流式执行引擎 API 完全指南:构建与运行执行计划

Apache Arrow C++ Acero 流式执行引擎 API 完全指南:构建与运行执行计划 Apache Arrow C Acero 流式执行引擎 API 完全指南构建与运行执行计划【免费下载链接】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/arrowAceroStreaming Execution是 Apache Arrow C 实现的流式查询引擎它将大数据集按小批次batch流式处理避免一次性加载全部数据从而支持任意规模输入的复杂计算。本文以 docs/source/cpp/api/acero.rst 这一 API 参考页为骨架系统讲解 Acero 的三大 API 分组——acero-api创建与运行执行计划、acero-nodes各执行节点的配置选项、acero-internals自定义节点所需的内部契约并结合仓库源码与官方用户指南帮助你掌握用Declaration声明计划、用DeclarationToXyz执行计划、为各类算子配置选项以及扩展自定义ExecNode的完整能力。注意Acero 目前仍处于实验阶段API 尚未稳定可能不经弃用周期直接变更官方用户指南在 docs/source/cpp/acero.rst 中对此有明确警告。一、Acero API 参考页的三层结构docs/source/cpp/api/acero.rst本身是一份精简的 Doxygen 参考页通过三个doxygengroup指令把散落在源码头部文件中的文档化注释聚合到同一页面参考页小节Doxygen 分组内容来源源码覆盖范围Creating and running execution plansacero-apicpp/src/arrow/acero/api.h含 exec_plan.hDeclaration、DeclarationToXyz方法族、QueryOptions等Configuration for execution nodesacero-nodescpp/src/arrow/acero/options.h全部*NodeOptions类Internals for creating custom nodesacero-internalscpp/src/arrow/acero/exec_plan.hExecPlan、ExecNode契约、ExecFactoryRegistry三个分组的实体入口都在 cpp/src/arrow/acero/api.h它声明了acero-api与acero-nodes两个defgroup并包含exec_plan.h与options.h两个头文件其中exec_plan.h内的ExecPlan、ExecNode等类又通过\addtogroup acero-internals归入内部组。理解这条“RST 指令 → Doxygen 组 → 头文件类”的映射链是阅读整个 API 参考页的前提。二、核心工作流Declaration → ExecPlan → 执行官方用户指南docs/source/cpp/acero/user_guide.rst给出了 Acero 的基本工作流创建一棵由Declaration对象组成的图graph描述整个执行计划调用某个DeclarationToXyz方法执行该Declaration从声明图创建新的ExecPlan每个Declaration对应计划中的一个ExecNode并根据所选方法自动追加一个 sink 节点执行ExecPlanDeclarationToReader例外它会在计划执行完成前就返回 reader计划执行完毕后被销毁。2.1 Declaration尚未构建的执行节点Declaration定义于 cpp/src/arrow/acero/exec_plan.h表示一个尚未实例化的ExecNode由于它的输入也可以是Declaration因此一棵Declaration树天然就描述了一张完整的执行图。其核心字段为factory_name要构造的节点工厂名称必须已注册到 exec node registry如filter、project、aggregateinputs输入声明列表元素类型为ExecNode*与Declaration的 variantoptions控制节点行为的选项对象必须使用与工厂匹配的ExecNodeOptions子类例如工厂名是project选项就应是ProjectNodeOptionslabel给节点的可选标签用于在计划中区分同类型节点。Declaration提供了多种便捷构造方式最简单的场景单输入、无输入可以直接用Declaration{factory_name, options}构造对于线性节点序列推荐使用静态方法Declaration::Sequence它会把前一个声明依次追加为后一个声明的输入。相比手写多层嵌套声明Sequence的写法更简洁可读// 不使用 Sequence 时的显式嵌套写法 Declaration{n3, {Declaration{n2, {Declaration{n1, {Declaration{n0, N0Opts{}}}, N1Opts{}}}, N2Opts{}}}, N3Opts{}}; // 使用 Sequence 的等价写法 Declaration::Sequence({ {n0, N0Opts{}}, {n1, N1Opts{}}, {n2, N2Opts{}}, {n3, N3Opts{}}, });若需要对执行过程做更精细的控制可调用Declaration::AddToPlan(plan, registry)把声明递归地加入一个已存在的ExecPlan会递归构造所有输入Declaration::IsValid()可用于校验声明而DeclarationToSchema、DeclarationToString则可在不实际执行的情况下分别计算计划的输出 schema 与调试用的字符串表示字符串表示仅用于调试完整序列化请使用 Substrait见 format/substrait。2.2 DeclarationToXyz一组执行入口acero-api组提供了一套以DeclarationTo开头的方法族全部定义于 cpp/src/arrow/acero/exec_plan.h。它们各自把结果以不同形态返回可根据使用场景选用方法返回形态适用场景DeclarationToTablestd::shared_ptrarrow::Table最简单把全部结果累积进内存中的 Table代价是必须全量驻留内存DeclarationToExecBatchesBatchesWithCommonSchemaExecBatch向量 公共 schema需要以 ExecBatch 粒度处理结果DeclarationToBatchesstd::vectorstd::shared_ptrRecordBatch以 RecordBatch 向量收集结果DeclarationToReaderstd::unique_ptrRecordBatchReader流式迭代消费读取不及时会触发背压暂停关闭 reader 会取消执行DeclarationToStatusStatus只运行计划、丢弃结果适合计划本身有副作用如写入节点或基准测试场景每个方法都有对应的...Async异步版本返回Future与接受QueryOptions的重载。以DeclarationToTable为例其内部逻辑cpp/src/arrow/acero/exec_plan.cc是在声明末尾追加{table_sink, sink_options}节点后整体执行DeclarationToReader则追加{sink, SinkNodeOptions(...)}cpp/src/arrow/acero/exec_plan.ccsink 节点把批次压入 FIFO 队列并由 reader 消费队列积压超过阈值时即触发背压。线程与背压行为所有同步方法均接受use_threads参数默认true。若为false所有 CPU 计算都在调用线程上完成I/O 任务仍由 I/O executor 承担。对于DeclarationToReaderuse_threadstrue时 CPU 工作在线程池上并行执行若 reader 消费不及时背压队列填满后计划会暂停队列相关默认值kDefaultBackgroundMaxQ 32、kDefaultBackgroundQRestart 16定义于 cpp/src/arrow/acero/exec_plan.huse_threadsfalse时 CPU 工作全部发生在RecordBatchReader::Next调用期间。提前关闭返回的 reader 可取消后续批次的计算但此时只会报告计算中遇到的错误不会报告取消错误。2.3 QueryOptions计划级配置QueryOptionscpp/src/arrow/acero/exec_plan.h用于在创建ExecPlan或调用DeclarationToXyz时指定计划级的全局行为核心字段如下字段类型默认值说明use_legacy_batchingboolfalse是否使用传统批处理策略。仅为支持Scanner::ToTable保留该方法依赖 scanner 的 batch 索引保持一致这在会切分 batch 的 ExecPlan 中不现实sequence_outputstd::optionalboolstd::nullopt输出是否按有意义顺序输出。默认在有排序时按序输出显式true时若无序将报错可用于校验查询显式false时立即输出可能乱序但降低延迟use_threadsbooltrue是否使用多后台线程做 CPU 密集计算为false时 CPU 工作全部在调用线程完成。设置了custom_cpu_executor时被忽略custom_cpu_executorExecutor*nullptr自定义 CPU 执行器须在计划生命周期内有效custom_io_executorExecutor*nullptr自定义 I/O 执行器为nullptr时使用受ARROW_IO_THREADS环境变量控制的全局 I/O 线程池memory_poolMemoryPool*default_memory_pool()计划分配内存所用的内存池function_registryFunctionRegistry*GetFunctionRegistry()计划使用的函数注册表field_namesvectorstring空输出列名为空时根据输入列自动生成非空时数量必须等于输出列数unaligned_buffer_handlingoptionalUnalignedBufferHandlingkWarn源数据缓冲区未对齐时的处理策略其中unaligned_buffer_handling关联一个独立的枚举UnalignedBufferHandling { kWarn, kIgnore, kReallocate, kError }cpp/src/arrow/acero/exec_plan.h并可通过环境变量ACERO_ALIGNMENT_HANDLING取值warn/ignore/reallocate/error覆盖默认行为。背景是Acero 内部和 compute 函数会对缓冲区做类型双关例如把uint8_t*强转为int32_t*做加法若源缓冲区未按类型对齐如 int32 数组不在 4 字节边界严格来说是 C 未定义行为。各策略含义为kWarn检查并发出警告大多数编译器与 CPU 对这种轻微未对齐是宽容的通常只是小幅性能损失kIgnore完全不检查kReallocate重新分配对齐的缓冲区并拷贝内容所有 Acero 内部分配的缓冲区保证已对齐此选项只影响源数据kError直接优雅中止计划。三、Execution Nodes 配置acero-nodes 分组全解acero-nodes组cpp/src/arrow/acero/options.h定义了所有执行节点的选项类全部继承自基类ExecNodeOptions仅当节点无配置时直接使用它。每个选项类与一个工厂名对应工厂名在节点源码中通过registry-AddFactory(...)注册实际名称与节点实现的对应关系可从源码确认类别工厂名选项类源码注册位置SourcesourceSourceNodeOptionssource_node.ccSourcetable_sourceTableSourceNodeOptionssource_node.ccSourcerecord_batch_sourceRecordBatchSourceNodeOptionssource_node.ccSourcerecord_batch_reader_sourceRecordBatchReaderSourceNodeOptionssource_node.ccSourceexec_batch_sourceExecBatchSourceNodeOptionssource_node.ccSourcearray_vector_sourceArrayVectorSourceNodeOptionssource_node.ccSourcenamed_tableNamedTableNodeOptionssource_node.ccComputefilterFilterNodeOptionsfilter_node.ccComputeprojectProjectNodeOptionsproject_node.ccComputeaggregateAggregateNodeOptionsaggregate_internal.ccComputepivot_longerPivotLongerNodeOptionspivot_longer_node.ccArrangementhashjoinHashJoinNodeOptionshash_join_node.ccArrangementasofjoinAsofJoinNodeOptionsasof_join_node.ccArrangementunion无N/Aunion_node.ccArrangementorder_byOrderByNodeOptionsorder_by_node.ccArrangementfetchFetchNodeOptionsfetch_node.ccSinksinkSinkNodeOptionssink_node.ccSinktable_sinkTableSinkNodeOptionssink_node.ccSinkconsuming_sinkConsumingSinkNodeOptionssink_node.ccSinkorder_by_sinkOrderBySinkNodeOptions已废弃sink_node.ccSinkselect_k_sinkSelectKSinkNodeOptions已废弃sink_node.cc官方用户指南在 docs/source/cpp/acero/user_guide.rst 的 “Available ExecNode Implementations” 一节以同样分类方式列出了这些算子下文按类别详解各选项类的配置语义。3.1 Source 节点把数据送入计划source/SourceNodeOptions最通用也最灵活的数据入口同时“相当难以配置”。它要求提供一个异步生成器函数无参数、返回arrow::Futurestd::optionalarrow::ExecBatch返回std::nullopt表示流结束这类函数即 Arrow 所称的AsyncGenerator。另外必须预先提供输出 schema——Acero 要求在开始处理前就知道执行图每一阶段的 schema。三个构造字段为output_schema生成的 batch 的 schemagenerator异步批次流以std::nullopt结尾ordering数据顺序默认Ordering::Unordered()。SourceNode在StartProducing期间开始调用生成器收到每个批次都会创建新任务向下游推送并把大于ExecPlan::kMaxBatchSize1 15 32768 行见 cpp/src/arrow/acero/exec_plan.h的大批次切分成小块再调用InputReceived。默认情况下它为输出批次赋予隐式顺序implicit ordering前提是生成器以确定性的方式产出数据。table_source/TableSourceNodeOptions数据已在内存中时的首选比手写异步生成器更简单高效。字段为tablestd::shared_ptrarrow::Table与max_batch_size默认kDefaultMaxBatchSize 1 20 1048576 行。节点按max_batch_size把大表切块以便并行处理但注意若源表本身的批次更小不会合并成更大的批次。基于迭代器的源SchemaSourceNodeOptionsItMaker是一个模板类接受“迭代器制造器”并可通过requires_io/io_executor把阻塞式迭代放到 I/O 线程上轮询避免占用 CPU 线程池。其具体实例包括RecordBatchSourceNodeOptions从RecordBatch迭代器取数ExecBatchSourceNodeOptions从ExecBatch迭代器取数ArrayVectorSourceNodeOptions从ArrayVector迭代器取数RecordBatchReaderSourceNodeOptions从RecordBatchReader读取每次迭代在 I/O 线程池新建任务执行。named_table/NamedTableNodeOptions定义一个惰性解析的 Arrow 表由名称唯一标识通常可在计划被消费时才解析该节点仅用于序列化目的永远不能实际执行字段为names序列化计划中的名称与schema表输出 schema。3.2 Compute 节点过滤、投影、聚合与透视filter/FilterNodeOptions按表达式过滤行filter_expression返回类型必须是布尔型只保留求值为true的行。若全部行被过滤节点会输出空批次以避免排序出现空洞。project/ProjectNodeOptions对每个输入批次执行一组表达式产出同长度的新列。字段expressions要执行的表达式列表输出每个表达式一列要保留输入列需为该列构造简单的field_ref表达式names输出列名为空时使用表达式ToString()的结果非空时长度必须与expressions相同。aggregate/AggregateNodeOptions聚合节点计算汇总统计量。构造参数为aggregatesstd::vectorcompute::Aggregate、keys分组键、segment_keys分段键。关键语义默认是 pipeline breaker必须累积全部输入后才产生输出keys非空时每个 aggregate 应为 HashAggregate 函数keys为空时假定为 ScalarAggregate 函数segment_keys是性能优化若输入已按一列或多列分区可指定为分段键每遇到分段键变化即输出截至当前已见数据的部分聚合结果从而支持流式输出如有序聚合分段键目前仅限单线程模式keys与segment_keys必须互斥disjoint未提供任何 measure 时直接输出唯一键列表输出列顺序为分段键 → 常规键 → 每个聚合一列。源码层面工厂aggregate会依据keys是否为空在ScalarAggregateNode与GroupByNode之间分派aggregate_internal.cc。pivot_longer/PivotLongerNodeOptions把若干列“翻转”为更多行即 UNPIVOT 操作常用于把每行含多个观测的宽表转成每行一个观测的长表。它通过row_templates行模板由特征值feature_values与度量列引用measurement_values组成度量缺失时用nullopt表示插入 null、feature_field_names特征列名、measurement_field_names度量列名三个字段描述变换。官方文档给出了把(time, left_temp, right_temp)变换为(time, location, temp)的示例输出行数为输入行数 × 行模板数。3.3 Arrangement 节点连接、排序与切片hashjoin/HashJoinNodeOptions基于哈希表的连接操作。JoinType枚举取值LEFT_SEMI、RIGHT_SEMI、LEFT_ANTI、RIGHT_ANTI、INNER、LEFT_OUTER、RIGHT_OUTER、FULL_OUTERoptions.h。提供多组构造函数核心字段join_type连接类型默认INNERleft_keys/right_keys左右输入的连接键FieldRef向量两者长度与类型必须一致output_all为true时输出左右输入的全部有效字段此时left_output/right_output被忽略left_output/right_outputoutput_allfalse时指定要输出的字段key_cmp键比较方式JoinKeyCmp { EQ, IS }决定 null 键之间是否相等默认全部为EQoutput_suffix_for_left/output_suffix_for_right左右输出字段名的后缀用于消除同名字段冲突默认均为空字符串filter残差过滤表达式作用于拼接后的输入 schema左字段在前、右字段在后可引用未参与输出的字段默认literal(true)disable_bloom_filter是否禁用连接中的 Bloom filter默认false即启用。asofjoin/AsofJoinNodeOptions实验性 APIasof 连接接受一个左表与任意多个右表按 “on” 键把各输入连接起来为左表的每一行输出一行。要求各输入批次必须已按 “on” 键排序。Keys结构体包含on_key单个字段各表须同类型同单位当前限整数、date 或 timestamp 类型使用非精确匹配与by_key字段列表使用精确相等匹配当前限整数、date、timestamp 或 base-binary 类型。tolerance决定匹配窗口行匹配当且仅当right.on - left.on ∈ [min(0, tolerance), max(0, tolerance)]负数表示 past-asof-join正数表示 future-asof-join零表示精确匹配单位与 “on” 键一致。构造时input_keys至少给出两个键第一个对应左表其余对应右表。union把 schema 完全相同的两个输入合并无配置选项。order_by/OrderByNodeOptions对数据施加新的排序。当前实现方式是累积全部数据、排序、再以更新后的 batch 索引发出暂不支持超内存larger-than-memory排序。构造参数为Ordering。fetch/FetchNodeOptions从流中切片一段行范围字段为offset跳过的行数与count保留的行数。kName常量为fetch。3.4 Sink 节点终止计划并消费结果sink/SinkNodeOptions把批次收集进 FIFO 队列generator指向一个std::functionFuturestd::optionalExecBatch()节点加入计划时被设置用于消费计划数据若消费不及时将累积并可能触发背压。schema是指向std::shared_ptrSchema的“出参”指针节点加入计划后、StartProducing之前被写入输出 schema。背压通过BackpressureOptions配置默认无背压resume_if_below0, pause_if_above0should_apply_backpressure()返回falseBackpressureOptions::DefaultBackpressure()使用默认阈值kDefaultBackpressureHighBytes 1 301 GiB达到即暂停生产与kDefaultBackpressureLowBytes 1 28256 MiB低于即恢复生产两常量定义于 options.h当pause_if_above 0时启用背压BackpressureMonitor接口bytes_in_use()/is_paused()可查询当前队列字节数与暂停状态。table_sink/TableSinkNodeOptions把全部输出累积进一个arrow::Table。output_table是“出参”指针非空且须在计划执行期间有效计划完成后被设置为指向结果表names可自定义输出列名须覆盖全部字段目前仅支持扁平 schema见 GH-31875。consuming_sink/ConsumingSinkNodeOptions在计划内部通过回调消费数据接受一个std::shared_ptrSinkNodeConsumer。SinkNodeConsumer接口包含三个阶段Init(schema, backpressure_control, plan)schema 定稿后、任何 Consume 之前调用一次常用于保存 schema、Consume(ExecBatch)消费一个批次、Finish()最后一个批次送达后调用返回的 future 应在所有未完成任务完成时结束计划提前终止或出错时不会被调用。配合BackpressureControlPause()/Resume()必须保证每次Pause()最终跟一次Resume()否则死锁可在消费端控制背压。已废弃的 sinkorder_by_sinkOrderBySinkNodeOptions对SortOptions排序后送入生成器与select_k_sinkSelectKSinkNodeOptions用SelectKOptions选出 top_k/bottom_k 行在官方节点表中被标记为 Deprecated见 docs/source/cpp/acero/user_guide.rst 的 Sink Nodes 表。四、InternalsExecPlan、ExecNode 契约与自定义节点4.1 ExecPlan计划的容器与生命周期ExecPlancpp/src/arrow/acero/exec_plan.h是执行计划的容器通过ExecPlan::Make(...)创建可传入QueryOptions、ExecContext与可选的KeyValueMetadata默认执行上下文为threaded_exec_context()提供以下关键操作nodes()获取计划中所有节点AddNode/EmplaceNode加入节点后者直接就地构造并返回指针Validate()校验计划StartProducing()按逆拓扑序启动所有节点保证任何节点都在其所有输入之前启动StopProducing()触发所有源停止生产新数据但计划仍会运行进行中的任务调用者仍需等待finished()完成后再销毁计划finished()所有任务完成后被标记完成的FutureHasMetadata()/metadata()计划元数据ToString()计划的字符串表示。当DeclarationToXyz方法无法满足需求例如自定义 sink 节点、或多个输出的计划时可直接手工运行计划。官方指南给出的步骤为创建ExecPlan→ 给Declaration图挂上 sink 声明 → 用Declaration::AddToPlan加入计划多输出时需逐个添加节点→Validate()→StartProducing()→ 等待finished()完成。注意学术文献与多数系统中默认执行计划至多一个输出DeclarationToXyz也依赖这一点但 Acero 的设计并不严格禁止多个 sink 节点。4.2 ExecNode节点的契约ExecNodecpp/src/arrow/acero/exec_plan.h是自定义节点的核心抽象Acero 通过实现它来扩展。其公开契约包括拓扑信息kind_name()节点种类名、num_inputs()/inputs()输入节点列表、input_labels()每个输入的功能标签、output()下游节点、output_schema()输出批次的数据类型、is_sink()无输出 schema 即为 sink、plan()、label()/SetLabel()。排序契约ordering()声明输出批次的顺序这是 Acero 中较微妙的概念。Ordering可以是Unordered顺序不确定如 hash-join 无可预测输出序或Implicit存在有意义的顺序但不体现在任何数据列中典型如从内存表读取数据时的隐式“行序”。filter/project 不改变排序只需保证输出批次的ExecBatch::index与对应输入批次一致order_by 产生全新排序hash-join 或聚合可能破坏排序。fetch 节点要求输入至少是隐式有序的asof join 则要求显式且与 on 键兼容的排序。维护排序的节点应注意避免 batch 索引出现空洞必要时发出空批次保持连续性。上游 API数据流入InputReceived(input, batch)接收输入批次节点处理后通常再调用下游的InputReceived转发结果需要累积的节点则把批次放入内存累积队列InputFinished(input, total_batches)标记某输入的批次总数可提前调用节点据此无论顺序如何都能判断输入是否收齐。这些方法可在StartProducing()成功后随时被并发调用并允许回调PauseProducing()/ResumeProducing()/StopProducing()。生命周期 APIInit()在ExecPlan创建与StartProducing之间执行任意初始化例如 Bloom filter 下推执行顺序未定义但同步进行、StartProducing()只能调用一次通常由ExecPlan::StartProducing()自动调用不应递归进入输入、PauseProducing(output, counter)/ResumeProducing(output, counter)背压提示可任意次数调用通过 counter 序号解决并发 pause/resume 的乱序问题——只有最高 counter 的调用有效、StopProducing()必须幂等且必须转发给输入因为 limit/top-k 等场景可能只停止部分计划。由于同步调用双向发生输入→输出、输出→输入节点实现必须小心重入与多线程并发最稳妥的做法是先更新内部状态、最后再通知输出。4.3 ExecFactoryRegistry注册与构造节点ExecFactoryRegistrycpp/src/arrow/acero/exec_plan.h是节点工厂的可扩展注册表工厂签名为std::functionResultExecNode*(ExecPlan*, std::vectorExecNode*, const ExecNodeOptions)。通过AddFactory(name, factory)注册重名会报错、GetFactory(name)获取未找到会报错。default_exec_factory_registry()返回内置工厂的默认注册表MakeExecNode(factory_name, plan, inputs, options, registry)则基于命名工厂构造节点。第三节表格中列出的所有AddFactory调用即是内置节点的注册证据自定义节点同样可以把自己的工厂注册进来从而用Declaration以factory_name引用。4.4 辅助工具Generator 与 reader 互转acero-api还包含若干生成器工具cpp/src/arrow/acero/exec_plan.hMakeGeneratorReader把一个ExecBatch生成器包装成RecordBatchReader不对发出批次施加顺序MakeReaderGenerator把RecordBatchReader转成可用于源节点的生成器接受io_executor与max_q/q_restart背压队列参数这对把现有 reader 接入执行计划非常实用。五、从声明到执行一个完整的端到端示例综合官方用户指南docs/source/cpp/acero/user_guide.rst中的示例模式与上文各节点语义一个典型的 Acero 程序化计划如下数据已在内存表table中#include arrow/acero/api.h #include arrow/compute/api_scalar.h #include arrow/compute/expression.h using namespace arrow::acero; using arrow::compute::call; using arrow::compute::field_ref; using arrow::compute::literal; // 1) 用 table_source 提供输入filter 过滤project 计算新列 Declaration source{table_source, TableSourceNodeOptions{table}}; Declaration filtered{filter, FilterNodeOptions{ call(greater, {field_ref(b), literal(3)})}}; Declaration projected{project, ProjectNodeOptions{ {field_ref(a), call(add, {field_ref(b), literal(1)})}, {a, b_plus_1}}}; // 2) 用 Sequence 把各步串成线性计划并执行结果累积为 Table auto result DeclarationToTable( Declaration::Sequence({source, filtered, projected}));要点回顾filter的表达式必须返回布尔类型project若不显式给出names将用表达式的ToString()作为列名DeclarationToTable会自动追加table_sink并阻塞直到计划完成。若改为DeclarationToReader则返回的 reader 支持流式消费与背压暂停若计划以写文件等副作用结尾且无需消费结果可用DeclarationToStatus。六、结语与进一步阅读Acero 的 API 参考页docs/source/cpp/api/acero.rst虽简短但其背后是三个组织严密的 Doxygen 分组acero-api负责计划声明与执行入口、acero-nodes提供每个算子的配置选项、acero-internals定义计划与节点的内部契约。掌握这三层结构你既能通过DeclarationDeclarationToXyz快速编写流式查询也能通过自定义ExecNode与ExecFactoryRegistry扩展 Acero 的能力边界。想深入了解可继续阅读仓库中的以下资源用户指南docs/source/cpp/acero.rst 与 docs/source/cpp/acero/user_guide.rst含全部节点示例概念总览docs/source/cpp/acero/overview.rstSubstrait 计划docs/source/cpp/acero/substrait.rstAcero 官方推荐的跨引擎计划描述方式可执行示例代码cpp/examples/arrow/execution_plan_documentation_examples.ccsource、filter、project、aggregate、hashjoin等节点示例的 literalinclude 来源节点源码与测试如 cpp/src/arrow/acero/filter_node.cc、cpp/src/arrow/acero/hash_join_node.cc、cpp/src/arrow/acero/aggregate_internal.cc 及各*_test.cc测试文件。【免费下载链接】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),仅供参考
返回列表