Apache Arrow C++ 行式数据与列式数据互转实战:从固定 Schema 到动态 Schema 的完整指南
2026/9/14 15:49:18 网站建设 项目流程

Apache Arrow C++ 行式数据与列式数据互转实战:从固定 Schema 到动态 Schema 的完整指南

【免费下载链接】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 以列式内存布局为核心优势,但现实世界中我们常常从数据库、日志、消息队列等系统拿到的是行式结构的数据(如结构体数组、JSON 文档)。本文基于 Apache Arrow 官方 C++ 文档 Row to columnar conversion,系统讲解两类行列转换场景:固定 Schema(编译期已知)动态 Schema(运行期才确定)。读完本文,你将掌握arrow::ArrayBuilder家族的手工组装技巧、arrow::RecordBatchBuilder的批量构建方式、基于arrow::VisitTypeInline/arrow::VisitArrayInline的访问者模式,以及arrow::TableBatchReader的分批零拷贝读取方法,并能在自己的 C++ 项目中直接复现完整可运行的转换代码。


一、为什么需要行列转换:场景与动机

Arrow 的列式内存格式(arrow::Tablearrow::RecordBatcharrow::Array)非常适合向量化计算与压缩,但外部系统(数据库驱动、REST API、日志采集器)交付数据的形态通常是逐行的。正如示例源码 row_wise_conversion_example.cc 的注释所述:

虽然我们想使用列式数据结构来构建高效运算,但我们经常从其他系统以行式方式接收数据。

因此,掌握"行 → 列"(写入 Arrow 数据结构)与"列 → 行"(从 Arrow 数据结构读回业务对象)的转换能力,是几乎所有 Arrow 接入方(连接器、ETL、数据导入导出工具)的基础功。本文按原文档结构分为两大层次:

层次Schema 状态典型代表核心手段
固定 Schema编译期已知结构体数组 ↔arrow::Table手写各类型 Builder
动态 Schema运行期才确定JSON 文档 ↔arrow::RecordBatchRecordBatchBuilder+ 访问者模式

二、固定 Schema:结构体数组与 arrow::Table 互转

原文档的第一部分演示了"将一个结构体数组转换为arrow::Table实例,再转回原始结构体数组",完整代码见 row_wise_conversion_example.cc。

2.1 数据模型:产品组件成本表

示例使用一个包含产品 ID、组件数量、每个组件成本的结构体作为行式载体:

struct data_row { int64_t id; int64_t components; std::vector<double> component_cost; };

对应到列式世界,目标arrow::Table由三列组成:

  • idint64)—— 产品 ID;
  • componentsint64)—— 组件数量;
  • component_costlist(float64))—— 每个组件的成本,是一个嵌套的列表列

2.2 行 → 列:手工组装 ArrayBuilder

Arrow 为每个数据类型提供了专门的 Builder 类,用于增量构造arrow::Array。对上述三列,需要三个 Builder,其中列表列比较特殊——它需要两层 Builder:顶层arrow::ListBuilder负责构建偏移量数组,嵌套的arrow::DoubleBuilder构建被偏移引用的底层值数组:

arrow::MemoryPool* pool = arrow::default_memory_pool(); Int64Builder id_builder(pool); Int64Builder components_builder(pool); ListBuilder component_cost_builder(pool, std::make_shared<DoubleBuilder>(pool)); // 该 builder 由 component_cost_builder 拥有 DoubleBuilder* component_item_cost_builder = (static_cast<DoubleBuilder*>(component_cost_builder.value_builder()));

示例源码 row_wise_conversion_example.cc 还给出了一个内存池使用建议:使用arrow::default_memory_pool()可以让底层内存区域原地扩容,效率更高(源码注释提到 jemalloc 池目前仅支持 Unix,不支持 Windows)。

随后逐行遍历数据并追加:

