Red Arrow Format:Apache Arrow 纯 Ruby 序列化实现的文件格式解析与实战指南
2026/9/14 5:52:59 网站建设 项目流程

Red Arrow Format:Apache Arrow 纯 Ruby 序列化实现的文件格式解析与实战指南

【免费下载链接】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

Red Arrow Format(red-arrow-formatgem)是 Apache Arrow 生态中一个纯 Ruby 实现的格式层库,提供 Arrow IPC(File 与 Streaming)两种形态的序列化/反序列化能力,但不包含计算、扫描等数据处理功能。读完本文,你将掌握它的安装方式、FileReader/FileWriter/StreamingReader/StreamingWriter的实际用法、Arrow IPC 文件在字节层的真实布局(魔数、continuation token、Footer、MetadataVersion V5),以及如何在仓库中运行其测试与基准测试、重新生成 FlatBuffers 绑定代码。

定位:只序列化,不处理数据

README 开宗明义:

Red Arrow Format is the pure Ruby Apache Arrow format serializer and deserializer implementation. This provides only serialize/deserialize features. If you want to process Apache Arrow data not only serialize/deserialize Apache Arrow data, you should use Red Arrow not Red Arrow Format.

也就是说,它是 Red Arrow 家族(ruby/red-arrow/)的“格式层”:

  • Red Arrow Format(本模块):把 Arrow IPC 字节流与 Ruby 的Schema/RecordBatch相互转换,纯 Ruby 实现,不依赖 C++ 扩展;
  • Red Arrow(ruby/red-arrow:基于libarrow的完整绑定,负责真正处理 Arrow 数据。

gemspec 中也复述了这一定位(Red Arrow Format provides only serialize/deserialize features...,见第 37~42 行),并声明了唯一的运行时依赖red-flatbuffers >= 0.0.8(第 50 行)——这正是“纯 Ruby + 纯 Ruby FlatBuffers 解码”的实现基础。当前仓库中 版本文件 显示为26.0.0-SNAPSHOT

安装

通过 Bundler,在Gemfile中加入:

gem "red-arrow-format"

或直接使用 RubyGems:

$ gem install red-arrow-format

使用时的加载方式是一行 require:

require "arrow-format"

入口文件 会一次性引入四个核心类(第 18~22 行):FileReaderFileWriterStreamingReaderStreamingWriter以及版本常量,因此 require 一次即可使用全部读写 API。

读取 Arrow 文件:FileReader

README 给出的最小示例:

require "arrow-format" File.open("/dev/shm/data.arrow", "rb") do |input| reader = ArrowFormat::FileReader.new(input) reader.each do |record_batch| # Use record_batch end end

结合 file-reader.rb 的源码,可以补充几点 README 未明说的事实:

  • 输入形态:构造函数(第 39~55 行)接受三种输入——IO(会用IO::Buffer.map做只读内存映射,零拷贝)、String(整段字节串)、或已经是IO::Buffer的对象;
  • 构造即校验initialize内依次执行validate(检查最小尺寸与首尾ARROW1魔数)、read_footerread_schemaread_dictionaries,因此拿到reader对象时 schema 与字典都已就绪;
  • 随机访问:除each顺序迭代外,FileReader提供n_record_batches(批数量)与read(i)(按块偏移直接读取第 i 个 RecordBatch,第 57~74 行),适合需要随机跳批的场景;
  • 暴露的属性reader.schemareader.metadata(自定义元数据)可直接读取。

文件结构校验与 Footer 定位

validate(第 85~104 行)要求输入至少为8 字节起始魔数 + 4 字节 footer 长度 + 6 字节结尾魔数,并比对首尾两个ARROW1标记(MAGIC = "ARROW1",第 25 行)。read_footer(第 106~112 行)的定位逻辑是典型的 IPC File 布局:从文件末尾倒数,先取结尾魔数前的 4 字节小端 s32 作为 footer 大小,再向前切出 Footer 的 FlatBuffers 字节。这与 format/File.fbs 定义的File布局一致。

continuation token 与向后兼容

read_block(第 114~171 行)负责按 Footer 中记录的Block(offset / metadata_length / body_length)切出一个 RecordBatch 消息,其中有一段值得注意的兼容逻辑(第 134~147 行):

  1. 先读 4 字节:若为 continuation 串(0xFFFFFFFF)则跳过;
  2. 若该值不是负数(即不是合法的 metadata 长度前驱),则视为版本 0.15.0 之前、没有 continuation token 的旧数据,回退为直接把该 4 字节当作 metadata length 重新解析。

源码注释原文即:“For backward compatibility of data produced prior to version 0.15.0. It doesn't have continuation token.” 这说明该读取器能消费相当早期的 Arrow IPC 文件。

字典解码

对于带字典的列,read_dictionaries(第 173~229 行)在构造阶段就完成了解码,并且严格实现了 Arrow 规范中的两条约束:delta 字典批次必须跟在非 delta 批次之后(否则抛FileReadError),同一字典 ID 不允许出现多个非 delta 批次。解码出的Dictionary对象挂到DictionaryType上供后续 RecordBatch 还原索引列使用。

写入 Arrow 文件:FileWriter

README 只演示了读取,但同一 gem 提供了对称的写入 API。file-writer.rb 中FileWriter < StreamingWriter,其完整生命周期是三步:

require "arrow-format" output = File.open("/dev/shm/data.arrow", "wb") writer = ArrowFormat::FileWriter.new(output) writer.start(schema) # schema 为 ArrowFormat::Schema writer.write_record_batch(record_batch) # 可多次调用 writer.finish # 写 Footer 并收尾
  • start(第 25~30 行):先写魔数ARROW1加 2 字节零填充(MAGIC_PADDING = "\x00\x00",把 6 字节魔数补齐到 8 字节边界),再委托StreamingWriter#start写出 schema 消息;
  • write_record_batch:继承自 streaming-writer.rb(第 42~52 行),会先为该批中出现的每个字典列写DictionaryBatch消息(仅增量部分,见下文),再写批消息本体;
  • finish(第 32~37 行):调用write_footer序列化 Footer 并追加 4 字节 footer 长度,最后再写一次ARROW1结尾魔数,与读取端的validate精确对应。

元数据版本与消息编码细节

StreamingWriter中有几个关键常量值得记住(streaming-writer.rb):

ALIGNMENT_SIZE = IO::Buffer.size_of(:u64) # 8 字节对齐 CONTINUATION = "\xFF\xFF\xFF\xFF" # 0xFFFFFFFF EOS = "\xFF\xFF\xFF\xFF\x00\x00\x00\x00" # 流式数据的结束标记

每条消息的编码由build_metadata(第 79~91 行)与write_message(第 148~154 行)完成:0xFFFFFFFFcontinuation + 小端 s32 metadata 长度(pack("l<"))+ 对齐到 8 字节边界的 FlatBuffersMessageversion = MetadataVersion::V5)。Footer 同样固定写为V5(file-writer.rb 第 42 行),并记录每个 Block 的offsetmeta_data_lengthbody_length,这正是读取端read_block的依据。

