news 2026/9/24 7:16:27

Apache Arrow Ruby(Red Arrow)开发命名约定:Reader/Writer 与 Loader/Saver 双层 API 设计解析

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Arrow Ruby(Red Arrow)开发命名约定:Reader/Writer 与 Loader/Saver 双层 API 设计解析
  • 数据工程
  • 大数据
  • 序列化
  • 数据分析

【免费下载链接】arrow

Apache Arrow is a multi-language toolbox for accelerated data interchange and in-memory processing

项目地址:https://gitcode.com/gh_mirrors/arrow13/arrow
点击查看免费下载

本文以 ruby/red-arrow/doc/text/development.md 中定义的开发命名约定为核心,系统讲解 Red Arrow(Apache Arrow 的 Ruby 绑定)中Reader/WriterLoader/Saver两类 API 的职责划分、设计动机与源码实现。读完本文,你将理解为什么有的类需要你手动打开 IO 流、有的类只需传一个路径,并能依据实际场景在两类 API 之间做出正确选择,同时掌握Arrow::Table.load/Table#save底层完整的格式分发与选项解析链路。

一、命名约定总览:两种互补的 API 风格

Red Arrow 的 IO 相关类遵循一套清晰的命名约定,其核心规则只有两条,但贯穿了整个读写体系:

命名后缀构造前提定位示例类
Reader需要一个已打开的 IO 流(input stream)底层、可组合的读取组件RecordBatchFileReaderRecordBatchStreamReader
Writer需要一个已打开的 IO 流(output stream)底层、可组合的写入组件RecordBatchFileWriterRecordBatchStreamWriter
Loader只需要一个路径(path)面向用户的便捷门面,内部打开路径后用 Reader 读取Arrow::TableLoaderArrow::CSVLoader
Saver只需要一个路径(path)面向用户的便捷门面,内部打开路径后用 Writer 写入Arrow::TableSaver

原文档对这条约定的定义非常精炼:

  • Reader 和 Writer 需要一个已打开的 IO 流("Reader and Writer require an opened IO stream")。
  • Loader 和 Saver 只需要一个路径,是便捷类("Loader and Saver require a path. They are convenient classes")。
  • Loader 负责打开路径,并借助 Reader 读取数据("Loader opens the path and reads data by Reader")。
  • Saver 负责打开路径,并借助 Writer 写入数据("Writer opens the path and writes data by Writer";此处原文笔误为 "Writer",结合 table-saver.rb 的实现,实际指 Saver 委托 Writer 完成写入)。

也就是说,Reader/Writer是"流驱动"的底层构件,而Loader/Saver是"路径驱动"的高层门面——两者不是平行的两套实现,而是分层协作:上层负责路径解析与 IO 流打开,下层负责真正的序列化/反序列化。

二、Reader/Writer:面向已打开 IO 流的底层构件

2.1 类族与文件位置

Red Arrow 中的核心 Reader/Writer 类包括:

  • Arrow::RecordBatchFileReader——读取 Arrow文件格式(random access,带 footer/元数据,可随机读取任意 record batch),见 record-batch-file-reader.rb;
  • Arrow::RecordBatchStreamReader——读取 Arrow流格式(streaming,只能顺序读取),见 record-batch-stream-reader.rb;
  • Arrow::RecordBatchFileWriter/Arrow::RecordBatchStreamWriter——对应两种格式的写入器,在 Ruby 侧没有手工包装文件,而是由 GObject Introspection 在运行时从 c_glib/arrow-glib/writer.h 自动生成绑定(见下文 2.3)。

从源码结构看,record-batch-file-reader.rb 只对自动生成的类补充了Enumerable混入与each遍历逻辑:

class RecordBatchFileReader include Enumerable def each return to_enum(__method__) {n_record_batches} unless block_given? n_record_batches.times do |i| yield(get_record_batch(i)) end end end

这体现了 Red Arrow 的整体架构:绝大多数类是 C 层(Apache Arrow C++ → Apache Arrow GLib)通过 GObject Introspection 自动生成的,Ruby 源文件只负责补充 Ruby 惯用的便利方法。加载全部绑定与扩展库的入口在 loader.rb,它继承自GObjectIntrospection::Loader,在post_load中通过require_libraries逐个加载所有手工扩展、再require_extension_library加载arrow.so

2.2 用法:先开流,再构造 Reader

由于 Reader 需要"已打开的 IO 流",典型用法是先创建Arrow::MemoryMappedInputStream等流对象,再传入 Reader 构造器。参考 example/read-file.rb:

require "arrow" Arrow::MemoryMappedInputStream.open("/tmp/file.arrow") do |input| reader = Arrow::RecordBatchFileReader.new(input) fields = reader.schema.fields reader.each_with_index do |record_batch, i| puts("=" * 40) puts("record-batch[#{i}]:") fields.each do |field| field_name = field.name values = record_batch.collect do |record| record[field_name] end puts(" #{field_name}: #{values.inspect}") end end end

