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 行):FileReader、FileWriter、StreamingReader、StreamingWriter以及版本常量,因此 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_footer、read_schema、read_dictionaries,因此拿到reader对象时 schema 与字典都已就绪; - 随机访问:除
each顺序迭代外,FileReader提供n_record_batches(批数量)与read(i)(按块偏移直接读取第 i 个 RecordBatch,第 57~74 行),适合需要随机跳批的场景; - 暴露的属性:
reader.schema与reader.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 行):
- 先读 4 字节:若为 continuation 串(
0xFFFFFFFF)则跳过; - 若该值不是负数(即不是合法的 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 字节边界的 FlatBuffersMessage(version = MetadataVersion::V5)。Footer 同样固定写为V5(file-writer.rb 第 42 行),并记录每个 Block 的offset、meta_data_length、body_length,这正是读取端read_block的依据。
字典的增量写入
write_dictionary(streaming-writer.rb 第 115~146 行)用一个@written_dictionary_offsets表跟踪每个字典 ID 已写到的偏移:新批次的字典值只有超出已写部分才会被切出来(data.slice(written_offset - current_base_offset))并以delta = true的DictionaryBatch写出。这与读取端read_dictionaries的“delta 必须跟随非 delta”校验是一对镜像实现,保证多批数据间共享字典时不会重复序列化。
流式形态:StreamingReader / StreamingWriter
除 File 形态外,lib/arrow-format.rb 还导出StreamingReader与StreamingWriter,对应 Arrow IPC 的流式布局(无 Footer,以EOS标记结束)。
streaming-reader.rb 的读取是按需拉取的:StreamingReader#each(第 48~59 行)通过StreamingPullReader#next_required_size询问“我还需要多少字节才能交出下一条 RecordBatch”,再按该大小从输入读取并consume。因此它既支持整块输入(File会走内存映射,String走IO::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.rb、record_batch.rb、schema.rb、body_compression.rb……)并非手写,而是由 Rakefile 中的flat_buffers:generate任务生成:
- 用
flatc --binary --schema ...把仓库根目录的 format/File.fbs 与 format/Message.fbs 编译为二进制 schema(.bfbs); - 用
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 testRakefile 中task default: :test,:test实际执行ruby test/run.rb(设置RUNNER_DEBUG=1时会加-v)。
测试的组织方式很有参考性:test/下按类型拆分为test-int8-array.rb、test-list-array.rb、test-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.yaml、file-writer.yaml、streaming-reader.yaml、streaming-writer.yaml。Rakefile(第 44~64 行)通过benchmark-driver工具为每个 yaml 生成rake benchmark:<name>任务,例如:
$ bundle exec rake benchmark:file-readerbenchmark:benchmark任务会顺序跑完全部四项。
小结:适合谁使用
综合 README、gemspec 与源码结构,red-arrow-format的价值边界很清晰:
| 需求 | 选择 |
|---|---|
| 在 Ruby 进程内零 C++ 依赖地解析/生产 Arrow IPC 字节流 | red-arrow-format(本模块) |
| 过滤、连接、聚合等 Arrow 数据计算 | red-arrow(ruby/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),仅供参考