新闻详情

新闻详情

首页 / 资讯中心 / 详情

cuDF-Polars DaskEngine 完全指南:在 Dask 集群上运行 GPU 流式查询

发布时间:2026/9/25 4:13:45来源:尧图网络
cuDF-Polars DaskEngine 完全指南:在 Dask 集群上运行 GPU 流式查询
数据分析数据工程机器学习【免费下载链接】cudfcuDF - GPU DataFrame Library项目地址https://gitcode.com/gh_mirrors/cu/cudf点击查看免费下载cuDF-Polars 是 NVIDIA cuDF 仓库python/cudf_polars中面向 Polars 用户的 GPU 查询引擎而DaskEngine是其中基于 Dask distributed 的多 GPU 后端。本文以官方文档 dask_engine.md 为主体结合引擎源码与测试系统讲解 DaskEngine 的集群模型、配置方式、与dask-cuda等启动器的协同规则以及集群诊断与生命周期管理。读完本文你将能够在自己的 Dask 集群单机多卡或 HPC 多节点上跑起分布式 GPU 流式查询并正确处理硬件绑定、内存资源与 UCXX 通信的冲突问题。DaskEngine 的集群模型一 GPU 一 Worker共享 UCXX 通信DaskEngine源码位于 dask.py类定义见 L935 起在 Dask distributed 集群上运行 cuDF-Polars 的流式执行器每个 Dask worker 对应一块 GPU由一个客户端client进程统一协调。分区partition数据沿查询计划流式推进shuffle、allgather、join 等集体操作在所有 worker 之间通过共享的UCXX communicator完成。启动时每个 worker 会被绑定到与其 GPU 最近的 CPU 核心与 NUMA 节点详见下文硬件绑定。与 RayEngine、SPMDEngine 一样DaskEngine 底层共用同一套流式执行器四者的差异只在 GPU worker 的供给方式。DaskEngine 适合团队中已有 Dask 部署或偏好 Dask 启动器的场景参见 engines.md 中的引擎对比表。快速开始零参数启动本地集群DaskEngine()不带任何参数时会自动创建一个distributed.LocalCluster按每个可见 GPU 一个 worker的规格启动一个distributed.Client在所有 worker 间引导bootstrap一个 UCXX communicator。退出时引擎创建的所有资源都会被拆除。最小示例import polars as pl from cudf_polars.engine.dask import DaskEngine with DaskEngine() as engine: result ( pl.scan_parquet(/data/dataset/*.parquet) .filter(pl.col(amount) 100) .group_by(customer_id) .agg(pl.col(amount).sum()) .collect(engineengine) ) print(result)从源码看自动创建集群的逻辑在 dask.py引擎读取CUDA_VISIBLE_DEVICES未设置则通过 NVML 探测全部 GPU见_get_visible_gpu_ids然后构造distributed.SpecCluster为每个 GPU 启动一个distributed.Nannyworker并设置nthreads1、memory_limit为系统内存上限、CUDA_VISIBLE_DEVICES指向对应 GPU。因此单卡机器上跑这段代码就是单 worker 单 GPU多卡机器则会自动扩展到所有可见 GPU。关于.collect()的重要提醒.collect()会把完整结果拉回客户端进程。对于大规模分布式输出应优先使用.sink_*()让每个 rank 直接写入自己的分区文件或在查询内部先聚合/采样再.collect()。详见 engines.md 的结果收集一节。用 StreamingOptions 配置 DaskEngine需要对引擎做自定义配置时构建StreamingOptions定义于 options.py完整字段说明见 options.md再通过DaskEngine.from_options()传入import polars as pl from cudf_polars.engine.options import StreamingOptions from cudf_polars.engine.dask import DaskEngine opts StreamingOptions(num_streaming_threads8, fallback_modesilent) with DaskEngine.from_options(opts) as engine: result pl.scan_parquet(/data/dataset/*.parquet).collect(engineengine)StreamingOptions的字段分为三个类别每类都有对应的环境变量前缀类别作用域环境变量前缀executor查询执行与分区如max_rows_per_partition、fallback_modeCUDF_POLARS__EXECUTOR__enginepl.GPUEngine的 kwargs如 Parquet 配置、内存资源、CUDA 流、硬件绑定CUDF_POLARS__rapidsmpf流式运行时如线程数、CUDA 流、spill、pinned 内存、日志级别RAPIDSMPF_优先级规则每个选项都有对应的环境变量。选项未显式设置时先从环境变量读取环境变量也不存在时由底层库应用内置默认值。布尔型环境变量接受{1, true, yes, y}为真、{0, false, no, n}为假。StreamingOptions的_opt工厂options.py正是按显式值 → 环境变量 → 内置默认的优先级实现这一行为。常用配置速查executor类别的关键字段完整表格见 options.md字段说明默认值num_py_executors内部 PythonThreadPoolExecutor的 worker 数8fallback_mode遇到不支持的操作被迫回退到 CPU 执行时的行为warn/raise/silentwarnmax_rows_per_partition每个分区的最大行数仅内存 DataFrame 源使用磁盘 IO 与动态规划不使用1_000_000broadcast_limitbroadcast join 的最大字节数autotarget_partition_sizeIO 与动态规划的目标分区字节数0表示自动autodynamic_planning动态规划配置None禁用启用sink_to_directory.sink_*()是否写目录dask引擎恒为True传False会抛ValueErrorTrueengine类别的关键字段字段说明默认值raise_on_fail出错时抛异常而不是回退 CPUFalseparquet_optionsParquet 配置dict 或ParquetOptions—memory_resource_configRMM 内存资源配置—hardware_binding硬件绑定策略可传HardwareBindingPolicy精细控制HardwareBindingPolicy()allow_gpu_sharingFalse默认时多个 rank 共享同一块物理 GPU 会抛错Falserapidsmpf类别是流式运行时的底层旋钮绝大多数用户不需要直接触碰例如num_streaming_threads执行协程的线程数默认 1、num_streams并发 GPU 执行的 CUDA 流数默认 16、spill_device_limitspill 前的显存软上限默认80%、pinned_memorypinned 主机内存默认启用。从字典构建配置StreamingOptions.from_dict()接受扁平字典便于从配置文件或 CLI 读取选项未知键会抛TypeErrorNone值表示字段未指定opts StreamingOptions.from_dict({ num_streaming_threads: 8, fallback_mode: silent, })from_dict的实现options.py会把None值转换为UNSPECIFIED从而让字段继续走环境变量 → 内置默认的解析链。StreamingOptions.to_dict()与from_dict()可无损往返适合序列化保存配置。自带 Dask client接入已有集群如果你已有一个正在运行的 Dask 集群例如多节点部署的 scheduler通过dask_client传入现有的distributed.Client即可接入from distributed import Client import polars as pl from cudf_polars.engine.dask import DaskEngine with Client(scheduler-address:8786) as dc: with DaskEngine(dask_clientdc) as engine: result pl.scan_parquet(/data/*.parquet).collect(engineengine)传入 client 时DaskEngine在退出时不会关闭该 client 及其集群——生命周期完全由调用方管理对应 dask.py 中仅关闭自有的owned_client/owned_cluster的逻辑。这也是多节点 HPC 部署的标准接入方式。预配置 GPU 集群关闭内置硬件绑定某些 Dask 启动器最典型的是dask_cuda.LocalCUDACluster已经为每个 worker 固定了 CPU affinity 并设置了CUDA_VISIBLE_DEVICES。此时应通过HardwareBindingPolicy禁用DaskEngine的内置硬件绑定避免两层绑定互相打架——后执行者胜出会导致最终放置位置不确定from dask_cuda import LocalCUDACluster from distributed import Client from cudf_polars.engine.dask import DaskEngine from cudf_polars.engine.hardware_binding import ( HardwareBindingPolicy, ) with Client(LocalCUDACluster()) as dc, DaskEngine( dask_clientdc, engine_options{ hardware_binding: HardwareBindingPolicy(enabledFalse), }, ) as engine: ...HardwareBindingPolicy是定义于 hardware_binding.py 的冻结 dataclass字段如下字段说明默认值enabled是否启用绑定False禁用全部绑定Truecpu是否绑定 CPU 核心Truememory是否绑定 NUMA 内存节点Truenetwork是否绑定网络设备默认关闭因为 UCX 通常能自动识别合适的 NICFalseskip_under_rrun进程由rrun启动时跳过绑定rrun 启动时已绑定Trueenable_once每个进程至多绑定一次Trueraise_on_fail绑定失败如 affinity、NUMA、拓扑发现失败是否抛异常False绑定目标 GPU 由CUDA_VISIBLE_DEVICES解析未设置时回退到 GPU 0。每个前端Dask 通过 nanny preload 或 SpecCluster、Ray 通过num_gpus1、SPMD 通过rrun负责为每个 worker 设置该变量。bind_to_gpu是线程安全的包装hardware_binding.py当enable_onceTrue时通过双重检查锁保证整个进程内底层bind()至多执行一次。绑定行为有对应测试覆盖见 test_bind_to_gpu.py。手动启动 workernanny preload 分配 GPU在多节点 HPC 集群手动启动 worker 时使用内置的nanny preload为每个 worker 分配一块 GPU。preload 会在 worker 子进程 spawn之前设置CUDA_VISIBLE_DEVICES# 每个节点上为每块 GPU 启动一个 worker每个 worker 单线程 dask worker SCHEDULER_ADDRESS:8786 --nworkers N --nthreads 1 \ --preload-nanny cudf_polars.engine.dask然后在客户端连接import polars as pl from distributed import Client from cudf_polars.engine.dask import DaskEngine with Client(SCHEDULER_ADDRESS:8786) as dc: with DaskEngine(dask_clientdc) as engine: result pl.scan_parquet(/data/*.parquet).collect(engineengine)注意职责划分硬件绑定CPU affinity、NUMA、网络由DaskEngine自动处理nanny preload 只负责 GPU 分配。从源码看preload 入口是dask_setup函数dask.py它必须配合--preload-nanny使用与--preload不兼容否则抛TypeError。其工作机制是读取当前节点的可见 GPU 列表用模块级计数器做轮询分配把第 i 个 worker 的CUDA_VISIBLE_DEVICES设为第i % len(gpu_ids)块 GPU。由于函数运行在 Nanny 进程内、先于 worker 子进程 spawn设置的环境变量会被 worker 继承。若 worker 数超过 GPU 数多个 worker 会共享 GPU此时需注意allow_gpu_sharing选项。使用dask-cuda-worker作为 nanny preload 的替代方案也可以使用 dask-cuda 项目提供的dask-cuda-worker启动 worker。它为每块可见 GPU 启动一个 worker并给每个 worker 安装一组插件CPUAffinity将 worker 固定到其 GPU 所在 NUMA 节点、RMMSetup以及一个配置 UCX 的 nanny preload。由于DaskEngine会为自己的流式运行时设置同样的东西两者必须协调否则会冲突遵循以下三条规则CPU affinity 在dask-cuda-worker中是强制的CPUAffinity插件总是安装且没有 CLI 开关可关闭。因此必须给DaskEngine传hardware_bindingHardwareBindingPolicy(enabledFalse)避免在 dask-cuda 已绑定之后再重复固定 affinity。不要给dask-cuda-worker传--rmm-pool-size、--rmm-managed-memory等 RMM 标志让DaskEngine通过其memory_resource_config见 options.md独占内存资源的管理权否则同一 worker 上会安装两套不同的内存资源。不要传--enable-tcp-over-ucx、--enable-infiniband、--enable-nvlink、--enable-rdmacmDaskEngine会自行引导 UCXX communicator 并选择传输层两侧同时启用可能产生不一致的 UCX 配置。# 每个节点上仅做 GPU 分配 CPU affinity不传 RMM、不传 UCX 标志 dask-cuda-worker SCHEDULER_ADDRESS:8786import polars as pl from distributed import Client from cudf_polars.engine.dask import DaskEngine from cudf_polars.engine.hardware_binding import ( HardwareBindingPolicy, ) with Client(SCHEDULER_ADDRESS:8786) as dc: with DaskEngine( dask_clientdc, engine_options{ # dask-cuda-worker 总是固定 CPU affinity禁用 DaskEngine 的 # 绑定避免两者冲突。 hardware_binding: HardwareBindingPolicy(enabledFalse), }, ) as engine: result pl.scan_parquet(/data/*.parquet).collect(engineengine)集群诊断查看每个 worker 的放置信息DaskEngine.gather_cluster_info()返回每个 worker 的放置信息hostname、pid、CUDA_VISIBLE_DEVICESwith DaskEngine() as engine: print(fcluster has {engine.nranks} workers) for info in engine.gather_cluster_info(): print( fhostname{info[hostname]}, pid{info[pid]}, fCUDA_VISIBLE_DEVICES{info[cuda_visible_devices]} )gather_cluster_info()的实现dask.py通过client.run把_get_cluster_info广播到所有 worker每个 worker 返回(rank, ClusterInfo)对客户端按 rank 排序汇总。同类 API 还有gather_statistics(clearFalse)收集每个 rank 的 RapidsMPF 运行时统计gather_io_summary(clearFalse)收集每个 rank 的 KvikIO I/O 统计需启用kvikio_statistics选项。这两者内部同样走_run_by_rank的逐 worker 执行 按 rank 排序模式dask.py。另外注意在rrun集群内部创建DaskEngine会直接抛RuntimeErrordask.py。这是因为 rrun 场景应由SPMDEngine处理——正确做法是单独启动 rrun 集群再让客户端连接它的节点。引擎生命周期context manager 与手动 shutdown推荐在脚本中使用 context manager 形式它保证即使抛异常worker 也会被拆除。在 Jupyter 等交互环境中with块无法跨 cell可以构造一次、复用多次最后手动engine.shutdown()# Cell 1: 启动引擎 from cudf_polars.engine.dask import DaskEngine engine DaskEngine()# Cell 2: 执行查询可复用同一引擎跑多个查询 import polars as pl result pl.scan_parquet(/data/*.parquet).collect(engineengine)# 最后: 拆除所有资源 engine.shutdown()shutdown()dask.py会先通过client.run在所有 worker 上执行_teardown_worker停止 KvikIO 监控、关闭线程池、释放 streaming context 与 communicator、删除持久化分区然后关闭引擎自有的 client 与 cluster。它是幂等的——重复调用安全若 worker 拆除时出现异常会汇总为ExceptionGroup抛出。更精细的持续查询partition 留在各自 worker 的 GPU 上、可继续.lazy()链式查询可用engine.execute(lf)参见 dask.py。深入原理DaskEngine 的 UCXX 引导流程从源码可以还原DaskEngine.__init__的完整启动序列dask.py阶段一_setup_root只在首个 worker 上运行。创建 rank 0 的 UCXX communicator、RMM 内存资源与统计对象并把序列化的 root UCXX 地址返回给客户端。绑定硬件bind_to_gpu也在此完成。阶段二_setup_worker在所有 worker 上并发执行包括 root。每个 worker 用 root 地址创建自己的 communicator然后全体进入barrier()同步之后各自构建流式Context、创建ThreadPoolExecutor线程名前缀dask-executor并通过rmm.mr.set_current_device_resource让 libcudf 的临时分配也使用同一内存资源。查询执行时evaluate_pipeline_dask_modedask.py预降级的 IR 通过client.run分发到所有 workerrank 0 收集 scan 统计并 allgather所有 rank 独立降级查询图然后各自在本机 GPU 上执行流式管线最后在客户端pl.concat各 rank 输出并按 rank 排序。这套引导与执行流程由 test_dask.py 等测试直接覆盖包括查询执行、统计收集与多 rank IO 场景test_statistics.py、test_io_multirank.py。更多参考options.mdStreamingOptions全部字段与环境变量engines.md四种引擎的执行模式与集群后端对比、结果收集规则usage.mdRayEngine 端到端入门from_options()/StreamingOptions模式对 DaskEngine 同样适用other_engines.mdSPMDEngine、默认单例引擎与内存引擎导航引擎源码dask.py、hardware_binding.py、options.py赞分享数据分析数据工程机器学习【免费下载链接】cudfcuDF - GPU DataFrame Library项目地址https://gitcode.com/gh_mirrors/cu/cudf点击查看免费下载相关推荐cudf-polars 使用指南用 RayEngine 在 GPU 上运行 Polars 查询cudf polars 使用指南用 RayEngine 在 GPU 上运行 Polars 查询 cudf polars 是 cuDF 项目提供的 Polars数据分析数据工程机器学习cuDF-Polars SPMDEngine 实战指南用 SPMD 模式在单机多 GPU 上运行 Polars 流式查询cuDF Polars SPMDEngine 实战指南用 SPMD 模式在单机多 GPU 上运行 Polars 流式查询 导读 本文围绕 cuDF Polar数据分析数据工程机器学习cuDF cudf-polars在 GPU 上执行 Polars 查询的引擎架构、配置与调优指南cuDF cudf polars在 GPU 上执行 Polars 查询的引擎架构、配置与调优指南 cuDF 为 Polars Lazy API 的 Pytho数据分析数据工程机器学习上一篇Plano 基于用户偏好的 LLM 路由用 routing_preferences 实现按意图自动选模的实战指南下一篇LMCache Kubernetes Operator 端到端冒烟测试套件与 Kubebuilder 开发规范实战指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
网站建设高端定制企业官网
RELATED