读取流格式的对应写法见 example/read-stream.rb,仅把RecordBatchFileReader换成RecordBatchStreamReader。二者的差异在于文件格式可随机访问(get_record_batch(i)按索引取),流格式只能顺序each消费。

2.3 从 Ruby 普通 IO 到 Arrow 流的桥接

一个值得注意的细节是:Reader/Writer 需要的是Arrow 的InputStream/OutputStream抽象,而不是 Ruby 的IO对象。Red Arrow 提供了多级桥接:

  • 文件路径 →Arrow::MemoryMappedInputStream(读)、Arrow::FileOutputStream(写,第二个参数false表示不追加,见 example/write-file.rb);
  • 内存缓冲 →Arrow::BufferInputStream/Arrow::BufferOutputStream
  • Ruby 的IO/StringIOGio::RubyInputStreamArrow::GIOInputStream(读),Gio::RubyOutputStreamArrow::GIOOutputStream(写)。

管道(pipe)场景是 "Reader 需要已打开 IO 流" 的最典型例证。参考 example/write-pipe.rb 与 example/read-pipe.rb:父进程把Table写入管道一端,子进程从管道另一端读取。因为管道没有路径,只有流,所以这里必须走 Reader/Writer 而非 Loader/Saver:

# 写入端(write-pipe.rb) table = Arrow::Table.new(a: [1, 2, 3], b: ["a", "b", "c"]) IO.pipe do |input, output| pid = spawn(RbConfig.ruby, File.join(__dir__, "read-pipe.rb"), in: input) input.close output.singleton_class.__send__(:undef_method, :seek) Gio::RubyOutputStream.open(output) do |gio_output| Arrow::GIOOutputStream.open(gio_output) do |arrow_output| Arrow::RecordBatchStreamWriter.open(arrow_output, table.schema) do |writer| writer.write_table(table) end end end output.close Process.waitpid(pid) end
# 读取端(read-pipe.rb) Gio::RubyInputStream.open($stdin) do |gio_input| Arrow::GIOInputStream.open(gio_input) do |arrow_input| reader = Arrow::RecordBatchStreamReader.new(arrow_input) p reader.read_all end end

这里还解释了为什么命名约定要求 Reader/Writer 只认流:流抽象让同一个读取/写入组件可以无差别地服务文件、内存、管道、网络 socket 乃至任意 Ruby IO,路径只是流的一种来源。

三、Loader/Saver:面向路径的便捷门面

3.1 一句话的职责定义与入口

Loader/Saver 的定位是"便捷类":调用方无需关心如何打开流、无需了解底层 Reader/Writer 的格式差异,只需给出路径。其职责正是原文档所写——Loader 打开路径后交给 Reader 读,Saver 打开路径后交给 Writer 写。

最常用的两个门面 API 定义在 table.rb:

class << self def load(path, options={}) # 第 29 行 TableLoader.load(path, options) end end # 实例方法(第 445 行附近) def save(output, options={}) # 委托给 TableSaver end

也就是说,Arrow::Table.load("/path/to/data.arrow")table.save("/path/to/out.arrow")这两条 README 里最常用的用法(见 ruby/red-arrow/README.md),底层分别由TableLoaderTableSaver完成"打开路径 → 构造流 → 委托 Reader/Writer"的完整流程。

3.2 TableLoader:来源识别与格式分发

table-loader.rb 的加载流程分三层:

第一层:识别输入来源(path / URI / Buffer / directory)load方法根据输入类型确定候选加载方法:

def load if @input.is_a?(URI) custom_load_method_candidates = [] if @input.scheme custom_load_method_candidates << "load_from_uri_#{@input.scheme}" end custom_load_method_candidates << "load_from_uri" elsif @input.is_a?(String) and ::File.directory?(@input) custom_load_method_candidates = ["load_from_directory"] else custom_load_method_candidates = ["load_from_file"] end # ...按候选方法逐一 dispatch,找不到则抛出列出可用来源的 ArgumentError end

基类实现了load_from_fileload_from_uri_http/https(统一走load_by_reader);load_from_directory则是预留的扩展钩子——从源码结构看,基类并未实现该方法,若传入目录路径,会抛出ArgumentError并列出所有可用的load_from_*来源。

第二层:按格式分发(load_as_*)load_by_reader读取@options[:format],动态调用对应的load_as_#{format}私有方法,基类支持以下格式:

