☰
pylibcudf 的 ORC 读写 API 完全指南:从 read_orc 到分块写入
2026/9/25 3:55:33 网站建设 项目流程
  • 数据分析
  • 数据工程
  • 机器学习

【免费下载链接】cudf

cuDF - GPU DataFrame Library

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

本篇技术指南以 cuDF 仓库中 pylibcudf 的 ORC(Optimized Row Columnar)格式 I/O 模块 API 文档为核心,系统讲解 pylibcudf 提供的 ORC 读取、写入、统计信息解析与分块(chunked)写入的完整接口。读者读完本文将掌握:如何用OrcReaderOptions精确控制读取范围与列投影、如何用OrcWriterOptions调优 stripe 与压缩参数、如何解析 ORC 文件级与 stripe 级列统计,以及在大数据量场景下如何使用OrcChunkedWriter分批写出数据。

pylibcudf.io.orc 模块概览

在 cuDF 仓库中,pylibcudf 的 ORC I/O 能力集中在 pylibcudf.io.orc 模块内,对应的 API 文档入口即 docs/cudf/source/pylibcudf/api_docs/io/orc.rst,通过 Sphinx 的automodule指令自动收集该模块的全部公开成员。模块导出的完整符号清单(见 orc.pyx 中的__all__)包括:

  • 读取:read_orc、read_parsed_orc_statistics
  • 读取配置:OrcReaderOptions、OrcReaderOptionsBuilder
  • 写入:write_orc
  • 写入配置:OrcWriterOptions、OrcWriterOptionsBuilder
  • 分块写入:OrcChunkedWriter、ChunkedOrcWriterOptions、ChunkedOrcWriterOptionsBuilder
  • 统计对象:OrcColumnStatistics、ParsedOrcStatistics
  • 能力探测:is_supported_read_orc、is_supported_write_orc

该模块底层直接绑定 libcudf 的 C++ 实现(cudf::io::orc_reader_options、cudf::io::read_orc、cudf::io::write_orc等,见 cpp/include/cudf/io/orc.hpp),因此所有 GPU 加速的解析、解压、编码与列裁剪都在设备端完成,Python 侧仅负责传递配置与接收结果。

读取 ORC 文件:OrcReaderOptions 与 read_orc

从构建器开始

读取 ORC 的第一步是构造OrcReaderOptions。与大多数 pylibcudf I/O 接口一致,推荐通过OrcReaderOptions.builder(source)创建OrcReaderOptionsBuilder,再用链式方法配置行为,最后调用build()生成选项对象。source是SourceInfo,可以指向文件路径、主机字节缓冲区(HostBuffer)或设备缓冲区(DeviceBuffer)。

import pylibcudf as plc source_info = plc.io.types.SourceInfo(["dataset.orc"]) options = plc.io.orc.OrcReaderOptions.builder(source_info).build() result = plc.io.orc.read_orc(options) # result 是 TableWithMetadata,包含 GPU 上的 Table 与列名等元数据

read_orc的完整签名还支持显式传入 CUDA 流与内存资源(见 orc.pyx):

read_orc(options, stream=None, mr=None)
  • stream:用于设备内存操作与 kernel 启动的 CUDA 流,不传则使用默认流;
  • mr:DeviceMemoryResource,控制返回表设备内存的分配来源(如使用池化内存资源)。

读取选项全解

OrcReaderOptions提供一组 setter 方法,覆盖行范围、列投影、stripe 选择、类型转换等维度。下表汇总各选项的含义与约束(对应 C++ 侧实现见 orc.hpp 与 orc.pyx):

方法参数说明约束
set_num_rows(nrows)int64_t从读取起点开始读取的行数;不设置则读到文件末尾不能为负;与set_stripes互斥
set_skip_rows(skip_rows)int64_t从文件开头跳过的行数不能为负;与set_stripes互斥
set_stripes(stripes)list[list[int]]每个输入源要读取的 stripe 编号列表(外层列表与输入源一一对应)非空时不允许再设置skip_rows/num_rows
set_columns(col_names)list[str]只读取指定名称的列(列投影)列名必须是字符串
set_decimal128_columns(val)list[str]将指定列(使用完全限定名)读为 128 位 Decimal 类型列名必须是字符串
set_timestamp_type(type_)DataType将时间戳列统一转换为指定时间戳类型—
set_source(src)SourceInfo覆盖已有数据源—
use_index(use)(builder 方法)bool是否使用 ORC 行索引(row index)加速读取,默认开启—