for (const data_row& row : rows) { ARROW_RETURN_NOT_OK(id_builder.Append(row.id)); ARROW_RETURN_NOT_OK(components_builder.Append(row.components)); // 标记一个新列表行的开始,这会记录值 builder 中的当前偏移 ARROW_RETURN_NOT_OK(component_cost_builder.Append()); // 存入实际值 ARROW_RETURN_NOT_OK(component_item_cost_builder->AppendValues( row.component_cost.data(), row.component_cost.size())); }

关键点:

  • 所有Append都可能失败(例如内存不足),必须检查返回值,这与 Arrow 全库的arrow::Status错误处理哲学一致;
  • 对列表列,先调component_cost_builder.Append()记下当前偏移,再向值 builder 追加元素,二者顺序不能颠倒;
  • 调用component_cost_builder.Finish()时会隐式完成嵌套值 builder,因此无需再单独Finish值 builder(见源码 row_wise_conversion_example.cc)。

最后声明 Schema 并把各数组组装成arrow::Table

std::vector<std::shared_ptr<arrow::Field>> schema_vector = { arrow::field("id", arrow::int64()), arrow::field("components", arrow::int64()), arrow::field("component_cost", arrow::list(arrow::float64()))}; auto schema = std::make_shared<arrow::Schema>(schema_vector); std::shared_ptr<arrow::Table> table = arrow::Table::Make(schema, {id_array, components_array, component_cost_array});

组装完成后,arrow::Table持有所有引用数据的所有权,离开函数作用域也不必担心悬空引用。

2.3 列 → 行:Schema 校验与零拷贝切片细节

反向转换(ColumnarTableToVector)的第一步是校验 Schema 是否匹配预期,避免把结构不符的表强行解包:

if (!expected_schema->Equals(*table->schema())) { return arrow::Status::Invalid("Schemas are not matching!"); }

随后通过std::static_pointer_casttable->column(i)->chunk(0)转成具体数组类型,其中列表列需要先取component_cost->values()得到底层DoubleArray

auto ids = std::static_pointer_cast<arrow::Int64Array>(table->column(0)->chunk(0)); auto component_cost = std::static_pointer_cast<arrow::ListArray>(table->column(2)->chunk(0)); auto component_cost_values = std::static_pointer_cast<arrow::DoubleArray>(component_cost->values());

读取每行时,用列表列的偏移量切片出该行的值区间:

const double* ccv_ptr = component_cost_values->raw_values(); for (int64_t i = 0; i < table->num_rows(); i++) { int64_t id = ids->Value(i); const double* first = ccv_ptr + component_cost->value_offset(i); const double* last = ccv_ptr + component_cost->value_offset(i + 1); std::vector<double> components_vec(first, last); rows.push_back({id, component, components_vec}); }

示例源码在这里特别强调了一个零拷贝切片陷阱(见 row_wise_conversion_example.cc):Arrow 支持零拷贝切片,因此数组的raw_values()原生指针可能带有切片偏移,手动访问裸指针时必须把value_offset(i)加上;而高层函数如Value(i)内部已自动处理该偏移,无需关心。这也是为什么对id/components直接调Value(i),而对列表列的值则需手工叠加偏移。

2.4 运行与验证

main中构造了 3 行样例数据{1, 1, {10.0}}{2, 3, {11.0, 12.0, 13.0}}{3, 2, {15.0, 25.0}},经过"行 → 表 → 行"往返后断言行数一致,并打印表格。运行输出应为:

ID Components Component prices 1 1 10 2 3 11 12 13 3 2 15 25

三、动态 Schema:以 RapidJSON 为例的通用转换框架

很多场景下,行数据的 Schema 在编译期未知(例如 JSON 文档、动态配置、来自其他系统的运行时协议)。原文档指出,Arrow 为此提供了若干实用工具:

工具作用
arrow::RecordBatchBuilder为一个完整 RecordBatch 创建并管理所有字段的 ArrayBuilder
arrow::VisitTypeInline按具体数组类型分派到专门的访问函数
arrow::enable_if_primitive_ctype等类型谓词(type-traits 文档)把模板函数收窄到特定 Arrow 类型,常与访问者模式配合
arrow::TableBatchReader一次读取表的一个批次,每个批次是零拷贝切片

完整的 RapidJSON 转换示例位于 rapidjson_row_converter.cc,它演示了如何编写任意 Schema的行列转换器,可读入运行参数:./rapidjson_row_converter [num_rows] [batch_size]

3.1 写方向:行 → Arrow RecordBatch

顶层函数 ConvertToRecordBatch

转换入口为ConvertToRecordBatch(源码 rapidjson_row_converter.cc):

arrow::Result<std::shared_ptr<arrow::RecordBatch>> ConvertToRecordBatch( const std::vector<rapidjson::Document>& rows, std::shared_ptr<arrow::Schema> schema) { // RecordBatchBuilder 会为 schema 中的每个字段创建数组构建器。 // 传入输出行数 rows.size() 可预分配正确的数组大小, // 但 string、byte 和 list 数组长度是动态的,无法预分配。 std::unique_ptr<arrow::RecordBatchBuilder> batch_builder; ARROW_ASSIGN_OR_RAISE( batch_builder, arrow::RecordBatchBuilder::Make(schema, arrow::default_memory_pool(), rows.size())); // 内部转换器负责把行值追加到提供的数组构建器上 JsonValueConverter converter(rows); for (int i = 0; i < batch_builder->num_fields(); ++i) { std::shared_ptr<arrow::Field> field = schema->field(i); arrow::ArrayBuilder* builder = batch_builder->GetField(i); ARROW_RETURN_NOT_OK(converter.Convert(*field.get(), builder)); } std::shared_ptr<arrow::RecordBatch> batch; ARROW_ASSIGN_OR_RAISE(batch, batch_builder->Flush()); // 用 ValidateFull() 检查数组构造是否正确 ARROW_RETURN_NOT_OK(batch->ValidateFull()); return batch; }

从源码结构看,其工作流是:RecordBatchBuilder::Make建好全部字段构建器 → 逐字段取GetField(i)JsonValueConverter追加行值 →Flush()产出批次 →ValidateFull()校验完整性。其中ValidateFull()会检查数组完整性(包括嵌套结构的深度校验),对调试新写的转换实现非常有用。

RecordBatchBuilder在头文件 table_builder.h 中定义,除了Make/GetField/Flush,还提供GetFieldAs<T>直接取强类型构建器、SetInitialCapacity/initial_capacity()控制预分配容量、num_fields()/schema()查询元信息。其实现是"为 schema 中每个字段创建一个ArrayBuilder子类实例"(内部持有std::vector<std::unique_ptr<ArrayBuilder>>),这正是它能"一键建好整批字段构建器"的原因。

中层:JsonValueConverter 与访问者模式

JsonValueConverter的职责是"为指定字段,把各行中该字段的值追加到给定的数组构建器"。为了按数据类型特化逻辑,它实现了一组Visit方法,并借助arrow::VisitTypeInline完成分派(见 rapidjson_row_converter.cc)。

arrow::VisitTypeInline的实现在 visit_type_inline.h:它本质是一个基于type.id()switch,通过宏为所有 Arrow 类型生成case TYPE##Type::type_id: return visitor->Visit(checked_cast<const TYPE##Type&>(type), args...);分支,把抽象arrow::DataType安全地分派到具体类型类的Visit重载上。访问者只需实现感兴趣的Visit重载,未实现的类型会落到默认实现(示例中返回Status::NotImplemented)。

Visit(const arrow::Int64Type&)为例,它遍历该列各行的 JSON 值,兼容 JSON 中整数可能出现的多种形态,缺失/空值则追加 null:

arrow::Status Visit(const arrow::Int64Type& type) { arrow::Int64Builder* builder = static_cast<arrow::Int64Builder*>(builder_); for (const auto& maybe_value : FieldValues()) { ARROW_ASSIGN_OR_RAISE(auto value, maybe_value); if (value->IsNull()) { ARROW_RETURN_NOT_OK(builder->AppendNull()); } else { if (value->IsUint()) { ARROW_RETURN_NOT_OK(builder->Append(value->GetUint())); } else if (value->IsInt()) { ARROW_RETURN_NOT_OK(builder->Append(value->GetInt())); } else if (value->IsUint64()) { ARROW_RETURN_NOT_OK(builder->Append(value->GetUint64())); } else if (value->IsInt64()) { ARROW_RETURN_NOT_OK(builder->Append(value->GetInt64())); } else { return arrow::Status::Invalid("Value is not an integer"); } } } return arrow::Status::OK(); }

DoubleTypeStringTypeBooleanTypeVisit结构与此类似,分别调用对应 Builder 的Append嵌套类型的处理更有代表性:

  • StructType:为每个子字段创建一个带root_path的子JsonValueConverter,递归转换子构建器后,再用FieldValues()统一追加 null 位图(rapidjson_row_converter.cc);
  • ListType:因为ListBuilder要求"值"与"偏移"交错追加,示例先用一个临时值构建器把整个值数组转换出来并Finish,然后对每一行调用Append(!value->IsNull())并通过AppendArraySlice按偏移切片拷入值(rapidjson_row_converter.cc)。
底层:FieldValues 与 DocValuesIterator

JsonValueConverter末尾的私有方法FieldValues()返回"当前字段跨所有行"的列值迭代器。对扁平行结构(如值向量),这个迭代器实现起来微不足道;但对 JSON 这类嵌套行结构,需要专门的迭代器来穿越嵌套层级——即DocValuesIterator(rapidjson_row_converter.cc)。

DocValuesIterator用两个状态量寻址 JSON 中的每个字段:

  • path:进入文档的字段名路径;
  • array_levels:需要穿越的数组层数。

示例注释给出了一个直观的例子(rapidjson_row_converter.cc):

{ "x": 3, // path: ["x"], array_levels: 0 "files": [ // path: ["files"], array_levels: 0 { // path: ["files"], array_levels: 1 "path": "my_str", // path: ["files", "path"], array_levels: 1 "sizes": [ // path: ["files", "size"], array_levels: 1 20, // path: ["files", "size"], array_levels: 2 22 ] } ] }

Next()方法维护一个ArrayPosition栈来记录每个数组层的当前位置,支持在空数组时回溯到下一个数组或行;字段缺失时返回全局的kNullJsonSingleton哨兵值(rapidjson_row_converter.cc),从而把"缺失"统一映射为 null。

3.2 读方向:Arrow RecordBatch → 行

顶层:ArrowToDocumentConverter 与分批处理

ArrowToDocumentConverter提供从 Arrow 批次/表到行(JSON 文档)的转换 API(rapidjson_row_converter.cc)。原文档特别指出:转成行的操作按小批次进行往往比整表一次转换更优,因此它拆成两个方法:

  • ConvertToVector:转换单个 RecordBatch,内部创建RowBatchBuilder,逐列设置字段并调用arrow::VisitArrayInline访问数组;
  • ConvertToIterator:借助arrow::TableBatchReader把表切成小批次迭代,最终返回arrow::Iterator<rapidjson::Document>,行可以逐个消费,也可以收集到容器中。
arrow::Iterator<rapidjson::Document> ConvertToIterator( std::shared_ptr<arrow::Table> table, size_t batch_size) { // 用 TableBatchReader 把表切成更小的批次,这些批次是至多 batch_size 行的零拷贝切片 auto batch_reader = std::make_shared<arrow::TableBatchReader>(*table); batch_reader->set_chunksize(batch_size); auto read_batch = this -> arrow::Result<arrow::Iterator<rapidjson::Document>> { ARROW_ASSIGN_OR_RAISE(auto rows, ConvertToVector(batch)); return arrow::MakeVectorIterator(std::move(rows)); }; auto nested_iter = arrow::MakeMaybeMapIterator( read_batch, arrow::MakeIteratorFromReader(std::move(batch_reader))); return arrow::MakeFlattenIterator(std::move(nested_iter)); }

这里用到了三个 Arrow 迭代器工具:arrow::MakeIteratorFromReader把 reader 包成批次迭代器,arrow::MakeMaybeMapIterator把每个批次映射为行迭代器(惰性求值),arrow::MakeFlattenIterator再把嵌套迭代器展平为单层行迭代器。

arrow::TableBatchReader在头文件 table.h 中定义,文档注释明确其转换是零拷贝的——每个批次都是表列切片上的视图。注意set_chunksize设置的是"期望的最大行数",实际每批行数可能更小,取决于表各列自身的分块特性(table.h)。

中层:RowBatchBuilder 与 enable_if_primitive_ctype

RowBatchBuilder负责填充输出行(rapidjson_row_converter.cc)。构造时预先reserve全部行数并初始化rapidjson::Document对象,避免不必要的扩容。它同样实现Visit()方法,但为了节省代码,对有原生 C 等价类型的数组(布尔、整数、浮点)写了一个模板方法,用arrow::enable_if_primitive_ctype限定:

template <typename ArrayType, typename DataClass = typename ArrayType::TypeClass> arrow::enable_if_primitive_ctype<DataClass, arrow::Status> Visit( const ArrayType& array) { assert(static_cast<int64_t>(rows_.size()) == array.length()); for (int64_t i = 0; i < array.length(); ++i) { if (!array.IsNull(i)) { rapidjson::Value str_key(field_->name(), rows_[i].GetAllocator()); rows_[i].AddMember(str_key, array.Value(i), rows_[i].GetAllocator()); } } return arrow::Status::OK(); }

enable_if_primitive_ctype定义于 type_traits.h:

template <typename T> using is_primitive_ctype = std::is_base_of<PrimitiveCType, T>; template <typename T, typename R = void> using enable_if_primitive_ctype = enable_if_t<is_primitive_ctype<T>::value, R>;

其原理是借助std::enable_if做 SFINAE 约束:只有当模板参数TPrimitiveCType的派生类时该重载才参与重载决议。这正是类型谓词 + 访问者模式的典型用法,RowBatchBuilder因此只需一个模板Visit覆盖所有原始类型,再为非原始类型单独写StringArrayStructArrayListArrayVisit重载(其中StructArray/ListArray会递归创建子RowBatchBuilder处理嵌套)。需要其他类型谓词(如日期时间类型has_c_type,见 type_traits.h)时可查阅 type-traits 文档。

arrow::VisitArrayInlineVisitTypeInline是对偶的:后者按DataType分派,前者按Array分派,定义于 visit_array_inline.h,内部同样以switch+ 宏展开实现。

数据流回顾

至此,"读方向"的完整链路可概括为:

arrow::Table └─ TableBatchReader(零拷贝分批) └─ ConvertToVector(逐批次) ├─ RowBatchBuilder 预建 num_rows 个空 Document └─ 逐列 VisitArrayInline: ├─ 原始类型 → 模板 Visit(enable_if_primitive_ctype) ├─ 字符串 → Visit(StringArray) ├─ 结构体 → 递归子 RowBatchBuilder └─ 列表 → 递归值子批次 + 偏移切片

3.3 完整链路验证:示例自带的自检断言

DoRowConversion(rapidjson_row_converter.cc)演示了完整的往返过程,可直接作为转换器正确性的自检用例:

  1. 准备 3 条 JSON 记录模板,循环填充num_rows行(默认 100)并打印原始 JSON;
  2. 声明目标 Schema:pk: int64date_created: utf8data: struct(deleted: bool, metrics: list(struct(key: utf8, value: int64)))
  3. 调用ConvertToRecordBatch转批次,Table::FromRecordBatches组装表并打印table->ToString()
  4. ArrowToDocumentConverter::ConvertToIteratorbatch_size(默认 100)转回文档流;
  5. 对每个输出文档执行一连串assert:校验pk是 Int64、date_created是字符串、data.deleted是布尔、metrics是数组,且数组元素含字符串key与 Int64value(rapidjson_row_converter.cc)。

这些断言与ValidateFull()共同构成了双重保障:前者验证"转换结果符合预期行形态",后者验证"构造出的 Arrow 数组内部结构合法"。这正是把本文模式复用到自己项目时应保留的最佳实践。


四、两种方案的对比与选型建议

维度固定 Schema(手工 Builder)动态 Schema(Builder + 访问者)
Schema 来源编译期硬编码运行期传入(其他系统提供或推断)
新增类型支持每加一种类型手写代码为访问者新增一个Visit重载即可
嵌套类型手动串联ListBuilder/StructBuilder层级递归转换 + 路径迭代器自动处理
典型适用类型已知的批量导入导出JSON/动态协议类连接器、通用工具库

原文档还给出了一个重要的工程经验:在转换过程中推断 Schema 是很有挑战的,很多系统会"先检查前 N 行来推断 Schema"(如果事先没有 Schema 可用)。因此推荐的做法是:转换前确定目标 Schema(来自其他系统或单独的函数),转换器只负责忠实执行字段映射。

关于构建效率,RecordBatchBuilder::Make(schema, pool, initial_capacity)的第三个参数允许预分配初始容量,除 string/byte/list 等变长类型外,都能在转换前预留好内存(见 rapidjson_row_converter.cc);RowBatchBuilder构造时的reserve同理。这些细节在高吞吐场景下对减少重分配有明显收益。


五、进一步阅读

  • 本主题的原始文档:Row to columnar conversion
  • 固定 Schema 完整示例:row_wise_conversion_example.cc
  • 动态 Schema 完整示例:rapidjson_row_converter.cc
  • RecordBatchBuilder头文件:table_builder.h
  • TableBatchReader头文件:table.h
  • 类型分派实现:visit_type_inline.h、visit_array_inline.h
  • 类型谓词定义:type_traits.h
  • 构建示例的 CMake 配置:cpp/examples/arrow/CMakeLists.txt

【免费下载链接】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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询