- 数据分析
- 数据工程
- 机器学习
【免费下载链接】cudf
cuDF - GPU DataFrame Library
本篇技术指南以 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/sum | int64 范围与求和 |
| 浮点(double) | minimum/maximum/sum | double 范围与求和 |
| 字符串(string) | minimum/maximum/sum | 字典序最小/最大字符串,sum为总字符长度 |
| 布尔(bucket) | true_count/false_count | 通过true_count与number_of_values推算得到 |
| Decimal | minimum/maximum/sum | 以字符串形式保存 |
| 日期(date) | minimum/maximum | UTC 时区的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_compression | SNAPPY | 压缩算法;传入AUTO会被归一化为SNAPPY |
enable_statistics(val) | enable_statistics | ORC_STATISTICS_ROW_GROUP | 统计收集粒度,见下文 |
stripe_size_bytes(val) | set_stripe_size_bytes | 64 MiB(64 * 1024 * 1024) | 单个 stripe 最大字节数 |
stripe_size_rows(val) | set_stripe_size_rows | 1,000,000 行 | 单个 stripe 最大行数 |
row_index_stride(val) | set_row_index_stride | 10,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
相关推荐
pylibcudf Parquet 读写 API 完全指南:从 ParquetReaderOptions 到 ChunkedParquetWriter 的 GPU 加速数据管线
pylibcudf Parquet 读写 API 完全指南:从 ParquetReaderOptions 到 ChunkedParquetWriter 的 GP
数据分析数据工程机器学习PyArrow 读写 Apache ORC 格式完全指南:从单文件到云存储
PyArrow 读写 Apache ORC 格式完全指南:从单文件到云存储 Apache ORC(Optimized Row Columnar)是一种开源的列式
数据工程大数据序列化数据分析scrcpy 安卓投屏:35ms 延迟、1 秒出首帧、手机零 App 安装
scrcpy 安卓投屏:35ms 延迟、1 秒出首帧、手机零 App 安装 敲一行 scrcpy ,约 1 秒后手机屏幕出现在电脑窗口,键盘输入、鼠标点击直接落
音视频
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考