字典的增量写入

write_dictionary(streaming-writer.rb 第 115~146 行)用一个@written_dictionary_offsets表跟踪每个字典 ID 已写到的偏移:新批次的字典值只有超出已写部分才会被切出来(data.slice(written_offset - current_base_offset))并以delta = trueDictionaryBatch写出。这与读取端read_dictionaries的“delta 必须跟随非 delta”校验是一对镜像实现,保证多批数据间共享字典时不会重复序列化。

流式形态:StreamingReader / StreamingWriter

除 File 形态外,lib/arrow-format.rb 还导出StreamingReaderStreamingWriter,对应 Arrow IPC 的流式布局(无 Footer,以EOS标记结束)。

streaming-reader.rb 的读取是按需拉取的:StreamingReader#each(第 48~59 行)通过StreamingPullReader#next_required_size询问“我还需要多少字节才能交出下一条 RecordBatch”,再按该大小从输入读取并consume。因此它既支持整块输入(File会走内存映射,StringIO::Buffer),也支持真正的流(如 Socket),ensure_schema(第 80~85 行)保证构造函数返回前已拿到 schema。写入侧的finish(streaming-writer.rb 第 54~57 行)写出的EOS\xFF\xFF\xFF\xFF\x00\x00\x00\x00,即 continuation 后跟 metadata length 为 0 的终止消息。

FlatBuffers 绑定代码是如何生成的

lib/arrow-format/org/apache/arrow/flatbuf/下 60 余个文件(footer.rbrecord_batch.rbschema.rbbody_compression.rb……)并非手写,而是由 Rakefile 中的flat_buffers:generate任务生成:

  1. flatc --binary --schema ...把仓库根目录的 format/File.fbs 与 format/Message.fbs 编译为二进制 schema(.bfbs);
  2. rbflatc.bfbs生成 Ruby 绑定,外层命名空间固定为ArrowFormat

这意味着本 gem 与 C++/Python/Java 各语言实现共享同一份 FlatBuffers 模式定义,任何对Message.fbs的演进(如新增body_compression字段)都可以经由这两个 fbs 文件同步到 Ruby 侧。

测试与基准测试

运行测试

与 README 的 Development 一节一致:

$ cd ruby/red-arrow-format $ bundle install $ bundle exec rake test

Rakefile 中task default: :test:test实际执行ruby test/run.rb(设置RUNNER_DEBUG=1时会加-v)。

测试的组织方式很有参考性:test/下按类型拆分为test-int8-array.rbtest-list-array.rbtest-timestamp-array.rb等 40 余个文件,公共逻辑集中在 test-reader.rb 的ReaderTests模块。其核心是roundtrip助手(第 19~52 行):用Arrow(red-arrow,即 C++ 实现)把数据写成真实 Arrow 文件,再交给被测的纯 Ruby reader 读回,逐值比对。这实际上构成了一个纯 Ruby 实现与 libarrow 之间的交叉验证套件——test_custom_metadata_field/test_custom_metadata_schema(第 54~90 行)还专门验证了 field 级与 schema 级自定义元数据的往返保真。

基准测试

benchmark/目录提供四个基准配置:file-reader.yamlfile-writer.yamlstreaming-reader.yamlstreaming-writer.yaml。Rakefile(第 44~64 行)通过benchmark-driver工具为每个 yaml 生成rake benchmark:<name>任务,例如:

$ bundle exec rake benchmark:file-reader

benchmark:benchmark任务会顺序跑完全部四项。

小结:适合谁使用

综合 README、gemspec 与源码结构,red-arrow-format的价值边界很清晰:

需求选择
在 Ruby 进程内零 C++ 依赖地解析/生产 Arrow IPC 字节流red-arrow-format(本模块)
过滤、连接、聚合等 Arrow 数据计算red-arrowruby/red-arrow
验证纯 Ruby 输出与 libarrow 的字节级兼容复用test/test-reader.rb的 roundtrip 模式

需要留意的实现前提:写入端固定产出MetadataVersion::V5元数据、小端字节序(schema.rb 中endianness硬编码为LITTLE);读取端则对 0.15.0 之前的旧格式保留了兼容路径。若你在 Ruby 服务中需要轻量消费由其他语言栈产出的.arrow文件(例如 README 示例中的/dev/shm/data.arrow场景),这正是ArrowFormat::FileReader的典型用武之地。

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

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

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

立即咨询