相关资讯

更多精彩内容,欢迎继续阅读

较早相关资讯

最新相关资讯

小喵V2电机驱动快速入门:简单积木实现4路电机调速与正反转控制 2026/9/25 4:55:58

小喵V2电机驱动快速入门:简单积木实现4路电机调速与正反转控制

小喵V2电机驱动快速入门:简单积木实现4路电机调速与正反转控制 【免费下载链接】miaow-v2 源师兄扩展项目: 小喵V2 | 由源师兄组织创建 项目地址: https://gitcode.com/yuanshixiong/miaow-v2 小喵V2是源师兄推出的 KittenBot 开源扩展项目,通过配…

阅读更多 →
VirtualBox嵌套虚拟化灰色锁定终极解决方案 2026/9/25 4:55:52

VirtualBox嵌套虚拟化灰色锁定终极解决方案

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →
视频剪辑素材宝藏库:可商用高清晰素材网站推荐与工作流整合 2026/9/25 4:55:52

视频剪辑素材宝藏库:可商用高清晰素材网站推荐与工作流整合

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →
边缘AI芯片选型:从场景约束反推技术方案 2026/9/25 4:55:33

边缘AI芯片选型:从场景约束反推技术方案

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →
高频变压器三明治绕法:原理、实操与EMI/效率优化 2026/9/25 4:55:33

高频变压器三明治绕法:原理、实操与EMI/效率优化

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →
网盘搜索引擎原理与实战:找资源不再靠运气 2026/9/25 4:55:33

网盘搜索引擎原理与实战:找资源不再靠运气

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →

今日资讯

本周资讯

本月资讯

看完文章仍有疑问?

联系尧图顾问,获取一对一建站咨询

立即免费咨询 📞 400-888-8888
📞 ✉