ARTICLE DETAIL

资讯详情

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

Dask cuDF:基于 Dask 构建 GPU 并行 DataFrame 的配置、查询规划与多 GPU 实战指南

Dask cuDF:基于 Dask 构建 GPU 并行 DataFrame 的配置、查询规划与多 GPU 实战指南 数据分析数据工程机器学习【免费下载链接】cudfcuDF - GPU DataFrame Library项目地址https://gitcode.com/gh_mirrors/cu/cudf点击查看免费下载本篇技术指南以 Dask cuDF 官方文档入口为核心系统讲解 Dask cuDF 如何作为 Dask DataFrame 的cudf后端被自动注册、如何通过 Dask 配置系统一键切换 GPU 后端、自 24.06 版本起默认启用的自动查询规划基于 Dask-Expressions 的谓词/投影下推以及使用LocalCUDACluster部署多 GPU 集群时的关键参数与最佳实践。读完本文你将能够独立完成 Dask cuDF 的环境配置、后端切换、查询优化理解与单节点/多节点 GPU 集群部署。一、Dask cuDF 是什么Dask DataFrame 的 GPU 后端扩展Dask cuDF读音 DASK KOO-dee-eff是 Dask 并行计算框架的一个扩展库。安装后它会自动注册为 Dask DataFrame 的cudfDataFrame 后端。这一点在仓库源码中可以得到直接印证打包配置 pyproject.toml 中声明了两组 Dask 后端入口点entry point分别面向传统 Dask DataFrame 和基于 Dask-Expressions 的新 API[project.entry-points.dask.dataframe.backends] cudf dask_cudf.backends:CudfBackendEntrypoint [project.entry-points.dask_expr.dataframe.backends] cudf dask_cudf.backends:CudfBackendEntrypoint入口点指向的CudfBackendEntrypoint类定义在 backends.py 中它为read_parquet、read_csv、read_json、read_orc、from_dict、to_backend等创建方法提供了 cudf 专属实现供 Dask 在dataframe.backend配置为cudf时自动调用。官方文档docs/dask_cudf/source/index.rst同时给出一个重要限定Dask cuDF 本身不提供多 GPU 或多节点执行能力。若要利用多块 GPU必须另行部署dask.distributed集群文档强烈推荐配合 dask-cuda 使用以便充分发挥 GPU 与网络硬件的全部特性。适用版本前提当前仓库中 VERSION 文件显示 dask-cudf 版本为 26.12.00其依赖声明要求 Python 3.11、cudf26.12.*、pandas3.0.0,3.1.0、numpy2.0,3.0Dask 生态依赖通过rapids-dask-dependency对齐见 pyproject.toml。二、使用方式一推荐通过 Dask 配置系统启用 cudf 后端这是文档推荐的用法。只需通过 Dask 配置系统把dataframe.backend设为cudfimport dask dask.config.set({dataframe.backend: cudf})也可以在运行代码前通过环境变量DASK_DATAFRAME__BACKENDcudf达成同样效果。设置完成后当使用下列dask.dataframe函数从磁盘格式创建新的 DataFrame 集合时公共 Dask DataFrame API 会自动走 cuDF 路径dask.dataframe.read_parquetdask.dataframe.read_jsondask.dataframe.read_csvdask.dataframe.read_orcdask.dataframe.read_hdfdask.dataframe.DataFrame.from_dict文档给出的对照示例import dask.dataframe as dd # 默认得到 pandas 后端的数据框 df dd.read_parquet(data.parquet, ...) import dask dask.config.set({dataframe.backend: cudf}) # 现在得到 cuDF 后端的数据框 df dd.read_parquet(data.parquet, ...)从后端入口点的实现看CudfBackendEntrypoint.read_parquet会将调用转发到 io/parquet.py 中面向 Dask-Expressions 的read_parquet_expr见 backends.pyread_csv/read_json/read_orc同理分别转发到dask_cudf.io子包下的对应实现——这就是配置切换后端背后的实际调用路径。其他创建函数的后端取决于输入当使用dask.dataframe.from_map、dask.dataframe.from_pandas、dask.dataframe.from_delayed、dask.dataframe.from_array等函数创建新集合时新集合的后端取决于传入的输入对象而非全局配置import pandas as pd import cudf # pandas 输入 → pandas 后端数据框 dd.from_pandas(pd.DataFrame({a: range(10)})) # cuDF 输入 → cuDF 后端数据框 dd.from_pandas(cudf.DataFrame({a: range(10)}))用 to_backend 显式迁移后端任何已有集合都可以随时通过dask.dataframe.DataFrame.to_backend迁移到指定后端# 确保是 cuDF 后端的数据框 df df.to_backend(cudf) # 确保是 pandas 后端的数据框 df df.to_backend(pandas)在源码层面to_backend由CudfBackendEntrypoint.to_backend构造一个ToCudfBackend表达式节点backends.py实际的 cuDF 到 pandas / pandas 到 cuDF 的转换则通过to_pandas_dispatch与to_cudf_dispatch两组 dispatch 完成backends.py。官方最佳实践文档同时提醒频繁在 CPU 与 GPU 之间来回搬运数据会显著拖慢性能应尽可能让数据留在 GPU 上见 best_practices.rst。三、使用方式二显式的 dask_cudf API除作为 Dask DataFrame 的cudf后端外Dask cuDF 还提供显式的dask_cudf模块 APIimport dask_cudf # 总是得到 cuDF 后端的数据框 df dask_cudf.read_parquet(data.parquet, ...)查看 dask_cudf/init.py 可以发现dask_cudf.read_parquet、read_csv、read_json、read_orc这几个显式函数本质上就是在config.set({dataframe.backend: cudf})的上下文中调用对应的dd.*函数——也就是说显式 API 与配置后使用 dask.dataframe走的是同一条路径直接使用它不会带来任何性能收益。此外文档特别指出显式 API 的部分用法与自动查询规划不兼容详见下一节因此推荐优先使用 CPU/GPU 可移植的dask.dataframeAPI。当前显式 API 的公开符号为DataFrame、Index、Series、concat、from_cudf、from_delayedinit.pygroupby_agg、to_orc等旧接口已被标记为弃用_deprecated_api并建议改用DataFrame.to_orc等替代方式。模块导入时还会执行一次关键的运行时修补——_patch_dask_expr()init.py其实现位于 _expr/expr.py用于修补 Dask-Expressions 的若干表达式节点例如用 GPU 友好的VarCudf替换Expr.var的方差计算、修补 cumulative 归约的参数传递保证 cuDF 后端在表达式优化下的正确性。四、查询规划默认开启的自动优化24.06 起自 24.06 版本起Dask cuDF默认提供自动查询规划只要在首次导入dask.dataframe时dataframe.query-planning配置为True默认值底层就会使用 Dask-Expressionsdask-expr来构建和优化计算图。文档用一个投影下推predicate pushdown / column pruning的例子说明其效果df dd.read_parquet(/my/parquet/dataset/) result df.sort_values(B)[A]未优化的表达式图df.pprint()Projection: columnsA SortValues: by[B] shuffle_methodtasks options{} ReadParquetFSSpec: path/my/parquet/dataset/ ...简化后的表达式图df.simplify().pprint()Projection: columnsA SortValues: by[B] shuffle_methodtasks options{} ReadParquetFSSpec: path/my/parquet/dataset/ columns[A, B] ...对比可以看到优化器把投影所需列[A, B]下推到了 Parquet 读取节点从而避免读取无关列。文档同时提醒无需手动优化或简化图——Dask 会在结果转换为任务图时经由dask.compute或dask.persist内部位于dask.optimize流程中自动完成表达式图的简化。五、多 GPU 与多节点执行分区、LocalCUDACluster 与内存管理分区与 out-of-core 计算Dask cuDF即 Dask DataFrame会尽可能自动把数据切分到足够小的任务使其能舒适地放进单块 GPU 的显存。这意味着计算一个查询所需的任务可以流式地送往单个 GPU 进程实现超出显存容量out-of-core的计算同样的任务也可以在多 GPU 集群上并行执行。部署 GPU 感知集群的典型代码要在多 GPU 上执行 Dask 工作流通常需要用 dask-cuda 部署 distributed 集群并用distributed.Client定义客户端。文档给出的完整示例from dask_cuda import LocalCUDACluster from distributed import Client if __name__ __main__: client Client( LocalCUDACluster( CUDA_VISIBLE_DEVICES0,1, # 使用两个 worker设备 0 和 1 rmm_pool_size0.9, # 将 GPU 显存的 90% 用作内存池加速分配 enable_cudf_spillTrue, # 启用 cuDF 原生溢出提升设备内存稳定性 local_directory/fast/scratch/, # 使用快速本地存储做溢出 ) ) df dd.read_parquet(/my/parquet/dataset/) agg df.groupby(B).sum() agg.compute() # 使用上面定义的集群执行四个关键参数值得逐一理解参数作用CUDA_VISIBLE_DEVICES指定参与计算的 GPU 设备每个设备对应一个 workerrmm_pool_size在 worker 上预分配 RMM 内存池如 0.9 表示显存的 90%使 cuDF 内存分配显著更快enable_cudf_spill开启 cuDF 原生溢出spilling支持把临时中间结果溢出到主机内存提升设备内存稳定性local_directory指定快速本地存储路径作为溢出的落盘位置仓库中 python/dask_cudf/README.md 提供了与上述完全一致的单节点多 GPU 快速上手示例可作为可直接复制运行的参考。关于 dask.compute 的重要警告文档以醒目 note 强调上面的示例使用dask.compute会在本地内存中物化一个具体的cudf.DataFrame对象——绝不要对一个无法舒适放进单块 GPU 显存的大集合调用dask.compute。六、最佳实践要点文档 toctree 关联内容文档索引页的 toctree 关联了 best_practices.rst其核心指导与上述配置参数直接配套值得随主文档一并掌握部署与配置用 Dask-CUDA 部署集群单机上即使只有一块 GPU也建议使用LocalCUDACluster——它便于把 worker 固定到指定设备、配置内存溢出选项且 distributed 调度器提供实时仪表盘诊断信息云/HPC 环境则建议用系统专属的部署工具。善用诊断工具浏览器仪表盘可实时查看 worker 资源、计算进度以及GPU标签页下的显存与利用率指标。启用 cuDF 溢出经典 ETL 工作负载建议开启 cuDF 原生溢出使用LocalCUDACluster时即设置enable_cudf_spillTrue。使用 RMM在 worker 进程上初始化 RMM 池如rmm_pool_size0.9可让 cuDF 的内存分配更快更高效。使用 dask.dataframe API优先用可移植的dask.dataframeAPI 后端配置而不是显式dask_cudf模块需要跨后端转换时用to_backend。避免急执行eager executionDask DataFrame 集合默认是惰性的但以下操作会立即执行底层任务图compute()执行整个任务图并在客户端进程中拼接所有分区——大集合切勿调用persist()同样执行整个任务图但计算结果保留在分布式 worker 内存中而非拼接回客户端若总分区大小超过全部 GPU 显存之和会导致大量溢出甚至 OOMlen/head/tail往往会执行部分或全部任务图检查数据时需谨慎sort_values/set_index需要急切收集分位数信息用于全局排序调用set_index时若全局集合不需要按新索引排序务必传入sortFalse。读取数据的调优调分区大小理想分区大小通常是单 GPU 显存容量的 1/32 ~ 1/8。经验法则是shuffle 密集型工作负载大规模排序、连接从 1/32 ~ 1/16 起步其余场景 1/16 ~ 1/8数据分布严重倾斜时可能需要 1/64 或更小。最简单的方式是在创建集合时调参read_parquet和read_csv都暴露blocksize参数控制最大分区大小创建后再用repartition是最后手段。优先使用 Parquet列式存储让 Dask 能做列投影与谓词下推。两个最关键参数blocksize默认256 MiB显存大于 8 GiB 的 GPU 上更大的值通常更快。它决定每个输出分区分到多少 Parquet row-group只按未压缩存储大小估算通常小于对应的cudf.DataFrameaggregate_files默认False当数据集包含大量小于blocksize一半的小文件时设为True通常更快若已知文件本身已对应合理分区大小设blocksizeNone可禁止文件切分从远程存储如对象存储读取大量此类文件时blocksizeNone还能避免昂贵的元数据收集。严格 1:1 文件-分区映射用dask.dataframe.from_map配合cudf.read_parquet手动构造分区——因为使用read_parquet时查询规划优化可能自动把多个文件聚合进同一分区即使aggregate_filesFalse。使用from_map实现自定义创建逻辑相比from_delayed它支持真正的惰性执行并能在映射函数接受columns关键字参数时启用列投影。务必显式指定meta参数否则 Dask 会急切地物化第一个分区——若首个可见设备上启用了大 RMM 池客户端侧的这次急切执行可能直接 OOM。排序、连接与分组排序、连接、分组都可能触发跨分区的全局 shuffle。数据能放进全局 GPU 显存时瓶颈通常是 worker 间通信数据超出全局显存时瓶颈通常是设备到主机的溢出。推荐做法使用 Dask-CUDA worker 的分布式集群尽可能启用 cuDF 原生溢出尽量避免 shuffle低基数 groupby 聚合用split_out1当其中一个集合只有很少分区例如 ≤5时join 使用broadcastTrue通信成为瓶颈时启用 UCX使 worker 间通信走 NVLink、InfiniBand 等高性能传输而非 TCP 套接字。用户自定义函数map_partitions是跨分区应用自定义逻辑的常用手段但它会产生一个不透明的 DataFrame 表达式阻止查询规划优化器做投影/过滤下推。由于列投影下推往往是最有效的优化应在调用map_partitions之前先选列、调用之后再选列也可以显式追加过滤操作来弥补过滤下推的损失。七、API 参考与文档入口Dask cuDF 总体上力求与 Dask DataFrame 提供完全一致的 API差异主要源于cuDF 并未完全镜像 pandas API以及 cuDF 在数据读写接口上提供了一些额外配置标志。因此简单的工作流迁移成本很低但利用更多特性的复杂工作流可能需要少量调整。api.rst 进一步说明了dask_cudf显式 API 的部分内容推荐用法仍是设置后端后使用 Dask DataFrame APIdask.config.set({dataframe.backend: cudf})除 Dask 通用的创建/存储接口外Dask cuDF 额外提供from_cudf与read_text等 cuDF 专属方法automodule直接来自 dask_cudf 包对于 Dask cuDF 不直接支持的磁盘格式官方建议两条路径先用 Dask 通用读取设施读入再通过to_backend(cudf)转换或使用from_map/from_delayed自行构造。文档入口的 toctreeindex.rst最终指向best_practices与api两页配合本文即可覆盖 Dask cuDF 文档的完整知识骨架。八、小结Dask cuDF 的技术价值在于把 cuDF 的 GPU DataFrame 能力无缝嵌入 Dask 生态通过 entry point 自动注册cudf后端pyproject.toml用一行配置或一个环境变量即可切换dask.config.set({dataframe.backend: cudf})自 24.06 起默认获得基于 Dask-Expressions 的自动查询规划投影/谓词下推simplify可验证并通过LocalCUDACluster的rmm_pool_size、enable_cudf_spill等参数把 RMM 内存池与 cuDF 原生溢出整合进多 GPU 部署。掌握分区大小调优 避免急执行 尽量留在 GPU 上这三条主线就能把 Dask cuDF 用得既快又稳。赞分享数据分析数据工程机器学习【免费下载链接】cudfcuDF - GPU DataFrame Library项目地址https://gitcode.com/gh_mirrors/cu/cudf点击查看免费下载相关推荐从单 GPU 到多 GPU 集群dask-cudf 并行计算完整指南dask-cuda 实战从单 GPU 到多 GPU 集群dask cudf 并行计算完整指南dask cuda 实战 dask cudf 是 GPU DataFrame 并数据分析数据工程机器学习Dask GPU 计算实战指南用 Delayed/Futures、CuPy 与 cuDF 构建 GPU 并行任务Dask GPU 计算实战指南用 Delayed/Futures、CuPy 与 cuDF 构建 GPU 并行任务 本指南围绕 docs/source/gpu.大数据数据分析任务调度dask-cuDF API 参考以 Dask DataFrame 后端方式使用 cuDF 并行 DataFrame 的完整指南dask cuDF API 参考以 Dask DataFrame 后端方式使用 cuDF 并行 DataFrame 的完整指南 本文基于仓库中的 API re数据分析数据工程机器学习上一篇终极指南Cutter高级搜索功能如何快速定位二进制文件中的关键模式与特征下一篇MJML邮件模板国际化多语言响应式邮件开发终极指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表