典型的分页读取示例:

options = plc.io.orc.OrcReaderOptions.builder(source_info).build() options.set_skip_rows(1000) # 跳过前 1000 行 options.set_num_rows(500) # 只读接下来的 500 行 options.set_columns(["a", "b"]) # 只取 a、b 两列 options.set_timestamp_type(plc.DataType(plc.TypeId.TIMESTAMP_MICROSECONDS)) result = plc.io.orc.read_orc(options)

几点源码级注意(见 orc.hpp):

  • set_stripes、set_skip_rows、set_num_rows之间存在互斥校验:stripes非空时,C++ 侧会通过CUDF_EXPECTS抛出cudf::logic_error(如 "Can't set stripes along with skip_rows");
  • skip_rows、num_rows传入负值会直接抛错;
  • 底层 C++ 还有enable_use_np_dtypes(numpy 兼容 dtype)、enable_ignore_timezone_in_stripe_footer(忽略 stripe footer 中的写入方时区)等选项,pylibcudf 目前暴露的是上表中的核心子集。

返回值 TableWithMetadata

read_orc返回TableWithMetadata(定义见 types.pyi),同时携带tbl(GPU 上的Table)与column_names(含嵌套子列名的列名规格)。可通过result.tbl取表,用result.columns取列元组,或调用column_names(...)获取扁平或含子列的列名列表。

解析 ORC 统计信息:read_parsed_orc_statistics

ORC 格式天然在文件 footer 与各 stripe 中携带列级统计(最小值、最大值、总和、空值信息等)。pylibcudf 提供read_parsed_orc_statistics(source_info, stream=None)直接读取并解析这些统计(见 orc.pyx),返回ParsedOrcStatistics对象:

stats = plc.io.orc.read_parsed_orc_statistics(plc.io.types.SourceInfo(["dataset.orc"])) print(stats.column_names) # 每列的列名 for col_stats in stats.file_stats: # 文件级统计(每列一个) print(col_stats.number_of_values) # 非空值数量 print(col_stats.has_null) # 是否含空值 print(col_stats.get("minimum")) # 按类型安全取值

ParsedOrcStatistics的三个属性对应(见 orc.pyx):

  • column_names:list[str],每列的列名;
  • file_stats:list[OrcColumnStatistics],文件级每列一条统计;
  • stripes_stats:list[list[OrcColumnStatistics]],外层按 stripe、内层按列组织。

OrcColumnStatistics提供统一的字典式访问:__getitem__、__contains__、get(item, default),以及number_of_values、has_null属性(可能为None,表示该统计缺失)。类型特定统计会按列类型填充不同键(解析逻辑见 orc.pyx,C++ 结构定义见 cpp/include/cudf/io/orc_metadata.hpp):

列类型可用键说明
整数(integer)minimum/maximum/sumint64 范围与求和
浮点(double)minimum/maximum/sumdouble 范围与求和
字符串(string)minimum/maximum/sum字典序最小/最大字符串,sum为总字符长度
布尔(bucket)true_count/false_count通过true_count与number_of_values推算得到
Decimalminimum/maximum/sum以字符串形式保存
日期(date)minimum/maximumUTC 时区的datetime.datetime
二进制(binary)sum总字节数
时间戳(timestamp)minimum/maximum依据 ORC-135 规范读取minimumUtc/maximumUtc并转为 UTCdatetime(毫秒精度)

写入 ORC 文件:OrcWriterOptions 与 write_orc

基本写入流程

写入同样采用 builder 模式:OrcWriterOptions.builder(sink, table)同时接收目标SinkInfo(文件路径或缓冲区)与要写出的Table,构建完成后调用plc.io.orc.write_orc(options, stream=None)执行(见 orc.pyx):

