ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Arrow C++ 文件 I/O 实战:以 C++ 读取与写入 IPC、CSV、Parquet 文件

Apache Arrow C++ 文件 I/O 实战:以 C++ 读取与写入 IPC、CSV、Parquet 文件 Apache Arrow C 文件 I/O 实战以 C 读取与写入 IPC、CSV、Parquet 文件【免费下载链接】arrowApache Arrow is a multi-language toolbox for accelerated data interchange and in-memory processing项目地址: https://gitcode.com/gh_mirrors/arrow12/arrow本文基于 Apache Arrow 官方 C 教程《Arrow File I/O》系统讲解如何用 Arrow C 库完成最常见的三种文件格式的读写Arrow IPC 文件RecordBatch、CSV 文件Table与 Parquet 文件Table。通过本文你将掌握io::ReadableFile/io::FileOutputStream的用法理解ipc::RecordBatchFileReader/ipc::RecordBatchWriter、csv::TableReader与parquet::arrow::FileReader的完整调用链并能在自己的应用里直接复用文中给出的可运行代码。本文对应的原始文档位于 docs/source/cpp/tutorials/io_tutorial.rst文中所有代码片段均取自仓库中的完整示例 cpp/examples/tutorial_examples/file_access_example.cc你可以直接对照阅读。前置条件在继续之前请确保你已具备一份可用的 Arrow 安装搭建方式参见 构建系统指南对 Arrow 基础数据结构的理解可先阅读 基础数据结构教程一个用于运行最终程序的目录——本示例会在运行目录下生成若干文件请提前做好准备。准备工作头文件、main() 与测试数据文件在动手编写文件 I/O 代码之前需要先补齐三块内容必要的头文件、把一切串联起来的main()以及供程序读取的测试文件。头文件Includes编写 C 代码前先引入必要的头文件。除了用于输出的iostream之外还需要为本文涉及的每种文件类型导入 Arrow 的 I/O 功能#include arrow/api.h #include arrow/csv/api.h #include arrow/io/api.h #include arrow/ipc/api.h #include parquet/arrow/reader.h #include parquet/arrow/writer.h #include iostream这些头文件分别对应Arrow 核心数据结构arrow/api.h、CSV 读写arrow/csv/api.h、底层 I/O 接口arrow/io/api.h、IPC 文件读写arrow/ipc/api.h以及 Parquet 的 Arrow 集成读写器parquet/arrow/reader.h与parquet/arrow/writer.h。注意 Parquet 相关头文件属于parquet命名空间在使用 Parquet 功能时需要额外链接parquet库。main() 与 RunMain()main()采用基础数据结构教程中介绍过的模式真正的逻辑放在RunMain()中main()只负责调用它并检查返回的Statusint main() { arrow::Status st RunMain(); if (!st.ok()) { std::cerr st std::endl; return 1; } return 0; }与之前一样main()与RunMain()成对出现。RunMain()返回arrow::StatusArrow 中所有可能出错的函数几乎都以Status或ResultT作为返回值这种设计让错误可以在调用链上一层层传播最终在main()中统一处理arrow::Status RunMain() { // 后续所有文件 I/O 逻辑都在这里展开 ... }生成用于读取的测试文件实际应用中你通常已经有自己的输入数据。但为了纯粹地演示 I/O 流程我们定义一个辅助函数GenInitialFile()先生成测试文件。这里用到了前一教程中的年/月/日数据先用Int8Builder/Int16Builder构建三个数组再用 schema 组装成一张Tablearrow::Status GenInitialFile() { // 构建两个 8 位整数数组和一个 16 位整数数组 —— 与基础 Arrow 示例一致 arrow::Int8Builder int8builder; int8_t days_raw[5] {1, 12, 17, 23, 28}; ARROW_RETURN_NOT_OK(int8builder.AppendValues(days_raw, 5)); std::shared_ptrarrow::Array days; ARROW_ASSIGN_OR_RAISE(days, int8builder.Finish()); int8_t months_raw[5] {1, 3, 5, 7, 1}; ARROW_RETURN_NOT_OK(int8builder.AppendValues(months_raw, 5)); std::shared_ptrarrow::Array months; ARROW_ASSIGN_OR_RAISE(months, int8builder.Finish()); arrow::Int16Builder int16builder; int16_t years_raw[5] {1990, 2000, 1995, 2000, 1995}; ARROW_RETURN_NOT_OK(int16builder.AppendValues(years_raw, 5)); std::shared_ptrarrow::Array years; ARROW_ASSIGN_OR_RAISE(years, int16builder.Finish()); // 将所有数组收集到一个 vector std::vectorstd::shared_ptrarrow::Array columns {days, months, years}; // 定义 schema 以初始化 Table std::shared_ptrarrow::Field field_day, field_month, field_year; std::shared_ptrarrow::Schema schema; field_day arrow::field(Day, arrow::int8()); field_month arrow::field(Month, arrow::int8()); field_year arrow::field(Year, arrow::int16()); schema arrow::schema({field_day, field_month, field_year}); // 用 schema 和数据创建 Table std::shared_ptrarrow::Table table; table arrow::Table::Make(schema, columns); // 以 IPC、CSV、Parquet 三种格式写出测试文件供后续示例使用 std::shared_ptrarrow::io::FileOutputStream outfile; ARROW_ASSIGN_OR_RAISE(outfile, arrow::io::FileOutputStream::Open(test_in.arrow)); ARROW_ASSIGN_OR_RAISE(std::shared_ptrarrow::ipc::RecordBatchWriter ipc_writer, arrow::ipc::MakeFileWriter(outfile, schema)); ARROW_RETURN_NOT_OK(ipc_writer-WriteTable(*table)); ARROW_RETURN_NOT_OK(ipc_writer-Close()); ARROW_ASSIGN_OR_RAISE(outfile, arrow::io::FileOutputStream::Open(test_in.csv)); ARROW_ASSIGN_OR_RAISE(auto csv_writer, arrow::csv::MakeCSVWriter(outfile, table-schema())); ARROW_RETURN_NOT_OK(csv_writer-WriteTable(*table)); ARROW_RETURN_NOT_OK(csv_writer-Close()); ARROW_ASSIGN_OR_RAISE(outfile, arrow::io::FileOutputStream::Open(test_in.parquet)); PARQUET_THROW_NOT_OK( parquet::arrow::WriteTable(*table, arrow::default_memory_pool(), outfile, 5)); return arrow::Status::OK(); }这段代码里出现了两个 Arrow 惯用宏值得说明ARROW_RETURN_NOT_OK(expr)执行表达式若返回非 OK 的Status则直接返回该Status实现错误的快速传播ARROW_ASSIGN_OR_RAISE(lhs, expr)执行返回ResultT的表达式成功时把值赋给lhs失败时返回错误Status。另外parquet::arrow::WriteTable与 Arrow 自身的 API 不同它抛出的异常使用PARQUET_THROW_NOT_OK来转换为异常抛出而非返回Status。要让后续代码正常工作必须在RunMain()的第一行调用GenInitialFile()来初始化环境// 用辅助函数生成每种格式的初始文件 —— 别担心本示例后面还会写一次 Table ARROW_RETURN_NOT_OK(GenInitialFile());读取与写入 Arrow IPC 文件IPC进程间通信文件是 Arrow 的原生序列化格式。整个流程分步来看先读后写。读取流程打开文件将文件绑定到ipc::RecordBatchFileReader将文件内容读取为RecordBatch。写入流程获取io::FileOutputStream从RecordBatch写入文件。打开文件io::ReadableFile要读取文件首先需要一种指向文件的句柄。在 Arrow 中这对应io::ReadableFile对象——和ArrayBuilder可以清空并构建新数组类似我们可以把它重新绑定到新文件上因此在后面的例子中会一直复用这个实例// 首先我们需要设置一个 ReadableFile 对象它让我们可以把 reader 指向磁盘上的正确数据。 // 在整个示例中我们会复用这个对象并把它重新绑定到多个文件上。 std::shared_ptrarrow::io::ReadableFile infile;单独的io::ReadableFile能做的不多——我们需要调用io::ReadableFile::Open把它真正绑定到一个文件。对本教程而言使用默认参数即可// 把 test_in.arrow 绑定到我们的文件指针 ARROW_ASSIGN_OR_RAISE(infile, arrow::io::ReadableFile::Open( test_in.arrow, arrow::default_memory_pool()));io::ReadableFile::Open的完整签名见 cpp/src/arrow/io/file.h为static Resultstd::shared_ptrReadableFile Open( const std::string path, MemoryPool* pool default_memory_pool());即第二个参数MemoryPool可以省略默认使用全局默认内存池因此示例中显式传入arrow::default_memory_pool()只是为了让内存分配更明确。从源码结构看ReadableFile继承自internal::RandomAccessFileConcurrencyWrapperReadableFile意味着它提供随机访问能力其内部DoReadAt/DoSeek等实现保证了按位置读取是线程安全的同时源码注释提示ReadableFile的读取是无缓冲的如果需要大量小尺寸读取建议在它之上再加一层缓冲如BufferedInputStream以获得更好性能。打开 Arrow 文件读取器io::ReadableFile过于通用不足以提供读取 Arrow 文件的全部功能。我们需要借助它得到ipc::RecordBatchFileReader对象——该对象封装了按正确格式解析 Arrow 文件所需的全部逻辑。通过ipc::RecordBatchFileReader::Open获取// 用库的 IPC 功能打开文件得到一个 reader 对象。 ARROW_ASSIGN_OR_RAISE(auto ipc_reader, arrow::ipc::RecordBatchFileReader::Open(infile));从源码看RecordBatchFileReader定义于 cpp/src/arrow/ipc/reader.h它专门面向文件格式带 footer 的随机访问文件与之相对的是面向流式数据的RecordBatchStreamReader。选择File系列意味着我们可以按索引随机读取任意一个 RecordBatch。将打开的 Arrow 文件读取为 RecordBatch读取 Arrow 文件必须使用RecordBatch因此先声明一个RecordBatch指针。Arrow 文件可能包含多个RecordBatch所以必须传入索引本文件只有一个因此传入 0// 借助 reader我们可以读取 RecordBatch。 // 注意这一点是 IPC 特有的其他格式我们关注 Table而这里使用的是 RecordBatch。 std::shared_ptrarrow::RecordBatch rbatch; ARROW_ASSIGN_OR_RAISE(rbatch, ipc_reader-ReadRecordBatch(0));准备 FileOutputStream输出时需要一个io::FileOutputStream。和io::ReadableFile一样我们会复用它打开方式与读取时类似// 和输入一样我们为输出文件获取一个对象。 std::shared_ptrarrow::io::FileOutputStream outfile; // 把它绑定到 test_out.arrow ARROW_ASSIGN_OR_RAISE(outfile, arrow::io::FileOutputStream::Open(test_out.arrow));io::FileOutputStream::Open的完整签名见 cpp/src/arrow/io/file.h为static Resultstd::shared_ptrFileOutputStream Open(const std::string path, bool append false);注意它的默认行为打开新文件时会把已存在的同路径文件截断为 0 字节即默认覆盖写若想追加到已有文件需传入append true。此外该类还提供基于文件描述符的Open(int fd)重载打开后文件描述符的所有权会转移给输出流Close()或析构时自动关闭。FileOutputStream实现自OutputStream接口其Write方法是线程安全的。从 RecordBatch 写入 Arrow 文件现在取回之前读到的RecordBatch用它和目标文件一起创建ipc::RecordBatchWriter。ipc::RecordBatchWriter需要两样东西目标文件我们刚创建的输出流RecordBatch的Schema以便后续继续写入相同格式的更多RecordBatch。Schema直接来自现有的RecordBatch// 用输出文件和 schema 设置 writer —— 一切在这里定义好准备发射。 ARROW_ASSIGN_OR_RAISE(std::shared_ptrarrow::ipc::RecordBatchWriter ipc_writer, arrow::ipc::MakeFileWriter(outfile, rbatch-schema()));然后调用ipc::RecordBatchWriter::WriteRecordBatch把RecordBatch写入文件// 写入 record batch。 ARROW_RETURN_NOT_OK(ipc_writer-WriteRecordBatch(*rbatch));对于 IPC 而言writer必须显式关闭因为它预期可能写入多个 batch需要在文件末尾写 footer 元数据// 对 IPC 来说writer 需要显式关闭。 ARROW_RETURN_NOT_OK(ipc_writer-Close());至此我们已经完成了一个 IPC 文件的读写闭环读取与写入 CSV 文件CSV 的处理流程同样是先读后写读取流程打开文件准备Table用csv::TableReader读取文件。写入流程获取io::FileOutputStream从Table写入文件。打开 CSV 文件与 Arrow 文件一样CSV 文件也需要io::ReadableFile直接复用之前的infile对象重新绑定即可// 把输入文件绑定到 test_in.csv ARROW_ASSIGN_OR_RAISE(infile, arrow::io::ReadableFile::Open(test_in.csv));准备 TableCSV 可以读入Table因此声明一个Table指针std::shared_ptrarrow::Table csv_table;读取 CSV 文件到 TableCSV reader 需要传入几个选项结构体——好在这些选项都有默认值可以直接传默认实例。本示例的 CSV 文件没有特殊分隔符且体积小所以直接用默认选项构造 reader// CSV reader 有几个用于各种选项的对象。这里我们先使用默认值。 ARROW_ASSIGN_OR_RAISE( auto csv_reader, arrow::csv::TableReader::Make( arrow::io::default_io_context(), infile, arrow::csv::ReadOptions::Defaults(), arrow::csv::ParseOptions::Defaults(), arrow::csv::ConvertOptions::Defaults()));csv::TableReader::Make的签名见 cpp/src/arrow/csv/reader.h要求传入五类参数io::IOContextI/O 上下文示例用default_io_context()std::shared_ptrio::InputStream输入流这里就是复用的infileReadOptions控制读取行为如use_threads、block_size、skip_rows、是否跳过表头行等ParseOptions控制 CSV 语法解析如分隔符、引号字符、是否允许错误等ConvertOptions控制类型推断与转换如column_types、null_values、strings_can_be_null等。更完整的选项说明可参考 格式化 API 参考。准备好 CSV reader 后用其csv::TableReader::Read方法填充Table// 读取 table。 ARROW_ASSIGN_OR_RAISE(csv_table, csv_reader-Read())从源码结构看TableReader是负责把整个 CSV 文件读成一个 ArrowTable的类cpp/src/arrow/csv/reader.h它同样提供异步版本ReadAsync()。如果需要增量、流式地处理大文件仓库中还提供csv::StreamingReader每次返回一个RecordBatch且类型推断只在首块完成、之后冻结相关细节见同文件 cpp/src/arrow/csv/reader.h。从 Table 写入 CSV 文件CSV 的写入与 IPC 从RecordBatch写入的流程几乎一致区别在于对象是Table且调用的是ipc::RecordBatchWriter::WriteTable而非WriteRecordBatch。注意这里复用了同一个 writer 类——因为我们持有的是Table所以调用WriteTable。先指定目标文件、使用Table的Schema然后写入Table// 把输出文件绑定到 test_out.csv ARROW_ASSIGN_OR_RAISE(outfile, arrow::io::FileOutputStream::Open(test_out.csv)); // CSV writer 有更简单的默认值更复杂的用法请参考 API 文档。 ARROW_ASSIGN_OR_RAISE(auto csv_writer, arrow::csv::MakeCSVWriter(outfile, csv_table-schema())); ARROW_RETURN_NOT_OK(csv_writer-WriteTable(*csv_table)); // 不是必须的但是一种安全的做法。 ARROW_RETURN_NOT_OK(csv_writer-Close());至此我们完成了 CSV 文件的读写读取与写入 Parquet 文件Parquet 是列式存储的事实标准格式Arrow 通过parquet命名空间提供深度集成。读取流程打开文件准备parquet::arrow::FileReader将文件读取为Table。写入流程将Table写入文件。打开 Parquet 文件Parquet 同样需要io::ReadableFile——我们已有这个对象直接对文件调用io::ReadableFile::Open// 把输入文件绑定到 test_in.parquet ARROW_ASSIGN_OR_RAISE(infile, arrow::io::ReadableFile::Open(test_in.parquet));设置 Parquet Reader和前面一样真正读取文件需要一个 Reader。此前我们都是从arrow命名空间获取各种格式的 Reader这次则进入parquet命名空间获取parquet::arrow::FileReaderstd::unique_ptrparquet::arrow::FileReader reader;接着调用parquet::arrow::OpenFile设置 reader。注意这一步是必要的——即便我们已经调用过io::ReadableFile::Open。还要注意这里把parquet::arrow::FileReader以引用指针形式传入而不是用返回值赋给它// 注意 Parquet 的 OpenFile() 以引用方式接收 reader而不是返回一个 reader。 PARQUET_THROW_NOT_OK( parquet::arrow::OpenFile(infile, arrow::default_memory_pool(), reader));这与之前 Arrow 命名空间中RecordBatchFileReader::Open、TableReader::Make的返回值风格形成了鲜明对比——Parquet 的 C API 沿用了 Parquet 项目自身输出参数 抛异常的接口风格这也是为什么这里用PARQUET_THROW_NOT_OK而不是ARROW_RETURN_NOT_OK。将 Parquet 文件读取为 Table拿到准备好的parquet::arrow::FileReader后就可以读取到Table——同样地Table必须按引用传入而不是作为返回值std::shared_ptrarrow::Table parquet_table; // 读取 table。 PARQUET_THROW_NOT_OK(reader-ReadTable(parquet_table));从 Table 写入 Parquet 文件对于一次性写入写 Parquet 文件不需要声明 writer 对象。只需把Table传进去、指定它进行内存分配时使用的内存池、告诉它写入位置以及需要把文件拆分时的块大小chunk size// 写 Parquet 不需要声明 writer 对象。只需绑定输出文件 // 然后传入 table、内存池、输出以及用于把 Table 拆块的 chunk size。 ARROW_ASSIGN_OR_RAISE(outfile, arrow::io::FileOutputStream::Open(test_out.parquet)); PARQUET_THROW_NOT_OK(parquet::arrow::WriteTable( *parquet_table, arrow::default_memory_pool(), outfile, 5));parquet::arrow::WriteTable的第四个参数chunk_size表示每个 RowGroup 大致包含的行数示例传 5Arrow 会按此把Table拆成多个 RowGroup 落盘对单次写入场景这是最简洁的 API。如果需要更精细的控制如压缩算法、字典编码、页大小等仓库中 Parquet 写入器提供了parquet::WriterProperties::Builder等更底层的配置入口可参考 docs/source/cpp/parquet.rst 进一步了解。结束程序最后RunMain()返回Status::OK让main()知道一切正常结束——与第一个教程中的做法一致return arrow::Status::OK(); }至此你已经在 Arrow C 中完成了 IPC、CSV 和 Parquet 三种格式的读取与写入能够正确地加载数据并写出结果了接下来就可以进入下一篇文章学习如何用 compute 函数处理数据。完整代码完整的可运行示例位于 cpp/examples/tutorial_examples/file_access_example.cc。该文件将上文所有代码片段按main()→RunMain()→GenInitialFile()的组织方式整合在一起其中三种格式共用同一个io::ReadableFileinfile与同一个io::FileOutputStreamoutfile演示了对象的复用与重新绑定三种格式的读写调用链含ARROW_ASSIGN_OR_RAISE/PARQUET_THROW_NOT_OK的错误处理均与上文一一对应示例会依次生成test_in.arrow、test_in.csv、test_in.parquet三个输入文件并在验证读写后产生test_out.arrow、test_out.csv、test_out.parquet三个输出文件。建议在你的工作目录中编译并运行该示例观察生成的六个文件再结合本文的调用链说明即可快速建立对 Arrow C 文件 I/O 的整体认知。【免费下载链接】arrowApache Arrow is a multi-language toolbox for accelerated data interchange and in-memory processing项目地址: https://gitcode.com/gh_mirrors/arrow12/arrow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表