ARTICLE DETAIL

资讯详情

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

Daft Parquet 基准测试实战:本地与 S3 读取、Row Group、谓词下推与编解码的对比评测方法

Daft Parquet 基准测试实战:本地与 S3 读取、Row Group、谓词下推与编解码的对比评测方法 Daft Parquet 基准测试实战本地与 S3 读取、Row Group、谓词下推与编解码的对比评测方法【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft本文基于 Daft 的 Parquet 基准测试套件 展开讲解如何搭建基准测试环境、如何运行并解读三类 pytest 基准命令定位 Daft 短板、跨框架对比、峰值内存分析并深入 conftest.py 与各测试文件说明每类基准场景本地读取、S3 批量读取、Row Group 数量、过滤下推、类型与压缩编解码背后的数据构造方式与源码级读取路径。读完后你可以独立复现这套评测流程并理解 Daft 与 PyArrow/boto3 在同一 Parquet 工作负载下的对比方法。基准目标与模块结构benchmarking/parquet/README.md 明确了这套基准的两个目标Find Parquet features that Daft underperforms on——找出 Daft 在 Parquet 特性上表现欠佳的地方Compare Daft against other frameworks——将 Daft 与其他框架PyArrow、boto3 直取进行对比。围绕这两个目标目录下的测试文件各司其职文件场景对比对象test_local.py本地 Parquet 读取覆盖不同行列规模与 Row Group 配置仅 Daft DataFrametest_bulk_reads.pyS3 上多文件批量读取Daft bulk API、PyArrow、boto3test_num_rowgroups.py不同 Row Group 数1/8/64下的列裁剪读取5 种读取函数全对比test_filter_pushdown.py不同选择率1%~90%下的谓词下推Daft vs PyArrowtest_types_and_codecs.py嵌套类型List/Struct/Map× 压缩编解码snappy/zstd/gzip/noneDaft vs PyArrow环境准备README 给出的完整搭建步骤如下先创建独立虚拟环境并安装锁定的依赖版本python -m venv .venv source .venv/bin/activate pip install -r benchmark-requirements.txtbenchmark-requirements.txt 中锁定了精确版本保证基准可复现pytest7.4.0 pytest-benchmark4.0.0 pytest-memray1.4.1 pyarrow19.0.1 boto31.28.3其中pytest-benchmark提供benchmarkfixture 与--benchmark-group-by报告分组pytest-memray提供--memray峰值内存采集pyarrow和boto3则是对比基线本身。随后安装待评测的 Daft——可以是发布版 wheel也可以是本地构建的版本pip install daft对比基线conftest 中的 5 种读取函数基准的对比质量取决于基线函数。conftest.py 定义了一个参数化的read_fnfixture它把每个用例展开为 5 个被测实现daft_native_read调用daft.recordbatch.MicroPartition.read_parquet(path, columnscolumns)后经to_arrow()返回走 Daft 原生 MicroPartition 路径daft_native_read_to_arrow调用daft.recordbatch.read_parquet_into_pyarrow(path, columnscolumns)这是 Daft 提供的「Rust 读取、直接产出 PyArrow Table」的低层接口见 recordbatch.pypyarrow_readpyarrow.parquet.read_tableS3 路径会先构造pafs.S3FileSystemboto3_get_object_read用 boto3get_object把整个对象下载到内存流再用 PyArrow 解析——代表「先全量下载再解析」的经典模式daft_dataframe_read完整的daft.read_parquet(path)select(*columns)to_arrow()代表用户日常使用的 DataFrame API。另有面向批量场景的bulk_read_fnfixture包含daft_into_pyarrow_bulk_read底层为daft.recordbatch.read_parquet_into_pyarrow_bulk、pyarrow_bulk_read逐文件 PyArrow 读取和boto3_bulk_read逐文件 boto3 下载读取。值得注意的是源码中两个低层 API 的默认参数read_parquet_into_pyarrow 支持start_offset、num_rows、row_groups、multithreaded_io、coerce_int96_timestamp_unit、string_encoding等精细控制file_timeout_ms默认 15 分钟而 read_parquet_into_pyarrow_bulk 额外提供row_groups_per_path与num_parallel_tasks默认 128 个并行任务这正是批量读取场景下 Daft 相比逐文件 PyArrow 调用的并行优势来源。场景一本地文件读取test_local.pytest_local.py 会先用 Fakerseed0 保证可复现生成合成数据每列按 5 种模式循环取整型、浮点、字符串、布尔、时间戳类型再由 session 级parquet_filefixture 按如下矩阵生成 9 个本地文件SIZE_CONFIGS(10_000_000, 1)千万行单列、(1_000_000, 32)均衡、(10_000, 1024)万行千列ROWGROUP_CONFIGS[1, 8, 64]个 Row Group。生成逻辑会把row_group_size设为num_rows // num_rowgroups并落盘到benchmarking/parquet/local_data/目录若文件已存在则直接复用避免重复生成。基准函数本身非常简洁def test_read_parquet(parquet_file, benchmark): import daft def read_parquet(): df daft.read_parquet(parquet_file) return df.to_arrow() benchmark(read_parquet)测试打上pytest.mark.benchmark(groupread_parquet_local)标记配合-m benchmark只选中基准用例。该场景用于观察「行数×列数×Row Group 数」三个维度对 Daft 本地读取的单独影响。场景二S3 多文件批量读取test_bulk_reads.pytest_bulk_reads.py 固定使用一个约 200MB、2 个 Row Group 的 TPC-H lineitem 分片s3://eventual-dev-benchmarking-fixtures/parquet-benchmarking/tpch/200MB-2RG/...每文件 5,515,199 行通过bulk_read_fnfixture 对三种实现横向对比test_read_parquet_num_files_single_column复制 1/2/4/8 次同一路径只读L_ORDERKEY单列验证返回行数与列名——考察「文件数增长时单列裁剪读取」的吞吐test_read_parquet_num_files_all_columns1/2/4 次路径、读取全部 16 列——考察全量读取的吞吐。这里断言了确切行数5,515,199与列数16保证基准跑的同时也在做数据正确性校验。场景三Row Group 数量对列裁剪的影响test_num_rowgroups.pytest_num_rowgroups.py 使用 100G 级 TPC-H lineitem 数据分别以 1、8、64 个 Row Group 打包1RG/8RG/64RG每个文件 18,751,674 行并设计 4 个读取组每个组都通过read_fnfixture 展开为 5 种实现num_rowgroups_single_column只读L_ORDERKEYnum_rowgroups_multi_contiguous_columns读相邻的L_ORDERKEY, L_PARTKEY, L_SUPPKEYnum_rowgroups_multi_sparse_columns读相距较远的L_ORDERKEY, L_TAXnum_rowgroups_all_columns读取全部 16 列。Row Group 越碎元数据与寻址开销越大而「连续列」与「稀疏列」的差异则考察 Parquet 文件内列物理布局对 IO 的影响。该场景直接服务于 README 的目标一——找出 Daft 在特定 Parquet 组织方式下可能落后的区域。注意S3 上的 200MB 分片与 100G 大文件需要可访问eventual-dev-benchmarking-fixtures存储桶的凭据这是运行该两组场景的前提。场景四谓词下推与选择率test_filter_pushdown.pytest_filter_pushdown.py 构造了分布已知的 2,000,000 行、8 个 Row Group 的本地文件seed42目录名内嵌参数local_data/r2000000_rg8/以便参数变更时强制重新生成filter_int[0, 100)均匀分布filter_int X即选择约 X% 的行filter_float[0.0, 1.0)均匀分布filter_strcat_00…cat_99均匀分布三个payload_*列用于模拟过滤后仍需物化的负载列。基准在选择率SELECTIVITIES [0.01, 0.10, 0.50, 0.90]上对整型、浮点、字符串谓词及「int AND float」复合谓词各跑一组 Daft / PyArrow 双实现例如整型谓词def test_daft_filter_int(filter_file, selectivity, benchmark): benchmark.group ffilter_int_{_sel_label(selectivity)} threshold int(selectivity * 100) benchmark(lambda: daft.read_parquet(filter_file).where(daft.col(filter_int) threshold).collect())PyArrow 侧则传papq.read_table(filter_file, filters[(filter_int, , threshold)])两者语义等价。此外还有一组无过滤的基线用例no_filter_baseline作为参照。低选择率1%时统计信息剪枝的收益最明显高选择率90%时几乎退化为全量读取——这组对比能清晰呈现两种引擎在不同选择率下的下推行为差异。场景五嵌套类型与压缩编解码test_types_and_codecs.pytest_types_and_codecs.py 关注解码性能取 500,000 行、4 个 Row Group 的数据在 3 种嵌套类型 × 4 种压缩编解码共 12 种组合下各生成一个独立文件types_codecs_{type}_{codec}.parquet类型ListInt641~9 个元素、Struct{a: Int64, b: String}、MapString, Int641~5 个条目编解码snappy、zstd、gzip、none不压缩。每组同时跑daft.read_parquet(filepath).collect()与papq.read_table(filepath)benchmark group 命名为{type}_{codec}便于直接对比同一组合下两个引擎的解码耗时。运行基准README 中的三条命令环境就绪后README 给出三条对应目标的运行命令均在仓库根目录执行目标一定位 Daft 短板的 Parquet 特性pytest benchmarking/parquet/ -m benchmark --benchmark-group-bygroup -k daft-k daft只选中函数名/id 中含daft的用例--benchmark-group-bygroup按 benchmark group 聚合报表聚焦 Daft 各条路径自身的绝对表现。目标二跨框架横向对比pytest benchmarking/parquet/ -m benchmark --benchmark-group-byparam:path按「参数:路径」分组同一 Parquet 文件下不同读取函数/引擎的结果会归到同一组中方便在报告中直接对照 Daft、PyArrow、boto3 的表现。峰值内存分析pytest benchmarking/parquet/ -m benchmark --memray需要已安装pytest-memray已在 benchmark-requirements.txt 中锁定 1.4.1--memray会为每个用例附加内存剖析用于观察各实现读取期间的峰值内存占用——例如「boto3 先全量下载进内存」与 Daft 流式读取之间的内存曲线差异。底层读取路径基准到底在测什么结合源码可以确认这条基准链最终落在 Daft 的 Rust 读取内核上daft_native_read_to_arrow调用的read_parquet_into_pyarrow在 recordbatch.py 中封装了_read_parquet_into_pyarrow支持按row_groups、start_offset、num_rows精确控制读取范围默认string_encodingutf-8、file_timeout_ms900_000读完直接把 Rust 侧 chunk 组装为pa.chunked_array列表构造 Table——省去了经 Daft 类型系统的转换开销完整 DataFrame 路径daft_dataframe_read则经过逻辑计划与 MicroPartition 层如 recordbatch_io.py 中read_parquet所示最终调用MicroPartition.read_parquet透传columns、num_rows、io_config、coerce_int96_timestamp_unit、multithreaded_io等参数再用_cast_table_to_schema把结果对齐到目标 schema 与列序。因此同一目录下的不同 fixture 实际上在测量三个不同层次纯 Rust 读内核daft_native_read_to_arrow、MicroPartition 中间层daft_native_read、以及含优化器/计划层的完整引擎路径daft_dataframe_read。这正是「找出 Daft 在哪些 Parquet 特性上落后」这一问题能够被精确定位到具体层次的原因。小结这套 Parquet 基准以 pytest pytest-benchmark 为骨架用锁版本的依赖清单保证可复现用带断言的读取函数兼顾正确性与性能并用 session 级 fixture 缓存本地合成数据、以 S3 上的 TPC-H 分片覆盖真实云对象存储场景。运行任何一组前只需记住本地场景开箱即用test_bulk_reads.py 与 test_num_rowgroups.py 依赖s3://eventual-dev-benchmarking-fixtures的可读凭据内存分析需保留--memray参数。按上述三条 README 命令分别执行即可同时获得 Daft 各读取层次的自画像与跨引擎的横向对照。【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表