import pylibcudf as plc import pyarrow as pa pa_table = pa.table({"a": [1.0, 2.0, None], "b": [True, None, False]}) plc_table = plc.Table.from_arrow(pa_table) sink = plc.io.types.SinkInfo(["output.orc"]) # 可选的列级元数据与 footer 键值元数据 tbl_meta = plc.io.types.TableInputMetadata(plc_table) user_data = {"source": "pylibcudf-demo"} options = ( plc.io.orc.OrcWriterOptions.builder(sink, plc_table) .metadata(tbl_meta) .key_value_metadata(user_data) .compression(plc.io.types.CompressionType.SNAPPY) .enable_statistics(plc.io.types.StatisticsFreq.STATISTICS_ROWGROUP) .build() ) plc.io.orc.write_orc(options)

写入选项与默认值

OrcWriterOptions的 setter 与 builder 方法对应关系及默认值如下(默认值常量见 orc.hpp,成员初始化见 orc.hpp):

builder 方法setter默认值说明
compression(comp)set_compressionSNAPPY压缩算法;传入AUTO会被归一化为SNAPPY
enable_statistics(val)enable_statisticsORC_STATISTICS_ROW_GROUP统计收集粒度,见下文
stripe_size_bytes(val)set_stripe_size_bytes64 MiB(64 * 1024 * 1024)单个 stripe 最大字节数
stripe_size_rows(val)set_stripe_size_rows1,000,000 行单个 stripe 最大行数
row_index_stride(val)set_row_index_stride10,000 行行索引(row index)跨度,即每个 row group 的最大行数
metadata(meta)set_metadata无TableInputMetadata,写入列级元数据
key_value_metadata(kvm)set_key_value_metadata空footer 的键值元数据(dict[str, str])

统计粒度取值来自StatisticsFreq枚举(见 types.pyi):STATISTICS_NONE(不收集)、STATISTICS_ROWGROUP/STATISTICS_PAGE/STATISTICS_COLUMN。为消除术语歧义,libcudf 专门定义了ORC_STATISTICS_STRIPE = STATISTICS_ROWGROUP与ORC_STATISTICS_ROW_GROUP = STATISTICS_PAGE两个常量——ORC 的 "stripe" 对应 Parquet 的 "row group",ORC 的 "row group" 对应 Parquet 的 "page"(见 orc.hpp)。

关键取值约束

写配置并非任意取值,C++ 侧有硬性校验(见 orc.hpp):

  • set_stripe_size_bytes:最小值 64 KiB(64 << 10),否则抛logic_error("64KB is the minimum stripe size");
  • set_stripe_size_rows:最小值 512 行("Maximum stripe size cannot be smaller than 512");
  • set_row_index_stride:最小值 512,且实际生效时向下取整到 8 的倍数(get_row_index_stride中unaligned_stride - unaligned_stride % 8);
  • 当 stripe 行数小于 row group 行数时,row group 大小会被自动缩减以适配 stripe 大小。

分块写入:OrcChunkedWriter 与 ChunkedOrcWriterOptions

当数据无法一次性整体驻留 GPU 内存、或需要流式地分批写出大量数据时,使用OrcChunkedWriter。流程为:先构造ChunkedOrcWriterOptions.builder(sink),配置压缩与统计等参数并build(),再通过OrcChunkedWriter.from_options(options, stream=None)创建写入器,之后循环调用writer.write(table),最后必须调用writer.close()收尾(见 orc.pyx):

sink = plc.io.types.SinkInfo(["chunked.orc"]) chunked_options = ( plc.io.orc.ChunkedOrcWriterOptions.builder(sink) .compression(plc.io.types.CompressionType.SNAPPY) .enable_statistics(plc.io.types.StatisticsFreq.STATISTICS_ROWGROUP) .build() ) writer = plc.io.orc.OrcChunkedWriter.from_options(chunked_options) for chunk_table in chunk_iterator: # 分批产出 pylibcudf Table writer.write(chunk_table) writer.close()

ChunkedOrcWriterOptions与OrcWriterOptions共享同一组调优参数(set_stripe_size_bytes、set_stripe_size_rows、set_row_index_stride及 builder 的compression/enable_statistics/key_value_metadata/metadata),默认值与校验规则完全一致。区别在于:普通写入要求构造时传入完整Table,而分块写入只需SinkInfo,Table在每次write调用时才提供。

