多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

Apache Arrow C++ Datasets 完整实战:从多文件读取到分区写入与谓词下推

Apache Arrow C++ Datasets 完整实战:从多文件读取到分区写入与谓词下推 大数据数据分析数据工程序列化【免费下载链接】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 Datasets 是 Arrow C 库中面向超内存larger than memory、多文件、表格式数据的查询与读写组件它把一个目录 多种文件格式 不同文件系统统一成单个可扫描的数据集对象并提供数据集发现、分区裁剪、列投影与谓词下推等优化能力。本文以仓库中 cpp/examples/arrow/dataset_documentation_example.cc 这一官方示例及其配套教程 docs/source/cpp/dataset.rst 为核心完整讲解如何用 C 创建、发现、扫描、过滤、投影以及读写分区数据集并深入剖析每个 API 背后的源码实现与性能取舍。读完本文你将掌握 Datasets API 的全链路用法能够直接在自己的项目中落地多文件即查、分区即裁剪的实战方案。一、先认识 Arrow Datasets设计目标与支持范围Datasets 库要解决的核心问题是当数据量超过内存、且分布在多个文件甚至多个文件系统中时如何像操作一张表一样高效地读写它们。从源码注释与文档看它提供了三方面能力见 docs/source/cpp/dataset.rst统一接口同一套 API 可对接不同数据源、不同文件格式本地磁盘、云存储等不同文件系统。数据源发现自动爬取目录、识别分区布局、支持多种分区方案、做基础 schema 归一化。优化读取支持谓词下推过滤行、投影选取与派生列并可选并行读取。当前支持的格式为Parquet、Feather / Arrow IPC、CSV、ORC其中 ORC 数据集目前只能读取、尚不能写入见 docs/source/cpp/dataset.rst。后续计划扩展到更多格式与数据源如数据库连接。在开始示例前文档先创建了一个含两个 Parquet 文件的小数据集作为后续所有示例的输入对应的源码函数为CreateTable与CreateExampleParquetDataset见 cpp/examples/arrow/dataset_documentation_example.cc。核心逻辑如下auto schema arrow::schema({arrow::field(a, arrow::int64()), arrow::field(b, arrow::int64()), arrow::field(c, arrow::int64())}); // 用 NumericBuilderInt64Type 生成三列数据 // a [0..9]b [9..0]c [1,2,1,2,...] ... return arrow::Table::Make(schema, {array_a, array_b, array_c});随后通过filesystem-CreateDir(base_path)创建parquet_dataset目录用parquet::arrow::WriteTable将表的Slice(0,5)与Slice(5)分别写成data1.parquet与data2.parquetARROW_ASSIGN_OR_RAISE(auto output, filesystem-OpenOutputStream(base_path /data1.parquet)); ARROW_RETURN_NOT_OK(parquet::arrow::WriteTable( *table-Slice(0, 5), arrow::default_memory_pool(), output, /*chunk_size*/2048));这里体现了 Datasets 的一贯思路数据集Dataset只是一个目录 格式约定的逻辑视图真正的读写都发生在 Fragment单个文件与 Scanner扫描器层面。二、数据集发现Dataset Discovery从目录爬取到 Dataset 对象2.1 FileSystemDatasetFactory入口与参数创建arrow::dataset::Dataset对象的标准方式是通过各类DatasetFactory文件系统场景使用的是FileSystemDatasetFactory。其声明位于 cpp/src/arrow/dataset/discovery.h支持三种构造途径显式文件列表、fs::FileSelector目录选择器、以及直接传入包含文件系统的 URI。示例中采用目录选择器方式// 创建数据集扫描文件系统中的文件 fs::FileSelector selector; selector.base_dir base_dir; ARROW_ASSIGN_OR_RAISE( auto factory, ds::FileSystemDatasetFactory::Make(filesystem, selector, format, ds::FileSystemFactoryOptions())); ARROW_ASSIGN_OR_RAISE(auto dataset, factory-Finish());关键参数说明filesystem读取使用的文件系统例如本地LocalFileSystem或 Amazon S3 文件系统这也是文档强调可以自由选择本地文件或 S3、Parquet 或 CSV的机制见 docs/source/cpp/dataset.rst。示例中文件系统由fs::FileSystemFromUri(uri, root_path)解析得到支持file://等 URI 形式见 cpp/src/arrow/filesystem/filesystem.h。selectorfs::FileSelector其recursive字段默认为falsemax_recursion默认为INT32_MAX见 cpp/src/arrow/filesystem/filesystem.h。读取分区数据集时需显式打开递归。formatstd::shared_ptrds::FileFormat决定文件如何被解析。optionsFileSystemFactoryOptions可配置分区方案等。2.2 发现阶段不读数据爬取、schema 推断与文件清单值得强调的一点是创建 Dataset 并不会开始读数据它只做两件事爬取目录必要时得到文件清单可通过FileSystemDataset::files()打印见 docs/source/cpp/dataset.rst// 仅对 FileSystemDataset 有效 for (const auto filename : dataset-files()) { std::cout filename std::endl; }推断数据集 schema默认取第一个文件的 schemastd::cout dataset-schema()-ToString() std::endl;完成发现后通过Dataset::NewScan()构建Scanner再调用Scanner::ToTable()把整个数据集读成一张arrow::TableARROW_ASSIGN_OR_RAISE(auto scan_builder, dataset-NewScan()); ARROW_ASSIGN_OR_RAISE(auto scanner, scan_builder-Finish()); return scanner-ToTable();在完整示例中ScanWholeDataset还会先打印每个 Fragment 的信息见 cpp/examples/arrow/dataset_documentation_example.ccARROW_ASSIGN_OR_RAISE(auto fragments, dataset-GetFragments()) for (const auto fragment : fragments) { std::cout Found fragment: (*fragment)-ToString() std::endl; }注意ToTable()会把整个数据集载入内存数据量大时可能消耗大量内存。文档明确提示应根据数据集规模配合下文过滤/投影使用见 docs/source/cpp/dataset.rst。2.3 读取不同文件格式一套代码多种格式示例用同样的CreateTable数据通过 Arrow IPC 写入器写成两个 Feather 文件CreateExampleFeatherDataset见 cpp/examples/arrow/dataset_documentation_example.ccARROW_ASSIGN_OR_RAISE(auto writer, arrow::ipc::MakeFileWriter(output.get(), table-schema())); ARROW_RETURN_NOT_OK(writer-WriteTable(*table-Slice(0, 5))); ARROW_RETURN_NOT_OK(writer-Close());读取时只需把 format 换成IpcFileFormat见 docs/source/cpp/dataset.rstauto format std::make_sharedds::IpcFileFormat(); // ... auto factory ds::FileSystemDatasetFactory::Make(filesystem, selector, format, options) .ValueOrDie();也就是说从本地 Parquet切换到S3 上的 CSV/Feather只改两个参数filesystem 与 format这正是统一接口的落地体现。示例入口RunDatasetDocumentation中按format_name分发feather→IpcFileFormatparquet→ParquetFileFormatparquet_hive→ParquetFileFormat 分区写入见 cpp/examples/arrow/dataset_documentation_example.cc。三、自定义文件格式reader_options、parse_options 与 FragmentScanOptionsFileFormat对象带有一组控制读取行为的属性。文档给出的典型例子是让 Parquet 的某一列以字典编码读入见 docs/source/cpp/dataset.rstauto format std::make_sharedds::ParquetFileFormat(); format-reader_options.dict_columns.insert(a);源码中ParquetFileFormat::ReaderOptions确实包含std::unordered_setstd::string dict_columns字段见 cpp/src/arrow/dataset/file_parquet.h。同理CSV 格式通过CsvFileFormat::parse_options类型为csv::ParseOptions见 cpp/src/arrow/dataset/file_csv.h可调整分隔符等解析行为。更进一步ScannerBuilder::FragmentScanOptions(...)可传入FragmentScanOptions对扫描做细粒度控制——例如 CSV 扫描时自定义什么字符串被转换为布尔 true/false见 docs/source/cpp/dataset.rst。该方法的签名在 cpp/src/arrow/dataset/scanner.h 中Status FragmentScanOptions(std::shared_ptrFragmentScanOptions fragment_scan_options);FileFormat的基类FileFormat在构造时接受一个default_fragment_scan_options见 cpp/src/arrow/dataset/file_base.h并提供了纯虚方法DefaultWriteOptions()见 cpp/src/arrow/dataset/file_base.h供写入时生成默认写选项。四、过滤数据列裁剪 谓词下推读取整个数据集常常是浪费的——只关心部分行列时ScannerBuilder提供了两个控制手段见 docs/source/cpp/dataset.rst。4.1 用 Project 选取列ARROW_RETURN_NOT_OK(scan_builder-Project({b}));Project的重载之一接收列名字符串向量见 cpp/src/arrow/dataset/scanner.hStatus Project(std::vectorstd::string columns);像Parquet 这类列式格式可以只从文件系统读取指定列显著降低 I/O 成本见 docs/source/cpp/dataset.rst。4.2 用 Filter 过滤行ARROW_RETURN_NOT_OK(scan_builder-Filter(cp::less(cp::field_ref(b), cp::literal(4))));Filter接收一个compute::Expression见 cpp/src/arrow/dataset/scanner.h。这里cp::field_ref(b)引用列 bcp::literal(4)是常量 4cp::less构造b 4的表达式。对于 Parquet该谓词还可借助行组统计信息减少读取的数据量见 docs/source/cpp/dataset.rst。完整示例中的FilterAndSelectDataset把两者组合只读列 b、只保留b 4的行见 cpp/examples/arrow/dataset_documentation_example.cc。ARROW_RETURN_NOT_OK(scan_builder-Project({b})); ARROW_RETURN_NOT_OK(scan_builder-Filter(cp::less(cp::field_ref(b), cp::literal(4))));五、投影列重命名、类型转换与派生新列Project的另一个重载接收表达式向量 列名向量见 cpp/src/arrow/dataset/scanner.h从而支持比选列更复杂的能力重命名列、类型转换、基于表达式派生新列见 docs/source/cpp/dataset.rst。Status Project(std::vectorcompute::Expression exprs, std::vectorstd::string names);示例ProjectDataset演示了三种典型用法见 cpp/examples/arrow/dataset_documentation_example.ccARROW_RETURN_NOT_OK(scan_builder-Project( { // 原样保留列 a cp::field_ref(a), // 将列 b 转成 float32 cp::call(cast, {cp::field_ref(b)}, arrow::compute::CastOptions::Safe(arrow::float32())), // 由 c 派生布尔列c 1 cp::equal(cp::field_ref(c), cp::literal(1)), }, {a_renamed, b_as_float32, c_1}));这里cp::call调用 Arrow compute 内核如castCastOptions::Safe(arrow::float32())指定安全转换到 float32。5.1 在保留原列的同时追加派生列注意传入表达式后结果表只包含列名向量中给定的列。若想原列 派生列都要可以从 dataset schema 出发动态构造表达式列表。示例SelectAndProjectDataset的做法是见 cpp/examples/arrow/dataset_documentation_example.ccstd::vectorstd::string names; std::vectorcp::Expression exprs; // 读入全部原始列 for (const auto field : dataset-schema()-fields()) { names.push_back(field-name()); exprs.push_back(cp::field_ref(field-name())); } // 再派生一个新列 names.emplace_back(b_large); exprs.push_back(cp::greater(cp::field_ref(b), cp::literal(1))); ARROW_RETURN_NOT_OK(scan_builder-Project(exprs, names));5.2 过滤与投影的协同文档特别提醒同时使用过滤与投影时Arrow 会自动计算所需读取的全部列。例如过滤条件用到了某个最终未选中的列Arrow 仍会读取该列用于求值过滤条件见 docs/source/cpp/dataset.rst保证结果正确性。六、读写分区数据集Partitioned Datasets6.1 什么是分区当数据集有一个或多个经常被过滤的列时与其全读再滤不如把文件按该列的取值组织进嵌套目录用目录名携带数据子集信息扫描时直接跳过不匹配的文件见 docs/source/cpp/dataset.rst。文档给出的 Hive 风格布局如下dataset_name/ year2007/ month01/ data0.parquet data1.parquet ... month02/ ... year2008/ month01/ ...这种/keyvalue/目录命名源自 Apache Hive文件dataset_name/year2007/month01/data0.parquet只包含year 2007 month 01的数据。6.2 用 Datasets 写入分区数据示例CreateExampleParquetHivePartitionedDataset展示了如何用 Datasets 写分区数据集见 cpp/examples/arrow/dataset_documentation_example.cc。先造一张含字符串分区列part的表再用InMemoryDataset包装并扫描auto dataset std::make_sharedds::InMemoryDataset(table); ARROW_ASSIGN_OR_RAISE(auto scanner_builder, dataset-NewScan()); ARROW_ASSIGN_OR_RAISE(auto scanner, scanner_builder-Finish());随后声明分区 schema、Hive 分区器与写入格式// 分区 schema 决定哪些字段属于分区键 auto partition_schema arrow::schema({arrow::field(part, arrow::utf8())}); // Hive 风格分区目录名为 keyvalue auto partitioning std::make_sharedds::HivePartitioning(partition_schema); // 写入 Parquet auto format std::make_sharedds::ParquetFileFormat(); ds::FileSystemDatasetWriteOptions write_options; write_options.file_write_options format-DefaultWriteOptions(); write_options.filesystem filesystem; write_options.base_dir base_path; write_options.partitioning partitioning; write_options.basename_template part{i}.parquet; ARROW_RETURN_NOT_OK(ds::FileSystemDataset::Write(write_options, scanner));FileSystemDataset::Write的静态签名在 cpp/src/arrow/dataset/file_base.hstatic Status Write(const FileSystemDatasetWriteOptions write_options, std::shared_ptrScanner scanner);FileSystemDatasetWriteOptions中还有一批值得了解的字段见 cpp/src/arrow/dataset/file_base.h字段默认值说明preserve_orderfalse多线程写入时是否保持行顺序开启可能显著降低性能max_partitions1024单个批次最多可写入的分区数basename_template—片段文件名模板{i}会被自增整数替换max_open_files900最多同时打开的文件数超出后关闭最久未用的文件max_rows_per_file0不限单个文件最多行数min_rows_per_group0行组最小行数攒够才落盘max_rows_per_group1 20行组最大行数existing_data_behaviorkError输出目录已存在时的行为create_dirtrue是否创建目录S3 等不需要目录的文件系统可关闭写入完成后目录结构为parta/part0.parquet、partb/part1.parquetParquet 文件中不再包含part列见 docs/source/cpp/dataset.rst。6.3 读取分区数据集让分区列回来读取时在FileSystemFactoryOptions里声明 Hive 分区工厂并开启递归扫描见 cpp/examples/arrow/dataset_documentation_example.ccfs::FileSelector selector; selector.base_dir base_dir; selector.recursive true; // 必须搜索子目录 ds::FileSystemFactoryOptions options; // 使用 Hive 风格分区让 Arrow 自动推断分区 schema options.partitioning ds::HivePartitioning::MakeFactory(); ARROW_ASSIGN_OR_RAISE(auto factory, ds::FileSystemDatasetFactory::Make( filesystem, selector, format, options)); ARROW_ASSIGN_OR_RAISE(auto dataset, factory-Finish());虽然分区字段不在实际文件中扫描时它们会被加回到结果表。文档给出了真实运行输出见 docs/source/cpp/dataset.rst$ ./debug/dataset_documentation_example file:///tmp parquet_hive partitioned Found fragment: /tmp/parquet_dataset/parta/part0.parquet Partition expression: (part a) Found fragment: /tmp/parquet_dataset/partb/part1.parquet Partition expression: (part b) Read 20 rows a: int64 -- field metadata -- PARQUET:field_id: 1 b: double ... part: string ---- # snip...每个 Fragment 都带一个分区表达式partition expression这是分区裁剪的关键数据结构。6.4 在分区键上过滤跳过整个文件对分区键做过滤时不匹配的文件根本不会被加载。示例FilterPartitionedDataset见 cpp/examples/arrow/dataset_documentation_example.cc// 基于分区值过滤分区表达式不匹配的文件不会参与读取 ARROW_RETURN_NOT_OK( scan_builder-Filter(cp::equal(cp::field_ref(part), cp::literal(b))));FileSystemDataset::GetFragmentsImpl会结合谓词做子树剪枝SetupSubtreePruning见 cpp/src/arrow/dataset/file_base.h这正是按目录信息避免加载不匹配文件的底层实现。6.5 其他分区方案显式 schema 与目录分区Hive 分区的分区键类型默认从文件路径推断HivePartitioning::MakeFactory()也可以直接构造HivePartitioning并显式定义分区键 schema见 docs/source/cpp/dataset.rstauto part std::make_sharedds::HivePartitioning(arrow::schema({ arrow::field(year, arrow::int16()), arrow::field(month, arrow::int8()), arrow::field(day, arrow::int32()) }));目录分区DirectoryPartitioning路径段直接是分区值、不含键名字段名由段的下标隐含例如/2019/11/15对应year2019, month11, day15。因为名字不在路径里必须显式传入auto part ds::DirectoryPartitioning::MakeFactory({year, month, day});DirectoryPartitioning::MakeFactory的原型位于 cpp/src/arrow/dataset/partition.h同样支持传入完整 schema 代替类型推断。另外Hive 分区对空值有约定占位符__HIVE_DEFAULT_PARTITION__kDefaultHiveNullFallback见 cpp/src/arrow/dataset/partition.h。6.6 分区性能考量粒度是把双刃剑文档专门用一节讨论分区性能见 docs/source/cpp/dataset.rst核心结论如下文件数量与并行度分区把数据集拆成多个文件读写都可并行但每个文件都带来文件系统交互开销与共享元数据如每个 Parquet 文件的 schema、行组统计信息使总体积变大。分区数是文件数的下限按日期分区一年的数据至少 365 个文件再加一个 1000 个取值维度的分区最多可达 36.5 万个文件容易产生几乎全是元数据的小文件。目录结构与发现开销嵌套目录带来裁剪收益但发现数据文件需要递归 list 目录过细分区会放大该成本365 个日期分区需要 365 次 list再加 1000 维度就是 36.5 万次。实操指南避免小于 20MB、大于 2GB 的文件避免超过 10,000 个分区对 Parquet 这类有行组概念的文件格式行组过小会让元数据占比过大Arrow 的写入器在多数场景下提供了合理的行组大小默认值。七、其他数据源内存数据与云存储7.1 内存数据集 InMemoryDataset如果数据已在内存中可以包装成InMemoryDataset以复用 Datasets 的过滤/投影/写出能力见 docs/source/cpp/dataset.rstauto table arrow::Table::FromRecordBatches(...); auto dataset std::make_sharedarrow::dataset::InMemoryDataset(std::move(table)); // 扫描、过滤…… auto scanner_builder dataset-NewScan();InMemoryDataset的定义位于 cpp/src/arrow/dataset/dataset.h。示例正是用它将生成的表写入本地磁盘供后续演示使用见 cpp/examples/arrow/dataset_documentation_example.cc。7.2 云存储Datasets 同样支持从云存储如 Amazon S3读取只需传入不同 filesystem 即可见 docs/source/cpp/dataset.rst。文件系统相关细节可查阅 docs/source/cpp/filesystems.rst。八、事务与 ACID 保证Datasets 的边界文档明确说明见 docs/source/cpp/dataset.rstDatasets API 不提供任何事务支持或 ACID 保证读写皆然。并发读是安全的并发写、或写与读并发可能产生意外行为。规避手段包括每个写入者使用唯一 basename template、使用临时目录存放新文件、或者显式维护文件列表而不依赖目录发现。写入过程中进程被意外杀掉可能留下不一致状态写调用通常在字节完全交付到 OS 页缓存后就返回随后的突然断电仍可能丢失部分文件。多数文件格式Parquet、IPC 等在文件末尾写 magic number可安全检测并丢弃半写文件CSV 没有这类机制部分写入的 CSV 可能被误判为合法文件。九、编译与运行完整示例示例的完整源码是 cpp/examples/arrow/dataset_documentation_example.cc约 400 行涵盖本文介绍的全部场景。它在 CMake 中按目标名dataset-documentation-example构建要求同时启用ARROW_PARQUET与ARROW_DATASET见 cpp/examples/arrow/CMakeLists.txt并显式声明依赖 parquet。命令行用法源码注释与文档一致见 cpp/examples/arrow/dataset_documentation_example.cc./debug/dataset-documentation-example file:///some_path/some_directory format [mode]formatfeather、parquet、parquet_hive三选一modeno_filter默认、filter、project、select_project、partitioned、filter_partitioned六选一见 cpp/examples/arrow/dataset_documentation_example.cc。程序会先按 format 在指定根目录下生成示例数据集再按 mode 执行对应的扫描策略最后打印读取行数与表内容std::cout Read table-num_rows() rows std::endl; std::cout table-ToString() std::endl;若参数不足 3 个程序为 CI 目的直接返回成功见 cpp/examples/arrow/dataset_documentation_example.cc。十、总结Datasets 编程心智模型综合示例与源码可以把 Datasets 的编程模型归纳为一条清晰的主线数据源抽象FileSystemFileFormat决定从哪读、怎么解析同一套DatasetFactory适配本地/S3 与 Parquet/CSV/Feather/ORC发现与逻辑视图Dataset只负责爬取目录、推断 schema、组织 Fragment不预读数据分区信息通过Partitioning注入扫描与优化ScannerBuilder提供Project选列/投影/派生与Filter谓词Parquet 等格式支持列裁剪与统计下推分区键过滤可跳过整个文件写出FileSystemDataset::WriteFileSystemDatasetWriteOptions把扫描结果按分区布局写回磁盘边界意识无事务/ACID 保证生产环境需自行设计写入互斥与部分文件清理策略。从文档与源码的对应关系看ScannerBuilder::Project/Filtercpp/src/arrow/dataset/scanner.h、FileSystemDatasetWriteOptionscpp/src/arrow/dataset/file_base.h、HivePartitioning/DirectoryPartitioningcpp/src/arrow/dataset/partition.h共同构成了 Datasets 的核心能力面。建议读者把本文与 docs/source/cpp/api/dataset.rst 的 API 参考对照阅读并在实际项目中以示例为起点逐步替换为自己的文件系统、格式与分区方案。赞分享大数据数据分析数据工程序列化【免费下载链接】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 C 数据集库Tabular Datasets实践指南多文件、分区与谓词下推Apache Arrow C 数据集库Tabular Datasets实践指南多文件、分区与谓词下推 导读 本文围绕 Apache Arrow C数据工程大数据序列化数据分析Apache Arrow C Datasets 实战指南从读取、过滤、投影到分区读写Apache Arrow C Datasets 实战指南从读取、过滤、投影到分区读写 本文以 docs/source/cpp/examples/datas数据工程数据分析大数据Apache Arrow C Datasets 实战从 Table 到多文件分区数据集的读写全流程Apache Arrow C Datasets 实战从 Table 到多文件分区数据集的读写全流程 Arrow C 的 arrow::dataset数据工程大数据序列化数据分析上一篇Super Productivity免费开源的时间追踪与任务管理终极解决方案下一篇Kata Containers 安装实践前提条件、Helm 部署、发布包与 Docker 集成全解创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表