尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

multidplyr 数据分区终极技巧:partition() 与 party_df 两大核心组件,把百万行数据框送上多核

multidplyr 数据分区终极技巧:partition() 与 party_df 两大核心组件,把百万行数据框送上多核 multidplyr 数据分区终极技巧partition() 与 party_df 两大核心组件把百万行数据框送上多核【免费下载链接】multidplyrA dplyr backend that partitions a data frame over multiple processes项目地址: https://gitcode.com/gh_mirrors/mu/multidplyrmultidplyr 是 R 语言生态中的多进程 dplyr 后端它把百万行级数据框自动分区到多个计算核心并行处理让你用熟悉的 dplyr 语法获得接近线性的加速。本文带你快速掌握两大核心组件partition()与party_df()覆盖集群创建、分组分区、直接加载与性能调优的完整实操路线帮你的数据分析任务真正跑满多核。 30 秒认识 multidplyrmultidplyr 由 tidyverse 团队开发思路可以一句话概括把一张大表拆成 N 块交给 N 个独立的 R 工作进程worker同时计算算完再拼回来。维度传统 dplyrmultidplyr计算位置单个 R 进程多个独立 worker 进程数据存放全部在主会话内存分散在各 worker语法filter()/mutate()/summarise()完全相同无需改写法最佳场景中小数据量百万行以上 较耗时的计算它的设计灵感来自 partools 与 distributedR数据落座后尽量少搬动把通信开销压到最低把并行收益做到最高。 两大核心组件partition() 与 party_df组件一partition() —— 内存中的数据一键分区如果你的数据已经读进主会话内存partition()是最省心的入口。它把数据框切分到每个 worker返回一个特殊的party_dfpartitioned data frame分区数据框之后所有 dplyr 动词照常使用library(multidplyr) cluster - new_cluster(4) cluster_library(cluster, dplyr) flight_dest - flights %% group_by(dest) %% partition(cluster) flight_dest %% summarise(delay mean(dep_delay, na.rm TRUE), n n()) %% collect()组件二party_df() —— 让每个 worker 各读各的文件数据太大根本装不进主会话那就让每个 worker 直接读取自己那份文件主进程全程不碰原始数据这是效率最高的路线cluster_assign_partition(cluster, files files) # 文件均分给 worker cluster_send(cluster, my_data - vroom::vroom(files)) # 各自并行读取 my_data - party_df(cluster, my_data) # 声明为分区数据框 如何选择 partition 与 party_df对比项partition()party_df()数据位置已在主会话内存分散在文件由各 worker 自行加载数据传输需要一次性分发到各 worker几乎零传输最快适用场景已加载好的中等数据框超大数据、按文件切好的数据典型代码df %% partition(cluster)party_df(cluster, my_data)两者最终都得到同一类对象——party_df后续操作方式完全一致。 三步上手建集群 → 分区 → 回收建集群new_cluster(n)创建 n 个 worker 进程操作系统会自动把它们铺满多核分区数据用partition()或直读文件 party_df()把数据放到各 worker回收结果计算完成后用collect()把结果拉回主会话。library(multidplyr) cluster - new_cluster(4) # 1. 建集群 cluster_library(cluster, dplyr) # 各 worker 加载依赖包 # 2. 分区partition 或 party_df 路线 # 3. 计算并回收 df %% summarise(...) %% collect() 一个会话里建议只建一个集群反复使用建集群的成本不低而 worker 常驻可以省下大量初始化时间。⚡ partition() 进阶技巧分组决定分区质量partition()内置了一个贪心均衡算法见 R/partydf.R 中worker_id()先group_by()再分区同一分组的所有行一定落在同一个 worker 上group_by()/summarise()等分组计算无需跨进程协调自动负载均衡算法逐个分组始终把下一个分组分给当前行数最少的 worker让各分片行数尽量接近输出中Shards: 4 [81,594--86,548 rows]就是每片行数区间未分组数据则按行顺序轮转分配同样保持均衡。一句话经验分区前先想好分组键分区质量直接决定并行效率。 party_df 高效路线数据不过主进程走party_df()路线时有几个实用细节用cluster_assign_each()/cluster_assign_partition()给每个 worker 分配不同的文件列表用cluster_send()下发读文件指令各 worker 并行执行partition()与party_df()创建的中间分片支持auto_rm自动清理计算过程中产生的临时对象会在下次调用前自动回收不占 worker 内存打印party_df对象会直接显示总维度、分组和各分片行数区间排查数据分布一目了然。 性能调优什么时候值得上多核worker 数建议比核心数少 1~2 个可用parallel::detectCores()查看核心数留出核心给操作系统和其他任务小于约 1000 万行的简单操作filter、mutate等进程间通信开销会吃掉收益此时用 data.table 系后端更划算multidplyr 真正的杀手锏是复杂计算官方示例中用do()给每个分组拟合广义可加模型GAM并行相比本地计算从约 5 秒降到更快任务越耗时越接近线性加速依赖包用cluster_library(cluster, mgcv)一次性批量加载到所有 worker中间步骤尽量留在集群内计算只在最后一步collect()回收减少往返。 延伸阅读项目内核心资料资料路径说明分区核心源码R/partydf.Rpartition()、party_df()、打印与回收逻辑单表动词实现R/dplyr-single.Rfilter/mutate/summarise等并行版双表动词实现R/dplyr-dual.R各类 join 与集合运算集群管理R/cluster.R、R/cluster-utils.Rnew_cluster()与数据传输工具官方入门教程vignettes/multidplyr.Rmd完整示例 性能对比测试用例tests/testthat/各功能的用法参考功能更新记录NEWS.md版本演进总结数据在内存就选partition()数据在文件就选party_df()配合少而稳的 worker 数和复杂任务的并行化multidplyr 能帮你把百万行数据框的计算稳稳送上多核。【免费下载链接】multidplyrA dplyr backend that partitions a data frame over multiple processes项目地址: https://gitcode.com/gh_mirrors/mu/multidplyr创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表