- 数据工程
- 大数据
- 序列化
- 数据分析
【免费下载链接】arrow
Apache Arrow is a multi-language toolbox for accelerated data interchange and in-memory processing
本文以 ruby/red-arrow/doc/text/development.md 中定义的开发命名约定为核心,系统讲解 Red Arrow(Apache Arrow 的 Ruby 绑定)中Reader/Writer与Loader/Saver两类 API 的职责划分、设计动机与源码实现。读完本文,你将理解为什么有的类需要你手动打开 IO 流、有的类只需传一个路径,并能依据实际场景在两类 API 之间做出正确选择,同时掌握Arrow::Table.load/Table#save底层完整的格式分发与选项解析链路。
一、命名约定总览:两种互补的 API 风格
Red Arrow 的 IO 相关类遵循一套清晰的命名约定,其核心规则只有两条,但贯穿了整个读写体系:
| 命名后缀 | 构造前提 | 定位 | 示例类 |
|---|---|---|---|
| Reader | 需要一个已打开的 IO 流(input stream) | 底层、可组合的读取组件 | RecordBatchFileReader、RecordBatchStreamReader |
| Writer | 需要一个已打开的 IO 流(output stream) | 底层、可组合的写入组件 | RecordBatchFileWriter、RecordBatchStreamWriter |
| Loader | 只需要一个路径(path) | 面向用户的便捷门面,内部打开路径后用 Reader 读取 | Arrow::TableLoader、Arrow::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/StringIO→Gio::RubyInputStream→Arrow::GIOInputStream(读),Gio::RubyOutputStream→Arrow::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),底层分别由TableLoader与TableSaver完成"打开路径 → 构造流 → 委托 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_file与load_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_file | load_as_arrow_file | RecordBatchFileReader | 自 1.0.0 起 |
:arrows/:arrow_streaming | load_as_arrows | RecordBatchStreamReader | 自 7.0.0 起 |
:orc | load_as_orc | ORCFileReader(若可用) | 支持:field_indexes选项 |
:csv | load_as_csv | CSVLoader | 见第四节 |
:tsv | load_as_tsv | CSVLoader(delimiter: "\t") | |
:feather | load_as_feather | FeatherFileReader | |
:json | load_as_json | JSONReader | 支持把选项映射到JSONReadOptions |
其中load_as_batch/load_as_stream是已废弃的旧格式名(源码中以@deprecated Use format: :arrow_file ...标注),分发表在构建错误信息时会把这些废弃格式从可用列表中剔除。
第三层:打开输入流并委托 Reader。open_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 endload_raw正是原文档所述 "Loader opens the path and reads data by Reader" 的直接代码体现:先打开路径对应的流,再构造 Reader 遍历 record batch,最后组装成Arrow::Table。table.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: :csv、compression: :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_csv(csv_save,先写 schema 字段名作表头,再逐行写raw_records)、save_as_tsv(col_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 endCSVLoader内部采用"两条腿走路"的策略:优先尝试用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原生支持URI与Buffer来源 |
| 需要控制流的生命周期、复用已打开的流 | 直接构造RecordBatchFileReader/RecordBatchStreamReader等 | Reader/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/Writer与Loader/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
相关推荐
Red Arrow 开发指南:Reader/Writer 与 Loader/Saver 命名约定及其 IO 流设计解析
Red Arrow 开发指南:Reader/Writer 与 Loader/Saver 命名约定及其 IO 流设计解析 Red Arrow(Apache Arr
数据工程数据分析大数据五步让旧 Mac 再战五年:OpenCore Legacy Patcher 免费安装新版 macOS 完整指南
五步让旧 Mac 再战五年:OpenCore Legacy Patcher 免费安装新版 macOS 完整指南 你的旧 Mac 跑起来还流畅,系统却告诉你"无法
操作系统固件驱动开发Red Arrow:基于 GObject Introspection 的 Apache Arrow Ruby 绑定入门与实战
Red Arrow:基于 GObject Introspection 的 Apache Arrow Ruby 绑定入门与实战 Red Arrow 是 Apach
数据工程数据分析大数据
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考