format 选项对应方法底层 Reader备注
:arrow(默认)load_as_arrow先尝试RecordBatchFileReader,失败则回退RecordBatchStreamReader扩展名无法识别时的兜底
:arrow_fileload_as_arrow_fileRecordBatchFileReader自 1.0.0 起
:arrows/:arrow_streamingload_as_arrowsRecordBatchStreamReader自 7.0.0 起
:orcload_as_orcORCFileReader(若可用)支持:field_indexes选项
:csvload_as_csvCSVLoader见第四节
:tsvload_as_tsvCSVLoaderdelimiter: "\t"
:featherload_as_featherFeatherFileReader
:jsonload_as_jsonJSONReader支持把选项映射到JSONReadOptions

其中load_as_batch/load_as_stream已废弃的旧格式名(源码中以@deprecated Use format: :arrow_file ...标注),分发表在构建错误信息时会把这些废弃格式从可用列表中剔除。

第三层:打开输入流并委托 Readeropen_input_stream按来源类型打开流,并统一交给load_raw

def open_input_stream case @input when Buffer yield(BufferInputStream.new(@input)) when URI @input.open do |ruby_input| # :stream 格式用 Gio::RubyInputStream + GIOInputStream; # 其他格式先整段读入 Buffer 再用 BufferInputStream(规避 GVL 问题) end else yield(MemoryMappedInputStream.new(@input)) end end def load_raw(input, reader) schema = reader.schema record_batches = [] reader.each do |record_batch| record_batches << record_batch end table = Table.new(schema, record_batches) table.refer_input(input) # 让 Table 引用输入流,保证底层数据存活 table end

load_raw正是原文档所述 "Loader opens the path and reads data by Reader" 的直接代码体现:先打开路径对应的流,再构造 Reader 遍历 record batch,最后组装成Arrow::Tabletable.refer_input(input)是一个容易被忽略但很关键的内存语义:Table 的数据可能直接引用输入流的底层内存(尤其 MemoryMappedInputStream),必须让 Table 持有流的引用,防止流被 GC 回收导致数据失效

3.3 扩展名自动识别与压缩选项

TableLoader/TableSaver另一个"便捷"之处在于自动从路径扩展名推断格式与压缩方式fill_options利用Arrow::PathExtension解析输入/输出路径:

def fill_options if @options[:format] and @options.key?(:compression) return end case @input when Buffer info = {} when URI extension = PathExtension.new(@input.path) info = extension.extract else extension = PathExtension.new(@input) info = extension.extract end format = info[:format] @options = @options.dup if format @options[:format] ||= format.to_sym else @options[:format] ||= :arrow # 无法识别时默认按 Arrow 文件格式 end unless @options.key?(:compression) @options[:compression] = info[:compression] end end

因此Table.load("data.csv.gz")无需显式传参——扩展名会解析出format: :csvcompression: :gzip,Loader 自动用CompressedInputStream解压后再交给CSVLoader。显式传入的:format/:compression选项优先级更高,会跳过自动推断。

四、Saver:打开路径后委托 Writer

table-saver.rb 与 Loader 完全对称:TableSaver.save(table, output, options={})先识别输出目标(URI 或文件路径),再按format分发到save_as_*方法。save_raw是核心委托点:

def save_raw(writer_class) open_output_stream do |output| writer_class.open(output, @table.schema) do |writer| writer.write_table(@table) end end end def save_as_arrow_file save_raw(RecordBatchFileWriter) # Arrow 文件格式 end def save_as_arrows save_raw(RecordBatchStreamWriter) # Arrow 流格式 end

这正是原文档 "Saver opens the path and writes data by Writer" 的代码化表达:open_output_stream负责打开(文件 →FileOutputStream.open(output, false),Buffer →BufferOutputStream),并按:compression选项决定是否包一层CompressedOutputStream(用Codec.new(compression)创建编解码器),最后把已打开的流和 schema 交给对应 Writer。写入选项同样由fill_options从扩展名自动推断(如.arrow:arrow_file)。

Saver 支持的格式与 Loader 对应:save_as_arrow_file(默认,等价于旧save_as_batch)、save_as_arrows(等价于旧save_as_stream)、save_as_csvcsv_save,先写 schema 字段名作表头,再逐行写raw_records)、save_as_tsvcol_sep: "\t")、save_as_feather(通过FeatherWriteProperties把选项映射为写属性后调用table.write_as_feather)。写入完成后save返回@table本身,便于链式调用。

五、CSVLoader:Loader 约定的第二个实例

arrow/csv-loader.rb 展示了 Loader 约定的另一种形态——输入既可以是路径(Pathname或以.csv结尾的字符串),也可以是内存中的 CSV 字符串数据

def load case @path_or_data when Pathname load_from_path(@path_or_data.to_path) when /\A.+\.csv\z/i load_from_path(@path_or_data) else load_data(@path_or_data) # 直接把字符串当作 CSV 数据解析 end end

