拓冰建站拓冰建站
首页 / 资讯中心 / 正文

Modin pandas on Dask 使用指南:引擎配置、集群运行与 DataFrame 互转

数据分析数据工程大数据【免费下载链接】modinModin: Scale your Pandas workflows by changing a single line of code项目地址https://gitcode.com/gh_mirrors/mo/modin点击查看免费下载本篇技术指南以仓库文档 docs/development/using_pandas_on_dask.rst 为核心骨架系统讲解在 Modin 中启用 Dask 执行引擎的完整路径从环境变量与代码配置的两种开启方式到单机本地运行、集群部署运行再到 Modin DataFrame 与 Dask DataFrame 之间的双向无拷贝转换。读完本文你将掌握 pandas on Dask 的完整配置姿势、to_dask/from_dask互转 API 的底层实现原理与使用限制并能在自己的机器或 Dask 集群上直接复现。背景Modin 的多引擎架构与 pandas 存储格式Modin 是一个通过修改一行 import 代码即可横向扩展 pandas 工作流的分布式 DataFrame 框架。其核心设计将执行引擎Execution Engine与存储格式Storage Format解耦存储格式Storage Format决定底层分区的内存表示。Modin 默认使用 pandas 作为底层分区的主内存格式并对从 API 层进入的查询做针对性优化因此绝大多数场景下你无需关心如何选择但也可以显式指定执行引擎Engine决定查询被分发到哪类分布式运行时上执行可选项包括 Ray、Dask、Unidist、Python 与 Native。本指南聚焦的pandas on Dask即pandas 存储格式 Dask 执行引擎的组合每个分区内部是一块 pandas DataFrame而分区的调度、物化与合并交给 Dask 的 Task Graph 与分布式调度器完成。仓库中该组合的实现位于 modin/core/execution/dask/implementations/pandas_on_dask/其中PandasOnDaskDataframe、PandasOnDaskDataframePartition与PandasOnDaskIO分别对应分区结构、分区粒度和 IO 入口。启用 pandas on Dask两种配置方式开启 pandas on Dask 执行需要同时设置引擎与存储格式两个配置项官方文档给出了两种等价方式。方式一环境变量在启动 Python 进程之前通过 shell 导出两个环境变量export MODIN_ENGINEdask export MODIN_STORAGE_FORMATpandas从源码看这两个变量分别对应 modin/config/envvars.py 中的两个配置类Engine类的varname MODIN_ENGINE其合法取值choices为(Ray, Dask, Python, Unidist, Native)且大小写不敏感配置框架会做normalizeStorageFormat类的varname MODIN_STORAGE_FORMAT用于声明底层的分区存储格式。方式二源代码内配置在 Python 代码中通过modin.config模块的put方法在运行时切换import modin.config as cfg cfg.Engine.put(dask) cfg.StorageFormat.put(pandas)值得注意的是Engine.put内部会联动设置Backend它会根据执行引擎 存储格式的组合推导出对应的后端实现再完成分发器的绑定见 envvars.py 中Engine.put的实现。因此实践中通常只需保证上述两个配置项一致地指向 pandas on Dask 组合即可。补充默认引擎的选择逻辑如果你既不设置环境变量也不调用putModin 会按一定优先级自动探测可用的执行引擎。从 envvars.py 的Engine._get_default实现可以看出探测顺序为若启用了MODIN_DEBUG或存在自定义引擎则回退到Python依次尝试导入ray、dask distributed、unidist首个满足最低版本要求的引擎即为默认值MIN_RAY_VERSION、MIN_DASK_VERSION、MIN_UNIDIST_VERSION定义于 modin/utils.py若三者均不可用则抛出提示安装对应引擎的ImportError例如pip install modin[dask]。因此只要你安装了 Dask 且未安装 Ray/UnidistMODIN_ENGINEdask甚至可能不是必须的但为了可读性与确定性显式声明依然是推荐做法。本地单机运行一行代码切换即可在单节点本地环境中使用 pandas on Dask 非常简单设置引擎后继续像使用 pandas 一样操作 Modin DataFrame 即可。文档给出的最小示例import modin.pandas as pd import modin.config as modin_cfg modin_cfg.Engine.put(dask) df pd.DataFrame(...) # 之后的操作照常进行本地模式下有两种 Dask 客户端初始化策略由 Modin 自动初始化不显式创建 ClientModin 会在首次需要调度时自行启动本地 Dask 集群LocalCluster/默认 Client自行初始化 Client你可以先from distributed import Client; client Client()建立本地集群Modin 会自动检测并复用当前进程的默认 Client而不是另起炉灶。这一点与PandasOnDaskIO中大量使用default_client()的实现细节一致详见下文转换一节。集群运行连接已部署的 Dask 集群当数据规模超出单机、需要横向扩展时可以对接一个独立部署的 Dask 集群Scheduler 多台 Worker。核心思路是先建立 Dask Client再启用 Modin 的 Dask 引擎此后 Modin 的所有分布式任务都会提交到你指定的集群上执行。from distributed import Client import modin.pandas as pd import modin.config as modin_cfg # 在这里定义你的集群例如 LocalCluster、SSHCluster、Kubernetes/KubeCluster 等 cluster ... client Client(cluster) modin_cfg.Engine.put(dask) df pd.DataFrame(...)要点说明Client(cluster)一旦创建便成为当前进程的默认客户端Modin 的 Dask 后端会通过distributed.client.default_client()拿到该客户端并把任务图提交给它见 io.py 的from_dask实现。集群的部署与运维方式SSH、Kubernetes、YARN、Cloud 等属于 Dask 生态的范畴官方文档原文也指向了 Dask 官方的集群部署文档本文不再展开。相比本地模式集群模式下唯一多出的是先创建 Client这一步其余 API 使用完全一致这正是改一行代码设计哲学的体现。Modin 与 Dask DataFrame 双向转换这是 pandas on Dask 最具实用价值的能力之一Modin DataFrame 可以与 Dask DataFrame 做无拷贝no-copy的按分区转换从而在同一个任务流中同时利用两套库的优势——例如在 Modin 中完成 pandas 风格的高层查询再把结果交给 Dask 生态做进一步的数据管道编排。公开 API 与最小示例文档给出的转换示例import modin.pandas as pd import modin.config as modin_cfg from modin.pandas.io import to_dask, from_dask modin_cfg.Engine.put(dask) df pd.DataFrame(...) # Modin DataFrame - Dask DataFrame dask_df to_dask(df) # Dask DataFrame - Modin DataFrame modin_df from_dask(dask_df)除此之外Modin DataFrame 还提供了.modin.to_dask()访问器写法。仓库测试 modin/tests/pandas/test_io.py 中的用例即使用该写法dask_df modin_df.modin.to_dask() df_equals(dask_df.compute(), pandas_df)to_dask/from_dask的函数定义位于 modin/pandas/io.pyfrom_dask在 L1096to_dask在 L1218两者都经FactoryDispatcher分发到当前启用的后端实现上。底层实现无拷贝分区转换的原理Modin → Daskto_dask核心实现位于 PandasOnDaskIO.to_dask流程为通过unwrap_partitions(modin_obj, axis0)取出 Modin DataFrame 按行划分的各分区该函数定义于 modin/distributed/dataframe/pandas/partitions.py若转换的是Series还会用df_to_series将每个分区单列表提交为真正的 pandas Series并处理未命名 Series 标签MODIN_UNNAMED_SERIES_LABEL最后调用dask.dataframe.from_delayed(partitions)把每个 Modin 分区直接封装为 Dask 的 Delayed 分区——分区本身不复制数据只是换了一层调度包装这正是无拷贝转换的含义。Dask → Modinfrom_dask核心实现位于 PandasOnDaskIO.from_dask流程为通过default_client()获取当前 Dask Client调用dask_obj.to_delayed()将 Dask DataFrame 展开为 Delayed 列表再用client.compute(...)提前物化为一批分区引用调用from_partitions(dask_futures, axis0)同样定义于 partitions.py按行拼接为一个 Modin DataFrame 的查询编译器query_compiler并返回。整个链路是公开函数modin.pandas.io.from_dask→FactoryDispatcher.from_daskdispatcher.py→BaseFactory._from_daskfactories.py→PandasOnDaskIO.from_dask。重要限制仅在 Dask 引擎下可用to_dask/from_dask只有在 Modin 使用 Dask 引擎时才有效这一点有两层源码保障在基类 modin/core/io/io.py 中BaseIO.from_dask与BaseIO.to_dask默认实现直接raise RuntimeError错误信息明确写着Modin DataFrame can only be converted to a Dask DataFrame if Modin uses a Dask engine只有 Dask 后端的PandasOnDaskIO覆盖了这两个方法仓库测试 test_io.py 中的test_df_to_dask、test_series_to_dask、test_from_dask均带有pytest.mark.skipif(conditionEngine.get() ! Dask, ...)跳过标记进一步印证了这一前提。因此若你的进程使用 Ray 或 Python 引擎调用互转 API 会抛出RuntimeError请务必先执行modin_cfg.Engine.put(dask)或导出MODIN_ENGINEdask再使用转换功能。工程实践建议与小结综合官方文档与仓库实现使用 pandas on Dask 时可遵循以下实践路径安装依赖确保安装了满足最低版本要求的dask与distributedModin 会在导入时校验版本不满足则提示pip install modin[dask]配置引擎在 shell 中导出MODIN_ENGINEdask、MODIN_STORAGE_FORMATpandas或在代码入口处调用cfg.Engine.put(dask)本地调试不建 Client 直接使用 Modin API让 Modin 自动初始化 Dask集群扩展先Client(cluster)接入既有 Dask 集群再执行 pandas 风格的 DataFrame 操作生态互操作用to_dask/from_dask或.modin.to_dask()在 Modin 与 Dask 之间按分区无拷贝切换各取所长。其他执行引擎Ray、Python、Unidist、MPI的用法在 docs/development/ 目录下另有对应文档如 using_pandas_on_ray.rst、using_pandas_on_python.rst其配置模型与本指南完全同构掌握了 pandas on Dask 之后即可举一反三。赞分享数据分析数据工程大数据【免费下载链接】modinModin: Scale your Pandas workflows by changing a single line of code项目地址https://gitcode.com/gh_mirrors/mo/modin点击查看免费下载相关推荐Modin pandas on Ray 使用指南以 Ray 为执行引擎的分布式 Pandas 工作负载Modin pandas on Ray 使用指南以 Ray 为执行引擎的分布式 Pandas 工作负载 Modin 是一款通过极简改动即可横向扩展 Panda数据分析数据工程大数据Modin 执行引擎初始化与配置完全指南Ray、Dask 与 MPIunidist实战Modin 执行引擎初始化与配置完全指南Ray、Dask 与 MPIunidist实战 Modin 作为一个通过只改一行 import即可将 pand数据分析数据工程大数据在 Ray 集群上运行 Dask 工作负载Dask-on-Ray 调度器使用与原理完全指南在 Ray 集群上运行 Dask 工作负载Dask on Ray 调度器使用与原理完全指南 本指南基于 Ray 官方文档《Using Dask on Ray》人工智能分布式训练强化学习任务调度模型推理服务后端上一篇core-js 支持的引擎与兼容性数据全解core-js-compat 的原理与实战用法下一篇Beads bd count 命令完整指南过滤、分组计数与底层实现创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门