说明:在底层 C++ 中,读取侧同样存在chunked_orc_reader(见 orc.hpp),用于把超大 ORC 文件按chunk_read_limit(输出字节上限)、pass_read_limit(临时内存上限)与output_row_granularity(行粒度)分块读回,pylibcudf 目前暴露的是写入方向的分块能力。

压缩支持探测:is_supported_read_orc / is_supported_write_orc

ORC 压缩算法是否可用取决于当前系统构建配置(例如某些压缩库未静态链接或运行时缺失)。CompressionType枚举提供NONE、AUTO、SNAPPY、GZIP、BZIP2、BROTLI、ZIP、XZ、ZLIB、LZ4、LZO、ZSTD等取值(见 types.pyi),但并非每种都必然可用。pylibcudf 提供两个运行时探测函数(见 orc.pyx):

if plc.io.orc.is_supported_write_orc(plc.io.types.CompressionType.ZSTD): # 仅在支持时才使用 ZSTD 写入 options = plc.io.orc.OrcWriterOptions.builder(sink, table) \ .compression(plc.io.types.CompressionType.ZSTD) \ .build() plc.io.orc.write_orc(options)

is_supported_read_orc(compression)与is_supported_write_orc(compression)分别检查读写方向的支持情况,C++ 侧文档明确标注这是"运行时检查"(runtime check),因此最佳实践是在选用非默认压缩算法前先探测。

测试验证与注意事项

pylibcudf 的 ORC 功能在 python/pylibcudf/tests/io/test_orc.py 中有系统性覆盖,可作为使用范式的参考:

  • test_read_orc_basic:参数化验证nrows/skiprows/columns组合,并通过set_source覆盖数据源,最终与 PyArrow 期望结果逐表比对;
  • test_read_orc_from_device_buffers:验证从DeviceBuffer构造SourceInfo直接读取;
  • test_roundtrip_pa_table:覆盖NONE/SNAPPY压缩、STATISTICS_NONE/STATISTICS_COLUMN统计粒度、以及 64 KiB stripe 与 512 行级参数的回环(round-trip)读写。

此外有几个写入侧事实值得留意(见 orc.hpp):

  • 若编码或压缩过程中抛出异常,则不会向 sink 写入任何数据(非部分写入);
  • UNIX 纪元前最后 999 毫秒内的时间戳在 ORC 中不可精确表示,读回时会晚一秒(与 Apache ORC 写入器行为一致,对应 ORC-763 / ORC-771);
  • 时间戳写入默认按 UTC 记录(writer_timezone默认"UTC"),如需与 Hive / Spark 等记录本地时区的写入器互操作,应显式设置写入方时区;空字符串或无法解析的时区名会在写入时被拒绝。

小结

pylibcudf.io.orc 是一套以 builder 模式贯穿读写两侧、参数语义与 libcudf C++ 实现一一对应的 GPU 加速 ORC 接口。读取侧通过OrcReaderOptions控制行范围、列投影、stripe 选择与 Decimal128 / 时间戳转换;统计侧用read_parsed_orc_statistics拿到文件级与 stripe 级列统计;写入侧用OrcWriterOptions调优压缩、stripe 与行索引参数;数据量超出单次内存承载时,OrcChunkedWriter提供分批写出方案。所有参数默认值与合法性校验均能在 cpp/include/cudf/io/orc.hpp 中找到源码级依据,建议在自定义压缩算法或极小 stripe 配置前,对照上文约束表并优先使用is_supported_write_orc做运行时探测。

  • 数据分析
  • 数据工程
  • 机器学习

【免费下载链接】cudf

cuDF - GPU DataFrame Library

项目地址:https://gitcode.com/gh_mirrors/cu/cudf
点击查看免费下载
上一篇:COLMAP三维重建实战:5种安装方案深度解析与性能优化
下一篇:Wasp 单命令自动化部署(Wasp Deploy)完整指南:从 `wasp deploy` 到 Fly.io 与 Railway 的生产落地

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

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询