CSVLoader内部采用"两条腿走路"的策略:优先尝试用Arrow::CSVReader(GLib 层高性能读取器)配合CSVReadOptions解析;若因选项不支持等原因失败(rescue Arrow::Error::Invalid, Gio::Error),则回退到 Ruby 标准库CSV逐行解析并构造Arrow::Table。它支持:headers(布尔值、列名数组或字符串三种形态)、:column_types:schema:encoding(用Gio::CharsetConverter做编码转换)、:delimiter(内部转成:col_sep)、:compression(用Codec+CompressedInputStream解压)等选项。

更值得一提的是它的列类型自动检测detect_robust_converters扫描全部行,按列归纳出:boolean/:integer/:float/:time/:date_time/:date/:string中的一种类型,若某列出现不一致则降级为:string,并为每列生成对应的selective_converter转换器(内置BOOLEAN_CONVERTER识别"true"/"false"ISO8601_CONVERTER识别 ISO 8601 时间字符串)。这保证了 CSV 加载后能得到类型正确的 Arrow 列。

六、如何选择:两条 API 路径的使用准则

结合以上源码分析,可以把原文档的命名约定转化为可操作的选择准则:

你的输入/输出形态推荐 API理由
有文件路径,希望一行搞定、自动识别格式Arrow::Table.load(path)/table.save(path)Loader/Saver 会完成开流、格式推断、压缩推断,代码量最小
需要从 HTTP(S) URI 或 Buffer 读取Arrow::Table.load(uri_or_buffer)TableLoader原生支持URIBuffer来源
需要控制流的生命周期、复用已打开的流直接构造RecordBatchFileReader/RecordBatchStreamReaderReader/Writer 只依赖流抽象,不关心流从哪来
管道、socket、任意 Ruby IO 上的实时读写Gio::RubyInputStream+GIOInputStream+ Reader(或对称的 Writer 链)没有路径,只能用流;见 example/read-pipe.rb
需要自定义加载来源(如目录、私有存储)继承TableLoader实现load_from_*/load_as_*方法分发机制按命名约定自动发现可用方法

一条实用经验:默认优先使用 Loader/Saver——它们不仅省去开流样板代码,还能从扩展名自动推导格式与压缩方式,并且对 URI、Buffer、文件三种来源统一了入口;只有当你需要对流的打开方式、生命周期或底层 Reader 行为做精细控制时,才下沉到 Reader/Writer 层

七、小结

Red Arrow 通过Reader/WriterLoader/Saver的命名约定,构建了一个清晰的读写分层:

  • Reader/Writer是流驱动的底层构件,可组合、可复用,服务于文件、内存、管道等一切"流"形态的数据源;
  • Loader/Saver是路径驱动的高层门面,内部完成"打开路径 → 构造流 → 委托 Reader/Writer"的完整链路,并额外提供格式与压缩的自动推断能力(见 table-loader.rb 与 table-saver.rb)。

理解这套约定,不仅能让你在阅读 Red Arrow 源码(从 example/ 下的读写示例到 lib 目录各组件)时快速定位类的职责,也能在开发中根据"有路径还是有流"瞬间选出正确的 API,并在此基础上通过实现load_from_*/save_as_*方法扩展自定义格式,这正是该命名约定作为开发指南的价值所在。

  • 数据工程
  • 大数据
  • 序列化
  • 数据分析

【免费下载链接】arrow

Apache Arrow is a multi-language toolbox for accelerated data interchange and in-memory processing

项目地址:https://gitcode.com/gh_mirrors/arrow13/arrow
点击查看免费下载

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/24 7:14:31

Modbus转MQTT数据采集全流程:从RS485到云端实战指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/24 7:06:04

在 CI/CD 流水线中使用 Regal 对 Rego 策略进行代码检查

后端认证鉴权云原生 【免费下载链接】opa Open Policy Agent (OPA) is an open source, general-purpose policy engine. 项目地址&#xff1a; https://gitcode.com/gh_mirrors/op/opa 点击查看 免费下载 Regal 是 Open Policy Agent 生态中专门用于 Rego 策略代码的 linter …

作者头像 李华
网站建设 2026/9/24 7:05:18

Qt工程打包为exe全流程:从windeployqt到安装包制作

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/24 6:58:06

代理IP团队化管理与选型实战:从API批量配IP到子账户权限

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/24 6:50:35

Qwen3.8-Flash 限时免费:9 月 30 日前在 Qoder 零 Credits 畅用

Qwen3.8-Flash 限时免费&#xff1a;9 月 30 日前在 Qoder 零 Credits 畅用 9 月 18 日&#xff0c;阿里 Agentic 编码平台 Qoder 官方宣布&#xff1a;Qwen3.8-Flash 模型限时免费开放&#xff0c;活动期为 2026 年 9 月 18 日 10:00 至 9 月 30 日 23:59:59。活动期间该模型…

作者头像 李华