拓冰建站拓冰建站
首页 / 资讯中心 / 正文

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

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::Table、arrow::RecordBatch、arrow::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::vectordouble component_cost; };对应到列式世界目标arrow::Table由三列组成idint64—— 产品 IDcomponentsint64—— 组件数量component_costlist(float64)—— 每个组件的成本是一个嵌套的列表列。2.2 行 → 列手工组装 ArrayBuilderArrow 为每个数据类型提供了专门的 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_sharedDoubleBuilder(pool)); // 该 builder 由 component_cost_builder 拥有 DoubleBuilder* component_item_cost_builder (static_castDoubleBuilder*(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::Tablestd::vectorstd::shared_ptrarrow::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_sharedarrow::Schema(schema_vector); std::shared_ptrarrow::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_cast把table-column(i)-chunk(0)转成具体数组类型其中列表列需要先取component_cost-values()得到底层DoubleArrayauto ids std::static_pointer_castarrow::Int64Array(table-column(0)-chunk(0)); auto component_cost std::static_pointer_castarrow::ListArray(table-column(2)-chunk(0)); auto component_cost_values std::static_pointer_castarrow::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::vectordouble components_vec(first, last); rows.push_back({id, component, components_vec}); }示例源码在这里特别强调了一个零拷贝切片陷阱见 row_wise_conversion_example.ccArrow 支持零拷贝切片因此数组的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 创建并管理所有字段的 ArrayBuilderarrow::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.ccarrow::Resultstd::shared_ptrarrow::RecordBatch ConvertToRecordBatch( const std::vectorrapidjson::Document rows, std::shared_ptrarrow::Schema schema) { // RecordBatchBuilder 会为 schema 中的每个字段创建数组构建器。 // 传入输出行数 rows.size() 可预分配正确的数组大小 // 但 string、byte 和 list 数组长度是动态的无法预分配。 std::unique_ptrarrow::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_ptrarrow::Field field schema-field(i); arrow::ArrayBuilder* builder batch_builder-GetField(i); ARROW_RETURN_NOT_OK(converter.Convert(*field.get(), builder)); } std::shared_ptrarrow::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还提供GetFieldAsT直接取强类型构建器、SetInitialCapacity/initial_capacity()控制预分配容量、num_fields()/schema()查询元信息。其实现是为 schema 中每个字段创建一个ArrayBuilder子类实例内部持有std::vectorstd::unique_ptrArrayBuilder这正是它能一键建好整批字段构建器的原因。中层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_castconst TYPE##Type(type), args...);分支把抽象arrow::DataType安全地分派到具体类型类的Visit重载上。访问者只需实现感兴趣的Visit重载未实现的类型会落到默认实现示例中返回Status::NotImplemented。以Visit(const arrow::Int64Type)为例它遍历该列各行的 JSON 值兼容 JSON 中整数可能出现的多种形态缺失/空值则追加 nullarrow::Status Visit(const arrow::Int64Type type) { arrow::Int64Builder* builder static_castarrow::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(); }DoubleType、StringType、BooleanType的Visit结构与此类似分别调用对应 Builder 的Append。嵌套类型的处理更有代表性StructType为每个子字段创建一个带root_path的子JsonValueConverter递归转换子构建器后再用FieldValues()统一追加 null 位图rapidjson_row_converter.ccListType因为ListBuilder要求值与偏移交错追加示例先用一个临时值构建器把整个值数组转换出来并Finish然后对每一行调用Append(!value-IsNull())并通过AppendArraySlice按偏移切片拷入值rapidjson_row_converter.cc。底层FieldValues 与 DocValuesIteratorJsonValueConverter末尾的私有方法FieldValues()返回当前字段跨所有行的列值迭代器。对扁平行结构如值向量这个迭代器实现起来微不足道但对 JSON 这类嵌套行结构需要专门的迭代器来穿越嵌套层级——即DocValuesIteratorrapidjson_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 文档的转换 APIrapidjson_row_converter.cc。原文档特别指出转成行的操作按小批次进行往往比整表一次转换更优因此它拆成两个方法ConvertToVector转换单个 RecordBatch内部创建RowBatchBuilder逐列设置字段并调用arrow::VisitArrayInline访问数组ConvertToIterator借助arrow::TableBatchReader把表切成小批次迭代最终返回arrow::Iteratorrapidjson::Document行可以逐个消费也可以收集到容器中。arrow::Iteratorrapidjson::Document ConvertToIterator( std::shared_ptrarrow::Table table, size_t batch_size) { // 用 TableBatchReader 把表切成更小的批次这些批次是至多 batch_size 行的零拷贝切片 auto batch_reader std::make_sharedarrow::TableBatchReader(*table); batch_reader-set_chunksize(batch_size); auto read_batch this - arrow::Resultarrow::Iteratorrapidjson::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_ctypeRowBatchBuilder负责填充输出行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_ctypeDataClass, arrow::Status Visit( const ArrayType array) { assert(static_castint64_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.htemplate typename T using is_primitive_ctype std::is_base_ofPrimitiveCType, T; template typename T, typename R void using enable_if_primitive_ctype enable_if_tis_primitive_ctypeT::value, R;其原理是借助std::enable_if做 SFINAE 约束只有当模板参数T是PrimitiveCType的派生类时该重载才参与重载决议。这正是类型谓词 访问者模式的典型用法RowBatchBuilder因此只需一个模板Visit覆盖所有原始类型再为非原始类型单独写StringArray、StructArray、ListArray的Visit重载其中StructArray/ListArray会递归创建子RowBatchBuilder处理嵌套。需要其他类型谓词如日期时间类型has_c_type见 type_traits.h时可查阅 type-traits 文档。arrow::VisitArrayInline与VisitTypeInline是对偶的后者按DataType分派前者按Array分派定义于 visit_array_inline.h内部同样以switch 宏展开实现。数据流回顾至此读方向的完整链路可概括为arrow::Table └─ TableBatchReader零拷贝分批 └─ ConvertToVector逐批次 ├─ RowBatchBuilder 预建 num_rows 个空 Document └─ 逐列 VisitArrayInline ├─ 原始类型 → 模板 Visitenable_if_primitive_ctype ├─ 字符串 → Visit(StringArray) ├─ 结构体 → 递归子 RowBatchBuilder └─ 列表 → 递归值子批次 偏移切片3.3 完整链路验证示例自带的自检断言DoRowConversionrapidjson_row_converter.cc演示了完整的往返过程可直接作为转换器正确性的自检用例准备 3 条 JSON 记录模板循环填充num_rows行默认 100并打印原始 JSON声明目标 Schemapk: int64、date_created: utf8、data: struct(deleted: bool, metrics: list(struct(key: utf8, value: int64)))调用ConvertToRecordBatch转批次Table::FromRecordBatches组装表并打印table-ToString()用ArrowToDocumentConverter::ConvertToIterator按batch_size默认 100转回文档流对每个输出文档执行一连串assert校验pk是 Int64、date_created是字符串、data.deleted是布尔、metrics是数组且数组元素含字符串key与 Int64valuerapidjson_row_converter.cc。这些断言与ValidateFull()共同构成了双重保障前者验证转换结果符合预期行形态后者验证构造出的 Arrow 数组内部结构合法。这正是把本文模式复用到自己项目时应保留的最佳实践。四、两种方案的对比与选型建议维度固定 Schema手工 Builder动态 SchemaBuilder 访问者Schema 来源编译期硬编码运行期传入其他系统提供或推断新增类型支持每加一种类型手写代码为访问者新增一个Visit重载即可嵌套类型手动串联ListBuilder/StructBuilder层级递归转换 路径迭代器自动处理典型适用类型已知的批量导入导出JSON/动态协议类连接器、通用工具库原文档还给出了一个重要的工程经验在转换过程中推断 Schema 是很有挑战的很多系统会先检查前 N 行来推断 Schema如果事先没有 Schema 可用。因此推荐的做法是转换前确定目标 Schema来自其他系统或单独的函数转换器只负责忠实执行字段映射。关于构建效率RecordBatchBuilder::Make(schema, pool, initial_capacity)的第三个参数允许预分配初始容量除 string/byte/list 等变长类型外都能在转换前预留好内存见 rapidjson_row_converter.ccRowBatchBuilder构造时的reserve同理。这些细节在高吞吐场景下对减少重分配有明显收益。五、进一步阅读本主题的原始文档Row to columnar conversion固定 Schema 完整示例row_wise_conversion_example.cc动态 Schema 完整示例rapidjson_row_converter.ccRecordBatchBuilder头文件table_builder.hTableBatchReader头文件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),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门