Apache Arrow GLib 官方示例精读:C / Lua / Vala 多语言列式数据读写实战指南
【免费下载链接】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
Apache Arrow GLib 是 Apache Arrow 的 C 语言封装库,它通过 GObject Introspection 机制让 Lua、Ruby、Vala 等多种语言可以在运行时或编译期获得与 C 一致的数据读写能力。本文以 c_glib/example/README.md 为骨架,逐一对官方示例代码进行精读:你将掌握 Arrow IPC 两种序列化格式(文件格式与流格式)的读写 API 调用链、RecordBatch 的构建与打印方法,并了解如何基于 Arrow GLib 在 C、Lua、Vala 三种语言中落地同样的列式数据处理逻辑,最后还会深入到扩展类型与网络传输两个高级示例的实现细节。
一、示例目录结构与定位
c_glib/example目录是 Arrow GLib 官方示例代码的集合。README 明确指出:C 语言的示例代码直接位于本目录,而其他语言绑定的示例代码位于各子目录中,例如 Lua 示例位于lua/子目录。
当前仓库中该目录的实际结构如下:
c_glib/example/ ├── README.md # C 示例说明文档 ├── read-file.c # 从文件格式中读取 Arrow 数据 ├── read-stream.c # 从流格式中读取 Arrow 数据 ├── extension-type.c # 自定义 UUID 扩展类型完整示例 ├── send-network.c # 通过网络发送 RecordBatch(客户端) ├── receive-network.c # 通过网络接收 RecordBatch(服务端) ├── lua/ # Lua 绑定示例(基于 LGI) │ ├── README.md │ ├── write-file.lua / read-file.lua │ └── write-stream.lua / read-stream.lua └── vala/ # Vala 绑定示例 ├── README.md ├── write-file.vala / read-file.vala └── write-stream.vala / read-stream.vala需要留意的是,README 的 C 示例清单中列出了build.c、write-file.c、write-stream.c,但这三项在源文件中处于 HTML 注释块(<!--- ... -->)内,且当前仓库的c_glib/example目录下并不存在对应的.c文件——这意味着 README 中处于注释状态的示例尚未落盘,实际可编译运行的 C 示例是read-file.c、read-stream.c、extension-type.c、send-network.c与receive-network.c五个文件。相比之下,Lua 与 Vala 子目录中的写文件、读文件、写流、读流四个示例均已完整提供。
所有 C 示例都通过统一的头文件入口引入 API:
#include <arrow-glib/arrow-glib.h>该头文件聚合了 Arrow GLib 的全部 C 接口(数组、数据类型、RecordBatch、读写器、内存映射输入流等),是使用 Arrow GLib 编写 C 程序的唯一前置依赖。
二、C 示例一:以文件格式读取 Arrow 数据(read-file.c)
read-file.c 演示了Arrow IPC 文件格式(file format)的读取流程。所谓文件格式,是指数据带有可随机访问的页脚(footer)元数据、可反复读取的持久化文件,典型应用场景是磁盘上的.arrow文件。
2.1 打开输入流
程序首先通过garrow_memory_mapped_input_stream_new()以内存映射方式打开文件,默认路径为/tmp/batch.arrow,也支持通过命令行参数argv[1]覆盖:
const char *input_path = "/tmp/batch.arrow"; GArrowMemoryMappedInputStream *input; if (argc > 1) input_path = argv[1]; input = garrow_memory_mapped_input_stream_new(input_path, &error);内存映射输入流使文件内容被映射到进程地址空间,读取时无需显式拷贝,适合需要随机访问的场景——这正是文件格式读取器的前提。
2.2 创建文件读取器
文件格式带有页脚,因此读取器要求输入流是可随机访问(seekable)的。示例通过类型宏GARROW_SEEKABLE_INPUT_STREAM(input)将内存映射流向上转型为GArrowSeekableInputStream,再交给文件读取器构造函数:
reader = garrow_record_batch_file_reader_new(GARROW_SEEKABLE_INPUT_STREAM(input), &error);这一点与 reader.h 中garrow_record_batch_file_reader_new(GArrowSeekableInputStream *file, GError **error)的函数签名严格对应:文件格式读取器只接受可寻址流。
2.3 按索引随机读取 RecordBatch
文件格式支持随机访问,因此示例用garrow_record_batch_file_reader_get_n_record_batches()拿到批次总数后,直接按下标循环读取任意一个 RecordBatch:
n = garrow_record_batch_file_reader_get_n_record_batches(reader); for (i = 0; i < n; i++) { record_batch = garrow_record_batch_file_reader_read_record_batch(reader, i, &error); ... print_record_batch(record_batch); g_object_unref(record_batch); }每个读取到的GArrowRecordBatch是GObject对象,用完后必须调用g_object_unref()释放引用,这是 GLib 内存管理的标准约定。
2.4 逐列打印 RecordBatch 内容
示例配套了两个打印函数,构成一条完整的"读出来、看得见"的验证链路:
print_record_batch():通过garrow_record_batch_get_n_columns()、garrow_record_batch_get_column_name()与garrow_record_batch_get_column_data()遍历每一列,输出列序号、列名与数据;print_array():通过garrow_array_get_value_type()获取数组的实际数据类型,再用宏ARRAY_CASE生成针对uint8~double共 10 种数值类型的switch分支,逐一调用对应类型数组的garrow_*_array_get_value()取值并格式化打印。
宏展开的核心是 GObject 的类型向下转型,例如:
case GARROW_TYPE_INT32: { GArrowInt32Array *real_array; real_array = GARROW_INT32_ARRAY(array); for (i = 0; i < n; i++) { g_print("%" G_GINT32_FORMAT, garrow_int32_array_get_value(real_array, i)); } } break;这展示了 Arrow GLib 的类型体系:运行时先用GArrowType枚举判别,再用GARROW_XXX_ARRAY()宏做类型安全的向下转换,最后调用具体类型的取值函数。
三、C 示例二:以流格式读取 Arrow 数据(read-stream.c)
read-stream.c 演示的是Arrow IPC 流格式(stream format)的读取。与文件格式不同,流格式是顺序、单向的消息序列:没有页脚、不支持随机跳转,只能从头到尾依次消费,天然适配管道、Socket、实时传输等场景。
3.1 创建流读取器
流读取器只需要顺序输入流,因此直接传入GARROW_INPUT_STREAM(input):
stream_reader = garrow_record_batch_stream_reader_new(GARROW_INPUT_STREAM(input), &error); reader = GARROW_RECORD_BATCH_READER(stream_reader);这里出现了第二个重要的继承关系:GArrowRecordBatchStreamReader是GArrowRecordBatchReader的子类,示例通过GARROW_RECORD_BATCH_READER(stream_reader)将其向上转型为通用读取器接口。
3.2 循环读取直到流结束
流格式没有"批次总数"的概念,示例采用读到 NULL 即结束的惯用循环:
while (TRUE) { record_batch = garrow_record_batch_reader_read_next(reader, &error); if (error) { /* 出错处理 */ } if (!record_batch) { break; /* 流已读完 */ } print_record_batch(record_batch); g_object_unref(record_batch); }garrow_record_batch_reader_read_next()返回NULL且error为空时表示正常到达流末尾;返回NULL且error非空时表示读取失败。这一"空指针 + GError"的双重约定,是 GLib 风格 I/O 的典型写法。
3.3 文件格式与流格式的选型对照
结合两个示例可以总结出选型要点:
| 维度 | 文件格式(read-file.c) | 流格式(read-stream.c) |
|---|---|---|
| 底层载体 | 带页脚元数据的随机访问文件 | 顺序消息序列 |
| 输入流要求 | 可寻址流(GArrowSeekableInputStream) | 仅需顺序流(GArrowInputStream) |
| 读取方式 | 按索引随机读取(read_record_batch(i)) | 顺序读取直到 NULL(read_next()) |
| 典型场景 | 磁盘持久化文件、反复读取 | 管道、网络传输、实时数据 |
四、C 示例三:自定义 UUID 扩展类型(extension-type.c)
extension-type.c 是 Arrow 扩展类型(Extension Type)机制的完整演示,它定义了一个以 16 字节fixed-size-binary为存储类型、名为uuid的自定义扩展类型,并完成注册、构建、序列化、反序列化、注销的全生命周期。
4.1 定义扩展数据类型与扩展数组
示例使用 GObject 的类型宏体系定义了两个类:
G_DECLARE_DERIVABLE_TYPE(ExampleUUIDArray, example_uuid_array, EXAMPLE, UUID_ARRAY, GArrowExtensionArray) G_DEFINE_TYPE(ExampleUUIDArray, example_uuid_array, GARROW_TYPE_EXTENSION_ARRAY) G_DECLARE_DERIVABLE_TYPE(ExampleUUIDDataType, example_uuid_data_type, EXAMPLE, UUID_DATA_TYPE, GArrowExtensionDataType) G_DEFINE_TYPE(ExampleUUIDDataType, example_uuid_data_type, GARROW_TYPE_EXTENSION_DATA_TYPE)扩展数据类型的父类是GArrowExtensionDataType,扩展数组的父类是GArrowExtensionArray。类型实现的关键在于覆写父类 vtable 中的四个虚函数:
extension_klass->get_extension_name = example_uuid_data_type_get_extension_name; extension_klass->equal = example_uuid_data_type_equal; extension_klass->deserialize = example_uuid_data_type_deserialize; extension_klass->serialize = example_uuid_data_type_serialize; extension_klass->get_array_gtype = example_uuid_data_type_get_array_gtype;get_extension_name():返回扩展名"uuid";serialize()/deserialize():扩展类型在序列化时必须附带一段标识数据,示例用常量字符串"uuid-serialized"作为序列化标识,反序列化时先校验标识是否匹配,再校验存储类型是否仍为fixed-size-binary(16);get_array_gtype():将扩展数据类型与对应的扩展数组类关联。
4.2 注册与构建数据
main()中首先取得全局扩展类型注册表并注册自定义类型:
GArrowExtensionDataTypeRegistry *registry = garrow_extension_data_type_registry_default(); ... garrow_extension_data_type_registry_register(registry, extension_data_type, &error);随后从扩展类型中取出存储类型,用GArrowFixedSizeBinaryArrayBuilder构造存储数组(两个 16 字节字符串加一个 NULL),再用garrow_extension_data_type_wrap_array()把存储数组"包装"成扩展数组,实现"存储数据 + 扩展语义"的解耦:
GArrowExtensionArray *extension_array = garrow_extension_data_type_wrap_array(extension_data_type, storage);4.3 序列化—反序列化闭环验证
示例随后把扩展数组装入 RecordBatch,用GArrowRecordBatchStreamWriter写入内存 Buffer(GArrowBufferOutputStream+GArrowResizableBuffer),再立刻用GArrowRecordBatchStreamReader读回,验证扩展类型在 IPC 往返后依然完整保留:
gchar *record_batch_content = garrow_record_batch_to_string(record_batch, &error); g_print("record batch:\n%s\n", record_batch_content); ... g_print("array: %s\n", G_OBJECT_TYPE_NAME(deserialized_array));最后在exit:标签处通过garrow_extension_data_type_registry_unregister()注销该扩展类型,体现了"注册—使用—注销"的完整资源管理闭环。该示例是理解 Arrow 扩展类型机制(自定义类型如何在 IPC 层无损传输)的最佳起点。
五、C 示例四:基于 Socket 的网络传输(send-network.c / receive-network.c)
send-network.c与receive-network.c是一对配套示例,演示如何通过 GLib 的GSocket体系在网络上实时传输 Arrow RecordBatch,本质上是把上一节的"流格式"应用到网络字节流上。
5.1 客户端:构建数据并写入 Socket
send-network.c 的运行方式是send-network PORT(例如send-network 2929)。客户端流程如下:
- 用
g_socket_client_new()+g_inet_socket_address_new_from_string("127.0.0.1", port)建立 TCP 连接到本机指定端口; - 用
build_schema()构建包含boolean与int32两列的 Schema; - 把 Socket 连接包装为 Arrow 输出流,创建流写入器:
GArrowGIOOutputStream *output = garrow_gio_output_stream_new(g_io_stream_get_output_stream(G_IO_STREAM(connection))); GArrowRecordBatchStreamWriter *writer = garrow_record_batch_stream_writer_new(GARROW_OUTPUT_STREAM(output), schema, &error);- 通过
GArrowRecordBatchBuilder构建 3 行 × 2 列的 RecordBatch(其中 boolean 列第 2 行为 NULL、int32 列第 1 行为 NULL),循环写入 5 个批次:
for (i = 0; i < n_record_batches; i++) { GArrowRecordBatch *record_batch = build_record_batch(); garrow_record_batch_writer_write_record_batch(GARROW_RECORD_BATCH_WRITER(writer), record_batch, &error); ... }这里体现了garrow_record_batch_builder_get_column_builder()配合garrow_*_array_builder_append_values()的批量构建方式:values 数组与 is_valids 数组同时传入,即可表达 NULL 语义。
5.2 服务端:监听端口并解析流
receive-network.c 使用GThreadedSocketService创建多线程 Socket 服务,监听任意可用端口(g_socket_listener_add_any_inet_port()),在event信号中打印监听地址,在incoming信号回调中处理每个连接:
GArrowGIOInputStream *input = garrow_gio_input_stream_new(g_io_stream_get_input_stream(G_IO_STREAM(connection))); GArrowRecordBatchStreamReader *reader = garrow_record_batch_stream_reader_new(GARROW_INPUT_STREAM(input), &error); while (TRUE) { record_batch = garrow_record_batch_reader_read_next(GARROW_RECORD_BATCH_READER(reader), &error); if (!record_batch) break; print_record_batch(record_batch); ... }服务端在 Unix 平台还注册了SIGINT/SIGTERM信号处理器(g_unix_signal_add())实现优雅退出。这对示例完整演示了 Arrow 流格式在网络传输中的应用:只要底层是可靠的字节流,Arrow 数据就可以在进程、机器之间无缝流动。
六、Lua 示例:基于 LGI 的动态绑定(lua/ 子目录)
Lua 子目录 的示例运行在LGI(Lua GObject Introspection)之上——LGI 在运行时读取 Arrow GLib 导出的 GIR 元数据,动态生成 Lua 绑定,因此 Lua 侧不需要编译任何绑定代码。
6.1 安装 LGI
README 给出了 Debian/Ubuntu 上的安装方式:
$ sudo apt install -y luarocks $ sudo luarocks install lgi6.2 加载绑定
所有 Lua 示例的开头都是一致的:
local lgi = require 'lgi' local Arrow = lgi.ArrowArrow命名空间即对应 C 侧全部以GArrow为前缀的类,命名规则为去掉G前缀:GArrowRecordBatchFileReader→Arrow.RecordBatchFileReader。
6.3 写入示例(write-file.lua / write-stream.lua)
write-file.lua 与 write-stream.lua 结构完全对称,仅写入器不同:前者用Arrow.RecordBatchFileWriter写文件格式到/tmp/batch.arrow,后者用Arrow.RecordBatchStreamWriter写流格式到/tmp/stream.arrow。
写入流程的核心步骤:
local output = Arrow.FileOutputStream.new(output_path, false) local writer = Arrow.RecordBatchStreamWriter.new(output, schema) local record_batch = Arrow.RecordBatch.new(schema, 4, columns) writer:write_record_batch(record_batch) -- 用 slice 构建第二个批次 local sliced_columns = {} for i, column in pairs(columns) do sliced_columns[i] = column:slice(1, 3) end record_batch = Arrow.RecordBatch.new(schema, 3, sliced_columns) writer:write_record_batch(record_batch) writer:close() output:close()两个值得注意的细节:
- Schema 由 10 个字段(
uint8到double)组成,每个字段由Arrow.Field.new(name, data_type)创建; - 第二个批次通过数组的
slice(1, 3)方法对同一份列数据做零拷贝切片(偏移 1、长度 3),展示了一个文件里写入多个不同形状批次的能力。
6.4 读取示例(read-file.lua / read-stream.lua)
read-file.lua 与 C 版read-file.c逻辑一致:Arrow.MemoryMappedInputStream打开文件,Arrow.RecordBatchFileReader按索引读取;区别在于 Lua 侧用record_batch:get_value(k)泛型取值,不再需要 C 版的 switch 分派——这正是动态语言的便利之处:
local input = Arrow.MemoryMappedInputStream.new(input_path) local reader = Arrow.RecordBatchFileReader.new(input) for i = 0, reader:get_n_record_batches() - 1 do local record_batch = reader:read_record_batch(i) for j = 0, record_batch:get_n_columns() - 1 do local column_data = record_batch:get_column_data(j) for k = 0, record_batch:get_n_rows() - 1 do io.write(column_data:get_value(k)) end end end input:close()read-stream.lua则使用Arrow.RecordBatchStreamReader顺序读取,运行方式均为lua 脚本名.lua [输入文件],文件参数缺省时使用/tmp下的默认路径。
七、Vala 示例:静态类型绑定(vala/ 子目录)
Vala 子目录 的示例通过valac编译器直接消费 Arrow GLib 的 GIR 元数据,在编译期生成强类型绑定。README 给出的构建命令是:
$ valac --pkg arrow-glib --pkg posix XXX.vala--pkg arrow-glib引入 Arrow GLib 的 Vala API 包,--pkg posix提供Posix.EXIT_SUCCESS等常量。
read-file.vala 与 C 版read-file.c一一对应,但类型系统更强:input as GArrow.SeekableInputStream是显式向下转型(因为MemoryMappedInputStream需要转换为SeekableInputStream才能传给文件读取器),数组取值则仍需要像 C 一样按类型分支:
switch (array.get_value_type()) { case GArrow.Type.UINT8: var concrete_array = array as GArrow.UInt8Array; for (var i = 0; i < n; i++) { stdout.printf("%hhu", concrete_array.get_value(i)); } break; ... }错误处理使用 Vala 的try/catch语法,捕获GError:
try { input = new GArrow.MemoryMappedInputStream(input_path); } catch (Error error) { stderr.printf("failed to open file: %s\n", error.message); return Posix.EXIT_FAILURE; }write-file.vala/write-stream.vala/read-stream.vala与 Lua 版本逻辑相同,仅语法不同。值得说明的是:Vala 子目录 README 的示例清单中提到的build.vala在当前仓库中并不存在,实际提供的是写/读文件与写/读流四个文件。
八、示例背后的实现原理:Arrow GLib 的封装层次
贯穿全部示例的一条主线是 Arrow GLib 的分层封装架构:
- 最底层是 Apache Arrow C++:
cpp/src中实现了列式内存布局、IPC 序列化(ipc/)、计算内核等核心能力; - 中间层是 Arrow GLib(
c_glib/arrow-glib/):将 C++ 对象包装为 GObject 类,例如 reader.h 中的GArrowRecordBatchFileReader/GArrowRecordBatchStreamReader,并统一暴露 C API,同时导出 GIR 元数据; - 最上层是各语言绑定:Lua 通过 LGI 运行时动态绑定,Vala 通过 valac 编译期静态绑定,Ruby 则通过 red-arrow / gobject-introspection gem(见 c_glib/README.md)。
这种"一次实现、多语言复用"的设计,正是示例目录中同一套读写逻辑以 C、Lua、Vala 三种语言各写一遍的根本原因——它们最终都调用同一份 C++ 实现。
关于构建与运行前提:要编译运行这些 C 示例,需要先构建并安装 Arrow C++ 与 Arrow GLib。Arrow GLib 使用 Meson + Ninja 构建,开发构建还需要 GTK-Doc 与 GObject Introspection;官方推荐直接使用发行版软件包,详见 c_glib/README.md 的 Install 章节。
九、小结:从示例到实战的迁移路径
| 场景 | 推荐示例 | 关键 API |
|---|---|---|
读取磁盘上的.arrow文件 | read-file.c | garrow_record_batch_file_reader_new+read_record_batch |
| 顺序消费管道/网络流 | read-stream.c | garrow_record_batch_stream_reader_new+read_next |
| 自定义业务数据类型 | extension-type.c | garrow_extension_data_type_registry_register+wrap_array |
| 跨进程/跨机器实时传输 | send-network.c 与 receive-network.c | GArrowGIOOutputStream/GArrowGIOInputStream+ StreamWriter/Reader |
| 脚本语言快速原型 | lua/ 子目录 | LGI +Arrow.RecordBatch* |
| 强类型语言生产代码 | vala/ 子目录 | valac --pkg arrow-glib |
无论选择哪种语言入口,其背后都是同一套 Arrow 列式格式与 IPC 协议。把官方示例跑通一遍,你就同时掌握了 Arrow 数据读写的最基本单元(Array、Schema、RecordBatch)、两种序列化载体(文件与流)以及它们在不同语言中的调用姿势,足以直接迁移到自己的数据处理管线中。
【免费下载链接】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),仅供参考