Apache Arrow C++ 示例解析:用 compute 比较列数据并写出 CSV 文件
【免费下载链接】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
本指南围绕 Apache Arrow 仓库中的官方 C++ 示例compute_and_write_csv_example.cc展开,完整演示了「构建数值列 → 两种方式比较大小 → 组装 Table → 写出 CSV」的端到端流程。读者将掌握NumericBuilder/BooleanBuilder构建数组、arrow::compute::CallFunction("greater", ...)调用比较内核、Table::Make组装多列数据,以及arrow::csv::WriteCSV落盘输出的全部技术细节,可直接迁移到自己的数据处理管线中。
示例文档与源码定位
该示例由文档 compute_and_write_example.rst 描述,并被收录在 C++ 示例文档目录 examples/index.rst 的 toctree 中。其完整可运行源码位于仓库 cpp/examples/arrow/compute_and_write_csv_example.cc,属于 Apache Arrow C++ 官方 examples 集合的一部分。
文档对示例的定位是:创建一张包含两个数值列的表,比较两列元素的数值大小,并把列数据及其比较结果写入 CSV 文件。示例的核心价值在于同时展示了两种比较实现路径——手写循环逐元素比较,以及调用 Arrow Compute 的比较内核——让读者直观对比「自行实现」与「使用引擎内核」的差异。
示例的构建条件与运行方式
在 cpp/examples/arrow/CMakeLists.txt 中,该示例仅在同时开启ARROW_COMPUTE与ARROW_CSV两个构建选项时才会被编译:
if(ARROW_COMPUTE AND ARROW_CSV) if(ARROW_BUILD_SHARED) set(COMPUTE_KERNELS_LINK_LIBS arrow_compute_shared) else() set(COMPUTE_KERNELS_LINK_LIBS arrow_compute_static) endif() add_arrow_example(compute_and_write_csv_example EXTRA_LINK_LIBS ${COMPUTE_KERNELS_LINK_LIBS}) endif()从源码结构看,add_arrow_example会将示例链接到 Arrow 核心库以及arrow_compute_shared(或静态版本arrow_compute_static)计算内核库。因此编译时需要确保:
ARROW_COMPUTE=ON(启用 Compute 模块)ARROW_CSV=ON(启用 CSV 读写模块)- 同时构建 Arrow C++ 库本体与 examples(构建 examples 的相关 CMake 选项,可在 cpp/CMakeLists.txt 中查看)
编译成功后可得到可执行文件compute_and_write_csv_example。源码头部注释明确说明了运行方式与输出位置:
./compute_and_write_csv_example程序运行后会在当前目录生成compute_and_write_output.csv文件,其中包含四个列:a、b、a>b? (self written)与a>b? (arrow)。
整体执行流程
示例的入口main只做一件事:调用RunMain并检查返回的arrow::Status,失败时打印错误信息并以非零码退出,成功则返回EXIT_SUCCESS。真正的工作集中在RunMain中,其流程可归纳为五个阶段:
- 初始化:调用
arrow::compute::Initialize()注册 Compute 内核; - 构建数组:用
NumericBuilder<Int64Type>构建两个 int64 列array_a、array_b; - 比较列:分别用显式循环和 compute 函数
greater得到两个布尔结果数组; - 组装 Table:通过
arrow::schema+Table::Make把 4 个数组合成一张表; - 写出 CSV:用
FileOutputStream打开文件,调用WriteCSV落盘。
下面逐阶段拆解。
第一阶段:构建 int64 数组
ARROW_RETURN_NOT_OK(arrow::compute::Initialize());ARROW_RETURN_NOT_OK宏在表达式返回非 OK 的Status时立即将该Status返回给调用者,这是 Arrow C++ 代码中标准的错误传播写法。arrow::compute::Initialize()负责注册内置的 compute 函数与内核(包括后文用到的greater),因此必须在使用任何CallFunction之前调用。
随后使用数值构建器创建两列数据:
arrow::NumericBuilder<arrow::Int64Type> int64_builder; arrow::BooleanBuilder boolean_builder; ARROW_RETURN_NOT_OK(int64_builder.Resize(8)); ARROW_RETURN_NOT_OK(boolean_builder.Resize(8)); std::vector<int64_t> int64_values = {1, 2, 3, 4, 5, 6, 7, 8}; ARROW_RETURN_NOT_OK(int64_builder.AppendValues(int64_values)); std::shared_ptr<arrow::Array> array_a; ARROW_RETURN_NOT_OK(int64_builder.Finish(&array_a)); int64_builder.Reset(); int64_values = {2, 5, 1, 3, 6, 2, 7, 4}; std::shared_ptr<arrow::Array> array_b; ARROW_RETURN_NOT_OK(int64_builder.AppendValues(int64_values)); ARROW_RETURN_NOT_OK(int64_builder.Finish(&array_b));关键点:
- Resize(8):预先为 8 个元素分配容量,避免追加过程中反复扩容,属于构建性能优化手段;
- AppendValues(向量):批量追加一组值。第一列
a为[1..8],第二列b为[2,5,1,3,6,2,7,4]; - Finish(&array):结束构建,产出
std::shared_ptr<arrow::Array>; - Reset():复用同一个 builder 构建第二列,避免重新实例化。
第二阶段:两种方式比较两列的大小
示例用两种途径计算a > b的布尔结果,并把它们都放进最终输出的表中,便于对照验证。
方式一:手写循环 + 数组 API
auto int64_array_a = std::static_pointer_cast<arrow::Int64Array>(array_a); auto int64_array_b = std::static_pointer_cast<arrow::Int64Array>(array_b); for (int64_t i = 0; i < 8; i++) { if ((!int64_array_a->IsNull(i)) && (!int64_array_b->IsNull(i))) { bool comparison_result = int64_array_a->Value(i) > int64_array_b->Value(i); boolean_builder.UnsafeAppend(comparison_result); } else { boolean_builder.UnsafeAppendNull(); } } std::shared_ptr<arrow::Array> array_a_gt_b_self; ARROW_RETURN_NOT_OK(boolean_builder.Finish(&array_a_gt_b_self));要点:
- 先把
std::shared_ptr<arrow::Array>向下转型为具体的arrow::Int64Array(static_pointer_cast),从而获得强类型的Value(i)访问器; IsNull(i)判空:只有两列对应位置都非空才比较;任一为空则追加UnsafeAppendNull(),保证输出列与输入列长度一致、空值语义正确;UnsafeAppend:由于前面已Resize(8),循环次数固定为 8,不会触发扩容,因此可以使用不检查容量的快速追加接口;- 该路径完全依赖数组 API,未使用 Compute 模块,代码量更大且需自行处理空值逻辑。
方式二:调用 compute 比较内核
ARROW_ASSIGN_OR_RAISE(arrow::Datum compared_datum, arrow::compute::CallFunction("greater", {array_a, array_b})); auto array_a_gt_b_compute = compared_datum.make_array();要点:
CallFunction("greater", {array_a, array_b})按函数名查找注册的比较内核并执行,返回arrow::Datum;ARROW_ASSIGN_OR_RAISE在失败时把Status作为错误返回,成功时把Datum赋给变量;greater是 Arrow Compute 内置的标量比较函数,其内核实现在 cpp/src/arrow/compute/kernels/scalar_compare.cc(MakeCompareFunction<Greater>("greater", greater_doc)),并配套有系统性的测试覆盖,见 scalar_compare_test.cc;- 从源码结构看,
CallFunction的字符串名与greater等比较操作在 expression_internal.h 等处也有对应映射,表达式 API 中同样使用同一套比较语义(参见 expression.cc 的call("greater", ...)); compared_datum.make_array()把Datum转换为std::shared_ptr<arrow::Array>,随后可直接放入 Table。
两种方式产出的布尔数组语义一致,示例正是通过最后 CSV 中两列比较结果逐行相同来验证 compute 内核与手写实现的等价性。
第三阶段:组装 Table
auto schema = arrow::schema({arrow::field("a", arrow::int64()), arrow::field("b", arrow::int64()), arrow::field("a>b? (self written)", arrow::boolean()), arrow::field("a>b? (arrow)", arrow::boolean())}); std::shared_ptr<arrow::Table> my_table = arrow::Table::Make( schema, {array_a, array_b, array_a_gt_b_self, array_a_gt_b_compute});arrow::schema({...})以字段列表构造Schema,字段由arrow::field(名称, 类型)描述:前两列为int64,后两列为boolean;Table::Make(schema, {array_a, ...})将 4 个等长数组按列组织成一张表,列顺序与 schema 中的字段顺序一一对应;- 列名特意区分
a>b? (self written)与a>b? (arrow),使 CSV 中能直观区分两种比较结果的来源。
第四阶段:写出 CSV 文件
auto csv_filename = "compute_and_write_output.csv"; ARROW_ASSIGN_OR_RAISE(auto outstream, arrow::io::FileOutputStream::Open(csv_filename)); ARROW_RETURN_NOT_OK(arrow::csv::WriteCSV( *my_table, arrow::csv::WriteOptions::Defaults(), outstream.get()));arrow::io::FileOutputStream::Open打开(不存在则创建)目标文件,返回std::shared_ptr<io::OutputStream>;arrow::csv::WriteCSV(*my_table, options, output)是 CSV 写入的高级入口,其声明位于 cpp/src/arrow/csv/writer.h。该头文件同时提供了针对RecordBatch与RecordBatchReader的重载,以及增量式写入的MakeCSVWriter工厂函数,适合流式场景;- writer.h 顶部注释明确了 CSV 序列化规则:非二进制类型不加引号、空值输出为空字符串;二进制类型非空数据加引号且内部引号以双引号转义,空值则保持空且不加引号;
WriteOptions::Defaults()返回默认写入配置,定义于 cpp/src/arrow/csv/options.cc,其结构体成员在 cpp/src/arrow/csv/options.h 中定义。
WriteOptions 关键字段
WriteOptions允许精细控制 CSV 输出格式,常用字段及默认值如下:
| 字段 | 默认值 | 含义 |
|---|---|---|
include_header | true | 是否写出首行列名头 |
batch_size | 1024 | 每次批量处理的最大行数,影响写入性能 |
delimiter | ',' | 字段分隔符 |
null_string | "" | 空值写出的字符串(不允许包含引号) |
eol | "\n" | 行结束符 |
quoting_style | QuotingStyle::Needed | 字段的引用(加引号)策略 |
quoting_header | QuotingStyle::Needed | 表头的引用策略(与Needed等效时会对所有列名加引号) |
示例使用Defaults(),因此输出文件包含以列名为首行的表头,字段以逗号分隔,行以\n结尾。需要自定义时,可修改上述字段后调用options.Validate()校验合法性(Validate()声明于 options.h 中同一结构体)。
预期输出结果
运行示例后,当前目录下的compute_and_write_output.csv内容如下(布尔值在 Arrow 的 CSV 序列化中输出为true/false,空值输出为空字符串):
a,b,a>b? (self written),a>b? (arrow) 1,2,false,false 2,5,false,false 3,1,true,true 4,3,true,true 5,6,false,false 6,2,true,true 7,7,false,false 8,4,true,true可以逐行核对:a>b? (self written)与a>b? (arrow)两列完全一致,证明 compute 内核greater与手写循环得到相同结果。该输出同时演示了 Arrow「列式数据组装 → 序列化落盘」的标准路径,行数与列数均与输入数组一一对应。
从示例到生产实践的延伸
该示例虽短,却串起了 Arrow C++ 中几个高频使用模式:
- 构建器模式:
NumericBuilder/BooleanBuilder与Resize/AppendValues/Finish/Reset的组合适用于批量构建列数据;数据量大时应先Resize以减少重分配开销; - Datum 统一抽象:
CallFunction的入参与返回值都是Datum,它可以包装 Array、Scalar 或 ChunkedArray,使同一个计算函数能作用于多种数据结构,这是将示例扩展到分块数据(ChunkedArray)或流式批次(RecordBatchReader+WriteCSV重载)的基础; - CSV 输出的两种形态:一次性写出整张表用
WriteCSV;若数据是持续产生的批次,可改用 writer.h 中提供的MakeCSVWriter增量写入接口; - 与 Dataset/文件系统生态衔接:本示例使用本地文件输出流,Arrow 的文件系统抽象同样支持将结果写到其他存储后端;仓库中 dataset_documentation_example.cc 等示例展示了更大规模数据集上的扫描与处理模式,可作为进一步学习的参照。
如需对照 Arrow C++ 的其他官方示例,可浏览 cpp/examples/arrow 目录,其中包含 Parquet 读写、Flight、Dataset 扫描、UDF 等多个主题的可运行样例。
【免费下载链接】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),仅供参考