ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Apache Arrow C++ IPC 读写 API 实战:流式格式与文件格式的完整解析

Apache Arrow C++ IPC 读写 API 实战:流式格式与文件格式的完整解析 Apache Arrow C IPC 读写 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/arrowApache Arrow 的 IPCInter-Process Communication是其在列式内存格式之上定义的一套高效序列化协议用于在不同进程、不同语言环境之间传输 RecordBatch 数据。本指南以 Apache Arrow C 的arrow::ipc命名空间为核心系统讲解 IPC 读写 API 的完整体系IpcReadOptions/IpcWriteOptions全部选项、阻塞式与事件驱动两种读取方式、写入器的工厂函数与终止语义以及 ReadStats / WriteStats 统计结构。读完本文你将掌握如何用 C 正确读取和写入 Arrow IPC Stream流式与 IPC File文件两种格式并理解压缩、字典、端序、对齐等关键参数对数据兼容性的影响。IPC 两种格式Stream 与 File 的本质区别在深入 API 之前先厘清 Arrow IPC 的两种格式。二者都由封装消息encapsulated message构成但布局与能力不同对应的读取类也不同。IPC Streaming Format流式格式数据是一系列顺序排列的消息。Schema 首先出现随后是字典批DictionaryBatch与记录批RecordBatch它们可以交错出现但任何字典键在 RecordBatch 中使用前必须先由 DictionaryBatch 定义。流式格式可以用 8 字节的 EOS 标记0xFFFFFFFF 0x00000000或直接关闭流来结束。这种格式只要求顺序输入流即可读取适合网络传输、管道等场景。格式细节见 Columnar.rst。IPC File Format文件格式是流式格式的超集支持随机访问。文件以魔数ARROW1开头结尾处写入一个包含 Schema 冗余副本、每个数据块内存偏移与大小的footer从而可以随机读取任意一个 RecordBatch。其布局为magic number ARROW1 empty padding bytes [to 8 byte boundary] STREAMING FORMAT with EOS FOOTER FOOTER SIZE: little-endian int32 magic number ARROW1文件格式要求随机访问文件io::RandomAccessFile扩展名建议为.arrow历史上也称 Feather V2。两种格式的结构示意同样见 Columnar.rst。IPC 官方文档页面 ipc.rst 正是围绕这两种格式的读写 API 展开的下面逐一深入。读写选项IpcWriteOptions 与 IpcReadOptionsIpcWriteOptions写侧的全部可调参数arrow::ipc::IpcWriteOptions定义在 options.h控制 IPC 消息序列化时的全部行为参数默认值说明allow_64bitfalse是否允许字段长度超出有符号 32 位整数范围。开启后部分实现可能无法解析该流max_recursion_depthkMaxNestingDepth64允许的最大 Schema 嵌套深度。对深嵌套 Schema需显式调大alignment8内存缓冲区写入时按该字节数补齐对齐到其倍数write_legacy_ipc_formatfalse是否写入 0.15.0 之前的旧 IPC 格式4 字节前缀而非 8 字节memory_pooldefault_memory_pool()写入过程中的内存池。IPC 以零拷贝为主但启用压缩等场景仍需分配内存codec空不压缩RecordBatch body 缓冲区的压缩编码器仅允许 UNCOMPRESSED、LZ4_FRAME、ZSTDmin_space_savings未设置应用压缩所需的最小空间节省比例按1.0 - compressed_size / uncompressed_size计算超出 [0,1] 范围视为错误。启用后对 Arrow C 12.0.0 之前的版本可能不可读use_threadstrue是否使用全局 CPU 线程池并行化压缩等计算任务emit_dictionary_deltasfalse字典变化时是否尝试输出字典增量delta而非完整替换。默认为 false 以最大化流兼容性嵌套字典从不输出增量unify_dictionariesfalse写 IPC 文件时是否统一各 chunk 的字典。IPC 文件格式不支持字典替换若为 false 且同一列字典不兼容则直接报错该选项对支持字典替换/增量的流式格式无效metadata_versionMetadataVersion::V5IPC 消息与元数据的格式版本。V5 需 1.0.0 及以上可读V4 需 0.8.0 及以上可读关于压缩格式规范还提醒若上层存储/传输协议已用 gzip 等压缩应避免再做 IPC 缓冲区压缩——双重压缩通常得不偿失见 Columnar.rst。IpcReadOptions读侧的全部可调参数arrow::ipc::IpcReadOptions定义在 options.h控制反序列化行为参数默认值说明max_recursion_depthkMaxNestingDepth64允许的最大 Schema 嵌套深度memory_pooldefault_memory_pool()读取过程中的内存池included_fields空反序列化 RecordBatch 时只包含的顶层 Schema 字段索引空表示全部字段use_threadstrue是否使用全局 CPU 线程池并行化解压缩等任务ensure_native_endiantrue是否将输入数据转换为平台原生端序。Arrow IPC 格式默认小端若收到的 Schema 端序与平台不同所有端序敏感缓冲区数值/时间/十进制类型的 value buffer以及变长 binary 和 list 类类型的 offset buffer都会被字节交换。该转换由 RecordBatchFileReader、RecordBatchStreamReader 和 StreamDecoder 完成ensure_alignmentAlignment::kAnyAlignment数据对齐策略。kAnyAlignment保持原样不复制kDataTypeSpecificAlignment按具体数据类型对齐k64ByteAlignment对齐到 64 字节边界。需要对齐时会通过配置的 memory_pool 复制数据到对齐内存pre_buffer_cache_optionsio::CacheOptions::LazyDefaults()预缓冲pre-buffering时的缓存行为控制读取 IPC阻塞式 API文档将读取 API 划分为阻塞式与事件驱动两类。阻塞式 API 对应两个类选择依据是目标格式reader.h。RecordBatchStreamReader读流式格式arrow::ipc::RecordBatchStreamReader继承自arrow::RecordBatchReader同步地从io::InputStream读取流。它在流的最前面读取 Schema以及字典随后依次读取 RecordBatch。三种Open重载分别接收MessageReader、裸指针io::InputStream*或shared_ptrio::InputStream均接受可选的IpcReadOptions#include arrow/ipc/reader.h auto input /* 一个 io::InputStream例如 FileInputStream */; arrow::ipc::IpcReadOptions read_options; read_options.use_threads true; read_options.included_fields {0, 1, 2}; // 只读前三个字段 ARROW_ASSIGN_OR_RAISE(auto reader, arrow::ipc::RecordBatchStreamReader::Open(input, read_options)); std::shared_ptrarrow::RecordBatch batch; while (true) { ARROW_ASSIGN_OR_RAISE(batch, reader-Next()); if (batch nullptr) break; // 读到 EOS 或流结束 // 处理 batch... }要点Next()返回nullptr即表示流已读完RecordBatchStreamReader::Open传入的流对象必须在 reader 生命周期内保持存活。RecordBatchFileReader读文件格式arrow::ipc::RecordBatchFileReader针对随机访问文件核心能力是按索引随机读取。除Open(io::RandomAccessFile*, options)外还支持传入footer_offset文件嵌入在更大文件/内存区域中时给出包含元数据 footer 的文件结尾绝对偏移。OpenAsync变体提供异步打开。常用成员schema()文件的 Schemanum_record_batches()文件中 RecordBatch 总数version()/metadata()元数据版本与 footer 中的自定义元数据ReadRecordBatch(int i)/ReadRecordBatchWithCustomMetadata(int i)按索引读取输入源支持零拷贝时不复制内存CountRows()计算文件总行数PreBufferMetadata(indices)预加载指定批次元数据与全部字典消息GetRecordBatchGenerator(...)获取可重入的 RecordBatch 异步生成器支持 I/O 合并coalesceToRecordBatches()/ToTable()一次性收集全部批次或拼接为 TableARROW_ASSIGN_OR_RAISE(auto file_reader, arrow::ipc::RecordBatchFileReader::Open(file, read_options)); std::cout batches: file_reader-num_record_batches() std::endl; ARROW_ASSIGN_OR_RAISE(auto batch, file_reader-ReadRecordBatch(0)); // 随机读取第 0 批低层零拷贝读取函数除了两个 Reader 类arrow::ipc还提供一组通用读取函数见 reader.hReadSchema()读取单独的 Schema 消息并填充 DictionaryMemo、ReadRecordBatch()从 Message 或流中读取单个批次。当输入支持零拷贝时这些函数不复制数据适合需要更细粒度控制的场景。读取 IPC事件驱动 APIListener StreamDecoder阻塞式 API 由调用方主动拉取数据事件驱动 API 则由用户把原始字节推给解码器解码器在解析出完整消息时回调监听器。文档明确指出要实现事件驱动读取必须实现arrow::ipc::Listener子类并传给arrow::ipc::StreamDecoder。该 API 自 0.17.0 引入目前标记为EXPERIMENTAL。Listener回调接口arrow::ipc::Listener见 reader.h是可覆写的回调集合回调默认行为说明OnEOS()返回 OK收到流结束标记时调用OnSchemaDecoded(schema)返回 OK解码出 Schema 时调用OnSchemaDecoded(schema, filtered_schema)转发给单参数版本额外提供只含读取字段的过滤后 Schema13.0.0 起OnRecordBatchDecoded(batch)返回 NotImplemented解码出 RecordBatch 时调用必须覆写OnRecordBatchWithMetadataDecoded(batch_with_metadata)转发给OnRecordBatchDecoded带自定义元数据的批次回调13.0.0 起仓库还提供了一个开箱即用的CollectListenerreader.h它收集解码出的 Schema 与全部批次并提供schema()、filtered_schema()、record_batches()、PopRecordBatch()等访问方法适合快速原型。StreamDecoder推式解码器arrow::ipc::StreamDecoderreader.h接收任意大小的数据块内部维护状态机class MyListener : public arrow::ipc::Listener { public: arrow::Status OnSchemaDecoded(std::shared_ptrarrow::Schema schema) override { schema_ schema; return arrow::Status::OK(); } arrow::Status OnRecordBatchDecoded( std::shared_ptrarrow::RecordBatch batch) override { batches_.push_back(std::move(batch)); return arrow::Status::OK(); } // ... OnEOS 等 }; auto listener std::make_sharedMyListener(); arrow::ipc::StreamDecoder decoder(listener, read_options); decoder.Consume(data, size); // 传入原始字节内部不拷贝调用方须保证内存存活 decoder.Consume(std::make_sharedarrow::Buffer(data, size)); // 或以 Buffer 传入StreamDecoder的两个关键性能接口Consume(const uint8_t*, int64_t)/Consume(shared_ptrBuffer)喂入数据。若数据足以解码出一个或多个批次会多次回调OnRecordBatchDecoded。next_required_size()返回推进解码器状态所需的最小字节数。若每次喂入恰好该大小的数据解码器不会使用内部缓冲避免拼接小分块的性能开销。官方注释给出的优化范式是循环get_data(decoder.next_required_size())再Consume见 reader.h。此外Reset()可复用解码器处理新流schema()返回流中批次的共享 Schema。读写统计ReadStats 与 WriteStats两个统计结构见 reader.h 与 writer.h字段一致帮助诊断与观测字段含义num_messages读/写的 IPC 消息总数num_record_batches读/写的 RecordBatch 数num_dictionary_batches字典批总数恒满足 num_dictionary_deltas num_replaced_dictionariesnum_dictionary_deltas字典增量批数num_replaced_dictionaries字典替换批数用无关新字典替换既有字典ReadStats额外含original_endianness记录 IPC 数据原始端序——只有当IpcReadOptions.ensure_native_endian设为 false 时才会在返回的 Schema 中体现。WriteStats额外含total_raw_body_size与total_serialized_body_size分别统计未压缩原始体大小与含 padding/压缩的序列化体大小。三个读取类与RecordBatchWriter均提供stats()访问当前统计。写入 IPC工厂函数与 RecordBatchWriter写侧 API 的核心是arrow::ipc::RecordBatchWriter抽象类writer.h提供WriteRecordBatch()、WriteTable()可指定max_chunksize传 -1 表示不限与Close()。实例由两个工厂函数创建record-batch-writer-factories组writer.hMakeStreamWriter(sink, schema, options)创建 IPC流式写入器输出到io::OutputStreamMakeFileWriter(sink, schema, options, metadata)创建 IPC文件写入器metadata参数为可选的 footer 自定义元数据。一个必须牢记的语义文件格式必须显式 Close()文档用加粗的篇幅强调了一个关键事实IPC 流式格式是可选终止的而 IPC 文件格式必须包含终止 footer。因此文件格式的写入器必须显式调用RecordBatchWriter::Close()否则产出的文件是损坏的corrupt。与之相对流式写入器不调用Close()也不会产生非法数据。#include arrow/ipc/writer.h #include arrow/io/file.h auto sink /* 一个 io::OutputStream例如 FileOutputStream */; arrow::ipc::IpcWriteOptions write_options; write_options.emit_dictionary_deltas true; // 写流式格式时允许字典增量 write_options.metadata_version arrow::ipc::MetadataVersion::V5; ARROW_ASSIGN_OR_RAISE(auto writer, arrow::ipc::MakeFileWriter(sink.get(), schema, write_options)); for (const auto batch : batches) { ARROW_RETURN_NOT_OK(writer-WriteRecordBatch(*batch)); } ARROW_RETURN_NOT_OK(writer-Close()); // 文件格式缺失将导致文件损坏WriteTable(table, max_chunksize)会把可能分块的 Table 自动切成一系列 RecordBatch 写入。仓库内真实用法stream-to-file 工具cpp/src/arrow/ipc/stream_to_file.cc完整演示了读流式、写文件的转换流程这也是对上述 API 组合的权威示例arrow::ipc::IpcWriteOptions write_options; write_options.emit_dictionary_deltas true; ARROW_ASSIGN_OR_RAISE(auto reader, RecordBatchStreamReader::Open(input)); ARROW_ASSIGN_OR_RAISE(auto writer, MakeFileWriter(sink, reader-schema(), write_options)); std::shared_ptrRecordBatch batch; while (true) { ARROW_ASSIGN_OR_RAISE(batch, reader-Next()); if (batch nullptr) break; RETURN_NOT_OK(writer-WriteRecordBatch(*batch)); } return writer-Close();该工具典型的用法是$ 流式输出程序 | stream-to-file file.arrow。仓库还提供了反向工具file_to_stream.cc两者共同构成读写两种格式的完整参照。低层序列化函数若需更精细的控制arrow::ipc还暴露了SerializeRecordBatch(batch, options)把单个批次序列化为封装消息 BufferSerializeSchema(schema, pool)序列化 SchemaWriteRecordBatchStream(batches, options, dst)一次性写多个同 Schema 批次GetRecordBatchSize(batch, options, size)计算含元数据的完整消息字节数便于预分配内存WriteIpcPayload/GetSchemaPayload/GetDictionaryPayload/GetRecordBatchPayload基于IpcPayload元数据头 零个或多个 body 缓冲区的中间结构见 writer.h的底层写路径。字典与压缩影响兼容性的高级考量结合前文选项与格式规范有两个容易踩坑的点值得单独说明字典Dictionary。流式格式支持字典替换replacement与增量delta而文件格式不支持字典替换——对同一字典 ID 只能发出一个非增量字典批增量字典批按其在 footer 中的顺序应用Columnar.rst。写文件时若各 chunk 字典不一致要么靠IpcWriteOptions.unify_dictionaries自动统一有运行时开销要么报错。这就是emit_dictionary_deltas默认关闭的原因保证最大兼容性。压缩。仅支持 UNCOMPRESSED、LZ4_FRAME、ZSTD 三种编码器min_space_savings可以避免压缩后反而更大的浪费但注意启用该选项后Arrow C 12.0.0 之前的版本无法读取见 options.h。端序。IPC 格式默认小端。跨平台交换数据时读侧默认开启ensure_native_endian自动字节交换若关闭该选项可通过ReadStats.original_endianness获知数据原始端序reader.h。小结Apache Arrow C 的arrow::ipc模块为两种 IPC 格式提供了层次分明、能力互补的完整读写路径读取阻塞式RecordBatchStreamReader/RecordBatchFileReader适合拉取模型后者支持按索引随机访问与异步生成器事件驱动ListenerStreamDecoder适合推送模型与流式管道写入MakeStreamWriter/MakeFileWriter创建写入器RecordBatchWriter统一承载批次写入文件格式务必Close()落 footer控制面IpcReadOptions/IpcWriteOptions覆盖嵌套深度、对齐、端序、压缩、字典与元数据版本等全部关键行为观测面ReadStats/WriteStats提供消息、批次、字典增量与体大小的量化统计。结合本文引用的源码options.h、reader.h、writer.h与仓库工具stream_to_file.cc、file_to_stream.cc你可以在此基础上直接落地跨进程、跨语言的数据交换或在管道式 ETL 与随机访问分析存储之间自由切换。【免费下载链接】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),仅供参考
返回列表